Skip to main content

semantic_memory/
journal.rs

1//! Transactional fact-create outbox for device-primary replication.
2//!
3//! V37 journal rows are retained as `legacy_unverified`. V38 introduces an
4//! explicit stream epoch, transaction-owned sequence/head state, and a
5//! domain-separated digest chain. Only exact typed payload bytes are canonical;
6//! derived embeddings and indexes are rebuilt by the replica.
7
8use crate::error::MemoryError;
9use rusqlite::{Connection, OptionalExtension, Transaction};
10use serde::{Deserialize, Serialize};
11use sha2::{Digest, Sha256};
12
13pub const FACT_CREATE_OPERATION: &str = "fact.create";
14pub const FACT_CREATE_PAYLOAD_SCHEMA: &str = "semantic_memory.fact.create.v1";
15pub const VERIFIED_RECORD_STATE: &str = "verified_v1";
16pub const LEGACY_RECORD_STATE: &str = "legacy_unverified";
17pub const GENESIS_PREDECESSOR: [u8; 32] = [0; 32];
18const COMPAT_STREAM_EPOCH: u64 = 1;
19
20/// Canonical fact-create authority payload. Field order is fixed by this struct.
21#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
22#[serde(deny_unknown_fields)]
23pub struct FactCreatePayloadV1 {
24    pub fact_id: String,
25    pub namespace: String,
26    pub content: String,
27    pub source: Option<String>,
28    pub metadata: Option<serde_json::Value>,
29}
30
31/// Closed receiver-side representation of one verified fact-create record.
32///
33/// This type contains only fields owned by semantic-memory's V38 journal
34/// contract. Transport authentication and signatures are deliberately owned by
35/// the caller (for example Mnemes); semantic-memory revalidates the canonical
36/// operation, schema, payload, digest chain, and stream ordering before apply.
37#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
38#[serde(deny_unknown_fields)]
39pub struct FactCreateReplicaEnvelopeV1 {
40    pub home_device_id: String,
41    pub store_id: String,
42    pub stream_epoch: u64,
43    pub sequence: i64,
44    pub operation_kind: String,
45    pub payload_schema: String,
46    pub payload: Vec<u8>,
47    pub payload_digest: [u8; 32],
48    pub predecessor_digest: [u8; 32],
49    pub envelope_digest: [u8; 32],
50}
51
52/// Durable receiver decision for a fact-create record.
53#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
54#[serde(rename_all = "snake_case")]
55pub enum ReplicaApplyOutcome {
56    Applied { sequence: i64, fact_id: String },
57    Duplicate { sequence: i64 },
58    Fork { sequence: i64 },
59    Gap { expected: i64, received: i64 },
60    EpochConflict { active: u64, received: u64 },
61}
62
63pub fn encode_fact_create_payload(payload: &FactCreatePayloadV1) -> Result<Vec<u8>, MemoryError> {
64    serde_json::to_vec(payload)
65        .map_err(|error| MemoryError::DigestError(format!("fact-create payload encoding: {error}")))
66}
67
68/// Domain-separated SHA-256 over length-prefixed exact fields.
69pub fn digest_fields(domain: &[u8], fields: &[&[u8]]) -> [u8; 32] {
70    let mut hasher = Sha256::new();
71    hasher.update((domain.len() as u64).to_be_bytes());
72    hasher.update(domain);
73    for field in fields {
74        hasher.update((field.len() as u64).to_be_bytes());
75        hasher.update(field);
76    }
77    hasher.finalize().into()
78}
79
80pub fn payload_digest(payload: &[u8]) -> [u8; 32] {
81    digest_fields(b"semantic-memory.payload.v1", &[payload])
82}
83
84#[allow(clippy::too_many_arguments)]
85pub fn envelope_digest(
86    home_device_id: &str,
87    store_id: &str,
88    stream_epoch: u64,
89    sequence: i64,
90    operation_kind: &str,
91    payload_schema: &str,
92    predecessor_digest: &[u8; 32],
93    payload_digest: &[u8; 32],
94) -> [u8; 32] {
95    digest_fields(
96        b"semantic-memory.envelope.v1",
97        &[
98            home_device_id.as_bytes(),
99            store_id.as_bytes(),
100            &stream_epoch.to_be_bytes(),
101            &sequence.to_be_bytes(),
102            operation_kind.as_bytes(),
103            payload_schema.as_bytes(),
104            predecessor_digest,
105            payload_digest,
106        ],
107    )
108}
109
110/// A verified outbox entry. V37 rows can still be inspected directly in SQLite
111/// but are never returned by the verified export API.
112#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
113pub struct JournalEntry {
114    pub journal_id: i64,
115    pub home_device_id: String,
116    pub store_id: String,
117    pub stream_epoch: u64,
118    pub sequence: i64,
119    pub operation_kind: String,
120    pub payload_schema: String,
121    pub payload: Vec<u8>,
122    pub payload_digest: [u8; 32],
123    pub predecessor_digest: [u8; 32],
124    pub envelope_digest: [u8; 32],
125    pub record_state: String,
126    pub created_at: String,
127}
128
129/// V37 migration retained for compatibility. These rows have no cryptographic
130/// chain and become explicitly `legacy_unverified` under V38.
131pub const MIGRATION_V37: &str = "\
132CREATE TABLE IF NOT EXISTS mutation_journal (
133    journal_id      INTEGER PRIMARY KEY AUTOINCREMENT,
134    home_device_id  TEXT NOT NULL,
135    store_id        TEXT NOT NULL,
136    sequence        INTEGER NOT NULL,
137    operation_kind  TEXT NOT NULL,
138    payload         BLOB NOT NULL,
139    created_at      TEXT NOT NULL DEFAULT (datetime('now'))
140);
141CREATE UNIQUE INDEX IF NOT EXISTS idx_journal_sequence
142    ON mutation_journal(home_device_id, store_id, sequence);
143";
144
145/// V38 verified stream state. SQLite ALTER defaults are literal constants so
146/// existing V37 databases migrate without fabricating verified metadata.
147pub const MIGRATION_V38: &str = r#"
148CREATE TABLE IF NOT EXISTS replication_streams (
149    home_device_id TEXT NOT NULL,
150    store_id TEXT NOT NULL,
151    stream_epoch INTEGER NOT NULL CHECK(stream_epoch > 0),
152    next_sequence INTEGER NOT NULL DEFAULT 1 CHECK(next_sequence > 0),
153    head_digest BLOB NOT NULL CHECK(length(head_digest) = 32),
154    PRIMARY KEY(home_device_id, store_id, stream_epoch)
155);
156ALTER TABLE mutation_journal ADD COLUMN stream_epoch INTEGER NOT NULL DEFAULT 0;
157ALTER TABLE mutation_journal ADD COLUMN payload_schema TEXT NOT NULL DEFAULT 'legacy.unverified';
158ALTER TABLE mutation_journal ADD COLUMN payload_digest BLOB NOT NULL DEFAULT X'0000000000000000000000000000000000000000000000000000000000000000';
159ALTER TABLE mutation_journal ADD COLUMN predecessor_digest BLOB NOT NULL DEFAULT X'0000000000000000000000000000000000000000000000000000000000000000';
160ALTER TABLE mutation_journal ADD COLUMN envelope_digest BLOB NOT NULL DEFAULT X'0000000000000000000000000000000000000000000000000000000000000000';
161ALTER TABLE mutation_journal ADD COLUMN record_state TEXT NOT NULL DEFAULT 'legacy_unverified';
162DROP INDEX IF EXISTS idx_journal_sequence;
163CREATE UNIQUE INDEX IF NOT EXISTS idx_journal_sequence_v38
164    ON mutation_journal(home_device_id, store_id, stream_epoch, sequence);
165"#;
166
167/// V39 receiver-side stream, inbox, and durable ACK projection.
168///
169/// These tables live in the semantic shard so fact state, stream advancement,
170/// inbox evidence, and the ACK decision commit in one SQLite transaction.
171pub const MIGRATION_V39: &str = r#"
172CREATE TABLE IF NOT EXISTS replication_inbox_streams (
173    home_device_id TEXT NOT NULL,
174    store_id TEXT NOT NULL,
175    stream_epoch INTEGER NOT NULL CHECK(stream_epoch > 0),
176    next_sequence INTEGER NOT NULL DEFAULT 1 CHECK(next_sequence > 0),
177    head_digest BLOB NOT NULL CHECK(length(head_digest) = 32),
178    updated_at TEXT NOT NULL DEFAULT (datetime('now')),
179    PRIMARY KEY(home_device_id, store_id)
180);
181CREATE TABLE IF NOT EXISTS replication_inbox (
182    home_device_id TEXT NOT NULL,
183    store_id TEXT NOT NULL,
184    stream_epoch INTEGER NOT NULL CHECK(stream_epoch > 0),
185    sequence INTEGER NOT NULL CHECK(sequence > 0),
186    operation_kind TEXT NOT NULL,
187    payload_schema TEXT NOT NULL,
188    payload BLOB NOT NULL,
189    payload_digest BLOB NOT NULL CHECK(length(payload_digest) = 32),
190    predecessor_digest BLOB NOT NULL CHECK(length(predecessor_digest) = 32),
191    envelope_digest BLOB NOT NULL CHECK(length(envelope_digest) = 32),
192    fact_id TEXT NOT NULL,
193    record_state TEXT NOT NULL DEFAULT 'applied_v1'
194        CHECK(record_state = 'applied_v1'),
195    applied_at TEXT NOT NULL DEFAULT (datetime('now')),
196    PRIMARY KEY(home_device_id, store_id, stream_epoch, sequence)
197);
198CREATE UNIQUE INDEX IF NOT EXISTS idx_replication_inbox_envelope
199    ON replication_inbox(home_device_id, store_id, envelope_digest);
200"#;
201
202fn validate_stream_identity(
203    home_device_id: &str,
204    store_id: &str,
205    stream_epoch: u64,
206) -> Result<(), MemoryError> {
207    if stream_epoch == 0 {
208        return Err(MemoryError::InvalidConfig {
209            field: "replication_stream_epoch",
210            reason: "must be positive".to_string(),
211        });
212    }
213    for (field, value) in [
214        ("journal_device_id", home_device_id),
215        ("journal_store_id", store_id),
216    ] {
217        if value.is_empty() || value.trim() != value || value.chars().any(char::is_whitespace) {
218            return Err(MemoryError::InvalidConfig {
219                field,
220                reason: "must be non-empty, trimmed, and contain no whitespace".to_string(),
221            });
222        }
223    }
224    Ok(())
225}
226
227/// Validate a receiver envelope and decode its strict canonical payload.
228pub fn validate_fact_create_replica_envelope(
229    envelope: &FactCreateReplicaEnvelopeV1,
230) -> Result<FactCreatePayloadV1, MemoryError> {
231    validate_stream_identity(
232        &envelope.home_device_id,
233        &envelope.store_id,
234        envelope.stream_epoch,
235    )?;
236    if envelope.sequence < 1 {
237        return Err(MemoryError::InvalidConfig {
238            field: "replication.sequence",
239            reason: "must be positive".to_string(),
240        });
241    }
242    if envelope.operation_kind != FACT_CREATE_OPERATION
243        || envelope.payload_schema != FACT_CREATE_PAYLOAD_SCHEMA
244    {
245        return Err(MemoryError::NotImplemented(format!(
246            "replication operation/schema not admitted: {}/{}",
247            envelope.operation_kind, envelope.payload_schema
248        )));
249    }
250    let expected_payload_digest = payload_digest(&envelope.payload);
251    if envelope.payload_digest != expected_payload_digest {
252        return Err(MemoryError::DigestError(
253            "fact-create replica payload digest mismatch".to_string(),
254        ));
255    }
256    let expected_envelope_digest = envelope_digest(
257        &envelope.home_device_id,
258        &envelope.store_id,
259        envelope.stream_epoch,
260        envelope.sequence,
261        &envelope.operation_kind,
262        &envelope.payload_schema,
263        &envelope.predecessor_digest,
264        &envelope.payload_digest,
265    );
266    if envelope.envelope_digest != expected_envelope_digest {
267        return Err(MemoryError::DigestError(
268            "fact-create replica envelope digest mismatch".to_string(),
269        ));
270    }
271    let payload: FactCreatePayloadV1 =
272        serde_json::from_slice(&envelope.payload).map_err(|error| MemoryError::CorruptData {
273            table: "replication_inbox",
274            row_id: envelope.sequence.to_string(),
275            detail: format!("invalid fact-create payload: {error}"),
276        })?;
277    if uuid::Uuid::parse_str(&payload.fact_id).is_err() {
278        return Err(MemoryError::CorruptData {
279            table: "replication_inbox",
280            row_id: envelope.sequence.to_string(),
281            detail: "fact-create payload fact_id must be a UUID".to_string(),
282        });
283    }
284    if payload.namespace.is_empty()
285        || payload.namespace.trim() != payload.namespace
286        || payload.namespace.chars().any(char::is_control)
287    {
288        return Err(MemoryError::CorruptData {
289            table: "replication_inbox",
290            row_id: envelope.sequence.to_string(),
291            detail: "fact-create namespace is empty, untrimmed, or contains control characters"
292                .to_string(),
293        });
294    }
295    if payload.content.is_empty() {
296        return Err(MemoryError::CorruptData {
297            table: "replication_inbox",
298            row_id: envelope.sequence.to_string(),
299            detail: "fact-create content is empty".to_string(),
300        });
301    }
302    Ok(payload)
303}
304
305fn digest_from_blob(column: usize, bytes: Vec<u8>) -> Result<[u8; 32], rusqlite::Error> {
306    bytes.try_into().map_err(|bytes: Vec<u8>| {
307        rusqlite::Error::FromSqlConversionFailure(
308            column,
309            rusqlite::types::Type::Blob,
310            Box::new(std::io::Error::new(
311                std::io::ErrorKind::InvalidData,
312                format!("expected 32-byte digest, got {} bytes", bytes.len()),
313            )),
314        )
315    })
316}
317
318fn row_to_entry(row: &rusqlite::Row<'_>) -> Result<JournalEntry, rusqlite::Error> {
319    let epoch: i64 = row.get(3)?;
320    let stream_epoch =
321        u64::try_from(epoch).map_err(|_| rusqlite::Error::IntegralValueOutOfRange(3, epoch))?;
322    Ok(JournalEntry {
323        journal_id: row.get(0)?,
324        home_device_id: row.get(1)?,
325        store_id: row.get(2)?,
326        stream_epoch,
327        sequence: row.get(4)?,
328        operation_kind: row.get(5)?,
329        payload_schema: row.get(6)?,
330        payload: row.get(7)?,
331        payload_digest: digest_from_blob(8, row.get(8)?)?,
332        predecessor_digest: digest_from_blob(9, row.get(9)?)?,
333        envelope_digest: digest_from_blob(10, row.get(10)?)?,
334        record_state: row.get(11)?,
335        created_at: row.get(12)?,
336    })
337}
338
339const ENTRY_SELECT: &str = "journal_id, home_device_id, store_id, stream_epoch, sequence, \
340operation_kind, payload_schema, payload, payload_digest, predecessor_digest, envelope_digest, \
341record_state, created_at";
342
343/// Append a verified record while the caller's semantic mutation transaction is active.
344/// The stream row is the allocator; no sequence is derived from journal contents.
345#[allow(clippy::too_many_arguments)]
346pub fn append_verified_in_tx(
347    tx: &Transaction<'_>,
348    home_device_id: &str,
349    store_id: &str,
350    stream_epoch: u64,
351    operation_kind: &str,
352    payload_schema: &str,
353    payload: &[u8],
354) -> Result<JournalEntry, MemoryError> {
355    validate_stream_identity(home_device_id, store_id, stream_epoch)?;
356    if operation_kind != FACT_CREATE_OPERATION || payload_schema != FACT_CREATE_PAYLOAD_SCHEMA {
357        return Err(MemoryError::NotImplemented(format!(
358            "replication operation/schema not admitted: {operation_kind}/{payload_schema}"
359        )));
360    }
361    let epoch = i64::try_from(stream_epoch).map_err(|_| MemoryError::InvalidConfig {
362        field: "replication_stream_epoch",
363        reason: "does not fit SQLite INTEGER".to_string(),
364    })?;
365
366    tx.execute(
367        "INSERT INTO replication_streams
368         (home_device_id, store_id, stream_epoch, next_sequence, head_digest)
369         VALUES (?1, ?2, ?3, 1, ?4)
370         ON CONFLICT(home_device_id, store_id, stream_epoch) DO NOTHING",
371        rusqlite::params![
372            home_device_id,
373            store_id,
374            epoch,
375            GENESIS_PREDECESSOR.as_slice()
376        ],
377    )?;
378
379    let (sequence, predecessor_bytes): (i64, Vec<u8>) = tx.query_row(
380        "SELECT next_sequence, head_digest FROM replication_streams
381         WHERE home_device_id = ?1 AND store_id = ?2 AND stream_epoch = ?3",
382        rusqlite::params![home_device_id, store_id, epoch],
383        |row| Ok((row.get(0)?, row.get(1)?)),
384    )?;
385    let predecessor = digest_from_blob(1, predecessor_bytes)?;
386    let payload_hash = payload_digest(payload);
387    let envelope_hash = envelope_digest(
388        home_device_id,
389        store_id,
390        stream_epoch,
391        sequence,
392        operation_kind,
393        payload_schema,
394        &predecessor,
395        &payload_hash,
396    );
397
398    tx.execute(
399        "INSERT INTO mutation_journal
400         (home_device_id, store_id, stream_epoch, sequence, operation_kind,
401          payload_schema, payload, payload_digest, predecessor_digest,
402          envelope_digest, record_state)
403         VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
404        rusqlite::params![
405            home_device_id,
406            store_id,
407            epoch,
408            sequence,
409            operation_kind,
410            payload_schema,
411            payload,
412            payload_hash.as_slice(),
413            predecessor.as_slice(),
414            envelope_hash.as_slice(),
415            VERIFIED_RECORD_STATE,
416        ],
417    )?;
418    let journal_id = tx.last_insert_rowid();
419
420    let advanced = tx.execute(
421        "UPDATE replication_streams
422         SET next_sequence = ?4, head_digest = ?5
423         WHERE home_device_id = ?1 AND store_id = ?2 AND stream_epoch = ?3
424           AND next_sequence = ?6 AND head_digest = ?7",
425        rusqlite::params![
426            home_device_id,
427            store_id,
428            epoch,
429            sequence + 1,
430            envelope_hash.as_slice(),
431            sequence,
432            predecessor.as_slice(),
433        ],
434    )?;
435    if advanced != 1 {
436        return Err(MemoryError::Other(
437            "replication stream allocator lost ownership".to_string(),
438        ));
439    }
440
441    tx.query_row(
442        &format!("SELECT {ENTRY_SELECT} FROM mutation_journal WHERE journal_id = ?1"),
443        [journal_id],
444        row_to_entry,
445    )
446    .map_err(MemoryError::Database)
447}
448
449/// Compatibility helper for fact-create tests and offline tools. It owns a
450/// transaction and therefore remains atomic, but callers should use the real
451/// MemoryStore fact mutation path instead.
452#[deprecated(note = "use the MemoryStore fact-create path or append_verified_in_tx")]
453pub fn append_journal_entry(
454    conn: &Connection,
455    home_device_id: &str,
456    store_id: &str,
457    operation_kind: &str,
458    payload: &[u8],
459) -> Result<i64, MemoryError> {
460    if operation_kind != "add_fact" && operation_kind != FACT_CREATE_OPERATION {
461        return Err(MemoryError::NotImplemented(format!(
462            "legacy journal operation not admitted: {operation_kind}"
463        )));
464    }
465    let tx = conn.unchecked_transaction()?;
466    let entry = append_verified_in_tx(
467        &tx,
468        home_device_id,
469        store_id,
470        COMPAT_STREAM_EPOCH,
471        FACT_CREATE_OPERATION,
472        FACT_CREATE_PAYLOAD_SCHEMA,
473        payload,
474    )?;
475    tx.commit()?;
476    Ok(entry.sequence)
477}
478
479/// Compatibility mutation wrapper. The semantic mutation and outbox record are
480/// committed together, and a failed closure consumes no sequence.
481#[deprecated(note = "use the typed MemoryStore mutation path")]
482pub fn mutate_and_journal<F, T>(
483    conn: &Connection,
484    home_device_id: &str,
485    store_id: &str,
486    operation_kind: &str,
487    payload: &[u8],
488    f: F,
489) -> Result<(i64, i64, T), MemoryError>
490where
491    F: FnOnce(&Connection) -> Result<T, MemoryError>,
492{
493    if operation_kind != "add_fact" && operation_kind != FACT_CREATE_OPERATION {
494        return Err(MemoryError::NotImplemented(format!(
495            "legacy journal operation not admitted: {operation_kind}"
496        )));
497    }
498    let tx = conn.unchecked_transaction()?;
499    let result = f(&tx)?;
500    let entry = append_verified_in_tx(
501        &tx,
502        home_device_id,
503        store_id,
504        COMPAT_STREAM_EPOCH,
505        FACT_CREATE_OPERATION,
506        FACT_CREATE_PAYLOAD_SCHEMA,
507        payload,
508    )?;
509    tx.commit()?;
510    Ok((entry.journal_id, entry.sequence, result))
511}
512
513#[derive(Debug, Clone, PartialEq, Eq)]
514pub enum ExportStatus {
515    /// The requested device/store/epoch has no verified stream. This is distinct
516    /// from [`ExportStatus::End`], which means an existing stream has no more
517    /// records at or after the requested sequence. Added in the V39 contract;
518    /// direct callers must handle it explicitly.
519    Empty,
520    /// The requested verified stream exists and has no more records at or after
521    /// the requested sequence.
522    End,
523    More,
524    Gap {
525        expected: i64,
526        found: Option<i64>,
527    },
528    Corrupt {
529        sequence: i64,
530        reason: String,
531    },
532}
533
534#[derive(Debug, Clone, PartialEq)]
535pub struct VerifiedExportBatch {
536    pub entries: Vec<JournalEntry>,
537    pub next_sequence: i64,
538    pub status: ExportStatus,
539}
540
541fn corrupt(
542    entries: Vec<JournalEntry>,
543    next_sequence: i64,
544    sequence: i64,
545    reason: impl Into<String>,
546) -> VerifiedExportBatch {
547    VerifiedExportBatch {
548        entries,
549        next_sequence,
550        status: ExportStatus::Corrupt {
551            sequence,
552            reason: reason.into(),
553        },
554    }
555}
556
557/// Export and verify a contiguous prefix. Legacy V37 rows are never promoted
558/// into this API and any stored digest/chain mismatch is typed as corruption.
559pub fn export_verified_contiguous(
560    conn: &Connection,
561    home_device_id: &str,
562    store_id: &str,
563    stream_epoch: u64,
564    start_sequence: i64,
565    limit: usize,
566) -> Result<VerifiedExportBatch, MemoryError> {
567    validate_stream_identity(home_device_id, store_id, stream_epoch)?;
568    if start_sequence < 1 || limit == 0 {
569        return Err(MemoryError::InvalidConfig {
570            field: "journal_export",
571            reason: "start_sequence and limit must be positive".to_string(),
572        });
573    }
574    let epoch = i64::try_from(stream_epoch).map_err(|_| MemoryError::InvalidConfig {
575        field: "replication_stream_epoch",
576        reason: "does not fit SQLite INTEGER".to_string(),
577    })?;
578    let stream_next: Option<i64> = conn
579        .query_row(
580            "SELECT next_sequence FROM replication_streams
581             WHERE home_device_id = ?1 AND store_id = ?2 AND stream_epoch = ?3",
582            rusqlite::params![home_device_id, store_id, epoch],
583            |row| row.get(0),
584        )
585        .optional()?;
586    let Some(stream_next) = stream_next else {
587        return Ok(VerifiedExportBatch {
588            entries: Vec::new(),
589            next_sequence: start_sequence,
590            status: ExportStatus::Empty,
591        });
592    };
593
594    let mut expected_predecessor = if start_sequence == 1 {
595        GENESIS_PREDECESSOR
596    } else {
597        let previous: Option<Vec<u8>> = conn
598            .query_row(
599                "SELECT envelope_digest FROM mutation_journal
600                 WHERE home_device_id = ?1 AND store_id = ?2 AND stream_epoch = ?3
601                   AND sequence = ?4",
602                rusqlite::params![home_device_id, store_id, epoch, start_sequence - 1],
603                |row| row.get(0),
604            )
605            .optional()?;
606        let Some(previous) = previous else {
607            return Ok(VerifiedExportBatch {
608                entries: Vec::new(),
609                next_sequence: start_sequence,
610                status: ExportStatus::Gap {
611                    expected: start_sequence - 1,
612                    found: None,
613                },
614            });
615        };
616        digest_from_blob(0, previous)?
617    };
618
619    let mut stmt = conn.prepare(&format!(
620        "SELECT {ENTRY_SELECT} FROM mutation_journal
621         WHERE home_device_id = ?1 AND store_id = ?2 AND stream_epoch = ?3
622           AND sequence >= ?4
623         ORDER BY sequence ASC LIMIT ?5"
624    ))?;
625    let rows = stmt.query_map(
626        rusqlite::params![
627            home_device_id,
628            store_id,
629            epoch,
630            start_sequence,
631            (limit + 1) as i64,
632        ],
633        row_to_entry,
634    )?;
635
636    let mut entries = Vec::new();
637    let mut expected = start_sequence;
638    let mut has_extra = false;
639    for row in rows {
640        let entry = row?;
641        if entries.len() == limit {
642            has_extra = true;
643            break;
644        }
645        if entry.sequence != expected {
646            return Ok(VerifiedExportBatch {
647                entries,
648                next_sequence: expected,
649                status: ExportStatus::Gap {
650                    expected,
651                    found: Some(entry.sequence),
652                },
653            });
654        }
655        if entry.operation_kind != FACT_CREATE_OPERATION
656            || entry.payload_schema != FACT_CREATE_PAYLOAD_SCHEMA
657            || entry.record_state != VERIFIED_RECORD_STATE
658        {
659            return Ok(corrupt(
660                entries,
661                expected,
662                entry.sequence,
663                "unadmitted operation, schema, or record state",
664            ));
665        }
666        let expected_payload = payload_digest(&entry.payload);
667        if entry.payload_digest != expected_payload {
668            return Ok(corrupt(
669                entries,
670                expected,
671                entry.sequence,
672                "payload digest mismatch",
673            ));
674        }
675        if entry.predecessor_digest != expected_predecessor {
676            return Ok(corrupt(
677                entries,
678                expected,
679                entry.sequence,
680                "predecessor digest mismatch",
681            ));
682        }
683        let expected_envelope = envelope_digest(
684            home_device_id,
685            store_id,
686            stream_epoch,
687            entry.sequence,
688            &entry.operation_kind,
689            &entry.payload_schema,
690            &entry.predecessor_digest,
691            &entry.payload_digest,
692        );
693        if entry.envelope_digest != expected_envelope {
694            return Ok(corrupt(
695                entries,
696                expected,
697                entry.sequence,
698                "envelope digest mismatch",
699            ));
700        }
701        expected_predecessor = entry.envelope_digest;
702        expected += 1;
703        entries.push(entry);
704    }
705
706    let status = if has_extra || expected < stream_next {
707        if has_extra {
708            ExportStatus::More
709        } else {
710            ExportStatus::Gap {
711                expected,
712                found: None,
713            }
714        }
715    } else if expected == stream_next {
716        ExportStatus::End
717    } else {
718        ExportStatus::Corrupt {
719            sequence: expected,
720            reason: format!(
721                "export advanced beyond stream allocator: export next {expected}, stream next {stream_next}"
722            ),
723        }
724    };
725    Ok(VerifiedExportBatch {
726        entries,
727        next_sequence: expected,
728        status,
729    })
730}
731
732/// Legacy shape retained for current callers. It exports only verified epoch-1
733/// records and does not hide corruption as successful completion.
734#[derive(Debug, Clone)]
735pub struct ExportedBatch {
736    pub entries: Vec<JournalEntry>,
737    pub next_seq: i64,
738    pub has_more: bool,
739}
740
741#[deprecated(note = "use export_verified_contiguous with an explicit epoch")]
742pub fn export_contiguous(
743    conn: &Connection,
744    home_device_id: &str,
745    store_id: &str,
746    start_seq: i64,
747    limit: usize,
748) -> Result<ExportedBatch, MemoryError> {
749    let batch = export_verified_contiguous(
750        conn,
751        home_device_id,
752        store_id,
753        COMPAT_STREAM_EPOCH,
754        start_seq,
755        limit,
756    )?;
757    if let ExportStatus::Corrupt { sequence, reason } = &batch.status {
758        return Err(MemoryError::CorruptData {
759            table: "mutation_journal",
760            row_id: sequence.to_string(),
761            detail: reason.clone(),
762        });
763    }
764    Ok(ExportedBatch {
765        entries: batch.entries,
766        next_seq: batch.next_sequence,
767        has_more: matches!(batch.status, ExportStatus::More),
768    })
769}
770
771/// Compatibility-only replay result. New remote admission must use the closed
772/// typed fact-create dispatcher, not this closure-based adapter.
773#[derive(Debug, Clone, PartialEq, Eq)]
774pub enum ReplayOutcome {
775    Applied { sequence: i64 },
776    AlreadyApplied { sequence: i64 },
777    Conflict { sequence: i64 },
778    Gap { expected: i64, received: i64 },
779}
780
781#[deprecated(note = "closure-based replay is compatibility-only; use closed typed admission")]
782pub fn replay_journal_entry<F>(
783    conn: &Connection,
784    home_device_id: &str,
785    store_id: &str,
786    sequence: i64,
787    operation_kind: &str,
788    payload: &[u8],
789    replay_fn: F,
790) -> Result<ReplayOutcome, MemoryError>
791where
792    F: FnOnce(&Connection) -> Result<(), MemoryError>,
793{
794    if operation_kind != "add_fact" && operation_kind != FACT_CREATE_OPERATION {
795        return Ok(ReplayOutcome::Conflict { sequence });
796    }
797    let tx = conn.unchecked_transaction()?;
798    let existing: Option<(String, Vec<u8>)> = tx
799        .query_row(
800            "SELECT operation_kind, payload FROM mutation_journal
801             WHERE home_device_id = ?1 AND store_id = ?2 AND stream_epoch = ?3 AND sequence = ?4",
802            rusqlite::params![
803                home_device_id,
804                store_id,
805                COMPAT_STREAM_EPOCH as i64,
806                sequence
807            ],
808            |row| Ok((row.get(0)?, row.get(1)?)),
809        )
810        .optional()?;
811    if let Some((stored_operation, stored_payload)) = existing {
812        if stored_operation == FACT_CREATE_OPERATION && stored_payload == payload {
813            return Ok(ReplayOutcome::AlreadyApplied { sequence });
814        }
815        return Ok(ReplayOutcome::Conflict { sequence });
816    }
817
818    let expected =
819        next_expected_sequence_for_epoch(&tx, home_device_id, store_id, COMPAT_STREAM_EPOCH)?;
820    if sequence != expected {
821        return Ok(ReplayOutcome::Gap {
822            expected,
823            received: sequence,
824        });
825    }
826    replay_fn(&tx)?;
827    let entry = append_verified_in_tx(
828        &tx,
829        home_device_id,
830        store_id,
831        COMPAT_STREAM_EPOCH,
832        FACT_CREATE_OPERATION,
833        FACT_CREATE_PAYLOAD_SCHEMA,
834        payload,
835    )?;
836    debug_assert_eq!(entry.sequence, sequence);
837    tx.commit()?;
838    Ok(ReplayOutcome::Applied { sequence })
839}
840
841pub fn next_expected_sequence_for_epoch(
842    conn: &Connection,
843    home_device_id: &str,
844    store_id: &str,
845    stream_epoch: u64,
846) -> Result<i64, MemoryError> {
847    validate_stream_identity(home_device_id, store_id, stream_epoch)?;
848    let epoch = i64::try_from(stream_epoch).map_err(|_| MemoryError::InvalidConfig {
849        field: "replication_stream_epoch",
850        reason: "does not fit SQLite INTEGER".to_string(),
851    })?;
852    Ok(conn
853        .query_row(
854            "SELECT next_sequence FROM replication_streams
855             WHERE home_device_id = ?1 AND store_id = ?2 AND stream_epoch = ?3",
856            rusqlite::params![home_device_id, store_id, epoch],
857            |row| row.get(0),
858        )
859        .optional()?
860        .unwrap_or(1))
861}
862
863#[deprecated(note = "use next_expected_sequence_for_epoch")]
864pub fn next_expected_sequence(
865    conn: &Connection,
866    home_device_id: &str,
867    store_id: &str,
868) -> Result<i64, MemoryError> {
869    next_expected_sequence_for_epoch(conn, home_device_id, store_id, COMPAT_STREAM_EPOCH)
870}
871
872#[cfg(test)]
873mod tests {
874    use super::*;
875    use std::sync::{Arc, Barrier};
876    use std::thread;
877    use std::time::Duration;
878
879    fn test_conn() -> Connection {
880        let conn = Connection::open_in_memory().unwrap();
881        conn.execute_batch(MIGRATION_V37).unwrap();
882        conn.execute_batch(MIGRATION_V38).unwrap();
883        conn
884    }
885
886    #[test]
887    #[allow(deprecated)]
888    fn first_record_uses_genesis_and_export_verifies_chain() {
889        let conn = test_conn();
890        append_journal_entry(&conn, "device-1", "store-1", "add_fact", b"payload-1").unwrap();
891        append_journal_entry(&conn, "device-1", "store-1", "add_fact", b"payload-2").unwrap();
892        let batch = export_verified_contiguous(&conn, "device-1", "store-1", 1, 1, 10).unwrap();
893        assert_eq!(batch.status, ExportStatus::End);
894        assert_eq!(batch.entries.len(), 2);
895        assert_eq!(batch.entries[0].predecessor_digest, GENESIS_PREDECESSOR);
896        assert_eq!(
897            batch.entries[1].predecessor_digest,
898            batch.entries[0].envelope_digest
899        );
900        assert_eq!(batch.next_sequence, 3);
901    }
902
903    #[test]
904    #[allow(deprecated)]
905    fn failed_mutation_consumes_no_sequence() {
906        let conn = test_conn();
907        let result = mutate_and_journal(
908            &conn,
909            "device-1",
910            "store-1",
911            "add_fact",
912            b"payload",
913            |_conn| Err::<(), _>(MemoryError::Database(rusqlite::Error::InvalidQuery)),
914        );
915        assert!(result.is_err());
916        assert_eq!(
917            next_expected_sequence_for_epoch(&conn, "device-1", "store-1", 1).unwrap(),
918            1
919        );
920    }
921
922    #[test]
923    #[allow(deprecated)]
924    fn same_sequence_changed_payload_is_conflict() {
925        let conn = test_conn();
926        conn.execute("CREATE TABLE replayed(value TEXT)", [])
927            .unwrap();
928        let first =
929            replay_journal_entry(&conn, "device-1", "store-1", 1, "add_fact", b"a", |conn| {
930                conn.execute("INSERT INTO replayed(value) VALUES ('a')", [])?;
931                Ok(())
932            })
933            .unwrap();
934        assert_eq!(first, ReplayOutcome::Applied { sequence: 1 });
935        let second =
936            replay_journal_entry(&conn, "device-1", "store-1", 1, "add_fact", b"b", |_conn| {
937                panic!("conflicting replay must not run")
938            })
939            .unwrap();
940        assert_eq!(second, ReplayOutcome::Conflict { sequence: 1 });
941        let value: String = conn
942            .query_row("SELECT value FROM replayed", [], |row| row.get(0))
943            .unwrap();
944        assert_eq!(value, "a");
945    }
946
947    #[test]
948    #[allow(deprecated)]
949    fn separate_connections_allocate_one_contiguous_stream() {
950        let temp = tempfile::TempDir::new().unwrap();
951        let path = temp.path().join("journal.db");
952        let setup = Connection::open(&path).unwrap();
953        setup.execute_batch(MIGRATION_V37).unwrap();
954        setup.execute_batch(MIGRATION_V38).unwrap();
955        drop(setup);
956
957        let workers = 8;
958        let barrier = Arc::new(Barrier::new(workers));
959        let handles: Vec<_> = (0..workers)
960            .map(|worker| {
961                let path = path.clone();
962                let barrier = Arc::clone(&barrier);
963                thread::spawn(move || {
964                    let conn = Connection::open(path).unwrap();
965                    conn.busy_timeout(Duration::from_secs(10)).unwrap();
966                    barrier.wait();
967                    append_journal_entry(
968                        &conn,
969                        "device-1",
970                        "store-1",
971                        "add_fact",
972                        format!("payload-{worker}").as_bytes(),
973                    )
974                    .unwrap()
975                })
976            })
977            .collect();
978        let mut sequences: Vec<_> = handles.into_iter().map(|h| h.join().unwrap()).collect();
979        sequences.sort_unstable();
980        assert_eq!(sequences, (1..=workers as i64).collect::<Vec<_>>());
981    }
982
983    #[test]
984    #[allow(deprecated)]
985    fn tampered_chain_is_reported_as_corrupt() {
986        let conn = test_conn();
987        append_journal_entry(&conn, "device-1", "store-1", "add_fact", b"payload").unwrap();
988        conn.execute(
989            "UPDATE mutation_journal SET predecessor_digest = ?1 WHERE sequence = 1",
990            [vec![7_u8; 32]],
991        )
992        .unwrap();
993        let batch = export_verified_contiguous(&conn, "device-1", "store-1", 1, 1, 10).unwrap();
994        assert!(matches!(
995            batch.status,
996            ExportStatus::Corrupt { sequence: 1, .. }
997        ));
998    }
999}