1use std::collections::BTreeSet;
8use std::fs::File;
9use std::path::{Path, PathBuf};
10
11use fs4::fs_std::FileExt;
12use graphforge_core::{GfError, ProjectErrorCode};
13use uuid::Uuid;
14
15use crate::project_checkpoints::checkpoint_retention_roots_after_writer_lock;
16use crate::project_failpoint;
17use crate::project_generation::{
18 resolve_project_generation, validated_generation_manifest_sha256, validated_generation_parent,
19};
20use crate::project_publication::{
21 ATTEMPTS_DIR, GENERATIONS_DIR, JournalPhase, LOCKS_DIR, TRANSACTIONS_DIR, WRITER_LOCK_FILE,
22 ensure_machine_directory, open_regular_lock, open_transaction_lock, read_journal,
23 sync_directory, write_journal,
24};
25
26const TRASH_DIR: &str = "trash";
27const MAX_RECOVERY_ENTRIES: usize = 10_000;
28const RETAINED_ANCESTORS: usize = 2;
29
30#[derive(Debug, Clone, PartialEq, Eq)]
32pub struct ProjectRecoveryReport {
33 pub selected_generation_uuid: Uuid,
35 pub repaired_journals: u64,
37 pub aborted_journals: u64,
39 pub removed_generations: u64,
41 pub preserved_unknown_entries: u64,
43}
44
45pub fn recover_project_transactions(
55 container_root: impl AsRef<Path>,
56) -> Result<ProjectRecoveryReport, GfError> {
57 let selected =
58 resolve_project_generation(container_root.as_ref()).map_err(map_recovery_resolution)?;
59 let root = selected.container_root().to_owned();
60 let writer_lock = acquire_recovery_lock(&root)?;
61
62 let selected = resolve_project_generation(&root).map_err(map_recovery_resolution)?;
65 let selected_uuid = selected.generation_uuid();
66 let checkpoint_roots = checkpoint_retention_roots_after_writer_lock(&root)?;
67 let retained = retained_generations(&root, &selected, &checkpoint_roots.roots)?;
68 let mut report = ProjectRecoveryReport {
69 selected_generation_uuid: selected_uuid,
70 repaired_journals: 0,
71 aborted_journals: 0,
72 removed_generations: 0,
73 preserved_unknown_entries: 0,
74 };
75
76 let transactions_root = root.join(TRANSACTIONS_DIR);
77 if transactions_root.exists() {
78 recover_journals(&root, &transactions_root, &selected, &retained, &mut report)?;
79 }
80 report.preserved_unknown_entries += count_unknown_generation_entries(&root, &retained)?;
81 drop(writer_lock);
82 Ok(report)
83}
84
85fn acquire_recovery_lock(root: &Path) -> Result<File, GfError> {
86 let lock_dir = ensure_machine_directory(root, Path::new(LOCKS_DIR))?;
87 sync_directory(root)?;
88 let writer_lock = open_regular_lock(&lock_dir.join(WRITER_LOCK_FILE))?;
89 if !FileExt::try_lock_exclusive(&writer_lock).map_err(storage_io)? {
90 return Err(project_error(
91 ProjectErrorCode::WriterBusy,
92 "phase=RECOVERY committed=false cause=live_writer_owns_kernel_lock",
93 ));
94 }
95 Ok(writer_lock)
96}
97
98fn recover_journals(
99 root: &Path,
100 transactions_root: &Path,
101 selected: &crate::ResolvedProjectGeneration,
102 retained: &BTreeSet<Uuid>,
103 report: &mut ProjectRecoveryReport,
104) -> Result<(), GfError> {
105 let mut journal_paths = bounded_directory_entries(transactions_root)?;
106 journal_paths.sort();
107 for journal_path in journal_paths {
108 if crate::project_publication::cleanup_atomicwrite_temp(&journal_path)? {
109 continue;
110 }
111 let Some(transaction_uuid) = journal_file_uuid(&journal_path) else {
112 return Err(recovery_corrupt(
113 "transaction directory contains a noncanonical journal entry",
114 ));
115 };
116 let transaction_lock = open_transaction_lock(root, transaction_uuid)?;
117 if !FileExt::try_lock_exclusive(&transaction_lock).map_err(storage_io)? {
118 continue;
122 }
123 let mut journal = read_journal(&journal_path)
124 .map_err(|_| recovery_corrupt("transaction journal is torn, invalid, or ambiguous"))?;
125 if parse_canonical_uuid(&journal.transaction_uuid) != Some(transaction_uuid) {
126 return Err(recovery_corrupt(
127 "transaction journal identity does not match its file name",
128 ));
129 }
130 let generation_uuid = parse_canonical_uuid(&journal.generation_uuid)
131 .ok_or_else(|| recovery_corrupt("transaction generation UUID is invalid"))?;
132
133 if generation_uuid == selected.generation_uuid() {
134 repair_reachable_journal(
135 &journal_path,
136 &mut journal,
137 selected.manifest_sha256(),
138 report,
139 )?;
140 } else if retained.contains(&generation_uuid) {
141 repair_reachable_journal(
142 &journal_path,
143 &mut journal,
144 validated_generation_manifest_sha256(root, generation_uuid)?,
145 report,
146 )?;
147 } else {
148 if journal.phase != JournalPhase::Published && journal.phase != JournalPhase::Aborted {
149 journal.phase = JournalPhase::Aborted;
150 write_journal(&journal_path, &journal)?;
151 report.aborted_journals += 1;
152 }
153 report.removed_generations += cleanup_abandoned_generation(
154 root,
155 transaction_uuid,
156 generation_uuid,
157 &journal.request_fingerprint,
158 retained,
159 )?;
160 }
161 }
162 Ok(())
163}
164
165fn repair_reachable_journal(
166 path: &Path,
167 journal: &mut crate::project_publication::JournalRecord,
168 manifest_sha256: [u8; 32],
169 report: &mut ProjectRecoveryReport,
170) -> Result<(), GfError> {
171 let expected_digest = digest_hex(manifest_sha256);
172 if journal.generation_manifest_sha256.as_deref() != Some(expected_digest.as_str()) {
173 return Err(recovery_corrupt(
174 "journal for committed generation does not match CURRENT manifest digest",
175 ));
176 }
177 if journal.phase != JournalPhase::Published {
178 journal.phase = JournalPhase::Published;
179 write_journal(path, journal)?;
180 report.repaired_journals += 1;
181 }
182 Ok(())
183}
184
185fn retained_generations(
186 root: &Path,
187 selected: &crate::ResolvedProjectGeneration,
188 checkpoint_roots: &[(Uuid, [u8; 32])],
189) -> Result<BTreeSet<Uuid>, GfError> {
190 let mut retained = BTreeSet::new();
191 retained.insert(selected.generation_uuid());
192 let mut parent = selected.parent_generation_uuid();
193 for _ in 0..RETAINED_ANCESTORS {
194 let Some(uuid) = parent else {
195 break;
196 };
197 if !retained.insert(uuid) {
198 return Err(recovery_corrupt("generation ancestry contains a cycle"));
199 }
200 parent = validated_generation_parent(root, uuid)?;
201 }
202 for (uuid, expected_digest) in checkpoint_roots {
203 let actual_digest = validated_generation_manifest_sha256(root, *uuid)?;
204 if actual_digest != *expected_digest {
205 return Err(recovery_corrupt(
206 "checkpoint generation manifest digest does not match registry",
207 ));
208 }
209 retained.insert(*uuid);
210 }
211 Ok(retained)
212}
213
214fn cleanup_abandoned_generation(
215 root: &Path,
216 transaction_uuid: Uuid,
217 generation_uuid: Uuid,
218 request_fingerprint: &str,
219 retained: &BTreeSet<Uuid>,
220) -> Result<u64, GfError> {
221 if retained.contains(&generation_uuid) {
222 return Ok(0);
223 }
224 let generation_name = generation_uuid.hyphenated().to_string();
225 let generation_path = root.join(GENERATIONS_DIR).join(&generation_name);
226 let attempt_path = root
227 .join(ATTEMPTS_DIR)
228 .join(transaction_uuid.hyphenated().to_string())
229 .join(request_fingerprint);
230 let mut removed = 0;
231 if attempt_path.exists() {
232 reject_real_directory(&attempt_path)?;
233 std::fs::remove_dir_all(&attempt_path).map_err(storage_io)?;
234 sync_directory(
235 attempt_path
236 .parent()
237 .expect("machine attempt path has a parent"),
238 )?;
239 let transaction_attempt_root = attempt_path
240 .parent()
241 .expect("machine attempt path has a parent");
242 match std::fs::remove_dir(transaction_attempt_root) {
243 Ok(()) => sync_directory(&root.join(ATTEMPTS_DIR))?,
244 Err(error) if error.kind() == std::io::ErrorKind::DirectoryNotEmpty => {}
245 Err(error) => return Err(storage_io(error)),
246 }
247 removed = 1;
248 }
249 let trash_root = ensure_machine_directory(root, Path::new(TRASH_DIR))?;
250 let trash_path = trash_root.join(&generation_name);
251
252 if trash_path.exists() {
253 remove_trash_entry(
254 root,
255 &trash_root,
256 &trash_path,
257 transaction_uuid,
258 generation_uuid,
259 )?;
260 return Ok(1);
261 }
262 if !generation_path.exists() {
263 return Ok(removed);
264 }
265 reject_real_directory(&generation_path)?;
266 let lease_path = generation_path.join("lease.lock");
267 let _lease = if lease_path.exists() {
268 let lease = open_regular_lock(&lease_path)?;
269 if !FileExt::try_lock_exclusive(&lease).map_err(storage_io)? {
270 return Ok(0);
271 }
272 Some(lease)
273 } else {
274 None
275 };
276
277 let current = resolve_project_generation(root).map_err(map_recovery_resolution)?;
278 if current.generation_uuid() == generation_uuid {
279 return Err(recovery_corrupt(
280 "cleanup candidate became the committed generation",
281 ));
282 }
283 std::fs::rename(&generation_path, &trash_path).map_err(storage_io)?;
284 sync_directory(&root.join(GENERATIONS_DIR))?;
285 sync_directory(&trash_root)?;
286 project_failpoint::hit(
287 "project.after_gc_move",
288 Some(transaction_uuid),
289 Some(generation_uuid),
290 "GC",
291 true,
292 )?;
293 remove_trash_entry(
294 root,
295 &trash_root,
296 &trash_path,
297 transaction_uuid,
298 generation_uuid,
299 )?;
300 Ok(1)
301}
302
303fn remove_trash_entry(
304 root: &Path,
305 trash_root: &Path,
306 trash_path: &Path,
307 transaction_uuid: Uuid,
308 generation_uuid: Uuid,
309) -> Result<(), GfError> {
310 reject_real_directory(trash_path)?;
311 if resolve_project_generation(root)
312 .map_err(map_recovery_resolution)?
313 .generation_uuid()
314 == generation_uuid
315 {
316 return Err(recovery_corrupt(
317 "trash entry is reachable from the committed pointer",
318 ));
319 }
320 std::fs::remove_dir_all(trash_path).map_err(storage_io)?;
321 sync_directory(trash_root)?;
322 project_failpoint::hit(
323 "project.after_gc_delete",
324 Some(transaction_uuid),
325 Some(generation_uuid),
326 "GC",
327 true,
328 )
329}
330
331fn count_unknown_generation_entries(
332 root: &Path,
333 retained: &BTreeSet<Uuid>,
334) -> Result<u64, GfError> {
335 let generations_root = root.join(GENERATIONS_DIR);
336 let entries = bounded_directory_entries(&generations_root)?;
337 let mut unknown = 0_u64;
338 for path in entries {
339 let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
340 unknown += 1;
341 continue;
342 };
343 let Some(uuid) = parse_canonical_uuid(name) else {
344 unknown += 1;
345 continue;
346 };
347 if !retained.contains(&uuid) {
348 unknown += 1;
351 }
352 }
353 Ok(unknown)
354}
355
356fn bounded_directory_entries(root: &Path) -> Result<Vec<PathBuf>, GfError> {
357 reject_real_directory(root)?;
358 let mut entries = Vec::new();
359 for entry in std::fs::read_dir(root).map_err(storage_io)? {
360 if entries.len() >= MAX_RECOVERY_ENTRIES {
361 return Err(recovery_corrupt("recovery entry limit exceeded"));
362 }
363 entries.push(entry.map_err(storage_io)?.path());
364 }
365 Ok(entries)
366}
367
368fn journal_file_uuid(path: &Path) -> Option<Uuid> {
369 let file_name = path.file_name()?.to_str()?;
370 let stem = file_name.strip_suffix(".json")?;
371 parse_canonical_uuid(stem)
372}
373
374fn parse_canonical_uuid(value: &str) -> Option<Uuid> {
375 let uuid = Uuid::parse_str(value).ok()?;
376 (uuid.hyphenated().to_string() == value).then_some(uuid)
377}
378
379fn reject_real_directory(path: &Path) -> Result<(), GfError> {
380 let metadata = std::fs::symlink_metadata(path).map_err(storage_io)?;
381 if !metadata.is_dir() || metadata.file_type().is_symlink() {
382 return Err(recovery_corrupt(
383 "recovery path is linked or not a directory",
384 ));
385 }
386 Ok(())
387}
388
389fn digest_hex(digest: [u8; 32]) -> String {
390 use std::fmt::Write as _;
391
392 digest
393 .iter()
394 .fold(String::with_capacity(64), |mut output, byte| {
395 write!(output, "{byte:02x}").expect("writing to String cannot fail");
396 output
397 })
398}
399
400fn recovery_corrupt(cause: &str) -> GfError {
401 project_error(
402 ProjectErrorCode::ProjectCorrupt,
403 format!(
404 "phase=RECOVERY committed=unknown cause={cause}; preserve the project and restore \
405 CURRENT plus its exact committed generation from a verified backup"
406 ),
407 )
408}
409
410fn map_recovery_resolution(error: GfError) -> GfError {
411 if error.code() == "GF_PROJECT_CORRUPT" {
412 recovery_corrupt("committed publication record is invalid or ambiguous")
413 } else {
414 error
415 }
416}
417
418fn storage_io(error: impl std::fmt::Display) -> GfError {
419 GfError::Storage(error.to_string())
420}
421
422fn project_error(code: ProjectErrorCode, message: impl Into<String>) -> GfError {
423 GfError::Project {
424 code,
425 message: message.into(),
426 }
427}
428
429#[cfg(test)]
430mod tests {
431 use std::process::Command;
432 use std::time::{Duration, Instant};
433
434 use sha2::{Digest, Sha256};
435
436 use super::*;
437 use crate::{
438 ProjectCapability, ProjectGenerationRequest, ProjectParticipant,
439 ProjectParticipantEncoding, ProjectStageOutcome, open_or_initialize_project,
440 stage_project_generation, stage_project_generation_optimistic,
441 };
442
443 const ENABLE_COOKIE: &str = "graphforge-internal-subprocess-v1";
444 const WRITER_HELPER: &str = "project_recovery::tests::subprocess_publication_writer";
445 const RECOVERY_HELPER: &str = "project_recovery::tests::subprocess_recovery_runner";
446 const INITIALIZER_HELPER: &str = "project_recovery::tests::subprocess_initializer";
447 const PRE_COMMIT_FAILPOINTS: &[&str] = &[
448 "project.after_writer_lock",
449 "project.after_journal_preparing",
450 "project.after_participant_write",
451 "project.after_participant_fsync",
452 "project.after_participant_dir_fsync",
453 "project.after_journal_staged",
454 "project.after_domain_validation",
455 "project.after_composite_validation",
456 "project.after_journal_validated",
457 "project.after_manifest_write",
458 "project.after_manifest_fsync",
459 "project.after_generation_dir_fsync",
460 "project.after_journal_durable",
461 "project.after_current_temp_write",
462 "project.after_current_temp_fsync",
463 "project.before_current_replace",
464 ];
465 const POST_COMMIT_FAILPOINTS: &[&str] = &[
466 "project.after_current_replace",
467 "project.after_root_fsync",
468 "project.after_journal_published",
469 ];
470
471 fn wait_for_writer_lock_release(root: &Path) {
472 let lock = open_regular_lock(&root.join(LOCKS_DIR).join(WRITER_LOCK_FILE)).unwrap();
473 let deadline = Instant::now() + Duration::from_secs(1);
474 loop {
475 if FileExt::try_lock_exclusive(&lock).unwrap() {
476 let acquired_at = Instant::now();
477 FileExt::unlock(&lock).unwrap();
478 assert!(
479 acquired_at < deadline,
480 "writer.lock remained owned after recovery completed"
481 );
482 return;
483 }
484 assert!(
485 Instant::now() < deadline,
486 "writer.lock remained owned after recovery completed"
487 );
488 std::thread::sleep(Duration::from_millis(1));
489 }
490 }
491
492 fn participant(capability: &str, family: &str) -> ProjectParticipant {
493 let bytes = format!("{capability}:{family}").into_bytes();
494 ProjectParticipant {
495 capability_id: capability.into(),
496 capability_version: 1,
497 record_family_id: family.into(),
498 record_version: 1,
499 encoding: ProjectParticipantEncoding::Parquet,
500 schema_fingerprint: Sha256::digest(format!("{capability}/{family}")).into(),
501 row_count: 1,
502 bytes,
503 }
504 }
505
506 fn participants(set: &str) -> Vec<ProjectParticipant> {
507 match set {
508 "graph" => vec![participant("graph", "nodes")],
509 "provenance" => vec![
510 participant("graph", "nodes"),
511 participant("provenance", "events"),
512 ],
513 "knowledge" => vec![
514 participant("graph", "nodes"),
515 participant("provenance", "events"),
516 participant("knowledge", "assertions"),
517 ],
518 other => panic!("unknown test participant set {other}"),
519 }
520 }
521
522 fn capabilities(set: &str) -> Vec<ProjectCapability> {
523 let mut capabilities = vec![ProjectCapability {
524 capability_id: "graph".into(),
525 capability_version: 1,
526 }];
527 if matches!(set, "provenance" | "knowledge") {
528 capabilities.push(ProjectCapability {
529 capability_id: "provenance".into(),
530 capability_version: 1,
531 });
532 }
533 if set == "knowledge" {
534 capabilities.push(ProjectCapability {
535 capability_id: "knowledge".into(),
536 capability_version: 1,
537 });
538 }
539 capabilities
540 }
541
542 fn spawn_writer(
543 root: &Path,
544 transaction_uuid: Uuid,
545 generation_uuid: Uuid,
546 set: &str,
547 failpoint: &str,
548 ) -> std::process::ExitStatus {
549 Command::new(std::env::current_exe().unwrap())
550 .arg("--exact")
551 .arg(WRITER_HELPER)
552 .arg("--nocapture")
553 .env("GRAPHFORGE_TEST_PROJECT_ROOT", root)
554 .env(
555 "GRAPHFORGE_TEST_TRANSACTION_UUID",
556 transaction_uuid.hyphenated().to_string(),
557 )
558 .env(
559 "GRAPHFORGE_TEST_GENERATION_UUID",
560 generation_uuid.hyphenated().to_string(),
561 )
562 .env("GRAPHFORGE_TEST_PARTICIPANT_SET", set)
563 .env("GRAPHFORGE_PROJECT_FAILPOINTS", ENABLE_COOKIE)
564 .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
565 .status()
566 .unwrap()
567 }
568
569 fn spawn_optimistic_writer(
570 root: &Path,
571 transaction_uuid: Uuid,
572 generation_uuid: Uuid,
573 failpoint: &str,
574 ) -> std::process::ExitStatus {
575 Command::new(std::env::current_exe().unwrap())
576 .arg("--exact")
577 .arg(WRITER_HELPER)
578 .arg("--nocapture")
579 .env("GRAPHFORGE_TEST_PROJECT_ROOT", root)
580 .env(
581 "GRAPHFORGE_TEST_TRANSACTION_UUID",
582 transaction_uuid.hyphenated().to_string(),
583 )
584 .env(
585 "GRAPHFORGE_TEST_GENERATION_UUID",
586 generation_uuid.hyphenated().to_string(),
587 )
588 .env("GRAPHFORGE_TEST_PARTICIPANT_SET", "graph")
589 .env("GRAPHFORGE_TEST_OPTIMISTIC", "1")
590 .env("GRAPHFORGE_PROJECT_FAILPOINTS", ENABLE_COOKIE)
591 .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
592 .status()
593 .unwrap()
594 }
595
596 fn assert_reopen(root: &Path, expected: Uuid, set: &str, expect_child: bool) {
597 let before_recovery = resolve_project_generation(root).unwrap();
598 assert_eq!(before_recovery.generation_uuid(), expected);
599 let expected_manifest_digest = before_recovery.manifest_sha256();
600 let first = recover_project_transactions(root).unwrap();
601 assert_eq!(first.selected_generation_uuid, expected);
602 let resolved = resolve_project_generation(root).unwrap();
603 assert_eq!(resolved.generation_uuid(), expected);
604 assert_eq!(resolved.manifest_sha256(), expected_manifest_digest);
605 let manifest_bytes =
606 std::fs::read(resolved.generation_root().join("manifest.json")).unwrap();
607 let manifest_digest: [u8; 32] = Sha256::digest(&manifest_bytes).into();
608 assert_eq!(manifest_digest, expected_manifest_digest);
609 let manifest: serde_json::Value = serde_json::from_slice(&manifest_bytes).unwrap();
610 if expect_child {
611 for participant in participants(set) {
612 let path = resolved
613 .participant_path(&participant.capability_id, &participant.record_family_id)
614 .unwrap();
615 let bytes = std::fs::read(path).unwrap();
616 assert_eq!(bytes, participant.bytes);
617 let persisted = manifest["participants"]
618 .as_array()
619 .unwrap()
620 .iter()
621 .find(|entry| {
622 entry["capability_id"] == participant.capability_id
623 && entry["record_family_id"] == participant.record_family_id
624 })
625 .unwrap();
626 let content_digest: [u8; 32] = Sha256::digest(&bytes).into();
627 assert_eq!(
628 persisted["content_sha256"].as_str().unwrap(),
629 digest_hex(content_digest)
630 );
631 }
632 }
633 wait_for_writer_lock_release(root);
634 let second = recover_project_transactions(root).unwrap();
635 assert_eq!(second.selected_generation_uuid, expected);
636 assert_eq!(second.repaired_journals, 0);
637 assert_eq!(second.aborted_journals, 0);
638 }
639
640 #[test]
641 fn subprocess_kill_matrix_never_exposes_a_partial_generation() {
642 for (failpoint, committed) in PRE_COMMIT_FAILPOINTS
643 .iter()
644 .map(|name| (*name, false))
645 .chain(POST_COMMIT_FAILPOINTS.iter().map(|name| (*name, true)))
646 {
647 let root = tempfile::tempdir().unwrap();
648 let parent = open_or_initialize_project(root.path())
649 .unwrap()
650 .generation_uuid();
651 let transaction_uuid = Uuid::now_v7();
652 let generation_uuid = Uuid::now_v7();
653 let status = spawn_writer(
654 root.path(),
655 transaction_uuid,
656 generation_uuid,
657 "graph",
658 failpoint,
659 );
660 assert_eq!(
661 status.code(),
662 Some(crate::project_failpoint::exit_code()),
663 "{failpoint} did not terminate at the named boundary"
664 );
665 assert_reopen(
666 root.path(),
667 if committed { generation_uuid } else { parent },
668 "graph",
669 committed,
670 );
671 }
672 }
673
674 #[test]
675 fn optimistic_commit_failpoints_never_expose_a_partial_generation() {
676 for failpoint in [
677 "project.after_optimistic_commit_lock",
678 "project.after_optimistic_promotion",
679 ] {
680 let root = tempfile::tempdir().unwrap();
681 let parent = open_or_initialize_project(root.path())
682 .unwrap()
683 .generation_uuid();
684 let status =
685 spawn_optimistic_writer(root.path(), Uuid::now_v7(), Uuid::now_v7(), failpoint);
686 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
687 assert_reopen(root.path(), parent, "graph", false);
688 }
689 }
690
691 #[test]
692 fn container_creation_failpoints_resume_only_exact_current_format() {
693 for failpoint in [
694 "project.after_format_fsync",
695 "project.after_container_dir_fsync",
696 ] {
697 for (active, killed) in [
698 (failpoint.to_owned(), true),
699 (format!("{failpoint}.error"), false),
700 ] {
701 let root = tempfile::tempdir().unwrap();
702 let status = Command::new(std::env::current_exe().unwrap())
703 .arg("--exact")
704 .arg(INITIALIZER_HELPER)
705 .arg("--nocapture")
706 .env("GRAPHFORGE_TEST_PROJECT_ROOT", root.path())
707 .env("GRAPHFORGE_PROJECT_FAILPOINTS", ENABLE_COOKIE)
708 .env("GRAPHFORGE_PROJECT_FAILPOINT", active)
709 .status()
710 .unwrap();
711 if killed {
712 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
713 } else {
714 assert!(status.success());
715 }
716
717 let reopened = open_or_initialize_project(root.path()).unwrap();
718 assert_eq!(
719 resolve_project_generation(root.path())
720 .unwrap()
721 .generation_uuid(),
722 reopened.generation_uuid()
723 );
724 }
725 }
726 }
727
728 #[test]
729 fn exact_format_does_not_authorize_mutating_an_unknown_layout() {
730 let root = tempfile::tempdir().unwrap();
731 std::fs::write(
732 root.path().join(crate::FORMAT_FILE),
733 crate::PROJECT_FORMAT_BYTES,
734 )
735 .unwrap();
736 std::fs::create_dir(root.path().join(GENERATIONS_DIR)).unwrap();
737 std::fs::write(root.path().join("unknown.db"), b"do-not-touch").unwrap();
738 let before = std::fs::read(root.path().join("unknown.db")).unwrap();
739
740 let error = open_or_initialize_project(root.path()).unwrap_err();
741
742 assert_eq!(error.code(), "GF_UNSUPPORTED_PROJECT_FORMAT");
743 assert_eq!(
744 std::fs::read(root.path().join("unknown.db")).unwrap(),
745 before
746 );
747 assert!(!root.path().join(crate::CURRENT_FILE).exists());
748 }
749
750 #[test]
751 fn recovery_removes_only_validated_atomic_journal_temps() {
752 let root = tempfile::tempdir().unwrap();
753 open_or_initialize_project(root.path()).unwrap();
754 let transactions = root.path().join(TRANSACTIONS_DIR);
755 std::fs::create_dir(&transactions).unwrap();
756 let writer_temp = transactions.join(".atomicwriteZ9y8X7");
757 std::fs::create_dir(&writer_temp).unwrap();
758 std::fs::write(writer_temp.join("tmpfile.tmp"), b"partial journal").unwrap();
759
760 recover_project_transactions(root.path()).unwrap();
761
762 assert!(!writer_temp.exists());
763 }
764
765 #[test]
766 fn recovery_rejects_spoofed_atomic_journal_temp() {
767 let root = tempfile::tempdir().unwrap();
768 open_or_initialize_project(root.path()).unwrap();
769 let transactions = root.path().join(TRANSACTIONS_DIR);
770 std::fs::create_dir(&transactions).unwrap();
771 let spoofed = transactions.join(".atomicwriteZ9y8X7");
772 std::fs::create_dir(&spoofed).unwrap();
773 std::fs::write(spoofed.join("unexpected"), b"preserve").unwrap();
774
775 let error = recover_project_transactions(root.path()).unwrap_err();
776
777 assert_eq!(error.code(), "GF_PROJECT_CORRUPT");
778 assert_eq!(
779 std::fs::read(spoofed.join("unexpected")).unwrap(),
780 b"preserve"
781 );
782 }
783
784 #[test]
785 fn injected_operation_errors_report_exact_commit_state() {
786 for failpoint in PRE_COMMIT_FAILPOINTS
787 .iter()
788 .copied()
789 .chain(std::iter::once("project.after_current_replace"))
790 {
791 let root = tempfile::tempdir().unwrap();
792 let parent = open_or_initialize_project(root.path())
793 .unwrap()
794 .generation_uuid();
795 let transaction_uuid = Uuid::now_v7();
796 let generation_uuid = Uuid::now_v7();
797 let status = spawn_writer(
798 root.path(),
799 transaction_uuid,
800 generation_uuid,
801 "graph",
802 &format!("{failpoint}.error"),
803 );
804 assert!(status.success(), "{failpoint}.error helper failed");
805 let committed = failpoint == "project.after_current_replace";
806 assert_reopen(
807 root.path(),
808 if committed { generation_uuid } else { parent },
809 "graph",
810 committed,
811 );
812 }
813 }
814
815 #[test]
816 fn provenance_and_knowledge_sets_follow_the_same_commit_boundary() {
817 for set in ["provenance", "knowledge"] {
818 for (failpoint, committed) in [
819 ("project.after_participant_fsync", false),
820 ("project.after_current_replace", true),
821 ] {
822 let root = tempfile::tempdir().unwrap();
823 let parent = open_or_initialize_project(root.path())
824 .unwrap()
825 .generation_uuid();
826 let transaction_uuid = Uuid::now_v7();
827 let generation_uuid = Uuid::now_v7();
828 let status = spawn_writer(
829 root.path(),
830 transaction_uuid,
831 generation_uuid,
832 set,
833 failpoint,
834 );
835 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
836 assert_reopen(
837 root.path(),
838 if committed { generation_uuid } else { parent },
839 set,
840 committed,
841 );
842 }
843 }
844 }
845
846 #[test]
847 fn recovered_aborted_transaction_can_retry_with_identical_inputs() {
848 let root = tempfile::tempdir().unwrap();
849 open_or_initialize_project(root.path()).unwrap();
850 let transaction_uuid = Uuid::now_v7();
851 let generation_uuid = Uuid::now_v7();
852 let status = spawn_writer(
853 root.path(),
854 transaction_uuid,
855 generation_uuid,
856 "graph",
857 "project.after_journal_staged",
858 );
859 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
860 recover_project_transactions(root.path()).unwrap();
861 let request = ProjectGenerationRequest {
862 transaction_uuid,
863 generation_uuid,
864 capabilities: capabilities("graph"),
865 participants: participants("graph"),
866 };
867
868 let ProjectStageOutcome::Staged(staged) =
869 stage_project_generation(root.path(), &request).unwrap()
870 else {
871 panic!("aborted transaction unexpectedly replayed as published");
872 };
873 staged
874 .validate(|_| Ok(()), |_, _| Ok(()))
875 .unwrap()
876 .publish()
877 .unwrap();
878 assert_eq!(
879 resolve_project_generation(root.path())
880 .unwrap()
881 .generation_uuid(),
882 generation_uuid
883 );
884 }
885
886 #[test]
887 fn reachable_ancestor_with_stale_journal_is_repaired_not_deleted() {
888 let root = tempfile::tempdir().unwrap();
889 open_or_initialize_project(root.path()).unwrap();
890 let first = ProjectGenerationRequest {
891 transaction_uuid: Uuid::now_v7(),
892 generation_uuid: Uuid::now_v7(),
893 capabilities: capabilities("graph"),
894 participants: participants("graph"),
895 };
896 let ProjectStageOutcome::Staged(staged) =
897 stage_project_generation(root.path(), &first).unwrap()
898 else {
899 panic!("new transaction replayed");
900 };
901 staged
902 .validate(|_| Ok(()), |_, _| Ok(()))
903 .unwrap()
904 .publish()
905 .unwrap();
906 let first_journal_path = root
907 .path()
908 .join(TRANSACTIONS_DIR)
909 .join(format!("{}.json", first.transaction_uuid.hyphenated()));
910 let mut first_journal = read_journal(&first_journal_path).unwrap();
911 first_journal.phase = JournalPhase::Durable;
912 write_journal(&first_journal_path, &first_journal).unwrap();
913
914 let second = ProjectGenerationRequest {
915 transaction_uuid: Uuid::now_v7(),
916 generation_uuid: Uuid::now_v7(),
917 capabilities: capabilities("graph"),
918 participants: participants("graph"),
919 };
920 let ProjectStageOutcome::Staged(staged) =
921 stage_project_generation(root.path(), &second).unwrap()
922 else {
923 panic!("new transaction replayed");
924 };
925 staged
926 .validate(|_| Ok(()), |_, _| Ok(()))
927 .unwrap()
928 .publish()
929 .unwrap();
930
931 let report = recover_project_transactions(root.path()).unwrap();
932
933 assert_eq!(report.repaired_journals, 1);
934 assert_eq!(
935 read_journal(&first_journal_path).unwrap().phase,
936 JournalPhase::Published
937 );
938 assert!(
939 root.path()
940 .join(GENERATIONS_DIR)
941 .join(first.generation_uuid.hyphenated().to_string())
942 .exists()
943 );
944 }
945
946 #[test]
947 fn torn_journal_fails_closed_without_changing_current() {
948 let root = tempfile::tempdir().unwrap();
949 let parent = open_or_initialize_project(root.path())
950 .unwrap()
951 .generation_uuid();
952 let transaction_uuid = Uuid::now_v7();
953 let generation_uuid = Uuid::now_v7();
954 let status = spawn_writer(
955 root.path(),
956 transaction_uuid,
957 generation_uuid,
958 "graph",
959 "project.after_journal_staged",
960 );
961 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
962 std::fs::write(
963 root.path()
964 .join(TRANSACTIONS_DIR)
965 .join(format!("{}.json", transaction_uuid.hyphenated())),
966 b"{torn",
967 )
968 .unwrap();
969
970 let error = recover_project_transactions(root.path()).unwrap_err();
971 assert_eq!(error.code(), "GF_PROJECT_CORRUPT");
972 assert!(error.to_string().contains("verified backup"));
973 assert_eq!(
974 resolve_project_generation(root.path())
975 .unwrap()
976 .generation_uuid(),
977 parent
978 );
979 }
980
981 #[test]
982 fn invalid_current_returns_recovery_guidance_without_election() {
983 let root = tempfile::tempdir().unwrap();
984 let selected = open_or_initialize_project(root.path())
985 .unwrap()
986 .generation_uuid();
987 std::fs::write(root.path().join(crate::CURRENT_FILE), b"{invalid\n").unwrap();
988
989 let error = recover_project_transactions(root.path()).unwrap_err();
990
991 assert_eq!(error.code(), "GF_PROJECT_CORRUPT");
992 assert!(error.to_string().contains("verified backup"));
993 assert!(
994 root.path()
995 .join(GENERATIONS_DIR)
996 .join(selected.hyphenated().to_string())
997 .exists()
998 );
999 }
1000
1001 #[test]
1002 fn live_writer_lock_blocks_recovery_without_metadata_heuristics() {
1003 let root = tempfile::tempdir().unwrap();
1004 open_or_initialize_project(root.path()).unwrap();
1005 let lock_dir = ensure_machine_directory(root.path(), Path::new(LOCKS_DIR)).unwrap();
1006 let lock = open_regular_lock(&lock_dir.join(WRITER_LOCK_FILE)).unwrap();
1007 FileExt::lock_exclusive(&lock).unwrap();
1008
1009 let error = recover_project_transactions(root.path()).unwrap_err();
1010
1011 assert_eq!(error.code(), "GF_WRITER_BUSY");
1012 }
1013
1014 #[test]
1015 fn corrupt_knowledge_bytes_do_not_block_graph_only_reopen_or_recovery() {
1016 let root = tempfile::tempdir().unwrap();
1017 open_or_initialize_project(root.path()).unwrap();
1018 let request = ProjectGenerationRequest {
1019 transaction_uuid: Uuid::now_v7(),
1020 generation_uuid: Uuid::now_v7(),
1021 capabilities: capabilities("knowledge"),
1022 participants: participants("knowledge"),
1023 };
1024 let ProjectStageOutcome::Staged(staged) =
1025 stage_project_generation(root.path(), &request).unwrap()
1026 else {
1027 panic!("new transaction replayed");
1028 };
1029 staged
1030 .validate(|_| Ok(()), |_, _| Ok(()))
1031 .unwrap()
1032 .publish()
1033 .unwrap();
1034 let resolved = resolve_project_generation(root.path()).unwrap();
1035 std::fs::write(
1036 resolved
1037 .participant_path("knowledge", "assertions")
1038 .unwrap(),
1039 b"future-or-corrupt-knowledge",
1040 )
1041 .unwrap();
1042
1043 let report = recover_project_transactions(root.path()).unwrap();
1044 assert_eq!(report.selected_generation_uuid, request.generation_uuid);
1045 let graph = resolve_project_generation(root.path())
1046 .unwrap()
1047 .participant_path("graph", "nodes")
1048 .unwrap();
1049 assert_eq!(std::fs::read(graph).unwrap(), b"graph:nodes");
1050 }
1051
1052 #[test]
1053 fn gc_crash_points_preserve_current_and_recover_idempotently() {
1054 for failpoint in ["project.after_gc_move", "project.after_gc_delete"] {
1055 let root = tempfile::tempdir().unwrap();
1056 let current = open_or_initialize_project(root.path())
1057 .unwrap()
1058 .generation_uuid();
1059 let transaction_uuid = Uuid::now_v7();
1060 let generation_uuid = Uuid::now_v7();
1061 let status = spawn_writer(
1062 root.path(),
1063 transaction_uuid,
1064 generation_uuid,
1065 "graph",
1066 "project.after_journal_staged",
1067 );
1068 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
1069 let status = Command::new(std::env::current_exe().unwrap())
1070 .arg("--exact")
1071 .arg(RECOVERY_HELPER)
1072 .arg("--nocapture")
1073 .env("GRAPHFORGE_TEST_PROJECT_ROOT", root.path())
1074 .env("GRAPHFORGE_PROJECT_FAILPOINTS", ENABLE_COOKIE)
1075 .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
1076 .status()
1077 .unwrap();
1078 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
1079
1080 assert_reopen(root.path(), current, "graph", false);
1081 assert!(
1082 !root
1083 .path()
1084 .join(GENERATIONS_DIR)
1085 .join(generation_uuid.hyphenated().to_string())
1086 .exists()
1087 );
1088 }
1089 }
1090
1091 #[test]
1092 fn subprocess_publication_writer() {
1093 if std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT").is_err() {
1094 return;
1095 }
1096 let root = PathBuf::from(std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT").unwrap());
1097 let transaction_uuid =
1098 Uuid::parse_str(&std::env::var("GRAPHFORGE_TEST_TRANSACTION_UUID").unwrap()).unwrap();
1099 let generation_uuid =
1100 Uuid::parse_str(&std::env::var("GRAPHFORGE_TEST_GENERATION_UUID").unwrap()).unwrap();
1101 let set = std::env::var("GRAPHFORGE_TEST_PARTICIPANT_SET").unwrap();
1102 let active = std::env::var("GRAPHFORGE_PROJECT_FAILPOINT").unwrap();
1103 let request = ProjectGenerationRequest {
1104 transaction_uuid,
1105 generation_uuid,
1106 capabilities: capabilities(&set),
1107 participants: participants(&set),
1108 };
1109 let result = (|| {
1110 let outcome = if std::env::var("GRAPHFORGE_TEST_OPTIMISTIC").is_ok() {
1111 stage_project_generation_optimistic(
1112 &root,
1113 &request,
1114 Sha256::digest(b"optimistic-subprocess-operation").into(),
1115 )?
1116 } else {
1117 stage_project_generation(&root, &request)?
1118 };
1119 let ProjectStageOutcome::Staged(staged) = outcome else {
1120 panic!("new transaction replayed");
1121 };
1122 staged.validate(|_| Ok(()), |_, _| Ok(()))?.publish()
1123 })();
1124 let error = result.expect_err("configured failpoint did not fire");
1125 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
1126 assert!(
1127 error
1128 .to_string()
1129 .contains(if active == "project.after_current_replace.error" {
1130 "committed=true"
1131 } else {
1132 "committed=false"
1133 })
1134 );
1135 }
1136
1137 #[test]
1138 fn subprocess_recovery_runner() {
1139 if std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT").is_err() {
1140 return;
1141 }
1142 let root = PathBuf::from(std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT").unwrap());
1143 recover_project_transactions(root).unwrap();
1144 panic!("configured recovery failpoint did not fire");
1145 }
1146
1147 #[test]
1148 fn subprocess_initializer() {
1149 if std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT").is_err() {
1150 return;
1151 }
1152 let root = PathBuf::from(std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT").unwrap());
1153 let active = std::env::var("GRAPHFORGE_PROJECT_FAILPOINT").unwrap();
1154 let result = open_or_initialize_project(root);
1155 if active.ends_with(".error") {
1156 let error = result.expect_err("configured initialization error did not fire");
1157 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
1158 assert!(error.to_string().contains("committed=false"));
1159 return;
1160 }
1161 result.unwrap();
1162 panic!("configured initialization failpoint did not fire");
1163 }
1164}