1use crate::query_mapped_rows;
4use crate::store::store_session::StoreRecords;
5use std::collections::{BTreeMap, BTreeSet};
6
7use rusqlite::{types::Value, Connection};
8use serde::{Deserialize, Serialize};
9
10use crate::{
11 DbError, COVEN_INITIALIZED_STATE_KEY, COVEN_SCHEMA_MANIFEST_STATE_KEY,
12 STORE_DEVICE_GENESIS_STATE_KEY, SYNC_ROUTING_CONTRACT_STATE_KEY, SYNC_ROUTING_HASH_STATE_KEY,
13};
14use coven_protocol::membership::OWNER_PUBKEY_STATE_KEY;
15use coven_protocol::store_commit::{
16 CommitFrontier, ObjectHash, RetainedReplaySnapshotAuthority, StoreBatchCommitRef,
17 StoreDeviceRegistrationRef, StoreRootRef,
18};
19
20pub(crate) fn migrate_retained_replay_schema_on(
21 conn: &Connection,
22 store_dir: &coven_foundation::store_dir::StoreDir,
23 policy: crate::CovenMigrationPolicy,
24 migrations: &[crate::Migration],
25 synced_tables: &[coven_protocol::synced_schema::SyncedTable],
26) -> Result<(), crate::OpenError> {
27 let records = StoreRecords::new(conn, store_dir);
28 let Some(baseline) = load_replay_baseline_metadata_on(records)? else {
29 return Ok(());
30 };
31 let mut image = Connection::open_in_memory().map_err(DbError::from)?;
32 crate::connection_io::deserialize_database_image_into(
33 &mut image,
34 &baseline.image_bytes(conn, store_dir)?,
35 )
36 .map_err(|error| DbError::context("open retained replay database image", error))?;
37 let routing = crate::database_open::load_coven_metadata(&image)?;
38 let coven_schema_is_current = crate::database_open::initialized_coven_schema_is_current(
39 &image,
40 routing.has_scoped_graph(),
41 )?;
42 let host_schema_version = crate::ensure_schema_supported(&image, migrations)?;
43 let target_host_schema_version = crate::supported_version(migrations);
44 if coven_schema_is_current && host_schema_version == target_host_schema_version {
45 return Ok(());
46 }
47 let had_schema_version = crate::get_protocol_state_on(
48 &image,
49 crate::coven_migration::COVEN_SCHEMA_VERSION_STATE_KEY,
50 )?
51 .is_some();
52 let transaction = image.unchecked_transaction().map_err(DbError::from)?;
53 crate::run_coven_migrations_in_transaction(&transaction, routing.has_scoped_graph(), policy)?;
54 let migrated_host_schema_version =
55 crate::run_migrations_in_transaction(&transaction, migrations)?;
56 if !had_schema_version {
57 crate::delete_protocol_state_on(
58 &transaction,
59 crate::coven_migration::COVEN_SCHEMA_VERSION_STATE_KEY,
60 )?;
61 }
62 transaction.commit().map_err(DbError::from)?;
63 let image_bytes = crate::connection_io::serialize_database_image(&image)?;
64 let blob_decls = crate::BlobDecls::from_tables(&image, synced_tables).map_err(DbError::from)?;
65 records.replace_retained_replay_image(
66 &baseline,
67 migrated_host_schema_version,
68 &image_bytes,
69 &blob_decls,
70 )?;
71 tracing::info!(
72 previous_host_schema_version = baseline.schema_version,
73 migrated_host_schema_version,
74 coven_schema_migrated = !coven_schema_is_current,
75 "Migrated retained replay baseline schema"
76 );
77 Ok(())
78}
79
80pub(crate) fn load_replay_baseline_on(
81 records: StoreRecords<'_>,
82) -> Result<Option<RetainedReplayBaseline>, DbError> {
83 let Some(baseline) = load_replay_baseline_metadata_on(records)? else {
84 return Ok(None);
85 };
86 records.validate_replay_baseline_image(&baseline)?;
87 records.validate_replay_authority(&baseline)?;
88 Ok(Some(baseline))
89}
90
91pub(super) fn load_replay_baseline_metadata_on(
92 records: StoreRecords<'_>,
93) -> Result<Option<RetainedReplayBaseline>, DbError> {
94 let stored = records.retained_replay_baseline_row()?;
95 let Some(stored) = stored else {
96 return Ok(None);
97 };
98 let schema_version = u32::try_from(stored.schema_version)
99 .map_err(|_| DbError::Message("retained replay schema version exceeds u32".to_string()))?;
100 let authority_hash = stored
101 .authority_hash
102 .parse()
103 .map_err(|error| DbError::context("retained replay authority hash", error))?;
104 let authority_bytes = records
105 .verified_payload(authority_hash)
106 .map_err(|error| DbError::context("read retained replay authority", error))?;
107 let authority: RetainedReplayAuthority = serde_json::from_slice(&authority_bytes)
108 .map_err(|error| DbError::context("retained replay authority", error))?;
109 if serde_json::to_vec(&authority)
110 .map_err(|error| DbError::context("serialize retained replay authority", error))?
111 != authority_bytes
112 {
113 return Err(DbError::Message(
114 "retained replay baseline metadata is not canonical".to_string(),
115 ));
116 }
117 Ok(Some(RetainedReplayBaseline {
118 schema_version,
119 routing_hash: stored
120 .routing_hash
121 .parse()
122 .map_err(|error| DbError::context("retained replay routing hash", error))?,
123 image_payload_hash: stored
124 .image_payload_hash
125 .parse()
126 .map_err(|error| DbError::context("retained replay image hash", error))?,
127 authority,
128 }))
129}
130
131pub(crate) fn install_generation_zero_replay_baseline_on(
132 records: StoreRecords<'_>,
133 schema_version: u32,
134 routing_hash: ObjectHash,
135 authority: RetainedReplayGenesisAuthority,
136) -> Result<RetainedReplayBaseline, DbError> {
137 if load_replay_baseline_on(records)?.is_some() {
138 return Err(DbError::Message(
139 "retained replay baseline already exists before founder activation".to_string(),
140 ));
141 }
142 records.install_generation_zero_replay_baseline(schema_version, routing_hash, authority)
143}
144
145pub(crate) fn install_snapshot_replay_baseline_on(
146 records: StoreRecords<'_>,
147 schema_version: u32,
148 routing_hash: ObjectHash,
149 authority: RetainedReplaySnapshotAuthority,
150 blob_decls: &crate::BlobDecls,
151) -> Result<RetainedReplayBaseline, DbError> {
152 if load_replay_baseline_on(records)?.is_some() {
153 return Err(DbError::Message(
154 "retained replay baseline already exists before snapshot bootstrap".to_string(),
155 ));
156 }
157 records.install_snapshot_replay_baseline(schema_version, routing_hash, authority, blob_decls)
158}
159
160pub(crate) fn ensure_founder_replay_baseline_on(
161 records: StoreRecords<'_>,
162 schema_version: u32,
163 routing_hash: ObjectHash,
164 authority: RetainedReplayGenesisAuthority,
165) -> Result<RetainedReplayBaseline, DbError> {
166 if let Some(existing) = load_replay_baseline_on(records)? {
167 let authority_matches = match &existing.authority {
168 RetainedReplayAuthority::Genesis(existing) => existing == &authority,
169 RetainedReplayAuthority::InstalledSnapshot(existing) => {
170 existing.store_root == authority.store_root
171 && existing.founder_registration == authority.founder_registration
172 }
173 };
174 if existing.schema_version != schema_version {
175 return Err(DbError::Message(format!(
176 "retained replay baseline schema version {} differs from database schema version {schema_version}",
177 existing.schema_version,
178 )));
179 }
180 if existing.routing_hash != routing_hash {
181 return Err(DbError::Message(format!(
182 "retained replay baseline routing hash {} differs from database routing hash {routing_hash}",
183 existing.routing_hash,
184 )));
185 }
186 if !authority_matches {
187 return Err(DbError::Message(format!(
188 "retained replay baseline authority {:?} differs from installed founder authority {:?}",
189 existing.authority, authority,
190 )));
191 }
192 return Ok(existing);
193 }
194 let accepted_history = records.accepted_history_count()?;
195 if accepted_history != 0 {
196 return Err(DbError::Message(
197 "accepted Store history exists without a retained replay baseline".to_string(),
198 ));
199 }
200 install_generation_zero_replay_baseline_on(records, schema_version, routing_hash, authority)
201}
202
203const GENESIS_PRESERVED_TABLES: &[&str] = &[
204 "protocol_state",
205 "store_protocol_root_authority",
206 "store_device_registration_activations",
207];
208
209#[derive(Debug, Clone, Copy, PartialEq, Eq)]
210enum ReplayTableDisposition {
211 Replace,
212 ReplaceWhenRouting,
213 Preserve,
214 ExactTransition,
215}
216
217const REPLAY_TABLES: &[(&str, ReplayTableDisposition)] = &[
218 ("active_store_publication", ReplayTableDisposition::Preserve),
219 ("activated_circle_acks", ReplayTableDisposition::Replace),
220 ("activated_store_acks", ReplayTableDisposition::Replace),
221 ("blob_locators", ReplayTableDisposition::ExactTransition),
225 ("blob_make_remote_intents", ReplayTableDisposition::Preserve),
226 ("circle_access_cache", ReplayTableDisposition::Replace),
227 (
228 "circle_bootstrap_coverage",
229 ReplayTableDisposition::Preserve,
230 ),
231 ("circle_close_exclusions", ReplayTableDisposition::Preserve),
232 (
233 "circle_control_activations",
234 ReplayTableDisposition::Replace,
235 ),
236 ("circle_current_state", ReplayTableDisposition::Replace),
237 ("circle_operation_uploads", ReplayTableDisposition::Preserve),
238 ("circle_operations", ReplayTableDisposition::Preserve),
239 ("cloud_outbox", ReplayTableDisposition::Preserve),
240 ("local_blob_refs", ReplayTableDisposition::Preserve),
241 ("local_cleanup_intents", ReplayTableDisposition::Preserve),
242 (
243 "local_store_device_registration",
244 ReplayTableDisposition::Preserve,
245 ),
246 (
247 "local_owner_recovery_publication",
248 ReplayTableDisposition::Preserve,
249 ),
250 (
251 "local_store_founder_graph",
252 ReplayTableDisposition::Preserve,
253 ),
254 (
255 "local_store_protocol_root",
256 ReplayTableDisposition::Preserve,
257 ),
258 ("materialized_commits", ReplayTableDisposition::Preserve),
265 (
266 "outbound_membership_mutation",
267 ReplayTableDisposition::Preserve,
268 ),
269 ("outbound_circle_acks", ReplayTableDisposition::Preserve),
270 ("outbound_circle_snapshot", ReplayTableDisposition::Preserve),
271 ("outbound_store_acks", ReplayTableDisposition::Preserve),
272 (
273 "outbound_store_device_exclusion",
274 ReplayTableDisposition::Preserve,
275 ),
276 ("outbound_store_snapshot", ReplayTableDisposition::Preserve),
277 ("payload_cleanup", ReplayTableDisposition::Preserve),
278 ("payload_owners", ReplayTableDisposition::Preserve),
279 ("payload_storage", ReplayTableDisposition::Preserve),
280 (
281 "protocol_inert_objects",
282 ReplayTableDisposition::ExactTransition,
283 ),
284 ("protocol_state", ReplayTableDisposition::Preserve),
285 (
286 "published_blob_drop_intents",
287 ReplayTableDisposition::Preserve,
288 ),
289 ("published_circle_acks", ReplayTableDisposition::Preserve),
290 (
291 "published_circle_snapshot",
292 ReplayTableDisposition::Preserve,
293 ),
294 ("published_store_acks", ReplayTableDisposition::Preserve),
295 ("published_store_snapshot", ReplayTableDisposition::Preserve),
296 ("reclaimed_store_packages", ReplayTableDisposition::Preserve),
297 ("remote_objects", ReplayTableDisposition::ExactTransition),
298 (
299 "retained_merge_materializations",
300 ReplayTableDisposition::Preserve,
301 ),
302 (
303 "retained_replay_baselines",
304 ReplayTableDisposition::Preserve,
305 ),
306 (
307 "retained_replay_blob_leases",
308 ReplayTableDisposition::Preserve,
309 ),
310 ("retained_replay_objects", ReplayTableDisposition::Preserve),
311 ("row_blob_locators", ReplayTableDisposition::Replace),
312 ("snapshot_coverage", ReplayTableDisposition::Preserve),
313 (
314 "store_author_exclusion_activations",
315 ReplayTableDisposition::Replace,
316 ),
317 (
318 "store_device_registration_activations",
319 ReplayTableDisposition::Replace,
320 ),
321 ("store_device_states", ReplayTableDisposition::Preserve),
325 (
326 "store_device_state_snapshots",
327 ReplayTableDisposition::Preserve,
328 ),
329 (
330 "store_protocol_root_authority",
331 ReplayTableDisposition::Preserve,
332 ),
333 (
334 "store_publication_current",
335 ReplayTableDisposition::Preserve,
336 ),
337 (
338 "store_publication_entries",
339 ReplayTableDisposition::Preserve,
340 ),
341 ("store_reclaim_operations", ReplayTableDisposition::Preserve),
342 ("store_write_blob_leases", ReplayTableDisposition::Preserve),
343 ("store_write_blobs", ReplayTableDisposition::Preserve),
344 ("store_write_packages", ReplayTableDisposition::Preserve),
345 ("store_write_partitions", ReplayTableDisposition::Preserve),
346 ("store_writes", ReplayTableDisposition::Preserve),
347 ("stream_activations", ReplayTableDisposition::Replace),
348 (
349 "_coven_audience",
350 ReplayTableDisposition::ReplaceWhenRouting,
351 ),
352 (
353 "_coven_row_routes",
354 ReplayTableDisposition::ReplaceWhenRouting,
355 ),
356];
357
358pub fn projection_table_names(include_routing: bool) -> Vec<String> {
359 REPLAY_TABLES
360 .iter()
361 .filter_map(|(table, disposition)| match disposition {
362 ReplayTableDisposition::Replace => Some((*table).to_string()),
363 ReplayTableDisposition::ReplaceWhenRouting if include_routing => {
364 Some((*table).to_string())
365 }
366 ReplayTableDisposition::ReplaceWhenRouting
367 | ReplayTableDisposition::Preserve
368 | ReplayTableDisposition::ExactTransition => None,
369 })
370 .collect()
371}
372
373#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
374#[serde(deny_unknown_fields)]
375pub struct RetainedReplayGenesisAuthority {
376 pub store_root: StoreRootRef,
377 pub founder_registration: StoreDeviceRegistrationRef,
378}
379
380#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
381#[serde(rename_all = "snake_case", deny_unknown_fields)]
382pub enum RetainedReplayAuthority {
383 Genesis(RetainedReplayGenesisAuthority),
384 InstalledSnapshot(RetainedReplaySnapshotAuthority),
388}
389
390#[derive(Debug, Clone, PartialEq, Eq)]
397pub struct RetainedReplayBaseline {
398 pub schema_version: u32,
399 pub routing_hash: ObjectHash,
400 pub image_payload_hash: ObjectHash,
401 pub authority: RetainedReplayAuthority,
402}
403
404impl RetainedReplayBaseline {
405 pub fn coverage(&self) -> &CommitFrontier {
407 match &self.authority {
408 RetainedReplayAuthority::Genesis(_) => {
409 static EMPTY: CommitFrontier = CommitFrontier(BTreeMap::new());
410 &EMPTY
411 }
412 RetainedReplayAuthority::InstalledSnapshot(authority) => &authority.metadata.coverage,
413 }
414 }
415
416 pub fn canonical_authority_bytes(&self) -> Result<Vec<u8>, DbError> {
417 serde_json::to_vec(&self.authority)
418 .map_err(|error| DbError::context("serialize retained replay authority", error))
419 }
420
421 pub(super) fn image_bytes(
424 &self,
425 conn: &Connection,
426 store_dir: &coven_foundation::store_dir::StoreDir,
427 ) -> Result<Vec<u8>, DbError> {
428 crate::payload_store::read_verified_payload_blocking(
429 conn,
430 store_dir,
431 self.image_payload_hash,
432 )
433 .map_err(|error| DbError::context("read retained replay image", error))
434 }
435
436 pub(crate) fn validate_image(
437 &self,
438 conn: &Connection,
439 store_dir: &coven_foundation::store_dir::StoreDir,
440 ) -> Result<(), DbError> {
441 let mut image = Connection::open_in_memory().map_err(DbError::from)?;
442 crate::connection_io::deserialize_database_image_into(
443 &mut image,
444 &self.image_bytes(conn, store_dir)?,
445 )
446 .map_err(|error| DbError::context("open retained replay database image", error))?;
447 self.validate_open_image(&image, store_dir)
448 }
449
450 pub(super) fn validate_open_image(
451 &self,
452 image: &Connection,
453 store_dir: &coven_foundation::store_dir::StoreDir,
454 ) -> Result<(), DbError> {
455 match &self.authority {
456 RetainedReplayAuthority::Genesis(_) => {
457 self.validate_image_metadata(image)?;
458 let protocol_keys = protocol_state_keys(image)?;
459 let founder_membership_cursor = founder_membership_cursor_key(image)?;
460 if protocol_keys
461 .iter()
462 .any(|key| !generation_zero_protocol_key(&founder_membership_cursor, key))
463 || !required_generation_zero_protocol_keys()
464 .iter()
465 .all(|key| protocol_keys.contains(*key))
466 || !protocol_keys.contains(&founder_membership_cursor)
467 {
468 return Err(DbError::Message(
469 "retained replay image protocol state is not the generation-zero set"
470 .to_string(),
471 ));
472 }
473 for table in crate::user_table_names(image).map_err(DbError::from)? {
474 let count: i64 = image
475 .query_row(
476 &format!("SELECT COUNT(*) FROM {}", crate::quote_ident(&table)),
477 [],
478 |row| row.get(0),
479 )
480 .map_err(DbError::from)?;
481 let expected = match table.as_str() {
482 "protocol_state" => None,
483 "store_protocol_root_authority"
484 | "store_device_registration_activations" => Some(1),
485 _ => Some(0),
486 };
487 if expected.is_some_and(|expected| count != expected) {
488 return Err(DbError::Message(format!(
489 "generation-zero retained replay image table {table:?} has {count} rows"
490 )));
491 }
492 }
493 validate_replay_image_foreign_keys(image)?;
494 }
495 RetainedReplayAuthority::InstalledSnapshot(authority) => {
496 authority.validate()?;
497 self.validate_image_metadata(image)?;
498 let mut actual = BTreeMap::new();
499 let rows = query_mapped_rows(
500 image,
501 "SELECT device_id, seq, commit_ref FROM snapshot_coverage ORDER BY device_id",
502 [],
503 |row| {
504 Ok((
505 row.get::<_, String>(0)?,
506 row.get::<_, i64>(1)?,
507 row.get::<_, String>(2)?,
508 ))
509 },
510 )?;
511 for (stream_id, sequence, encoded) in rows {
512 let reference: StoreBatchCommitRef =
513 serde_json::from_str(&encoded).map_err(|error| {
514 DbError::context("snapshot replay coverage reference", error)
515 })?;
516 if sequence < 0
517 || u64::try_from(sequence).ok() != Some(reference.coord.sequence())
518 {
519 return Err(DbError::Message(
520 "snapshot replay coverage sequence differs from its exact reference"
521 .to_string(),
522 ));
523 }
524 if actual.insert(stream_id, reference).is_some() {
525 return Err(DbError::Message(
526 "snapshot replay coverage repeats a Store stream".to_string(),
527 ));
528 }
529 }
530 if actual != self.coverage().clone().into_refs() {
531 return Err(DbError::Message(
532 "snapshot replay image coverage differs from its baseline".to_string(),
533 ));
534 }
535 let materialized_commits: i64 = image
536 .query_row("SELECT COUNT(*) FROM materialized_commits", [], |row| {
537 row.get(0)
538 })
539 .map_err(DbError::from)?;
540 if materialized_commits != 0 {
541 return Err(DbError::Message(
542 "snapshot replay baseline contains materialized_commits rows".to_string(),
543 ));
544 }
545 let mut verified_authority =
546 super::VerifiedStoreAuthority::for_replay_baseline(self.clone());
547 crate::StoreDatabase::validate_snapshot_retained_inputs_on(
548 StoreRecords::new(image, store_dir),
549 &mut verified_authority,
550 &authority.store_root,
551 self.coverage(),
552 )?;
553 validate_replay_image_foreign_keys(image)?;
554 }
555 }
556 let routing = crate::database_open::load_coven_metadata(image)?;
557 if routing.hash() != self.routing_hash {
558 return Err(DbError::Message(
559 "retained replay image routing contract differs from its baseline".to_string(),
560 ));
561 }
562 crate::database_open::validate_initialized_coven_schema(image, routing.has_scoped_graph())?;
563 crate::store_authority_records::validate_replay_authority_on(image, self)
564 }
565
566 fn validate_image_metadata(&self, image: &Connection) -> Result<(), DbError> {
567 let stored_schema_version: u32 = image
568 .pragma_query_value(None, "user_version", |row| row.get(0))
569 .map_err(DbError::from)?;
570 if stored_schema_version != self.schema_version {
571 return Err(DbError::Message(format!(
572 "retained replay image schema version is {stored_schema_version}, expected {}",
573 self.schema_version
574 )));
575 }
576 let stored_routing_hash =
577 crate::required_protocol_state_on(image, SYNC_ROUTING_HASH_STATE_KEY)?;
578 if stored_routing_hash != self.routing_hash.to_string() {
579 return Err(DbError::Message(
580 "retained replay image routing hash differs from its baseline".to_string(),
581 ));
582 }
583 Ok(())
584 }
585}
586
587pub(super) struct GenerationZeroReplayImage {
588 image: Connection,
589}
590
591impl GenerationZeroReplayImage {
592 pub(super) fn database_bytes(&self) -> Result<Vec<u8>, DbError> {
593 crate::connection_io::serialize_database_image(&self.image)
594 }
595
596 pub(super) fn validate(
597 &self,
598 baseline: &RetainedReplayBaseline,
599 store_dir: &coven_foundation::store_dir::StoreDir,
600 ) -> Result<(), DbError> {
601 baseline.validate_open_image(&self.image, store_dir)
602 }
603}
604
605pub(crate) fn copy_table_with_conflicts(
610 source: &Connection,
611 target: &Connection,
612 table: &str,
613 ignore_existing: bool,
614) -> Result<(), DbError> {
615 let pragma = format!("PRAGMA table_info({})", crate::quote_ident(table));
616 let columns = query_mapped_rows(source, &pragma, [], |row| row.get::<_, String>(1))?;
617 if columns.is_empty() {
618 return Err(DbError::Message(format!(
619 "retained replay projection table {table:?} is absent"
620 )));
621 }
622 let quoted_columns = columns
623 .iter()
624 .map(|column| crate::quote_ident(column))
625 .collect::<Vec<_>>()
626 .join(", ");
627 let select = format!("SELECT {quoted_columns} FROM {}", crate::quote_ident(table));
628 let rows = query_mapped_rows(source, &select, [], |row| {
629 (0..columns.len())
630 .map(|index| row.get::<_, Value>(index))
631 .collect::<rusqlite::Result<Vec<_>>>()
632 })?;
633 let placeholders = (1..=columns.len())
634 .map(|index| format!("?{index}"))
635 .collect::<Vec<_>>()
636 .join(", ");
637 let verb = if ignore_existing {
638 "INSERT OR IGNORE"
639 } else {
640 "INSERT"
641 };
642 let insert = format!(
643 "{verb} INTO {} ({quoted_columns}) VALUES ({placeholders})",
644 crate::quote_ident(table)
645 );
646 for values in rows {
647 target
648 .execute(&insert, rusqlite::params_from_iter(values))
649 .map_err(DbError::from)?;
650 }
651 Ok(())
652}
653
654pub(super) fn project_generation_zero_image(
655 source: &Connection,
656) -> Result<GenerationZeroReplayImage, DbError> {
657 let source_bytes = crate::connection_io::serialize_database_image(source)?;
658 let mut image = Connection::open_in_memory().map_err(DbError::from)?;
659 crate::connection_io::deserialize_database_image_into(&mut image, &source_bytes)
660 .map_err(|error| DbError::context("open retained replay database image", error))?;
661 image
662 .pragma_update(None, "foreign_keys", "OFF")
663 .map_err(DbError::from)?;
664 let transaction = image.unchecked_transaction().map_err(DbError::from)?;
665 let founder_membership_cursor = founder_membership_cursor_key(&transaction)?;
666 for table in crate::user_table_names(&transaction).map_err(DbError::from)? {
667 if GENESIS_PRESERVED_TABLES.contains(&table.as_str()) {
668 continue;
669 }
670 transaction
671 .execute_batch(&format!("DELETE FROM {}", crate::quote_ident(&table)))
672 .map_err(DbError::from)?;
673 }
674 let protocol_keys = protocol_state_keys(&transaction)?;
675 for key in protocol_keys {
676 if !generation_zero_protocol_key(&founder_membership_cursor, &key) {
677 crate::delete_protocol_state_on(&transaction, &key)?;
678 }
679 }
680 transaction
681 .execute("DELETE FROM sqlite_sequence", [])
682 .map_err(DbError::from)?;
683 transaction.commit().map_err(DbError::from)?;
684 image.execute_batch("VACUUM").map_err(DbError::from)?;
685 image
686 .pragma_update(None, "foreign_keys", "ON")
687 .map_err(DbError::from)?;
688 let violations: bool = image
689 .query_row(
690 "SELECT EXISTS(SELECT 1 FROM pragma_foreign_key_check)",
691 [],
692 |row| row.get(0),
693 )
694 .map_err(DbError::from)?;
695 if violations {
696 return Err(DbError::Message(
697 "generation-zero retained replay image violates foreign keys".to_string(),
698 ));
699 }
700 Ok(GenerationZeroReplayImage { image })
701}
702
703fn validate_replay_image_foreign_keys(image: &Connection) -> Result<(), DbError> {
704 let violations: bool = image
705 .query_row(
706 "SELECT EXISTS(SELECT 1 FROM pragma_foreign_key_check)",
707 [],
708 |row| row.get(0),
709 )
710 .map_err(DbError::from)?;
711 if violations {
712 return Err(DbError::Message(
713 "retained replay image violates foreign keys".to_string(),
714 ));
715 }
716 Ok(())
717}
718
719fn protocol_state_keys(connection: &Connection) -> Result<BTreeSet<String>, DbError> {
720 let mut statement = connection
721 .prepare("SELECT key FROM protocol_state ORDER BY key")
722 .map_err(DbError::from)?;
723 let keys = statement
724 .query_map([], |row| row.get::<_, String>(0))
725 .map_err(DbError::from)?
726 .collect::<rusqlite::Result<BTreeSet<_>>>()
727 .map_err(DbError::from)?;
728 Ok(keys)
729}
730
731fn required_generation_zero_protocol_keys() -> &'static [&'static str] {
732 &[
733 COVEN_INITIALIZED_STATE_KEY,
734 COVEN_SCHEMA_MANIFEST_STATE_KEY,
735 OWNER_PUBKEY_STATE_KEY,
736 STORE_DEVICE_GENESIS_STATE_KEY,
737 SYNC_ROUTING_CONTRACT_STATE_KEY,
738 SYNC_ROUTING_HASH_STATE_KEY,
739 ]
740}
741
742fn founder_membership_cursor_key(connection: &Connection) -> Result<String, DbError> {
743 let bytes: Vec<u8> = connection
744 .query_row(
745 "SELECT store_protocol_root_bytes
746 FROM store_protocol_root_authority WHERE singleton = 1",
747 [],
748 |row| row.get(0),
749 )
750 .map_err(DbError::from)?;
751 let root = coven_protocol::store_commit::StoreProtocolRoot::parse(&bytes)
752 .map_err(|error| DbError::context("retained replay Store root", error))?;
753 let stream = coven_protocol::membership::derive_founder_stream_id(
754 &root.descriptor.store_root_id().to_string(),
755 &root.descriptor.founder_pubkey,
756 );
757 Ok(
758 crate::InitialStoreMembershipAuthority::cursor_state_key_for_stream(
759 &root.descriptor.founder_grant,
760 stream,
761 ),
762 )
763}
764
765fn generation_zero_protocol_key(founder_membership_cursor: &str, key: &str) -> bool {
766 required_generation_zero_protocol_keys().contains(&key) || founder_membership_cursor == key
767}
768
769#[cfg(test)]
770#[path = "retained_replay_test.rs"]
771mod tests;