Skip to main content

coven_database/
snapshot_objects.rs

1use crate::blob_records::live_blob_row;
2use crate::blob_records::load_activated_registration_on;
3use crate::blob_records::validate_live_blob_locator;
4use crate::blob_records::validate_stored_row_binding_on;
5use crate::remote_object_records::load_remote_object_on;
6use crate::PreparedSnapshotBlob;
7
8use super::*;
9
10pub(crate) fn validate_snapshot_object_owners_on(
11    conn: &Connection,
12    root: &coven_protocol::store_commit::StoreRootRef,
13    reference: &StoreSnapshotRef,
14    meta: &SnapshotMeta,
15) -> Result<(), DbError> {
16    load_activated_registration_on(conn, root, &meta.author_registration)?;
17    let expected = coven_protocol::remote_object::SnapshotObjectOwner::Store {
18        metadata_slot: reference.object.slot().clone(),
19    };
20    validate_snapshot_object_owner_records_on(
21        conn,
22        &expected,
23        &meta.history_summary.pending_device_join_snapshot_slots(),
24    )
25}
26
27pub fn validate_snapshot_author(
28    author: &StoreDeviceRegistrationRef,
29    local: &StoreDeviceRegistrationRef,
30    label: &str,
31) -> Result<(), DbError> {
32    if author == local {
33        Ok(())
34    } else {
35        Err(DbError::Message(format!(
36            "staged {label} snapshot author differs from local activation"
37        )))
38    }
39}
40
41pub fn validate_snapshot_image(
42    image: &SnapshotImageRef,
43    prepared: &PreparedExactObject,
44    plaintext_hash: ObjectHash,
45    stored_hash: ObjectHash,
46    stored_size: u64,
47    expected_slot: String,
48    label: &str,
49) -> Result<(), DbError> {
50    if image.object == *prepared.reference()
51        && plaintext_hash == image.image_hash
52        && stored_hash == image.object.stored_hash()
53        && stored_size == image.object.stored_size()
54        && image.object.slot().logical_key() == expected_slot
55    {
56        Ok(())
57    } else {
58        Err(DbError::Message(format!(
59            "staged {label} snapshot image differs from its exact reference"
60        )))
61    }
62}
63
64/// Load the exact image named by authenticated snapshot metadata. Both Store
65/// and Circle publication retain the plaintext and upload bytes under this ref.
66pub(crate) fn load_snapshot_image_on(
67    conn: &Connection,
68    store_dir: &coven_foundation::store_dir::StoreDir,
69    reference: &SnapshotImageRef,
70    label: &str,
71) -> Result<PreparedProtocolObject<Vec<u8>>, DbError> {
72    let value = crate::payload_store::read_payload_blocking(conn, store_dir, reference.image_hash)
73        .map_err(|error| DbError::context(format!("outbound {label} snapshot image"), error))?;
74    if ObjectHash::digest(&value) != reference.image_hash {
75        return Err(DbError::Message(format!(
76            "outbound {label} snapshot image differs from its exact hash"
77        )));
78    }
79    let stored = crate::payload_store::read_payload_blocking(
80        conn,
81        store_dir,
82        reference.object.stored_hash(),
83    )
84    .map_err(|error| {
85        DbError::context(format!("outbound prepared {label} snapshot image"), error)
86    })?;
87    let prepared = PreparedExactObject::new(reference.object.clone(), stored).map_err(|error| {
88        DbError::context(format!("outbound prepared {label} snapshot image"), error)
89    })?;
90    Ok(PreparedProtocolObject { value, prepared })
91}
92
93pub(crate) fn validate_snapshot_blob_plans_on(
94    conn: &Connection,
95    gates: &Gates,
96    synced_tables: &[SyncedTable],
97    owner: &coven_protocol::remote_object::SnapshotObjectOwner,
98    blobs: &[PreparedSnapshotBlob],
99) -> Result<(), DbError> {
100    for blob in blobs {
101        blob.remote
102            .validate()
103            .map_err(|error| DbError::context("snapshot remote blob", error))?;
104        let owners = blob.remote.snapshot_owners().collect::<Vec<_>>();
105        if owners != [owner] {
106            return Err(DbError::Message(
107                "snapshot blob owner differs from the prepared snapshot".to_string(),
108            ));
109        }
110        if blob.bindings.is_empty()
111            || blob.bindings.iter().any(|binding| {
112                binding.blob().object() != blob.remote.object()
113                    || binding.blob().locator().audience() != blob.authority.remote_audience()
114            })
115        {
116            return Err(DbError::Message(
117                "snapshot blob plan has inconsistent exact references".to_string(),
118            ));
119        }
120        for binding in &blob.bindings {
121            let table = synced_tables
122                .iter()
123                .find(|table| table.name() == binding.table())
124                .ok_or_else(|| {
125                    DbError::Message(format!(
126                        "snapshot blob names undeclared table {:?}",
127                        binding.table()
128                    ))
129                })?;
130            let declaration = table.blob().ok_or_else(|| {
131                DbError::Message(format!(
132                    "snapshot blob names table {:?} without a blob declaration",
133                    table.name()
134                ))
135            })?;
136            let row = live_blob_row(conn, table.name(), binding.row_id(), declaration)?
137                .ok_or_else(|| {
138                    DbError::Message(format!(
139                        "snapshot blob row {:?}/{:?} is absent",
140                        table.name(),
141                        binding.row_id()
142                    ))
143                })?;
144            let audience = gate::live_row_audience(conn, gates, table.name(), binding.row_id())
145                .map_err(DbError::from)?;
146            let audience = RemoteAudience::try_from(audience).map_err(DbError::from)?;
147            validate_live_blob_locator(
148                binding.table(),
149                binding.row_id(),
150                binding.column(),
151                binding.row_stamp(),
152                binding.blob(),
153                declaration,
154                &row,
155                &audience,
156            )?;
157        }
158    }
159    Ok(())
160}
161
162pub(crate) fn persist_snapshot_image_on(
163    conn: &Connection,
164    store_dir: &coven_foundation::store_dir::StoreDir,
165    image: &SnapshotImageRef,
166    owner: coven_protocol::remote_object::SnapshotObjectOwner,
167    label: &str,
168) -> Result<(), DbError> {
169    let image = RemoteObjectRecord::snapshot_activated_image(image, owner)
170        .map_err(|error| DbError::context(format!("{label} ownership"), error))?;
171    persist_exact_remote_object_on(conn, store_dir, &image, label)
172}
173
174/// Persist this snapshot candidate's exact membership rollup ownership.
175pub(crate) fn persist_membership_rollup_on(
176    conn: &Connection,
177    store_dir: &coven_foundation::store_dir::StoreDir,
178    rollup: &coven_protocol::store_commit::MembershipRollupRef,
179    owner: coven_protocol::remote_object::SnapshotObjectOwner,
180    label: &str,
181) -> Result<(), DbError> {
182    let rollup = RemoteObjectRecord::snapshot_activated_membership_rollup(rollup, owner)
183        .map_err(|error| DbError::context(format!("{label} ownership"), error))?;
184    persist_exact_remote_object_on(conn, store_dir, &rollup, label)
185}
186
187pub fn snapshot_generation_as_i64(generation: u64, label: &str) -> Result<i64, DbError> {
188    i64::try_from(generation)
189        .map_err(|_| DbError::Message(format!("{label} generation exceeds SQLite INTEGER")))
190}
191
192pub(crate) fn validate_snapshot_object_owner_records_on(
193    conn: &Connection,
194    expected: &coven_protocol::remote_object::SnapshotObjectOwner,
195    pending_store_snapshots: &BTreeSet<coven_protocol::objects::ObjectSlot>,
196) -> Result<(), DbError> {
197    let mut statement = conn
198        .prepare("SELECT object_id FROM remote_objects ORDER BY object_id")
199        .map_err(DbError::from)?;
200    let object_ids = statement
201        .query_map([], |row| row.get::<_, String>(0))
202        .map_err(DbError::from)?
203        .collect::<Result<Vec<_>, _>>()
204        .map_err(DbError::from)?;
205    drop(statement);
206    for object_id in object_ids {
207        let parsed = object_id.parse().map_err(|error| {
208            DbError::context(format!("snapshot remote object id {object_id:?}"), error)
209        })?;
210        let remote = load_remote_object_on(conn, parsed)?;
211        for owner in remote.snapshot_owners() {
212            let matches = match (owner, expected) {
213                (
214                    coven_protocol::remote_object::SnapshotObjectOwner::Store { metadata_slot },
215                    coven_protocol::remote_object::SnapshotObjectOwner::Store { .. },
216                ) => owner == expected || pending_store_snapshots.contains(metadata_slot),
217                (
218                    coven_protocol::remote_object::SnapshotObjectOwner::Circle {
219                        activation,
220                        generation,
221                    },
222                    coven_protocol::remote_object::SnapshotObjectOwner::Circle {
223                        activation: expected_activation,
224                        generation: expected_generation,
225                    },
226                ) => activation == expected_activation && generation <= expected_generation,
227                _ => owner == expected,
228            };
229            if !matches {
230                return Err(DbError::Message(format!(
231                    "snapshot remote object {object_id} belongs to another snapshot"
232                )));
233            }
234        }
235    }
236    Ok(())
237}
238
239pub(crate) fn replace_snapshot_object_owners_on(
240    conn: &Connection,
241    owner: &coven_protocol::remote_object::SnapshotObjectOwner,
242    blobs: &[PreparedSnapshotBlob],
243    pending_store_snapshots: &BTreeSet<coven_protocol::objects::ObjectSlot>,
244) -> Result<(), DbError> {
245    let mut statement = conn.prepare("SELECT object_id FROM remote_objects ORDER BY object_id")?;
246    let object_ids = statement
247        .query_map([], |row| row.get::<_, String>(0))?
248        .collect::<Result<Vec<_>, _>>()?;
249    drop(statement);
250    for object_id in object_ids {
251        let object_id = object_id
252            .parse()
253            .map_err(|error| DbError::context("snapshot remote object id", error))?;
254        let mut remote = load_remote_object_on(conn, object_id)?;
255        let retained = blobs
256            .iter()
257            .any(|blob| blob.remote.object_id() == object_id);
258        remote
259            .replace_snapshot_owners_for_image(retained.then_some(owner), pending_store_snapshots)
260            .map_err(|error| DbError::context("replace exported snapshot ownership", error))?;
261        update_remote_object_on(conn, object_id, &remote)?;
262    }
263    Ok(())
264}
265
266fn retain_snapshot_blob_on(conn: &Connection, blob: &PreparedSnapshotBlob) -> Result<(), DbError> {
267    let object_id = blob.remote.object_id();
268    let exists: bool = conn
269        .query_row(
270            "SELECT EXISTS(SELECT 1 FROM remote_objects WHERE object_id = ?1)",
271            [object_id.to_string()],
272            |row| row.get(0),
273        )
274        .map_err(DbError::from)?;
275    let merged = if exists {
276        let mut existing = load_remote_object_on(conn, object_id)?;
277        for owner in blob.remote.snapshot_owners() {
278            existing
279                .merge_snapshot_owner(blob.bindings[0].blob(), owner.clone())
280                .map_err(|error| DbError::context("merge snapshot blob owner", error))?;
281        }
282        existing
283    } else {
284        blob.remote.clone()
285    };
286    let encoded = serde_json::to_string(&merged)
287        .map_err(|error| DbError::context("serialize snapshot blob", error))?;
288    conn.execute(
289        "INSERT INTO remote_objects (object_id, state) VALUES (?1, ?2)
290         ON CONFLICT(object_id) DO UPDATE SET state = excluded.state",
291        rusqlite::params![object_id.to_string(), encoded],
292    )
293    .map_err(DbError::from)?;
294    crate::blob_records::record_stored_locator_on(conn, blob.bindings[0].blob())?;
295    Ok(())
296}
297
298/// Install the snapshot's row bindings only into its own database image.
299pub(crate) fn install_snapshot_blob_plan_on(
300    conn: &Connection,
301    blob: &PreparedSnapshotBlob,
302) -> Result<(), DbError> {
303    retain_snapshot_blob_on(conn, blob)?;
304    let object_id = blob.remote.object_id();
305    let authority = serde_json::to_string(&blob.authority)
306        .map_err(|error| DbError::context("serialize snapshot blob authority", error))?;
307    for binding in &blob.bindings {
308        conn.execute(
309            "INSERT INTO row_blob_locators
310         (table_name, row_id, column_name, row_stamp, audience_authority, remote_object_id)
311         VALUES (?1, ?2, ?3, ?4, ?5, ?6)
312         ON CONFLICT(table_name, row_id, column_name, row_stamp) DO NOTHING",
313            rusqlite::params![
314                binding.table(),
315                binding.row_id(),
316                binding.column(),
317                binding.row_stamp(),
318                authority,
319                object_id.to_string(),
320            ],
321        )
322        .map_err(DbError::from)?;
323        validate_stored_row_binding_on(conn, binding, &blob.authority, object_id)?;
324    }
325    Ok(())
326}
327
328/// Publication retains the snapshot's payloads without changing current row
329/// bindings: unpublished edits may have replaced or removed those rows.
330pub(crate) fn install_snapshot_blob_plans_on(
331    conn: &Connection,
332    blobs: &[PreparedSnapshotBlob],
333) -> Result<(), DbError> {
334    for blob in blobs {
335        retain_snapshot_blob_on(conn, blob)?;
336    }
337    Ok(())
338}
339
340/// A locally reconstructed cut retains only its live Store blobs for the
341/// accepted snapshot. Current rows may include later edits, and importing the
342/// cut's entire inventory would restore objects reclaimed after that cut.
343pub(crate) fn retain_reconstructed_snapshot_blobs_on(
344    conn: &Connection,
345    image: &Connection,
346    tables: &[SyncedTable],
347    snapshot: &StoreSnapshotRef,
348) -> Result<(), DbError> {
349    let gates = Gates::from_tables(image, tables)?;
350    let owner = coven_protocol::remote_object::SnapshotObjectOwner::Store {
351        metadata_slot: snapshot.object.slot().clone(),
352    };
353    for encoded in crate::query_mapped_rows(
354        image,
355        "SELECT remote_object_id FROM blob_locators ORDER BY remote_object_id",
356        [],
357        |row| row.get::<_, String>(0),
358    )? {
359        let object_id = encoded.parse()?;
360        let captured = load_remote_object_on(image, object_id)?;
361        let locator = crate::blob_records::carried_blob_locator(
362            &captured,
363            "reconstructed snapshot inventory",
364        )?;
365        if locator.audience() != RemoteAudience::Store {
366            continue;
367        }
368        let stored =
369            coven_protocol::blob::locator::StoredBlobRef::new(locator, captured.object().clone())?;
370        match Database::stored_blob_reference_state_on(image, &gates, tables, &stored)? {
371            crate::StoredBlobReferenceState::NotLiveRemote => continue,
372            crate::StoredBlobReferenceState::Unresolved => {
373                return Err(DbError::Message(format!(
374                    "reconstructed snapshot blob {object_id} has unresolved locality"
375                )));
376            }
377            crate::StoredBlobReferenceState::LiveRemote => {}
378        }
379        let mut remote = load_remote_object_on(conn, object_id)?;
380        remote.merge_snapshot_owner(&stored, owner.clone())?;
381        update_remote_object_on(conn, object_id, &remote)?;
382    }
383    Ok(())
384}