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
64pub(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
174pub(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
298pub(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
328pub(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
340pub(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}