1use 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#[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#[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#[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
68pub 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#[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
129pub 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
145pub 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
167pub 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
227pub 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#[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#[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#[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 Empty,
520 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
557pub 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#[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#[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}