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 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    // A commit's position is a fact accepted at commit time, maintained
270    // incrementally by materialization and retraction — never rebuilt by a
271    // projection replay. A filtered replay (one that omits covered or
272    // beyond-cutoff Circle packages) would re-derive the commit's retained-input
273    // hash from its filtered packages and drift from the preserved
274    // `retained_merge_materializations` row, breaking the `materialized_commits`
275    // foreign key. Only explicit retraction removes these rows.
276    ("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    // Preserved for the same reason as `materialized_commits`: a commit's
346    // derived device state is accepted at commit time, never recomputed by a
347    // filtered replay; only explicit retraction removes it.
348    ("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    /// The device installed a snapshot image and replays forward from it. Named
405    /// for what happened, not for a property of the snapshot: installing one
406    /// never required the store's other devices to have acknowledged it.
407    InstalledSnapshot(RetainedReplaySnapshotAuthority),
408}
409
410/// The database image a replay starts from, and the authority that says which
411/// history it covers.
412///
413/// `image_payload_hash` names the image bytes in the payload store, as
414/// `authority_hash` on the row names the authority's canonical bytes. The row
415/// holds the facts without carrying either payload in its own columns.
416#[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    /// The verified bytes of the baseline image. Replays open a private,
433    /// writable connection from these bytes inside the replay capability.
434    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
630/// Copy `table` from `source` into `target`. With `ignore_existing`, a row whose
631/// unique key already exists in `target` is left untouched instead of failing —
632/// used when installing a Circle image onto a Store image that already carries the
633/// shared, deterministic audience-routing rows.
634pub(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;