1use std::collections::BTreeMap;
11use std::fs;
12use std::io::Write;
13use std::path::{Path, PathBuf};
14use std::sync::atomic::{AtomicU64, Ordering};
15
16use crate::commit_marker::{
17 CommitMarkerRecord, MARKER_SEGMENT_HEADER_BYTES, MarkerSegmentHeader, recover_valid_prefix,
18 segment_id_for_commit_seq,
19};
20use crate::symbol_log::scan_symbol_segment;
21use fsqlite_error::{FrankenError, Result};
22use fsqlite_types::{EpochId, ObjectId, SymbolRecord};
23use tracing::{debug, error, info, warn};
24
25const MASTER_KEY_DOMAIN: &[u8] = b"fsqlite:symbol-auth-master:v1";
31
32const EPOCH_KEY_DOMAIN: &[u8] = b"fsqlite:symbol-auth:epoch:v1";
36
37const ROOT_POINTER_AUTH_DOMAIN: &[u8] = b"fsqlite:ecs-root-auth:v1";
39
40const ROOT_BOOTSTRAP_BEAD_ID: &str = "bd-1hi.25";
42const ROOT_BOOTSTRAP_LOGGING_STANDARD: &str = "bd-1fpm";
44static ROOT_TMP_SUFFIX_COUNTER: AtomicU64 = AtomicU64::new(0);
46
47#[derive(Debug)]
54pub struct EpochClock {
55 current: AtomicU64,
56}
57
58impl EpochClock {
59 #[must_use]
61 pub fn new(initial: EpochId) -> Self {
62 Self {
63 current: AtomicU64::new(initial.get()),
64 }
65 }
66
67 #[must_use]
69 pub fn current(&self) -> EpochId {
70 EpochId::new(self.current.load(Ordering::Acquire))
71 }
72
73 pub fn increment(&self) -> Result<EpochId> {
82 loop {
83 let old = self.current.load(Ordering::Acquire);
84 let new = old.checked_add(1).ok_or_else(|| {
85 error!(
86 bead_id = "bd-3go.12",
87 old_epoch = old,
88 "epoch counter overflow — cannot increment past u64::MAX"
89 );
90 FrankenError::OutOfRange {
91 what: "ecs_epoch".to_owned(),
92 value: old.to_string(),
93 }
94 })?;
95 if self
96 .current
97 .compare_exchange_weak(old, new, Ordering::AcqRel, Ordering::Acquire)
98 .is_ok()
99 {
100 info!(
101 bead_id = "bd-3go.12",
102 old_epoch = old,
103 new_epoch = new,
104 "epoch incremented"
105 );
106 return Ok(EpochId::new(new));
107 }
108 }
109 }
110
111 pub fn store(&self, epoch: EpochId) {
116 self.current.store(epoch.get(), Ordering::Release);
117 }
118}
119
120#[derive(Debug, Clone, PartialEq, Eq)]
124pub struct EpochAuthKey([u8; 32]);
125
126impl EpochAuthKey {
127 #[must_use]
129 pub fn as_bytes(&self) -> &[u8; 32] {
130 &self.0
131 }
132}
133
134pub fn derive_master_key_from_dek(dek: &[u8; 32]) -> [u8; 32] {
140 let keyed_hasher = blake3::Hasher::new_keyed(dek);
141 let mut hasher = keyed_hasher;
142 hasher.update(MASTER_KEY_DOMAIN);
143 let hash = hasher.finalize();
144 debug!(
145 bead_id = "bd-3go.12",
146 domain = std::str::from_utf8(MASTER_KEY_DOMAIN).unwrap_or("<invalid>"),
147 "derived master key from DEK with domain separation"
148 );
149 *hash.as_bytes()
150}
151
152#[must_use]
158pub fn derive_epoch_auth_key(master_key: &[u8; 32], epoch: EpochId) -> EpochAuthKey {
159 let mut hasher = blake3::Hasher::new_keyed(master_key);
160 hasher.update(EPOCH_KEY_DOMAIN);
161 hasher.update(&epoch.get().to_le_bytes());
162 let hash = hasher.finalize();
163 debug!(
164 bead_id = "bd-3go.12",
165 epoch = epoch.get(),
166 domain = std::str::from_utf8(EPOCH_KEY_DOMAIN).unwrap_or("<invalid>"),
167 "derived epoch auth key (NOT logging key material)"
168 );
169 EpochAuthKey(*hash.as_bytes())
170}
171
172#[derive(Debug, Clone, Copy, PartialEq, Eq)]
176pub enum BarrierOutcome {
177 AllArrived {
179 new_epoch: EpochId,
181 },
182 Timeout {
184 arrived: usize,
186 expected: usize,
188 },
189 Cancelled,
191}
192
193#[derive(Debug)]
201pub struct EpochBarrier {
202 current_epoch: EpochId,
204 expected: usize,
206 arrived: AtomicU64,
208 cancelled: std::sync::atomic::AtomicBool,
210}
211
212impl EpochBarrier {
213 #[must_use]
218 pub fn new(current_epoch: EpochId, participants: usize) -> Self {
219 info!(
220 bead_id = "bd-3go.12",
221 epoch = current_epoch.get(),
222 participants,
223 "epoch barrier created"
224 );
225 Self {
226 current_epoch,
227 expected: participants,
228 arrived: AtomicU64::new(0),
229 cancelled: std::sync::atomic::AtomicBool::new(false),
230 }
231 }
232
233 #[must_use]
235 pub fn epoch(&self) -> EpochId {
236 self.current_epoch
237 }
238
239 #[must_use]
241 pub fn arrived_count(&self) -> usize {
242 let val = self.arrived.load(Ordering::Acquire);
243 usize::try_from(val).unwrap_or(usize::MAX)
244 }
245
246 #[must_use]
248 pub fn expected_count(&self) -> usize {
249 self.expected
250 }
251
252 pub fn arrive(&self, participant_name: &str) -> bool {
256 if self.cancelled.load(Ordering::Acquire) {
257 warn!(
258 bead_id = "bd-3go.12",
259 participant = participant_name,
260 "participant arrived at cancelled barrier — ignoring"
261 );
262 return false;
263 }
264 let prev = self.arrived.fetch_add(1, Ordering::AcqRel);
265 let new_count = usize::try_from(prev.saturating_add(1)).unwrap_or(usize::MAX);
266 debug!(
267 bead_id = "bd-3go.12",
268 participant = participant_name,
269 arrived = new_count,
270 expected = self.expected,
271 "barrier participant arrived"
272 );
273 new_count >= self.expected
274 }
275
276 #[must_use]
278 pub fn is_complete(&self) -> bool {
279 self.arrived_count() >= self.expected
280 }
281
282 pub fn cancel(&self) {
284 self.cancelled.store(true, Ordering::Release);
285 warn!(
286 bead_id = "bd-3go.12",
287 epoch = self.current_epoch.get(),
288 arrived = self.arrived_count(),
289 expected = self.expected,
290 "epoch barrier cancelled — epoch will NOT advance"
291 );
292 }
293
294 #[must_use]
296 pub fn is_cancelled(&self) -> bool {
297 self.cancelled.load(Ordering::Acquire)
298 }
299
300 pub fn resolve(&self, clock: &EpochClock) -> Result<BarrierOutcome> {
309 if self.is_cancelled() {
310 return Ok(BarrierOutcome::Cancelled);
311 }
312 if !self.is_complete() {
313 return Ok(BarrierOutcome::Timeout {
314 arrived: self.arrived_count(),
315 expected: self.expected,
316 });
317 }
318 let new_epoch = clock.increment()?;
319 info!(
320 bead_id = "bd-3go.12",
321 old_epoch = self.current_epoch.get(),
322 new_epoch = new_epoch.get(),
323 participants = self.expected,
324 "epoch transition completed — all participants arrived"
325 );
326 Ok(BarrierOutcome::AllArrived { new_epoch })
327 }
328}
329
330pub fn validate_symbol_epoch(
341 symbol_epoch: EpochId,
342 window: &fsqlite_types::SymbolValidityWindow,
343) -> Result<()> {
344 if window.contains(symbol_epoch) {
345 Ok(())
346 } else {
347 error!(
348 bead_id = "bd-3go.12",
349 symbol_epoch = symbol_epoch.get(),
350 window_from = window.from_epoch.get(),
351 window_to = window.to_epoch.get(),
352 "symbol epoch outside validity window — fail-closed rejection"
353 );
354 Err(FrankenError::DatabaseCorrupt {
355 detail: format!(
356 "symbol epoch {} outside validity window [{}, {}]",
357 symbol_epoch.get(),
358 window.from_epoch.get(),
359 window.to_epoch.get(),
360 ),
361 })
362 }
363}
364
365pub const ECS_ROOT_POINTER_MAGIC: [u8; 4] = *b"FSRT";
369pub const ECS_ROOT_POINTER_VERSION: u32 = 1;
371pub const ECS_ROOT_POINTER_BYTES: usize = 56;
373const ECS_ROOT_POINTER_CHECKSUM_INPUT_BYTES: usize = 32;
375const ECS_ROOT_POINTER_AUTH_INPUT_BYTES: usize = 40;
377
378pub const ROOT_MANIFEST_MAGIC: [u8; 8] = *b"FSQLROOT";
380pub const ROOT_MANIFEST_VERSION: u32 = 1;
382
383#[derive(Debug, Clone, Copy, PartialEq, Eq)]
388pub struct EcsRootPointer {
389 pub manifest_object_id: ObjectId,
391 pub ecs_epoch: EpochId,
393 pub root_auth_tag: [u8; 16],
395}
396
397impl EcsRootPointer {
398 #[must_use]
400 pub const fn unauthed(manifest_object_id: ObjectId, ecs_epoch: EpochId) -> Self {
401 Self {
402 manifest_object_id,
403 ecs_epoch,
404 root_auth_tag: [0_u8; 16],
405 }
406 }
407
408 #[must_use]
410 pub fn authed(manifest_object_id: ObjectId, ecs_epoch: EpochId, master_key: &[u8; 32]) -> Self {
411 let mut pointer = Self::unauthed(manifest_object_id, ecs_epoch);
412 let auth_input = pointer.auth_input_bytes();
413 pointer.root_auth_tag = compute_root_pointer_auth_tag(master_key, &auth_input);
414 pointer
415 }
416
417 #[must_use]
419 pub fn encode(&self) -> [u8; ECS_ROOT_POINTER_BYTES] {
420 let mut out = [0_u8; ECS_ROOT_POINTER_BYTES];
421 out[0..4].copy_from_slice(&ECS_ROOT_POINTER_MAGIC);
422 out[4..8].copy_from_slice(&ECS_ROOT_POINTER_VERSION.to_le_bytes());
423 out[8..24].copy_from_slice(self.manifest_object_id.as_bytes());
424 out[24..32].copy_from_slice(&self.ecs_epoch.get().to_le_bytes());
425 let checksum = xxhash_rust::xxh3::xxh3_64(&out[..ECS_ROOT_POINTER_CHECKSUM_INPUT_BYTES]);
426 out[32..40].copy_from_slice(&checksum.to_le_bytes());
427 out[40..56].copy_from_slice(&self.root_auth_tag);
428 out
429 }
430
431 pub fn decode(
437 bytes: &[u8],
438 symbol_auth_enabled: bool,
439 master_key: Option<&[u8; 32]>,
440 ) -> Result<Self> {
441 if bytes.len() != ECS_ROOT_POINTER_BYTES {
442 return Err(FrankenError::DatabaseCorrupt {
443 detail: format!(
444 "ecs/root size mismatch: expected {ECS_ROOT_POINTER_BYTES}, got {}",
445 bytes.len()
446 ),
447 });
448 }
449 if bytes[0..4] != ECS_ROOT_POINTER_MAGIC {
450 return Err(FrankenError::DatabaseCorrupt {
451 detail: format!(
452 "invalid ecs/root magic: {:02X?} (reason=bad_magic)",
453 &bytes[0..4]
454 ),
455 });
456 }
457 let version = read_u32_le_at(bytes, 4, "root.version")?;
458 if version != ECS_ROOT_POINTER_VERSION {
459 return Err(FrankenError::DatabaseCorrupt {
460 detail: format!(
461 "unsupported ecs/root version {version} (expected {ECS_ROOT_POINTER_VERSION})"
462 ),
463 });
464 }
465
466 let mut manifest_id = [0_u8; 16];
467 manifest_id.copy_from_slice(&bytes[8..24]);
468 let manifest_object_id = ObjectId::from_bytes(manifest_id);
469 let ecs_epoch_raw = read_u64_le_at(bytes, 24, "root.ecs_epoch")?;
470 let ecs_epoch = EpochId::new(ecs_epoch_raw);
471
472 let stored_checksum = read_u64_le_at(bytes, 32, "root.checksum")?;
473 let computed_checksum =
474 xxhash_rust::xxh3::xxh3_64(&bytes[..ECS_ROOT_POINTER_CHECKSUM_INPUT_BYTES]);
475 if stored_checksum != computed_checksum {
476 error!(
477 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
478 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
479 reason_code = "checksum_mismatch",
480 stored_checksum = stored_checksum,
481 computed_checksum = computed_checksum,
482 "ecs/root checksum verification failed"
483 );
484 return Err(FrankenError::DatabaseCorrupt {
485 detail: format!(
486 "ecs/root checksum mismatch (reason=checksum_mismatch): stored={stored_checksum:#018X}, computed={computed_checksum:#018X}"
487 ),
488 });
489 }
490
491 let mut root_auth_tag = [0_u8; 16];
492 root_auth_tag.copy_from_slice(&bytes[40..56]);
493
494 if symbol_auth_enabled {
495 let Some(master_key) = master_key else {
496 return Err(FrankenError::DatabaseCorrupt {
497 detail: "symbol_auth enabled but master key is missing (reason=auth_failed)"
498 .to_owned(),
499 });
500 };
501 let expected = compute_root_pointer_auth_tag(
502 master_key,
503 &bytes[..ECS_ROOT_POINTER_AUTH_INPUT_BYTES],
504 );
505 if root_auth_tag != expected {
506 error!(
507 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
508 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
509 reason_code = "auth_failed",
510 "ecs/root auth-tag verification failed"
511 );
512 return Err(FrankenError::DatabaseCorrupt {
513 detail: "ecs/root auth tag verification failed (reason=auth_failed)".to_owned(),
514 });
515 }
516 } else if root_auth_tag != [0_u8; 16] {
517 return Err(FrankenError::DatabaseCorrupt {
518 detail: "ecs/root auth tag must be all-zero when symbol_auth=off".to_owned(),
519 });
520 }
521
522 Ok(Self {
523 manifest_object_id,
524 ecs_epoch,
525 root_auth_tag,
526 })
527 }
528
529 #[must_use]
531 fn auth_input_bytes(&self) -> [u8; ECS_ROOT_POINTER_AUTH_INPUT_BYTES] {
532 let encoded = self.encode();
533 let mut out = [0_u8; ECS_ROOT_POINTER_AUTH_INPUT_BYTES];
534 out.copy_from_slice(&encoded[..ECS_ROOT_POINTER_AUTH_INPUT_BYTES]);
535 out
536 }
537}
538
539#[derive(Debug, Clone, PartialEq, Eq)]
541pub struct RootManifest {
542 pub database_name: String,
544 pub current_commit: ObjectId,
546 pub commit_seq: u64,
548 pub schema_snapshot: ObjectId,
550 pub schema_epoch: u64,
552 pub ecs_epoch: EpochId,
554 pub checkpoint_base: ObjectId,
556 pub gc_horizon: u64,
558 pub created_at: u64,
560 pub updated_at: u64,
562}
563
564impl RootManifest {
565 pub fn encode(&self) -> Result<Vec<u8>> {
571 let name_bytes = self.database_name.as_bytes();
572 let name_len = u32::try_from(name_bytes.len()).map_err(|_| FrankenError::OutOfRange {
573 what: "root_manifest.database_name_len".to_owned(),
574 value: name_bytes.len().to_string(),
575 })?;
576
577 let mut out = Vec::with_capacity(name_bytes.len().saturating_add(128));
578 out.extend_from_slice(&ROOT_MANIFEST_MAGIC);
579 out.extend_from_slice(&ROOT_MANIFEST_VERSION.to_le_bytes());
580 out.extend_from_slice(&name_len.to_le_bytes());
581 out.extend_from_slice(name_bytes);
582 out.extend_from_slice(self.current_commit.as_bytes());
583 out.extend_from_slice(&self.commit_seq.to_le_bytes());
584 out.extend_from_slice(self.schema_snapshot.as_bytes());
585 out.extend_from_slice(&self.schema_epoch.to_le_bytes());
586 out.extend_from_slice(&self.ecs_epoch.get().to_le_bytes());
587 out.extend_from_slice(self.checkpoint_base.as_bytes());
588 out.extend_from_slice(&self.gc_horizon.to_le_bytes());
589 out.extend_from_slice(&self.created_at.to_le_bytes());
590 out.extend_from_slice(&self.updated_at.to_le_bytes());
591 let checksum = xxhash_rust::xxh3::xxh3_64(&out);
592 out.extend_from_slice(&checksum.to_le_bytes());
593 Ok(out)
594 }
595
596 pub fn decode(bytes: &[u8]) -> Result<Self> {
602 if bytes.len() < 120 {
603 return Err(FrankenError::DatabaseCorrupt {
604 detail: format!(
605 "root manifest too short: expected >= 120 bytes, got {}",
606 bytes.len()
607 ),
608 });
609 }
610 if bytes[0..8] != ROOT_MANIFEST_MAGIC {
611 return Err(FrankenError::DatabaseCorrupt {
612 detail: format!("invalid root manifest magic: {:02X?}", &bytes[0..8]),
613 });
614 }
615 let version = read_u32_le_at(bytes, 8, "root_manifest.version")?;
616 if version != ROOT_MANIFEST_VERSION {
617 return Err(FrankenError::DatabaseCorrupt {
618 detail: format!(
619 "unsupported root manifest version {version} (expected {ROOT_MANIFEST_VERSION})"
620 ),
621 });
622 }
623
624 let name_len_u32 = read_u32_le_at(bytes, 12, "root_manifest.database_name_len")?;
625 let name_len = u32_to_usize(name_len_u32, "root_manifest.database_name_len")?;
626 let mut cursor = 16_usize;
627 let name_end = checked_add(cursor, name_len, "root_manifest.database_name_end")?;
628 if name_end > bytes.len() {
629 return Err(FrankenError::DatabaseCorrupt {
630 detail: format!(
631 "root manifest name out of bounds: end={name_end}, len={}",
632 bytes.len()
633 ),
634 });
635 }
636 let database_name = std::str::from_utf8(&bytes[cursor..name_end])
637 .map_err(|err| FrankenError::DatabaseCorrupt {
638 detail: format!("root manifest database_name is not UTF-8: {err}"),
639 })?
640 .to_owned();
641 cursor = name_end;
642
643 let current_commit = read_object_id_at(bytes, cursor, "root_manifest.current_commit")?;
644 cursor = checked_add(cursor, 16, "root_manifest.cursor.current_commit")?;
645 let commit_seq = read_u64_le_at(bytes, cursor, "root_manifest.commit_seq")?;
646 cursor = checked_add(cursor, 8, "root_manifest.cursor.commit_seq")?;
647 let schema_snapshot = read_object_id_at(bytes, cursor, "root_manifest.schema_snapshot")?;
648 cursor = checked_add(cursor, 16, "root_manifest.cursor.schema_snapshot")?;
649 let schema_epoch = read_u64_le_at(bytes, cursor, "root_manifest.schema_epoch")?;
650 cursor = checked_add(cursor, 8, "root_manifest.cursor.schema_epoch")?;
651 let ecs_epoch_raw = read_u64_le_at(bytes, cursor, "root_manifest.ecs_epoch")?;
652 let ecs_epoch = EpochId::new(ecs_epoch_raw);
653 cursor = checked_add(cursor, 8, "root_manifest.cursor.ecs_epoch")?;
654 let checkpoint_base = read_object_id_at(bytes, cursor, "root_manifest.checkpoint_base")?;
655 cursor = checked_add(cursor, 16, "root_manifest.cursor.checkpoint_base")?;
656 let gc_horizon = read_u64_le_at(bytes, cursor, "root_manifest.gc_horizon")?;
657 cursor = checked_add(cursor, 8, "root_manifest.cursor.gc_horizon")?;
658 let created_at = read_u64_le_at(bytes, cursor, "root_manifest.created_at")?;
659 cursor = checked_add(cursor, 8, "root_manifest.cursor.created_at")?;
660 let updated_at = read_u64_le_at(bytes, cursor, "root_manifest.updated_at")?;
661 cursor = checked_add(cursor, 8, "root_manifest.cursor.updated_at")?;
662
663 let checksum_end = checked_add(cursor, 8, "root_manifest.cursor.checksum_end")?;
664 if checksum_end != bytes.len() {
665 return Err(FrankenError::DatabaseCorrupt {
666 detail: format!(
667 "root manifest trailing bytes present: parsed_end={checksum_end}, actual_len={}",
668 bytes.len()
669 ),
670 });
671 }
672 let stored_checksum = read_u64_le_at(bytes, cursor, "root_manifest.checksum")?;
673 let computed_checksum = xxhash_rust::xxh3::xxh3_64(&bytes[..cursor]);
674 if stored_checksum != computed_checksum {
675 return Err(FrankenError::DatabaseCorrupt {
676 detail: format!(
677 "root manifest checksum mismatch: stored={stored_checksum:#018X}, computed={computed_checksum:#018X}"
678 ),
679 });
680 }
681
682 Ok(Self {
683 database_name,
684 current_commit,
685 commit_seq,
686 schema_snapshot,
687 schema_epoch,
688 ecs_epoch,
689 checkpoint_base,
690 gc_horizon,
691 created_at,
692 updated_at,
693 })
694 }
695}
696
697#[must_use]
699pub fn compute_root_pointer_auth_tag(master_key: &[u8; 32], magic_to_checksum: &[u8]) -> [u8; 16] {
700 let mut hasher = blake3::Hasher::new_keyed(master_key);
701 hasher.update(ROOT_POINTER_AUTH_DOMAIN);
702 hasher.update(magic_to_checksum);
703 let digest = hasher.finalize();
704 let mut out = [0_u8; 16];
705 out.copy_from_slice(&digest.as_bytes()[..16]);
706 out
707}
708
709#[derive(Debug, Clone, PartialEq, Eq)]
711pub struct NativeBootstrapLayout {
712 pub ecs_dir: PathBuf,
714}
715
716impl NativeBootstrapLayout {
717 #[must_use]
719 pub fn new(ecs_dir: impl Into<PathBuf>) -> Self {
720 Self {
721 ecs_dir: ecs_dir.into(),
722 }
723 }
724
725 #[must_use]
727 pub fn root_path(&self) -> PathBuf {
728 self.ecs_dir.join("root")
729 }
730
731 #[must_use]
733 pub fn symbols_dir(&self) -> PathBuf {
734 self.ecs_dir.join("symbols")
735 }
736
737 #[must_use]
739 pub fn markers_dir(&self) -> PathBuf {
740 self.ecs_dir.join("markers")
741 }
742}
743
744#[derive(Debug, Clone, PartialEq, Eq)]
746pub struct NativeBootstrapState {
747 pub root_pointer: EcsRootPointer,
749 pub manifest: RootManifest,
751 pub latest_marker: CommitMarkerRecord,
753 pub schema_snapshot_bytes: Vec<u8>,
755 pub checkpoint_base_bytes: Vec<u8>,
757}
758
759pub fn read_root_pointer(
765 root_path: &Path,
766 symbol_auth_enabled: bool,
767 master_key: Option<&[u8; 32]>,
768) -> Result<EcsRootPointer> {
769 let bytes = fs::read(root_path).map_err(|err| {
770 error!(
771 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
772 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
773 reason_code = "scan_failed",
774 path = %root_path.display(),
775 error = %err,
776 "failed reading ecs/root"
777 );
778 FrankenError::Io(err)
779 })?;
780 EcsRootPointer::decode(&bytes, symbol_auth_enabled, master_key)
781}
782
783pub fn write_root_pointer_atomic(root_path: &Path, pointer: EcsRootPointer) -> Result<()> {
789 let Some(parent) = root_path.parent() else {
790 return Err(FrankenError::DatabaseCorrupt {
791 detail: format!("ecs/root has no parent directory: {}", root_path.display()),
792 });
793 };
794 fs::create_dir_all(parent)?;
795
796 let pid = std::process::id();
797 let suffix = ROOT_TMP_SUFFIX_COUNTER.fetch_add(1, Ordering::SeqCst);
798 let tmp_name = format!(".root.tmp.{pid}.{suffix}");
799 let tmp_path = parent.join(tmp_name);
800
801 let bytes = pointer.encode();
802 let mut temp = fs::OpenOptions::new()
803 .write(true)
804 .create_new(true)
805 .open(&tmp_path)?;
806 temp.write_all(&bytes)?;
807 temp.sync_all()?;
808 fs::rename(&tmp_path, root_path)?;
809 let parent_dir = fs::File::open(parent)?;
810 parent_dir.sync_all()?;
811
812 info!(
813 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
814 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
815 path = %root_path.display(),
816 root_epoch = pointer.ecs_epoch.get(),
817 "wrote ecs/root atomically"
818 );
819
820 Ok(())
821}
822
823pub fn build_root_pointer(
829 manifest_object_id: ObjectId,
830 ecs_epoch: EpochId,
831 symbol_auth_enabled: bool,
832 master_key: Option<&[u8; 32]>,
833) -> Result<EcsRootPointer> {
834 if symbol_auth_enabled {
835 let Some(master_key) = master_key else {
836 return Err(FrankenError::DatabaseCorrupt {
837 detail: "symbol_auth enabled but master key is missing (reason=auth_failed)"
838 .to_owned(),
839 });
840 };
841 Ok(EcsRootPointer::authed(
842 manifest_object_id,
843 ecs_epoch,
844 master_key,
845 ))
846 } else {
847 Ok(EcsRootPointer::unauthed(manifest_object_id, ecs_epoch))
848 }
849}
850
851pub fn bootstrap_native_mode(
862 layout: &NativeBootstrapLayout,
863 symbol_auth_enabled: bool,
864 master_key: Option<&[u8; 32]>,
865) -> Result<NativeBootstrapState> {
866 debug!(
867 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
868 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
869 step = 1_u8,
870 root_path = %layout.root_path().display(),
871 symbol_auth_enabled = symbol_auth_enabled,
872 "bootstrap step 1: reading ecs/root"
873 );
874 let root_path = layout.root_path();
875 let root_pointer = read_root_pointer(&root_path, symbol_auth_enabled, master_key)?;
876 info!(
877 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
878 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
879 step = 1_u8,
880 duration_ms = 0_u64,
881 root_epoch = root_pointer.ecs_epoch.get(),
882 manifest_object_id = %root_pointer.manifest_object_id,
883 "bootstrap steps 1-3 complete"
884 );
885 bootstrap_from_root_pointer(layout, root_pointer)
886}
887
888pub fn bootstrap_native_mode_with_recovery(
894 layout: &NativeBootstrapLayout,
895 symbol_auth_enabled: bool,
896 master_key: Option<&[u8; 32]>,
897) -> Result<NativeBootstrapState> {
898 match bootstrap_native_mode(layout, symbol_auth_enabled, master_key) {
899 Ok(state) => Ok(state),
900 Err(initial_err) => {
901 debug!(
902 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
903 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
904 reason_code = "retry_scan_recovery",
905 error = %initial_err,
906 "bootstrap entering degraded scan-based recovery path"
907 );
908 warn!(
909 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
910 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
911 reason_code = "retry_scan_recovery",
912 error = %initial_err,
913 "bootstrap from ecs/root failed; attempting scan-based recovery"
914 );
915
916 let recovered_pointer =
917 recover_root_pointer_from_scan(layout, symbol_auth_enabled, master_key)?;
918
919 let state = bootstrap_from_root_pointer(layout, recovered_pointer)?;
924
925 write_root_pointer_atomic(&layout.root_path(), recovered_pointer)?;
927 Ok(state)
928 }
929 }
930}
931
932pub fn recover_root_pointer_from_scan(
938 layout: &NativeBootstrapLayout,
939 symbol_auth_enabled: bool,
940 master_key: Option<&[u8; 32]>,
941) -> Result<EcsRootPointer> {
942 let marker_tip = scan_latest_marker(layout.markers_dir().as_path())?;
943 let mut grouped: BTreeMap<ObjectId, Vec<SymbolRecord>> = BTreeMap::new();
944 let symbol_segments = sorted_segment_paths(layout.symbols_dir().as_path())?;
945 debug!(
946 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
947 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
948 segments = symbol_segments.len(),
949 marker_tip_commit_seq = marker_tip.as_ref().map_or(0_u64, |m| m.commit_seq),
950 "scan recovery started"
951 );
952
953 for (_, segment_path) in &symbol_segments {
954 debug!(
955 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
956 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
957 segment = %segment_path.display(),
958 "scan recovery inspecting symbol segment"
959 );
960 let scan = scan_symbol_segment(segment_path)?;
961 for row in scan.records {
962 grouped
963 .entry(row.record.object_id)
964 .or_default()
965 .push(row.record);
966 }
967 }
968
969 let mut best: Option<(ObjectId, RootManifest, bool)> = None;
970 for (object_id, records) in grouped {
971 let Ok(payload) = reconstruct_payload_from_source_symbols(records) else {
972 continue;
973 };
974 let Ok(manifest) = RootManifest::decode(&payload) else {
975 continue;
976 };
977
978 let marker_matches = marker_tip.as_ref().is_some_and(|tip| {
979 manifest.current_commit.as_bytes() == &tip.marker_id
980 && manifest.commit_seq == tip.commit_seq
981 });
982
983 match &best {
984 None => best = Some((object_id, manifest, marker_matches)),
985 Some((_, best_manifest, best_marker_matches)) => {
986 let better_marker_match = marker_matches && !best_marker_matches;
987 let better_commit = manifest.commit_seq > best_manifest.commit_seq;
988 let better_update = manifest.commit_seq == best_manifest.commit_seq
989 && manifest.updated_at > best_manifest.updated_at;
990 if better_marker_match || better_commit || better_update {
991 best = Some((object_id, manifest, marker_matches));
992 }
993 }
994 }
995 }
996
997 let Some((manifest_object_id, manifest, _)) = best else {
998 error!(
999 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1000 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1001 reason_code = "scan_failed",
1002 segments_scanned = symbol_segments.len(),
1003 "scan recovery could not find a valid RootManifest candidate"
1004 );
1005 return Err(FrankenError::DatabaseCorrupt {
1006 detail: "scan recovery failed: no valid RootManifest candidate (reason=scan_failed)"
1007 .to_owned(),
1008 });
1009 };
1010
1011 info!(
1012 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1013 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1014 segments_scanned = symbol_segments.len(),
1015 best_candidate_commit_seq = manifest.commit_seq,
1016 chosen_root_pointer = %manifest_object_id,
1017 "scan recovery selected root manifest candidate"
1018 );
1019
1020 build_root_pointer(
1021 manifest_object_id,
1022 manifest.ecs_epoch,
1023 symbol_auth_enabled,
1024 master_key,
1025 )
1026}
1027
1028#[allow(clippy::too_many_lines)]
1029fn bootstrap_from_root_pointer(
1030 layout: &NativeBootstrapLayout,
1031 root_pointer: EcsRootPointer,
1032) -> Result<NativeBootstrapState> {
1033 info!(
1034 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1035 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1036 step = 4_u8,
1037 root_epoch = root_pointer.ecs_epoch.get(),
1038 manifest_object_id = %root_pointer.manifest_object_id,
1039 "bootstrap step 4: loading root manifest object"
1040 );
1041 let manifest_bytes = fetch_object_payload(
1042 layout.symbols_dir().as_path(),
1043 root_pointer.manifest_object_id,
1044 root_pointer.ecs_epoch,
1045 )?;
1046 let manifest = RootManifest::decode(&manifest_bytes)?;
1047 info!(
1048 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1049 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1050 step = 4_u8,
1051 duration_ms = 0_u64,
1052 root_epoch = root_pointer.ecs_epoch.get(),
1053 object_id = %root_pointer.manifest_object_id,
1054 "bootstrap step 4 complete"
1055 );
1056
1057 if manifest.ecs_epoch != root_pointer.ecs_epoch {
1058 error!(
1059 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1060 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1061 reason_code = "epoch_mismatch",
1062 root_epoch = root_pointer.ecs_epoch.get(),
1063 manifest_epoch = manifest.ecs_epoch.get(),
1064 "bootstrap step 5 failed: root epoch != manifest epoch"
1065 );
1066 return Err(FrankenError::DatabaseCorrupt {
1067 detail: format!(
1068 "root/manifest epoch mismatch (reason=epoch_mismatch): root={}, manifest={}",
1069 root_pointer.ecs_epoch.get(),
1070 manifest.ecs_epoch.get()
1071 ),
1072 });
1073 }
1074 info!(
1075 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1076 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1077 step = 5_u8,
1078 duration_ms = 0_u64,
1079 root_epoch = root_pointer.ecs_epoch.get(),
1080 object_id = %root_pointer.manifest_object_id,
1081 "bootstrap step 5 complete"
1082 );
1083
1084 info!(
1085 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1086 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1087 step = 6_u8,
1088 commit_seq = manifest.commit_seq,
1089 "bootstrap step 6: verifying marker"
1090 );
1091 let latest_marker = fetch_marker_record(layout.markers_dir().as_path(), manifest.commit_seq)?;
1092 if latest_marker.marker_id != *manifest.current_commit.as_bytes() {
1093 error!(
1094 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1095 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1096 reason_code = "marker_mismatch",
1097 manifest_commit_seq = manifest.commit_seq,
1098 "bootstrap marker mismatch"
1099 );
1100 return Err(FrankenError::DatabaseCorrupt {
1101 detail:
1102 "root manifest current_commit does not match marker stream (reason=marker_mismatch)"
1103 .to_owned(),
1104 });
1105 }
1106 info!(
1107 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1108 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1109 step = 6_u8,
1110 duration_ms = 0_u64,
1111 root_epoch = root_pointer.ecs_epoch.get(),
1112 object_id = %manifest.current_commit,
1113 "bootstrap step 6 complete"
1114 );
1115
1116 let schema_snapshot_bytes = fetch_object_payload(
1117 layout.symbols_dir().as_path(),
1118 manifest.schema_snapshot,
1119 root_pointer.ecs_epoch,
1120 )?;
1121 info!(
1122 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1123 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1124 step = 7_u8,
1125 duration_ms = 0_u64,
1126 root_epoch = root_pointer.ecs_epoch.get(),
1127 object_id = %manifest.schema_snapshot,
1128 "bootstrap step 7 complete"
1129 );
1130 let checkpoint_base_bytes = fetch_object_payload(
1131 layout.symbols_dir().as_path(),
1132 manifest.checkpoint_base,
1133 root_pointer.ecs_epoch,
1134 )?;
1135 info!(
1136 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1137 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1138 step = 8_u8,
1139 duration_ms = 0_u64,
1140 root_epoch = root_pointer.ecs_epoch.get(),
1141 object_id = %manifest.checkpoint_base,
1142 "bootstrap step 8 complete"
1143 );
1144
1145 info!(
1146 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1147 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1148 step = 9_u8,
1149 duration_ms = 0_u64,
1150 root_epoch = root_pointer.ecs_epoch.get(),
1151 commit_seq = manifest.commit_seq,
1152 schema_epoch = manifest.schema_epoch,
1153 "bootstrap sequence completed"
1154 );
1155
1156 Ok(NativeBootstrapState {
1157 root_pointer,
1158 manifest,
1159 latest_marker,
1160 schema_snapshot_bytes,
1161 checkpoint_base_bytes,
1162 })
1163}
1164
1165fn fetch_object_payload(
1166 symbols_dir: &Path,
1167 object_id: ObjectId,
1168 root_epoch: EpochId,
1169) -> Result<Vec<u8>> {
1170 let mut records = Vec::new();
1171 let segments = sorted_segment_paths(symbols_dir)?;
1172 debug!(
1173 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1174 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1175 root_epoch = root_epoch.get(),
1176 object_id = %object_id,
1177 segment_count = segments.len(),
1178 "fetching bootstrap object payload from symbol log"
1179 );
1180
1181 for (_, segment_path) in segments {
1182 let scan = scan_symbol_segment(&segment_path)?;
1183 if scan.header.epoch_id > root_epoch.get() {
1184 error!(
1185 bead_id = ROOT_BOOTSTRAP_BEAD_ID,
1186 logging_standard = ROOT_BOOTSTRAP_LOGGING_STANDARD,
1187 reason_code = "future_epoch",
1188 segment = %segment_path.display(),
1189 segment_epoch = scan.header.epoch_id,
1190 root_epoch = root_epoch.get(),
1191 "bootstrap rejected future-epoch segment"
1192 );
1193 return Err(FrankenError::DatabaseCorrupt {
1194 detail: format!(
1195 "future-epoch segment rejected (reason=future_epoch): segment_epoch={}, root_epoch={}",
1196 scan.header.epoch_id,
1197 root_epoch.get()
1198 ),
1199 });
1200 }
1201 for row in scan.records {
1202 if row.record.object_id == object_id {
1203 records.push(row.record);
1204 }
1205 }
1206 }
1207
1208 if records.is_empty() {
1209 return Err(FrankenError::DatabaseCorrupt {
1210 detail: format!("object {object_id} not found in symbol logs"),
1211 });
1212 }
1213 reconstruct_payload_from_source_symbols(records)
1214}
1215
1216fn reconstruct_payload_from_source_symbols(mut records: Vec<SymbolRecord>) -> Result<Vec<u8>> {
1217 records.sort_by_key(|record| record.esi);
1218 let Some(first) = records.first() else {
1219 return Err(FrankenError::DatabaseCorrupt {
1220 detail: "cannot reconstruct payload from empty symbol set".to_owned(),
1221 });
1222 };
1223 let first_oti = first.oti;
1225 let symbol_size_u64 = u64::from(first_oti.t);
1226 if symbol_size_u64 == 0 {
1227 return Err(FrankenError::DatabaseCorrupt {
1228 detail: "symbol_size=0 in OTI".to_owned(),
1229 });
1230 }
1231
1232 let transfer_len_usize = u64_to_usize(first_oti.f, "oti.f")?;
1233 let source_symbols = first_oti.f.div_ceil(symbol_size_u64);
1234 let source_symbols_usize = u64_to_usize(source_symbols, "source_symbols")?;
1235 let symbol_size_usize = u32_to_usize(first_oti.t, "oti.t")?;
1236 let total_bytes = source_symbols_usize
1237 .checked_mul(symbol_size_usize)
1238 .ok_or_else(|| FrankenError::DatabaseCorrupt {
1239 detail: "reconstruction size overflow".to_owned(),
1240 })?;
1241 let mut out = vec![0_u8; total_bytes];
1242 let mut seen = vec![false; source_symbols_usize];
1243
1244 for record in records {
1245 if u64::from(record.esi) >= source_symbols {
1246 continue;
1247 }
1248 if record.oti != first_oti {
1249 return Err(FrankenError::DatabaseCorrupt {
1250 detail: "inconsistent OTI across object symbols".to_owned(),
1251 });
1252 }
1253 let idx = u32_to_usize(record.esi, "esi")?;
1254 let start =
1255 idx.checked_mul(symbol_size_usize)
1256 .ok_or_else(|| FrankenError::DatabaseCorrupt {
1257 detail: "symbol offset overflow".to_owned(),
1258 })?;
1259 let end = checked_add(start, symbol_size_usize, "symbol_end")?;
1260 if end > out.len() {
1261 return Err(FrankenError::DatabaseCorrupt {
1262 detail: "symbol write out of bounds during reconstruction".to_owned(),
1263 });
1264 }
1265 if record.symbol_data.len() != symbol_size_usize {
1266 return Err(FrankenError::DatabaseCorrupt {
1267 detail: "symbol size does not match OTI.t".to_owned(),
1268 });
1269 }
1270 out[start..end].copy_from_slice(&record.symbol_data);
1271 seen[idx] = true;
1272 }
1273
1274 if !seen.iter().all(|bit| *bit) {
1275 return Err(FrankenError::DatabaseCorrupt {
1276 detail: "insufficient source symbols to reconstruct object payload".to_owned(),
1277 });
1278 }
1279 out.truncate(transfer_len_usize);
1280 Ok(out)
1281}
1282
1283fn fetch_marker_record(markers_dir: &Path, commit_seq: u64) -> Result<CommitMarkerRecord> {
1284 let segment_id = segment_id_for_commit_seq(commit_seq);
1285 let segment_path = markers_dir.join(format!("segment-{segment_id:06}.log"));
1286 let bytes = fs::read(&segment_path)?;
1287 if bytes.len() < MARKER_SEGMENT_HEADER_BYTES {
1288 return Err(FrankenError::DatabaseCorrupt {
1289 detail: format!(
1290 "marker segment {} shorter than header: {} bytes",
1291 segment_path.display(),
1292 bytes.len()
1293 ),
1294 });
1295 }
1296 let header =
1297 MarkerSegmentHeader::decode(&bytes[..MARKER_SEGMENT_HEADER_BYTES]).map_err(|err| {
1298 FrankenError::DatabaseCorrupt {
1299 detail: format!(
1300 "marker header decode failed for {}: {err}",
1301 segment_path.display()
1302 ),
1303 }
1304 })?;
1305 let records = recover_valid_prefix(&bytes).map_err(|err| FrankenError::DatabaseCorrupt {
1306 detail: format!(
1307 "marker segment recover failed for {}: {err}",
1308 segment_path.display()
1309 ),
1310 })?;
1311
1312 if commit_seq < header.start_commit_seq {
1313 return Err(FrankenError::DatabaseCorrupt {
1314 detail: format!(
1315 "commit_seq {commit_seq} precedes segment start {}",
1316 header.start_commit_seq
1317 ),
1318 });
1319 }
1320 let index_u64 = commit_seq - header.start_commit_seq;
1321 let index = u64_to_usize(index_u64, "marker_index")?;
1322 let Some(record) = records.get(index) else {
1323 return Err(FrankenError::DatabaseCorrupt {
1324 detail: format!(
1325 "marker for commit_seq {commit_seq} missing in segment {}",
1326 segment_path.display()
1327 ),
1328 });
1329 };
1330 if !record.verify_marker_id() {
1331 return Err(FrankenError::DatabaseCorrupt {
1332 detail: "marker_id verification failed (reason=marker_mismatch)".to_owned(),
1333 });
1334 }
1335 if index > 0 {
1336 for i in 1..=index {
1337 if records[i].prev_marker_id != records[i - 1].marker_id {
1338 return Err(FrankenError::DatabaseCorrupt {
1339 detail: format!("marker hash chain gap at index {i} (reason=marker_chain_gap)"),
1340 });
1341 }
1342 }
1343 }
1344 Ok(record.clone())
1345}
1346
1347fn scan_latest_marker(markers_dir: &Path) -> Result<Option<CommitMarkerRecord>> {
1348 let segments = sorted_segment_paths(markers_dir)?;
1349 let mut best: Option<CommitMarkerRecord> = None;
1350 for (_, segment_path) in segments {
1351 let bytes = fs::read(&segment_path)?;
1352 if bytes.len() < MARKER_SEGMENT_HEADER_BYTES {
1353 continue;
1354 }
1355 let Ok(records) = recover_valid_prefix(&bytes) else {
1356 continue;
1357 };
1358 if let Some(last) = records.last() {
1359 let replace = best
1360 .as_ref()
1361 .is_none_or(|existing| last.commit_seq > existing.commit_seq);
1362 if replace {
1363 best = Some(last.clone());
1364 }
1365 }
1366 }
1367 Ok(best)
1368}
1369
1370fn sorted_segment_paths(dir: &Path) -> Result<Vec<(u64, PathBuf)>> {
1371 if !dir.exists() {
1372 return Ok(Vec::new());
1373 }
1374 let mut out = Vec::new();
1375 for entry in fs::read_dir(dir)? {
1376 let entry = entry?;
1377 if !entry.file_type()?.is_file() {
1378 continue;
1379 }
1380 let name_os = entry.file_name();
1381 let Some(name) = name_os.to_str() else {
1382 continue;
1383 };
1384 let Some(segment_id) = parse_segment_id(name) else {
1385 continue;
1386 };
1387 out.push((segment_id, entry.path()));
1388 }
1389 out.sort_by_key(|(segment_id, _)| *segment_id);
1390 Ok(out)
1391}
1392
1393fn parse_segment_id(name: &str) -> Option<u64> {
1394 let body = name.strip_prefix("segment-")?.strip_suffix(".log")?;
1395 body.parse::<u64>().ok()
1396}
1397
1398fn read_object_id_at(bytes: &[u8], offset: usize, field: &str) -> Result<ObjectId> {
1399 let end = checked_add(offset, 16, field)?;
1400 if end > bytes.len() {
1401 return Err(FrankenError::DatabaseCorrupt {
1402 detail: format!("{field} out of bounds: end={end}, len={}", bytes.len()),
1403 });
1404 }
1405 let mut raw = [0_u8; 16];
1406 raw.copy_from_slice(&bytes[offset..end]);
1407 Ok(ObjectId::from_bytes(raw))
1408}
1409
1410fn read_u32_le_at(bytes: &[u8], offset: usize, field: &str) -> Result<u32> {
1411 let end = checked_add(offset, 4, field)?;
1412 if end > bytes.len() {
1413 return Err(FrankenError::DatabaseCorrupt {
1414 detail: format!("{field} out of bounds: end={end}, len={}", bytes.len()),
1415 });
1416 }
1417 Ok(u32::from_le_bytes(
1418 bytes[offset..end].try_into().expect("fixed 4-byte field"),
1419 ))
1420}
1421
1422fn read_u64_le_at(bytes: &[u8], offset: usize, field: &str) -> Result<u64> {
1423 let end = checked_add(offset, 8, field)?;
1424 if end > bytes.len() {
1425 return Err(FrankenError::DatabaseCorrupt {
1426 detail: format!("{field} out of bounds: end={end}, len={}", bytes.len()),
1427 });
1428 }
1429 Ok(u64::from_le_bytes(
1430 bytes[offset..end].try_into().expect("fixed 8-byte field"),
1431 ))
1432}
1433
1434fn u32_to_usize(value: u32, field: &str) -> Result<usize> {
1435 usize::try_from(value).map_err(|_| FrankenError::OutOfRange {
1436 what: field.to_owned(),
1437 value: value.to_string(),
1438 })
1439}
1440
1441fn u64_to_usize(value: u64, field: &str) -> Result<usize> {
1442 usize::try_from(value).map_err(|_| FrankenError::OutOfRange {
1443 what: field.to_owned(),
1444 value: value.to_string(),
1445 })
1446}
1447
1448fn checked_add(lhs: usize, rhs: usize, field: &str) -> Result<usize> {
1449 lhs.checked_add(rhs)
1450 .ok_or_else(|| FrankenError::DatabaseCorrupt {
1451 detail: format!("{field} overflow"),
1452 })
1453}
1454
1455#[cfg(test)]
1456mod tests {
1457 use std::fs;
1458 use std::path::Path;
1459
1460 use crate::commit_marker::MarkerSegmentHeader;
1461 use crate::symbol_log::{SymbolSegmentHeader, append_symbol_record, ensure_symbol_segment};
1462 use fsqlite_types::{ObjectId, Oti, SymbolRecord, SymbolRecordFlags, SymbolValidityWindow};
1463 use tempfile::TempDir;
1464
1465 use super::*;
1466
1467 const BEAD_ID: &str = "bd-3go.12";
1468
1469 #[test]
1472 fn test_epoch_id_monotone() {
1473 let clock = EpochClock::new(EpochId::ZERO);
1474 let mut prev = clock.current();
1475 for i in 0..100 {
1476 let next_result = clock.increment();
1477 assert!(
1478 next_result.is_ok(),
1479 "bead_id={BEAD_ID} case=epoch_monotone_increment_{i} err={next_result:?}"
1480 );
1481 let Ok(next) = next_result else {
1482 return;
1483 };
1484 assert!(
1485 next > prev,
1486 "bead_id={BEAD_ID} case=epoch_monotone prev={} next={}",
1487 prev.get(),
1488 next.get()
1489 );
1490 prev = next;
1491 }
1492 assert_eq!(
1493 clock.current().get(),
1494 100,
1495 "bead_id={BEAD_ID} case=epoch_monotone_final"
1496 );
1497 }
1498
1499 #[test]
1502 fn test_symbol_validity_window_rejects_future() {
1503 let current = EpochId::new(5);
1504 let window = SymbolValidityWindow::default_window(current);
1505 let future = EpochId::new(6);
1506 assert!(
1507 !window.contains(future),
1508 "bead_id={BEAD_ID} case=validity_window_rejects_future"
1509 );
1510 let result = validate_symbol_epoch(future, &window);
1511 assert!(
1512 result.is_err(),
1513 "bead_id={BEAD_ID} case=validity_window_future_epoch_error"
1514 );
1515 }
1516
1517 #[test]
1520 fn test_symbol_validity_window_accepts_past() {
1521 let current = EpochId::new(10);
1522 let window = SymbolValidityWindow::default_window(current);
1523 for past in [0, 1, 5, 9, 10] {
1524 let epoch = EpochId::new(past);
1525 assert!(
1526 window.contains(epoch),
1527 "bead_id={BEAD_ID} case=validity_window_accepts_past epoch={past}"
1528 );
1529 let result = validate_symbol_epoch(epoch, &window);
1530 assert!(
1531 result.is_ok(),
1532 "bead_id={BEAD_ID} case=validity_window_past_epoch_ok epoch={past}"
1533 );
1534 }
1535 }
1536
1537 #[test]
1540 fn test_epoch_scoped_key_derivation() {
1541 let master_key = [0xAB_u8; 32];
1542 let key_5 = derive_epoch_auth_key(&master_key, EpochId::new(5));
1543 let key_6 = derive_epoch_auth_key(&master_key, EpochId::new(6));
1544 assert_ne!(
1545 key_5, key_6,
1546 "bead_id={BEAD_ID} case=epoch_keys_differ_across_epochs"
1547 );
1548 let key_5_again = derive_epoch_auth_key(&master_key, EpochId::new(5));
1550 assert_eq!(
1551 key_5, key_5_again,
1552 "bead_id={BEAD_ID} case=epoch_key_deterministic"
1553 );
1554 }
1555
1556 #[test]
1559 fn test_epoch_key_derivation_domain_separation() {
1560 let dek = [0x42_u8; 32];
1561 let master_key = derive_master_key_from_dek(&dek);
1562 assert_ne!(
1564 master_key, dek,
1565 "bead_id={BEAD_ID} case=master_key_differs_from_dek"
1566 );
1567 let auth_key = derive_epoch_auth_key(&master_key, EpochId::ZERO);
1569 assert_ne!(
1570 auth_key.as_bytes(),
1571 &master_key,
1572 "bead_id={BEAD_ID} case=auth_key_differs_from_master"
1573 );
1574 }
1575
1576 #[test]
1579 fn test_epoch_transition_barrier_all_arrive() {
1580 let clock = EpochClock::new(EpochId::new(5));
1581 let barrier = EpochBarrier::new(EpochId::new(5), 4);
1582
1583 assert!(!barrier.arrive("WriteCoordinator"));
1584 assert!(!barrier.arrive("SymbolStore"));
1585 assert!(!barrier.arrive("Replicator"));
1586 assert!(barrier.arrive("CheckpointGc"));
1587
1588 assert!(
1589 barrier.is_complete(),
1590 "bead_id={BEAD_ID} case=barrier_complete"
1591 );
1592
1593 let outcome = barrier.resolve(&clock).expect("resolve must succeed");
1594 assert_eq!(
1595 outcome,
1596 BarrierOutcome::AllArrived {
1597 new_epoch: EpochId::new(6),
1598 },
1599 "bead_id={BEAD_ID} case=barrier_all_arrived_epoch_incremented"
1600 );
1601 assert_eq!(
1602 clock.current().get(),
1603 6,
1604 "bead_id={BEAD_ID} case=clock_advanced_after_barrier"
1605 );
1606 }
1607
1608 #[test]
1611 fn test_epoch_transition_barrier_timeout() {
1612 let clock = EpochClock::new(EpochId::new(5));
1613 let barrier = EpochBarrier::new(EpochId::new(5), 4);
1614
1615 barrier.arrive("WriteCoordinator");
1616 barrier.arrive("SymbolStore");
1617 barrier.arrive("Replicator");
1618 let outcome = barrier.resolve(&clock).expect("resolve must succeed");
1621 assert_eq!(
1622 outcome,
1623 BarrierOutcome::Timeout {
1624 arrived: 3,
1625 expected: 4,
1626 },
1627 "bead_id={BEAD_ID} case=barrier_timeout_epoch_unchanged"
1628 );
1629 assert_eq!(
1630 clock.current().get(),
1631 5,
1632 "bead_id={BEAD_ID} case=clock_unchanged_after_timeout"
1633 );
1634 }
1635
1636 #[test]
1639 fn test_epoch_bootstrap_from_ecs_root() {
1640 let root_epoch = EpochId::new(7);
1644 let window = SymbolValidityWindow::default_window(root_epoch);
1645
1646 assert!(
1648 !window.contains(EpochId::new(8)),
1649 "bead_id={BEAD_ID} case=bootstrap_rejects_future"
1650 );
1651 assert!(
1653 window.contains(EpochId::new(7)),
1654 "bead_id={BEAD_ID} case=bootstrap_accepts_current"
1655 );
1656 assert!(
1658 window.contains(EpochId::ZERO),
1659 "bead_id={BEAD_ID} case=bootstrap_accepts_zero"
1660 );
1661 }
1662
1663 #[test]
1666 fn test_barrier_cancelled() {
1667 let clock = EpochClock::new(EpochId::new(3));
1668 let barrier = EpochBarrier::new(EpochId::new(3), 2);
1669
1670 barrier.arrive("WriteCoordinator");
1671 barrier.cancel();
1672
1673 let outcome = barrier.resolve(&clock).expect("resolve must succeed");
1674 assert_eq!(
1675 outcome,
1676 BarrierOutcome::Cancelled,
1677 "bead_id={BEAD_ID} case=barrier_cancelled_epoch_unchanged"
1678 );
1679 assert_eq!(
1680 clock.current().get(),
1681 3,
1682 "bead_id={BEAD_ID} case=clock_unchanged_after_cancel"
1683 );
1684 }
1685
1686 #[test]
1689 fn test_epoch_clock_store_and_recover() {
1690 let clock = EpochClock::new(EpochId::ZERO);
1691 clock.store(EpochId::new(42));
1692 assert_eq!(
1693 clock.current().get(),
1694 42,
1695 "bead_id={BEAD_ID} case=clock_store_recovery"
1696 );
1697 let next = clock.increment().expect("increment after store");
1698 assert_eq!(
1699 next.get(),
1700 43,
1701 "bead_id={BEAD_ID} case=clock_increment_after_store"
1702 );
1703 }
1704
1705 #[test]
1708 fn test_validity_window_boundary() {
1709 let window = SymbolValidityWindow::new(EpochId::new(3), EpochId::new(7));
1710 assert!(
1711 !window.contains(EpochId::new(2)),
1712 "bead_id={BEAD_ID} case=window_below_lower_bound"
1713 );
1714 assert!(
1715 window.contains(EpochId::new(3)),
1716 "bead_id={BEAD_ID} case=window_at_lower_bound"
1717 );
1718 assert!(
1719 window.contains(EpochId::new(5)),
1720 "bead_id={BEAD_ID} case=window_within_bounds"
1721 );
1722 assert!(
1723 window.contains(EpochId::new(7)),
1724 "bead_id={BEAD_ID} case=window_at_upper_bound"
1725 );
1726 assert!(
1727 !window.contains(EpochId::new(8)),
1728 "bead_id={BEAD_ID} case=window_above_upper_bound"
1729 );
1730 }
1731
1732 const ROOT_BEAD_ID: &str = "bd-1hi.25";
1733
1734 fn make_object_id(seed: u8) -> ObjectId {
1735 ObjectId::from_bytes([seed; 16])
1736 }
1737
1738 fn test_master_key() -> [u8; 32] {
1739 [0xA5; 32]
1740 }
1741
1742 fn create_layout() -> (TempDir, NativeBootstrapLayout) {
1743 let temp_dir = TempDir::new().expect("tempdir");
1744 let layout = NativeBootstrapLayout::new(temp_dir.path().join("ecs"));
1745 std::fs::create_dir_all(layout.symbols_dir()).expect("create symbols dir");
1746 std::fs::create_dir_all(layout.markers_dir()).expect("create markers dir");
1747 (temp_dir, layout)
1748 }
1749
1750 fn write_single_symbol_object(
1751 symbols_dir: &Path,
1752 segment_id: u64,
1753 epoch_id: EpochId,
1754 object_id: ObjectId,
1755 payload: &[u8],
1756 ) {
1757 let header =
1758 SymbolSegmentHeader::new(segment_id, epoch_id.get(), 1_700_000_000 + segment_id);
1759 let segment_path = symbols_dir.join(format!("segment-{segment_id:06}.log"));
1760 ensure_symbol_segment(&segment_path, header).expect("ensure symbol segment");
1761 let symbol_size = u32::try_from(payload.len()).expect("payload fits u32");
1762 let oti = Oti {
1763 f: u64::from(symbol_size),
1764 al: 1,
1765 t: symbol_size,
1766 z: 1,
1767 n: 1,
1768 };
1769 let record = SymbolRecord::new(
1770 object_id,
1771 oti,
1772 0,
1773 payload.to_vec(),
1774 SymbolRecordFlags::SYSTEMATIC_RUN_START,
1775 );
1776 append_symbol_record(symbols_dir, header, &record).expect("append symbol");
1777 }
1778
1779 fn write_marker_segment(
1780 markers_dir: &Path,
1781 start_commit_seq: u64,
1782 records: &[CommitMarkerRecord],
1783 ) {
1784 let segment_id = segment_id_for_commit_seq(start_commit_seq);
1785 let header = MarkerSegmentHeader::new(segment_id, start_commit_seq);
1786 let mut bytes = Vec::from(header.encode());
1787 for record in records {
1788 bytes.extend_from_slice(&record.encode());
1789 }
1790 let segment_path = markers_dir.join(format!("segment-{segment_id:06}.log"));
1791 std::fs::write(segment_path, bytes).expect("write marker segment");
1792 }
1793
1794 fn make_marker(commit_seq: u64, prev: [u8; 16], salt: u8) -> CommitMarkerRecord {
1795 CommitMarkerRecord::new(
1796 commit_seq,
1797 1_800_000_000_000_000_000 + commit_seq,
1798 [salt; 16],
1799 [salt.wrapping_add(1); 16],
1800 prev,
1801 )
1802 }
1803
1804 fn make_manifest(
1805 database_name: &str,
1806 current_commit: ObjectId,
1807 commit_seq: u64,
1808 schema_snapshot: ObjectId,
1809 schema_epoch: u64,
1810 ecs_epoch: EpochId,
1811 checkpoint_base: ObjectId,
1812 ) -> RootManifest {
1813 RootManifest {
1814 database_name: database_name.to_owned(),
1815 current_commit,
1816 commit_seq,
1817 schema_snapshot,
1818 schema_epoch,
1819 ecs_epoch,
1820 checkpoint_base,
1821 gc_horizon: commit_seq,
1822 created_at: 1_800_000_000,
1823 updated_at: 1_800_000_123,
1824 }
1825 }
1826
1827 fn must_err_contains<T: std::fmt::Debug>(result: Result<T>, needle: &str, case: &str) {
1828 let err = result.expect_err(case);
1829 let detail = err.to_string();
1830 assert!(
1831 detail.contains(needle),
1832 "bead_id={ROOT_BEAD_ID} case={case} expected_substring={needle} actual={detail}"
1833 );
1834 }
1835
1836 fn write_bootstrap_objects(
1837 layout: &NativeBootstrapLayout,
1838 root_epoch: EpochId,
1839 manifest_id: ObjectId,
1840 manifest: &RootManifest,
1841 schema_payload: &[u8],
1842 checkpoint_payload: &[u8],
1843 markers: &[CommitMarkerRecord],
1844 ) {
1845 let manifest_bytes = manifest.encode().expect("encode manifest");
1846 write_single_symbol_object(
1847 layout.symbols_dir().as_path(),
1848 1,
1849 root_epoch,
1850 manifest_id,
1851 &manifest_bytes,
1852 );
1853 write_single_symbol_object(
1854 layout.symbols_dir().as_path(),
1855 1,
1856 root_epoch,
1857 manifest.schema_snapshot,
1858 schema_payload,
1859 );
1860 write_single_symbol_object(
1861 layout.symbols_dir().as_path(),
1862 1,
1863 root_epoch,
1864 manifest.checkpoint_base,
1865 checkpoint_payload,
1866 );
1867 if let Some(first) = markers.first() {
1868 write_marker_segment(layout.markers_dir().as_path(), first.commit_seq, markers);
1869 }
1870 }
1871
1872 fn write_valid_bootstrap_fixture(
1873 layout: &NativeBootstrapLayout,
1874 root_epoch: EpochId,
1875 ) -> (ObjectId, RootManifest, CommitMarkerRecord, Vec<u8>, Vec<u8>) {
1876 let manifest_id = make_object_id(0x70);
1877 let schema_payload = b"schema-cache-v1".to_vec();
1878 let checkpoint_payload = b"checkpoint-cache-v1".to_vec();
1879 let marker = make_marker(0, [0_u8; 16], 0x71);
1880 let manifest = make_manifest(
1881 "db-valid",
1882 ObjectId::from_bytes(marker.marker_id),
1883 marker.commit_seq,
1884 make_object_id(0x72),
1885 1,
1886 root_epoch,
1887 make_object_id(0x73),
1888 );
1889 write_bootstrap_objects(
1890 layout,
1891 root_epoch,
1892 manifest_id,
1893 &manifest,
1894 &schema_payload,
1895 &checkpoint_payload,
1896 std::slice::from_ref(&marker),
1897 );
1898 (
1899 manifest_id,
1900 manifest,
1901 marker,
1902 schema_payload,
1903 checkpoint_payload,
1904 )
1905 }
1906
1907 #[test]
1908 fn test_ecs_root_pointer_encode_decode() {
1909 let pointer = EcsRootPointer::unauthed(make_object_id(0x11), EpochId::new(7));
1910 let encoded = pointer.encode();
1911 assert_eq!(encoded.len(), ECS_ROOT_POINTER_BYTES);
1912 let decoded = EcsRootPointer::decode(&encoded, false, None).expect("decode root pointer");
1913 assert_eq!(
1914 decoded, pointer,
1915 "bead_id={ROOT_BEAD_ID} case=root_roundtrip"
1916 );
1917 }
1918
1919 #[test]
1920 fn test_ecs_root_pointer_magic() {
1921 let pointer = EcsRootPointer::unauthed(make_object_id(0x22), EpochId::new(3));
1922 let mut encoded = pointer.encode();
1923 encoded[0] = b'X';
1924 let result = EcsRootPointer::decode(&encoded, false, None);
1925 assert!(
1926 result.is_err(),
1927 "bead_id={ROOT_BEAD_ID} case=root_bad_magic"
1928 );
1929 }
1930
1931 #[test]
1932 fn test_ecs_root_pointer_checksum_tamper() {
1933 let pointer = EcsRootPointer::unauthed(make_object_id(0x33), EpochId::new(9));
1934 let mut encoded = pointer.encode();
1935 encoded[9] ^= 0xFF;
1936 let result = EcsRootPointer::decode(&encoded, false, None);
1937 assert!(
1938 result.is_err(),
1939 "bead_id={ROOT_BEAD_ID} case=root_checksum_tamper"
1940 );
1941 }
1942
1943 #[test]
1944 fn test_root_auth_tag_verification() {
1945 let key = test_master_key();
1946 let pointer = EcsRootPointer::authed(make_object_id(0x44), EpochId::new(12), &key);
1947 let encoded = pointer.encode();
1948 let decoded =
1949 EcsRootPointer::decode(&encoded, true, Some(&key)).expect("auth decode succeeds");
1950 assert_eq!(decoded, pointer);
1951
1952 let mut tampered = encoded;
1953 tampered[40] ^= 0x01;
1954 let result = EcsRootPointer::decode(&tampered, true, Some(&key));
1955 assert!(
1956 result.is_err(),
1957 "bead_id={ROOT_BEAD_ID} case=root_auth_tamper"
1958 );
1959 }
1960
1961 #[test]
1962 fn test_root_auth_tag_zero_when_off() {
1963 let pointer = build_root_pointer(make_object_id(0x55), EpochId::new(2), false, None)
1964 .expect("build unauthed root");
1965 assert_eq!(pointer.root_auth_tag, [0_u8; 16]);
1966 let decoded =
1967 EcsRootPointer::decode(&pointer.encode(), false, None).expect("decode off mode");
1968 assert_eq!(decoded.root_auth_tag, [0_u8; 16]);
1969 }
1970
1971 #[test]
1972 fn test_root_manifest_encode_decode() {
1973 let manifest = make_manifest(
1974 "db-main",
1975 make_object_id(0x10),
1976 5,
1977 make_object_id(0x20),
1978 3,
1979 EpochId::new(7),
1980 make_object_id(0x30),
1981 );
1982 let encoded = manifest.encode().expect("encode manifest");
1983 let decoded = RootManifest::decode(&encoded).expect("decode manifest");
1984 assert_eq!(decoded, manifest);
1985 }
1986
1987 #[test]
1988 fn test_root_manifest_magic() {
1989 let manifest = make_manifest(
1990 "db-main",
1991 make_object_id(0x10),
1992 5,
1993 make_object_id(0x20),
1994 3,
1995 EpochId::new(7),
1996 make_object_id(0x30),
1997 );
1998 let mut encoded = manifest.encode().expect("encode manifest");
1999 encoded[0] = b'X';
2000 assert!(
2001 RootManifest::decode(&encoded).is_err(),
2002 "bead_id={ROOT_BEAD_ID} case=manifest_bad_magic"
2003 );
2004 }
2005
2006 #[test]
2007 fn test_bootstrap_step_4_epoch_guard() {
2008 let (_tmp, layout) = create_layout();
2009 let manifest_id = make_object_id(0x66);
2010 let schema_id = make_object_id(0x67);
2011 let checkpoint_id = make_object_id(0x68);
2012 let marker = make_marker(0, [0_u8; 16], 0x60);
2013 let manifest = make_manifest(
2014 "future-segment",
2015 ObjectId::from_bytes(marker.marker_id),
2016 marker.commit_seq,
2017 schema_id,
2018 1,
2019 EpochId::new(3),
2020 checkpoint_id,
2021 );
2022 let manifest_bytes = manifest.encode().expect("encode manifest");
2023 write_single_symbol_object(
2024 layout.symbols_dir().as_path(),
2025 1,
2026 EpochId::new(4),
2027 manifest_id,
2028 &manifest_bytes,
2029 );
2030 let pointer =
2031 build_root_pointer(manifest_id, EpochId::new(3), false, None).expect("build root");
2032 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2033
2034 let result = bootstrap_native_mode(&layout, false, None);
2035 assert!(
2036 result.is_err(),
2037 "bead_id={ROOT_BEAD_ID} case=future_epoch_guard"
2038 );
2039 }
2040
2041 #[test]
2042 fn test_bootstrap_step_5_epoch_invariant() {
2043 let (_tmp, layout) = create_layout();
2044 let root_epoch = EpochId::new(7);
2045 let manifest_id = make_object_id(0x74);
2046 let marker = make_marker(0, [0_u8; 16], 0x75);
2047 let manifest = make_manifest(
2048 "epoch-mismatch",
2049 ObjectId::from_bytes(marker.marker_id),
2050 marker.commit_seq,
2051 make_object_id(0x76),
2052 1,
2053 EpochId::new(8),
2054 make_object_id(0x77),
2055 );
2056 write_bootstrap_objects(
2057 &layout,
2058 root_epoch,
2059 manifest_id,
2060 &manifest,
2061 b"schema",
2062 b"checkpoint",
2063 &[marker],
2064 );
2065 let pointer = build_root_pointer(manifest_id, root_epoch, false, None).expect("root");
2066 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2067 must_err_contains(
2068 bootstrap_native_mode(&layout, false, None),
2069 "epoch_mismatch",
2070 "bootstrap_step_5_epoch_invariant",
2071 );
2072 }
2073
2074 #[test]
2075 fn test_bootstrap_step_6_marker_verification() {
2076 let (_tmp, layout) = create_layout();
2077 let root_epoch = EpochId::new(9);
2078 let manifest_id = make_object_id(0x78);
2079 let marker = make_marker(0, [0_u8; 16], 0x79);
2080 let manifest = make_manifest(
2081 "marker-mismatch",
2082 make_object_id(0x7A),
2083 marker.commit_seq,
2084 make_object_id(0x7B),
2085 1,
2086 root_epoch,
2087 make_object_id(0x7C),
2088 );
2089 write_bootstrap_objects(
2090 &layout,
2091 root_epoch,
2092 manifest_id,
2093 &manifest,
2094 b"schema",
2095 b"checkpoint",
2096 &[marker],
2097 );
2098 let pointer = build_root_pointer(manifest_id, root_epoch, false, None).expect("root");
2099 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2100 must_err_contains(
2101 bootstrap_native_mode(&layout, false, None),
2102 "marker_mismatch",
2103 "bootstrap_step_6_marker_verification",
2104 );
2105 }
2106
2107 #[test]
2108 fn test_bootstrap_full_sequence() {
2109 let (_tmp, layout) = create_layout();
2110 let root_epoch = EpochId::new(11);
2111 let (manifest_id, manifest, marker, schema_payload, checkpoint_payload) =
2112 write_valid_bootstrap_fixture(&layout, root_epoch);
2113 let pointer = build_root_pointer(manifest_id, root_epoch, false, None).expect("root");
2114 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2115
2116 let state = bootstrap_native_mode(&layout, false, None).expect("bootstrap ok");
2117 assert_eq!(state.root_pointer, pointer);
2118 assert_eq!(state.manifest, manifest);
2119 assert_eq!(state.latest_marker, marker);
2120 assert_eq!(state.schema_snapshot_bytes, schema_payload);
2121 assert_eq!(state.checkpoint_base_bytes, checkpoint_payload);
2122 }
2123
2124 #[test]
2125 fn test_bootstrap_corrupted_root_recovery() {
2126 let (_tmp, layout) = create_layout();
2127 let root_epoch = EpochId::new(12);
2128 let (_manifest_id, manifest, marker, schema_payload, checkpoint_payload) =
2129 write_valid_bootstrap_fixture(&layout, root_epoch);
2130 fs::write(layout.root_path(), [0xFF_u8; 7]).expect("write corrupt root");
2131
2132 let recovered = bootstrap_native_mode_with_recovery(&layout, false, None).expect("recover");
2133 assert_eq!(recovered.manifest, manifest);
2134 assert_eq!(recovered.latest_marker, marker);
2135 assert_eq!(recovered.schema_snapshot_bytes, schema_payload);
2136 assert_eq!(recovered.checkpoint_base_bytes, checkpoint_payload);
2137 let persisted =
2138 read_root_pointer(&layout.root_path(), false, None).expect("read recovered");
2139 assert_eq!(
2140 persisted.manifest_object_id,
2141 recovered.root_pointer.manifest_object_id
2142 );
2143 assert_eq!(persisted.ecs_epoch, recovered.root_pointer.ecs_epoch);
2144 }
2145
2146 #[test]
2147 fn test_crash_safe_root_update() {
2148 let (_tmp, layout) = create_layout();
2149 let pointer_a = EcsRootPointer::unauthed(make_object_id(0x80), EpochId::new(1));
2150 let pointer_b = EcsRootPointer::unauthed(make_object_id(0x81), EpochId::new(2));
2151 write_root_pointer_atomic(&layout.root_path(), pointer_a).expect("write A");
2152 write_root_pointer_atomic(&layout.root_path(), pointer_b).expect("write B");
2153 let decoded = read_root_pointer(&layout.root_path(), false, None).expect("decode");
2154 assert_eq!(
2155 decoded, pointer_b,
2156 "bead_id={ROOT_BEAD_ID} case=root_atomic_swap"
2157 );
2158 let entries = fs::read_dir(layout.ecs_dir.as_path()).expect("list ecs dir");
2159 for entry in entries {
2160 let entry = entry.expect("entry");
2161 let name = entry.file_name();
2162 let name = name.to_string_lossy();
2163 assert!(
2164 !name.starts_with(".root.tmp."),
2165 "bead_id={ROOT_BEAD_ID} case=temp_root_file_leaked file={name}"
2166 );
2167 }
2168 }
2169
2170 #[test]
2171 fn prop_root_pointer_roundtrip() {
2172 let key = test_master_key();
2173 for seed in [0_u8, 1, 17, 99, 255] {
2174 for epoch in [0_u64, 1, 2, 17, 255, 4_096, 1 << 20] {
2175 let id = make_object_id(seed);
2176 let plain = EcsRootPointer::unauthed(id, EpochId::new(epoch));
2177 let plain_roundtrip =
2178 EcsRootPointer::decode(&plain.encode(), false, None).expect("plain decode");
2179 assert_eq!(plain_roundtrip, plain);
2180 let authed = EcsRootPointer::authed(id, EpochId::new(epoch), &key);
2181 let authed_roundtrip = EcsRootPointer::decode(&authed.encode(), true, Some(&key))
2182 .expect("auth decode");
2183 assert_eq!(authed_roundtrip, authed);
2184 }
2185 }
2186 }
2187
2188 #[test]
2189 fn test_bootstrap_rejects_marker_chain_gap() {
2190 let (_tmp, layout) = create_layout();
2191 let root_epoch = EpochId::new(13);
2192 let manifest_id = make_object_id(0x82);
2193 let m0 = make_marker(0, [0_u8; 16], 0x83);
2194 let m1 = make_marker(1, m0.marker_id, 0x84);
2195 let m2 = make_marker(2, [0xEE_u8; 16], 0x85);
2196 let manifest = make_manifest(
2197 "marker-gap",
2198 ObjectId::from_bytes(m2.marker_id),
2199 m2.commit_seq,
2200 make_object_id(0x86),
2201 1,
2202 root_epoch,
2203 make_object_id(0x87),
2204 );
2205 write_bootstrap_objects(
2206 &layout,
2207 root_epoch,
2208 manifest_id,
2209 &manifest,
2210 b"schema-gap",
2211 b"checkpoint-gap",
2212 &[m0, m1, m2],
2213 );
2214 let pointer = build_root_pointer(manifest_id, root_epoch, false, None).expect("root");
2215 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2216 must_err_contains(
2217 bootstrap_native_mode(&layout, false, None),
2218 "marker_chain_gap",
2219 "bootstrap_rejects_marker_chain_gap",
2220 );
2221 }
2222
2223 #[test]
2224 fn test_bootstrap_schema_snapshot_loads() {
2225 let (_tmp, layout) = create_layout();
2226 let root_epoch = EpochId::new(14);
2227 let (manifest_id, _manifest, _marker, schema_payload, _checkpoint_payload) =
2228 write_valid_bootstrap_fixture(&layout, root_epoch);
2229 let pointer = build_root_pointer(manifest_id, root_epoch, false, None).expect("root");
2230 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2231 let state = bootstrap_native_mode(&layout, false, None).expect("bootstrap");
2232 assert_eq!(state.schema_snapshot_bytes, schema_payload);
2233 }
2234
2235 #[test]
2236 fn test_bootstrap_checkpoint_base_warms_cache() {
2237 let (_tmp, layout) = create_layout();
2238 let root_epoch = EpochId::new(15);
2239 let (manifest_id, _manifest, _marker, _schema_payload, checkpoint_payload) =
2240 write_valid_bootstrap_fixture(&layout, root_epoch);
2241 let pointer = build_root_pointer(manifest_id, root_epoch, false, None).expect("root");
2242 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2243 let state = bootstrap_native_mode(&layout, false, None).expect("bootstrap");
2244 assert_eq!(state.checkpoint_base_bytes, checkpoint_payload);
2245 }
2246
2247 #[test]
2248 fn test_bootstrap_happy_path_from_root() {
2249 test_bootstrap_full_sequence();
2250 }
2251
2252 #[test]
2253 fn test_bootstrap_corrupt_root_pointer_recovers_by_scan() {
2254 test_bootstrap_corrupted_root_recovery();
2255 }
2256
2257 #[test]
2258 fn test_bootstrap_root_auth_mismatch_fails() {
2259 let (_tmp, layout) = create_layout();
2260 let root_epoch = EpochId::new(16);
2261 let (manifest_id, _manifest, _marker, _schema_payload, _checkpoint_payload) =
2262 write_valid_bootstrap_fixture(&layout, root_epoch);
2263 let good_key = test_master_key();
2264 let bad_key = [0x5A_u8; 32];
2265 let pointer =
2266 build_root_pointer(manifest_id, root_epoch, true, Some(&good_key)).expect("root");
2267 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2268 must_err_contains(
2269 bootstrap_native_mode(&layout, true, Some(&bad_key)),
2270 "auth_failed",
2271 "bootstrap_root_auth_mismatch_fails",
2272 );
2273 }
2274
2275 #[test]
2276 fn test_bootstrap_root_pointer_corrupt_checksum_fails_then_scan() {
2277 let (_tmp, layout) = create_layout();
2278 let root_epoch = EpochId::new(17);
2279 let (manifest_id, _manifest, _marker, _schema_payload, _checkpoint_payload) =
2280 write_valid_bootstrap_fixture(&layout, root_epoch);
2281 let pointer = build_root_pointer(manifest_id, root_epoch, false, None).expect("root");
2282 let mut root_bytes = pointer.encode();
2283 root_bytes[10] ^= 0xAA;
2284 fs::write(layout.root_path(), root_bytes).expect("write corrupt root");
2285 let state = bootstrap_native_mode_with_recovery(&layout, false, None).expect("recovered");
2286 assert_eq!(state.root_pointer.manifest_object_id, manifest_id);
2287 assert_eq!(state.root_pointer.ecs_epoch, root_epoch);
2288 }
2289
2290 #[test]
2291 fn test_e2e_native_mode_open_close_reopen() {
2292 let (_tmp, layout) = create_layout();
2293 let root_epoch = EpochId::new(18);
2294 let (manifest_id, manifest, marker, _schema_payload, _checkpoint_payload) =
2295 write_valid_bootstrap_fixture(&layout, root_epoch);
2296 let pointer = build_root_pointer(manifest_id, root_epoch, false, None).expect("root");
2297 write_root_pointer_atomic(&layout.root_path(), pointer).expect("write root");
2298
2299 let first_open = bootstrap_native_mode(&layout, false, None).expect("first open");
2300 let second_open = bootstrap_native_mode(&layout, false, None).expect("second open");
2301 assert_eq!(first_open.manifest, manifest);
2302 assert_eq!(first_open.latest_marker, marker);
2303 assert_eq!(
2304 second_open.manifest.current_commit,
2305 first_open.manifest.current_commit
2306 );
2307 assert_eq!(second_open.root_pointer, first_open.root_pointer);
2308
2309 fs::write(layout.root_path(), [0_u8; 9]).expect("corrupt root");
2310 let recovered = bootstrap_native_mode_with_recovery(&layout, false, None).expect("reopen");
2311 assert_eq!(
2312 recovered.manifest.current_commit,
2313 first_open.manifest.current_commit
2314 );
2315 assert_eq!(
2316 recovered.manifest.commit_seq,
2317 first_open.manifest.commit_seq
2318 );
2319 }
2320
2321 #[test]
2322 fn test_e2e_bootstrap_cold_start() {
2323 test_bootstrap_full_sequence();
2324 }
2325
2326 #[test]
2327 fn test_e2e_bootstrap_after_crash() {
2328 test_bootstrap_root_pointer_corrupt_checksum_fails_then_scan();
2329 }
2330
2331 #[test]
2332 fn test_e2e_bootstrap_schema_migration() {
2333 let (_tmp, layout) = create_layout();
2334 let root_epoch = EpochId::new(19);
2335 let manifest_id_v1 = make_object_id(0x88);
2336 let marker0 = make_marker(0, [0_u8; 16], 0x89);
2337 let marker1 = make_marker(1, marker0.marker_id, 0x8A);
2338 write_marker_segment(
2339 layout.markers_dir().as_path(),
2340 0,
2341 &[marker0, marker1.clone()],
2342 );
2343
2344 let schema_v1 = make_object_id(0x8B);
2345 let schema_v2 = make_object_id(0x8C);
2346 let checkpoint_id = make_object_id(0x8D);
2347 let manifest_v1 = make_manifest(
2348 "schema-v1",
2349 ObjectId::from_bytes(marker1.marker_id),
2350 1,
2351 schema_v1,
2352 1,
2353 root_epoch,
2354 checkpoint_id,
2355 );
2356 let manifest_v2 = make_manifest(
2357 "schema-v2",
2358 ObjectId::from_bytes(marker1.marker_id),
2359 1,
2360 schema_v2,
2361 2,
2362 root_epoch,
2363 checkpoint_id,
2364 );
2365 let manifest_id_v2 = make_object_id(0x8E);
2366 write_single_symbol_object(
2367 layout.symbols_dir().as_path(),
2368 1,
2369 root_epoch,
2370 manifest_id_v1,
2371 &manifest_v1.encode().expect("manifest v1"),
2372 );
2373 write_single_symbol_object(
2374 layout.symbols_dir().as_path(),
2375 1,
2376 root_epoch,
2377 manifest_id_v2,
2378 &manifest_v2.encode().expect("manifest v2"),
2379 );
2380 write_single_symbol_object(
2381 layout.symbols_dir().as_path(),
2382 1,
2383 root_epoch,
2384 schema_v1,
2385 b"schema-v1",
2386 );
2387 write_single_symbol_object(
2388 layout.symbols_dir().as_path(),
2389 1,
2390 root_epoch,
2391 schema_v2,
2392 b"schema-v2",
2393 );
2394 write_single_symbol_object(
2395 layout.symbols_dir().as_path(),
2396 1,
2397 root_epoch,
2398 checkpoint_id,
2399 b"checkpoint",
2400 );
2401
2402 let pointer_v2 = build_root_pointer(manifest_id_v2, root_epoch, false, None).expect("root");
2403 write_root_pointer_atomic(&layout.root_path(), pointer_v2).expect("write root");
2404 let state = bootstrap_native_mode(&layout, false, None).expect("bootstrap");
2405 assert_eq!(state.manifest.schema_epoch, 2);
2406 assert_eq!(state.schema_snapshot_bytes, b"schema-v2".to_vec());
2407 }
2408
2409 #[test]
2410 fn test_ecs_root_pointer_checksum_roundtrip() {
2411 test_ecs_root_pointer_encode_decode();
2412 }
2413
2414 #[test]
2415 fn test_ecs_root_pointer_auth_tag_verifies() {
2416 test_root_auth_tag_verification();
2417 }
2418
2419 #[test]
2420 fn test_bootstrap_future_epoch_guard() {
2421 test_bootstrap_step_4_epoch_guard();
2422 }
2423
2424 #[test]
2425 fn test_root_manifest_epoch_must_match_root_pointer() {
2426 test_bootstrap_step_5_epoch_invariant();
2427 }
2428
2429 #[test]
2430 fn test_bootstrap_commit_marker_matches_current_commit() {
2431 test_bootstrap_step_6_marker_verification();
2432 }
2433
2434 #[test]
2435 fn test_bd_1hi_25_unit_compliance_gate() {
2436 assert_eq!(ROOT_BOOTSTRAP_BEAD_ID, "bd-1hi.25");
2437 assert_eq!(ROOT_BOOTSTRAP_LOGGING_STANDARD, "bd-1fpm");
2438 assert_eq!(ECS_ROOT_POINTER_MAGIC, *b"FSRT");
2439 assert_eq!(ROOT_MANIFEST_MAGIC, *b"FSQLROOT");
2440 }
2441
2442 #[test]
2443 fn prop_bd_1hi_25_structure_compliance() {
2444 for name in ["segment-000001.log", "segment-999999.log"] {
2445 assert!(
2446 parse_segment_id(name).is_some(),
2447 "bead_id={ROOT_BEAD_ID} case=parse_segment_id_valid name={name}"
2448 );
2449 }
2450 for name in ["segment.log", "segment-aa.log", "other-000001.log"] {
2451 assert!(
2452 parse_segment_id(name).is_none(),
2453 "bead_id={ROOT_BEAD_ID} case=parse_segment_id_invalid name={name}"
2454 );
2455 }
2456 }
2457
2458 #[test]
2459 fn test_e2e_bd_1hi_25_compliance() {
2460 test_e2e_native_mode_open_close_reopen();
2461 }
2462}