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