Skip to main content

coven_database/store/store_session/
retained_replay.rs

1//! Private accepted-history baselines and deterministic retained replay.
2
3use 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 provenance belongs to accepted object lifetime, including orphaned
222    // ciphertext whose row version a replay omits. Exact activation and reclaim
223    // transitions update this index alongside remote_objects.
224    ("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    // A commit's position is a fact accepted at commit time and is never rebuilt by a
259    // projection replay. A filtered replay (one that omits covered or
260    // beyond-cutoff Circle packages) would re-derive the commit's retained-input
261    // hash from its filtered packages and drift from the preserved
262    // `retained_merge_materializations` row, breaking the `materialized_commits`
263    // foreign key.
264    ("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    // Preserved for the same reason as `materialized_commits`: a commit's
322    // derived device state is accepted at commit time, never recomputed by a
323    // filtered replay; only explicit retraction removes it.
324    ("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    /// The device installed a snapshot image and replays forward from it. Named
385    /// for what happened, not for a property of the snapshot: installing one
386    /// never required the store's other devices to have acknowledged it.
387    InstalledSnapshot(RetainedReplaySnapshotAuthority),
388}
389
390/// The database image a replay starts from, and the authority that says which
391/// history it covers.
392///
393/// `image_payload_hash` names the image bytes in the payload store, as
394/// `authority_hash` on the row names the authority's canonical bytes. The row
395/// holds the facts without carrying either payload in its own columns.
396#[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    /// Accepted history covered by the authority this image starts from.
406    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    /// The verified bytes of the baseline image. Replays open a private,
422    /// writable connection from these bytes inside the replay capability.
423    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
605/// Copy `table` from `source` into `target`. With `ignore_existing`, a row whose
606/// unique key already exists in `target` is left untouched instead of failing —
607/// used when installing a Circle image onto a Store image that already carries the
608/// shared, deterministic audience-routing rows.
609pub(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;