1use std::collections::BTreeSet;
4use std::fmt::Write as _;
5use std::fs::{self, File, OpenOptions};
6use std::io::{Read, Write};
7use std::path::{Path, PathBuf};
8use std::sync::Arc;
9use std::time::{SystemTime, UNIX_EPOCH};
10
11use arrow::array::{
12 ArrayRef, FixedSizeBinaryBuilder, StringArray, TimestampMicrosecondArray, UInt32Array,
13};
14use arrow::datatypes::{DataType, Field, Schema, TimeUnit};
15use arrow::record_batch::RecordBatch;
16use fs4::fs_std::FileExt;
17use graphforge_core::canonical::{CANONICAL_CONTRACT_VERSION, CanonicalDomain, fingerprint};
18use graphforge_core::{GfError, ProjectErrorCode};
19use parquet::arrow::ArrowWriter;
20use parquet::file::properties::WriterProperties;
21use serde::{Deserialize, Serialize};
22use sha2::{Digest, Sha256};
23use unicode_normalization::UnicodeNormalization;
24use uuid::Uuid;
25
26use crate::project_failpoint;
27use crate::project_generation::resolve_verified_generation;
28use crate::project_publication::{
29 LOCKS_DIR, ProjectCapability, ProjectGenerationRequest, ProjectParticipant,
30 ProjectParticipantEncoding, ProjectStageOutcome, RevertJournalExtension, WRITER_LOCK_FILE,
31 ensure_machine_directory, load_published_revert, load_revert_journal_extension,
32 open_regular_lock, stage_project_generation_with_lock, sync_directory,
33};
34use crate::resolve_project_generation;
35
36const CHECKPOINTS_DIR: &str = "checkpoints";
37const REGISTRY_FILE: &str = "registry.json";
38const CHECKSUM_FILE: &str = "registry.json.sha256";
39const INTENT_FILE: &str = "registry.txn.json";
40const CHECKPOINT_LOCK_FILE: &str = "checkpoints.lock";
41const MAX_REGISTRY_BYTES: u64 = 8 * 1024 * 1024;
42const MAX_ACTIVE: usize = 1_024;
43const MAX_TOMBSTONES: usize = 4_096;
44const MAX_NAME_BYTES: usize = 128;
45const MAX_DESCRIPTION_BYTES: usize = 1_024;
46const MAX_REASON_BYTES: usize = 1_024;
47const RESTORATION_FAMILY: &str = "restoration_transition";
48const RESTORATION_CONTRACT_VERSION: u32 = 1;
49
50#[derive(Debug, Clone)]
52pub struct CheckpointCreateRequest {
53 pub operation_uuid: Uuid,
55 pub name: String,
57 pub description: Option<String>,
59 pub actor_uuid: Option<Uuid>,
61}
62
63#[derive(Debug, Clone)]
65pub struct CheckpointDeleteRequest {
66 pub operation_uuid: Uuid,
68 pub name: String,
70 pub actor_uuid: Option<Uuid>,
72}
73
74#[derive(Debug, Clone)]
76pub struct CheckpointRevertRequest {
77 pub operation_uuid: Uuid,
79 pub name: String,
81 pub reason: String,
83 pub actor_uuid: Option<Uuid>,
85}
86
87#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(deny_unknown_fields)]
90pub struct CheckpointRecord {
91 pub checkpoint_uuid: Uuid,
93 pub name: String,
95 pub generation_uuid: Uuid,
97 pub generation_manifest_sha256: String,
99 pub description: Option<String>,
101 pub created_at: i64,
103 pub created_by: Option<Uuid>,
105 pub create_operation_uuid: Uuid,
107 pub create_request_sha256: String,
109 pub created_revision: u64,
111}
112
113#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
114#[serde(deny_unknown_fields)]
115struct CheckpointTombstone {
116 checkpoint_uuid: Uuid,
117 name: String,
118 generation_uuid: Uuid,
119 generation_manifest_sha256: String,
120 description: Option<String>,
121 created_at: i64,
122 created_by: Option<Uuid>,
123 create_operation_uuid: Uuid,
124 create_request_sha256: String,
125 created_revision: u64,
126 deleted_at: i64,
127 deleted_by: Option<Uuid>,
128 delete_operation_uuid: Uuid,
129 delete_request_sha256: String,
130 deleted_revision: u64,
131}
132
133#[derive(Debug, Clone, PartialEq, Eq)]
135pub struct CheckpointReceipt {
136 pub operation: &'static str,
138 pub operation_uuid: Uuid,
140 pub checkpoint_uuid: Uuid,
142 pub name: String,
144 pub source_generation_uuid: Uuid,
146 pub prior_current_generation_uuid: Option<Uuid>,
148 pub result_generation_uuid: Option<Uuid>,
150 pub registry_revision: u64,
152 pub committed_at: i64,
154}
155
156#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
157#[serde(deny_unknown_fields)]
158struct Registry {
159 format: String,
160 format_version: u32,
161 revision: u64,
162 active: Vec<CheckpointRecord>,
163 tombstones: Vec<CheckpointTombstone>,
164}
165
166impl Registry {
167 fn empty() -> Self {
168 Self {
169 format: "graphforge-checkpoints".into(),
170 format_version: 1,
171 revision: 0,
172 active: Vec::new(),
173 tombstones: Vec::new(),
174 }
175 }
176
177 fn canonical_bytes(&self) -> Result<Vec<u8>, GfError> {
178 validate_registry(self)?;
179 let mut bytes = serde_json::to_vec(self).map_err(registry_serde)?;
180 bytes.push(b'\n');
181 if bytes.len() as u64 > MAX_REGISTRY_BYTES {
182 return Err(project_error(
183 ProjectErrorCode::ResourceLimit,
184 "checkpoint registry exceeds 8 MiB",
185 ));
186 }
187 Ok(bytes)
188 }
189}
190
191#[derive(Debug, Clone, Serialize, Deserialize)]
192#[serde(deny_unknown_fields)]
193struct RegistryIntent {
194 transaction_uuid: Uuid,
195 previous_revision: Option<u64>,
196 previous_sha256: Option<String>,
197 next_revision: u64,
198 next_sha256: String,
199 registry_temp: String,
200 checksum_temp: String,
201}
202
203struct MutationLocks {
204 writer: Option<File>,
205 checkpoint: Option<File>,
206}
207
208impl MutationLocks {
209 fn transfer_writer_for_revert_publication(&mut self) -> File {
210 self.writer
211 .take()
212 .expect("writer lock must be present until revert publication")
213 }
214
215 fn release_revert_replay(mut self) -> Result<(), GfError> {
216 let checkpoint = self
217 .checkpoint
218 .take()
219 .expect("checkpoint lock must be present");
220 let writer = self.writer.take().expect("writer lock must be present");
221 release_revert_replay_locks(&checkpoint, &writer)
222 }
223}
224
225impl Drop for MutationLocks {
226 fn drop(&mut self) {
227 if let Some(checkpoint) = &self.checkpoint {
228 let _ = FileExt::unlock(checkpoint);
229 }
230 if let Some(writer) = &self.writer {
231 let _ = FileExt::unlock(writer);
232 }
233 }
234}
235
236pub fn create_checkpoint(
238 container_root: impl AsRef<Path>,
239 request: &CheckpointCreateRequest,
240) -> Result<CheckpointReceipt, GfError> {
241 let name = validate_name(&request.name)?;
242 validate_description(request.description.as_deref())?;
243 let root = canonical_project_root(container_root.as_ref())?;
244 let _locks = acquire_mutation_locks(&root)?;
245 let checkpoint_root = checkpoint_root(&root)?;
246 recover_pair(&checkpoint_root)?;
247 let mut registry = read_registry(&checkpoint_root)?;
248 let request_digest = create_request_digest(request, &name);
249 let request_hex = hex(&request_digest);
250
251 if let Some(row) = registry
252 .active
253 .iter()
254 .find(|row| row.create_operation_uuid == request.operation_uuid)
255 {
256 if row.create_request_sha256 == request_hex {
257 return Ok(create_receipt(row));
258 }
259 return Err(project_error(
260 ProjectErrorCode::TransactionConflict,
261 "checkpoint operation UUID was reused with different canonical request bytes",
262 ));
263 }
264 if registry
265 .tombstones
266 .iter()
267 .any(|row| row.delete_operation_uuid == request.operation_uuid)
268 {
269 return Err(project_error(
270 ProjectErrorCode::TransactionConflict,
271 "checkpoint operation UUID was already used by delete_checkpoint",
272 ));
273 }
274 if let Some(row) = registry
275 .tombstones
276 .iter()
277 .find(|row| row.create_operation_uuid == request.operation_uuid)
278 {
279 if row.create_request_sha256 == request_hex {
280 return Ok(create_tombstone_receipt(row));
281 }
282 return Err(project_error(
283 ProjectErrorCode::TransactionConflict,
284 "checkpoint create operation UUID was reused with different canonical request bytes",
285 ));
286 }
287 if registry.active.iter().any(|row| row.name == name) {
288 return Err(project_error(
289 ProjectErrorCode::CheckpointExists,
290 "checkpoint name already exists",
291 ));
292 }
293 if registry.active.len() >= MAX_ACTIVE {
294 return Err(project_error(
295 ProjectErrorCode::ResourceLimit,
296 "active checkpoint limit is 1024",
297 ));
298 }
299
300 let selected = resolve_project_generation(&root)?;
301 let now = utc_micros()?;
302 let checkpoint_uuid = checkpoint_uuid(request.operation_uuid, request_digest);
303 let revision = registry.revision.checked_add(1).ok_or_else(|| {
304 project_error(
305 ProjectErrorCode::ResourceLimit,
306 "checkpoint registry revision overflow",
307 )
308 })?;
309 let row = CheckpointRecord {
310 checkpoint_uuid,
311 name,
312 generation_uuid: selected.generation_uuid(),
313 generation_manifest_sha256: hex(&selected.manifest_sha256()),
314 description: request.description.clone(),
315 created_at: now,
316 created_by: request.actor_uuid,
317 create_operation_uuid: request.operation_uuid,
318 create_request_sha256: request_hex,
319 created_revision: revision,
320 };
321 registry.revision = revision;
322 registry.active.push(row.clone());
323 registry.active.sort_by(|left, right| {
324 (&left.name, left.checkpoint_uuid).cmp(&(&right.name, right.checkpoint_uuid))
325 });
326 commit_registry(&checkpoint_root, ®istry, request.operation_uuid)?;
327 Ok(create_receipt(&row))
328}
329
330pub fn delete_checkpoint(
332 container_root: impl AsRef<Path>,
333 request: &CheckpointDeleteRequest,
334) -> Result<CheckpointReceipt, GfError> {
335 let name = validate_name(&request.name)?;
336 let root = canonical_project_root(container_root.as_ref())?;
337 let _locks = acquire_mutation_locks(&root)?;
338 let checkpoint_root = checkpoint_root(&root)?;
339 recover_pair(&checkpoint_root)?;
340 let mut registry = read_registry(&checkpoint_root)?;
341 let digest = delete_request_digest(request, &name);
342 let digest_hex = hex(&digest);
343 if let Some(row) = registry
344 .tombstones
345 .iter()
346 .find(|row| row.delete_operation_uuid == request.operation_uuid)
347 {
348 if row.delete_request_sha256 == digest_hex {
349 return Ok(delete_receipt(row));
350 }
351 return Err(project_error(
352 ProjectErrorCode::TransactionConflict,
353 "checkpoint delete operation UUID was reused with different canonical request bytes",
354 ));
355 }
356 if registry
357 .active
358 .iter()
359 .any(|row| row.create_operation_uuid == request.operation_uuid)
360 || registry
361 .tombstones
362 .iter()
363 .any(|row| row.create_operation_uuid == request.operation_uuid)
364 {
365 return Err(project_error(
366 ProjectErrorCode::TransactionConflict,
367 "checkpoint operation UUID was already used by checkpoint",
368 ));
369 }
370 let index = registry
371 .active
372 .iter()
373 .position(|row| row.name == name)
374 .ok_or_else(|| {
375 project_error(
376 ProjectErrorCode::CheckpointNotFound,
377 "checkpoint name does not exist",
378 )
379 })?;
380 let row = registry.active.remove(index);
381 let now = utc_micros()?;
382 registry.revision = registry.revision.checked_add(1).ok_or_else(|| {
383 project_error(
384 ProjectErrorCode::ResourceLimit,
385 "checkpoint registry revision overflow",
386 )
387 })?;
388 let tombstone = CheckpointTombstone {
389 checkpoint_uuid: row.checkpoint_uuid,
390 name: row.name,
391 generation_uuid: row.generation_uuid,
392 generation_manifest_sha256: row.generation_manifest_sha256,
393 description: row.description,
394 created_at: row.created_at,
395 created_by: row.created_by,
396 create_operation_uuid: row.create_operation_uuid,
397 create_request_sha256: row.create_request_sha256,
398 created_revision: row.created_revision,
399 deleted_at: now,
400 deleted_by: request.actor_uuid,
401 delete_operation_uuid: request.operation_uuid,
402 delete_request_sha256: digest_hex,
403 deleted_revision: registry.revision,
404 };
405 registry.tombstones.push(tombstone.clone());
406 registry
407 .tombstones
408 .sort_by_key(|row| (row.deleted_revision, row.checkpoint_uuid));
409 if registry.tombstones.len() > MAX_TOMBSTONES {
410 registry
411 .tombstones
412 .drain(..registry.tombstones.len() - MAX_TOMBSTONES);
413 }
414 commit_registry(&checkpoint_root, ®istry, request.operation_uuid)?;
415 Ok(delete_receipt(&tombstone))
416}
417
418pub fn list_checkpoints(
420 container_root: impl AsRef<Path>,
421) -> Result<Vec<CheckpointRecord>, GfError> {
422 let root = canonical_project_root(container_root.as_ref())?;
423 let checkpoint_root = checkpoint_root(&root)?;
424 let (_checkpoint_lock, registry) = read_registry_for_read(&root, &checkpoint_root)?;
425 Ok(registry.active)
426}
427
428pub fn open_checkpoint_generation(
430 container_root: impl AsRef<Path>,
431 name: &str,
432) -> Result<(CheckpointRecord, crate::ResolvedProjectGeneration), GfError> {
433 let name = validate_name(name)?;
434 let root = canonical_project_root(container_root.as_ref())?;
435 let checkpoint_root = checkpoint_root(&root)?;
436 let (_checkpoint_lock, registry) = read_registry_for_read(&root, &checkpoint_root)?;
437 let row = registry
438 .active
439 .iter()
440 .find(|row| row.name == name)
441 .cloned()
442 .ok_or_else(|| {
443 project_error(
444 ProjectErrorCode::CheckpointNotFound,
445 "checkpoint name does not exist",
446 )
447 })?;
448 let generation = resolve_verified_generation(
449 &root,
450 row.generation_uuid,
451 decode_digest(&row.generation_manifest_sha256)?,
452 )?;
453 let after = read_registry(&checkpoint_root)?;
454 if after.revision != registry.revision
455 || !after.active.iter().any(|candidate| candidate == &row)
456 {
457 return Err(project_error(
458 ProjectErrorCode::CheckpointNotFound,
459 "checkpoint changed while its generation was being pinned",
460 ));
461 }
462 Ok((row, generation))
463}
464
465#[expect(
467 clippy::too_many_lines,
468 reason = "the revert transaction is intentionally linear so lock ownership and publication order remain auditable"
469)]
470pub fn revert_checkpoint<T, V>(
471 container_root: impl AsRef<Path>,
472 request: &CheckpointRevertRequest,
473 select_timestamp: T,
474 validate_source: V,
475) -> Result<(CheckpointReceipt, crate::ResolvedProjectGeneration), GfError>
476where
477 T: FnOnce() -> Result<i64, GfError>,
478 V: FnOnce(&crate::ResolvedProjectGeneration) -> Result<(), GfError>,
479{
480 let requested_name = validate_name(&request.name)?;
481 let requested_reason = validate_reason(&request.reason)?;
482 let root = canonical_project_root(container_root.as_ref())?;
483 let transaction_uuid = revert_transaction_uuid(request.operation_uuid);
484 let mut locks = acquire_mutation_locks(&root)?;
485 let checkpoint_root = checkpoint_root(&root)?;
486 recover_pair(&checkpoint_root)?;
487 let registry = read_registry(&checkpoint_root)?;
488 let prior_current = resolve_project_generation(&root)?;
489
490 if let Some((extension, receipt)) = load_published_revert(&root, transaction_uuid)? {
491 validate_revert_replay_request(request, &requested_name, &requested_reason, &extension)?;
492 let resolved = resolve_verified_generation(
493 &root,
494 receipt.generation_uuid,
495 receipt.generation_manifest_sha256,
496 )?;
497 validate_source(&resolved)?;
498 let replay = revert_receipt(
499 request,
500 &requested_name,
501 &extension,
502 receipt.generation_uuid,
503 )?;
504 locks.release_revert_replay()?;
505 return Ok((replay, resolved));
506 }
507
508 let prior_extension = load_revert_journal_extension(&root, transaction_uuid)?;
509 let (checkpoint, source, restored_at, registry_revision) =
510 if let Some(extension) = prior_extension.as_ref() {
511 let checkpoint_uuid = parse_uuid(&extension.checkpoint_uuid)?;
512 let source_uuid = parse_uuid(&extension.source_generation_uuid)?;
513 let source_digest = decode_digest(&extension.source_manifest_sha256)?;
514 let source = resolve_verified_generation(&root, source_uuid, source_digest)?;
515 let row = CheckpointRecord {
516 checkpoint_uuid,
517 name: extension.checkpoint_name.clone(),
518 generation_uuid: source_uuid,
519 generation_manifest_sha256: extension.source_manifest_sha256.clone(),
520 description: None,
521 created_at: 0,
522 created_by: None,
523 create_operation_uuid: Uuid::nil(),
524 create_request_sha256: "0".repeat(64),
525 created_revision: extension.registry_revision,
526 };
527 (
528 row,
529 source,
530 extension.restored_at,
531 extension.registry_revision,
532 )
533 } else {
534 let row = registry
535 .active
536 .iter()
537 .find(|row| row.name == requested_name)
538 .cloned()
539 .ok_or_else(|| {
540 project_error(
541 ProjectErrorCode::CheckpointNotFound,
542 "checkpoint name does not exist",
543 )
544 })?;
545 let source = resolve_verified_generation(
546 &root,
547 row.generation_uuid,
548 decode_digest(&row.generation_manifest_sha256)?,
549 )?;
550 (row, source, select_timestamp()?, registry.revision)
551 };
552
553 let request_digest = revert_request_digest(
554 request.operation_uuid,
555 &requested_name,
556 checkpoint.checkpoint_uuid,
557 source.generation_uuid(),
558 source.manifest_sha256(),
559 &requested_reason,
560 request.actor_uuid,
561 );
562 let request_hex = hex(&request_digest);
563 let restoration_uuid = restoration_uuid(request.operation_uuid, request_digest);
564 let original_prior_uuid = prior_extension.as_ref().map_or_else(
565 || Ok(prior_current.generation_uuid()),
566 |value| parse_uuid(&value.prior_current_generation_uuid),
567 )?;
568 let generation_uuid = restored_generation_uuid(
569 transaction_uuid,
570 checkpoint.checkpoint_uuid,
571 source.generation_uuid(),
572 source.manifest_sha256(),
573 original_prior_uuid,
574 restored_at,
575 request_digest,
576 );
577 let expected_extension = RevertJournalExtension {
578 operation_uuid: request.operation_uuid.to_string(),
579 request_sha256: request_hex,
580 checkpoint_uuid: checkpoint.checkpoint_uuid.to_string(),
581 checkpoint_name: requested_name.clone(),
582 source_generation_uuid: source.generation_uuid().to_string(),
583 source_manifest_sha256: hex(&source.manifest_sha256()),
584 prior_current_generation_uuid: original_prior_uuid.to_string(),
585 restored_at,
586 reason: requested_reason.clone(),
587 actor_uuid: request.actor_uuid.map(|value| value.to_string()),
588 restoration_uuid: restoration_uuid.to_string(),
589 registry_revision,
590 };
591 if prior_extension
592 .as_ref()
593 .is_some_and(|value| value != &expected_extension)
594 {
595 return Err(project_error(
596 ProjectErrorCode::TransactionConflict,
597 "revert operation UUID was reused with different canonical request bytes",
598 ));
599 }
600
601 validate_source(&source)?;
602 let mut participants = source
603 .participant_snapshots()?
604 .into_iter()
605 .filter(|snapshot| {
606 !(snapshot.capability_id == crate::WORKSPACE_CAPABILITY_ID
607 && snapshot.record_family_id == RESTORATION_FAMILY)
608 })
609 .map(snapshot_to_participant)
610 .collect::<Result<Vec<_>, GfError>>()?;
611 participants.push(restoration_participant(
612 restoration_uuid,
613 checkpoint.checkpoint_uuid,
614 source.generation_uuid(),
615 source.manifest_sha256(),
616 parse_uuid(&expected_extension.prior_current_generation_uuid)?,
617 generation_uuid,
618 request.operation_uuid,
619 request.actor_uuid,
620 &requested_reason,
621 restored_at,
622 )?);
623 let capabilities = source
624 .capabilities()
625 .into_iter()
626 .map(|value| ProjectCapability {
627 capability_id: value.capability_id,
628 capability_version: value.capability_version,
629 })
630 .collect();
631 let publication = ProjectGenerationRequest {
632 transaction_uuid,
633 generation_uuid,
634 capabilities,
635 participants,
636 };
637 let expected_parent_uuid = prior_current.generation_uuid();
638 let expected_participants = publication
639 .participants
640 .iter()
641 .map(|row| {
642 (
643 row.capability_id.clone(),
644 row.record_family_id.clone(),
645 row.record_version,
646 row.row_count,
647 )
648 })
649 .collect::<BTreeSet<_>>();
650 let writer = locks.transfer_writer_for_revert_publication();
651 let receipt = match stage_project_generation_with_lock(
652 root.clone(),
653 writer,
654 prior_current,
655 &publication,
656 Some(expected_extension),
657 )? {
658 ProjectStageOutcome::AlreadyPublished(receipt) => receipt,
659 ProjectStageOutcome::Staged(staged) => {
660 staged
661 .validate(
662 |rows| {
663 let actual = rows
664 .iter()
665 .map(|row| {
666 (
667 row.capability_id.clone(),
668 row.record_family_id.clone(),
669 row.record_version,
670 row.row_count,
671 )
672 })
673 .collect::<BTreeSet<_>>();
674 if rows.len() != expected_participants.len()
675 || actual != expected_participants
676 || rows.iter().filter(|row| {
677 row.capability_id == crate::WORKSPACE_CAPABILITY_ID
678 && row.record_family_id == RESTORATION_FAMILY
679 && row.encoding == "parquet"
680 && row.record_version == RESTORATION_CONTRACT_VERSION
681 && row.row_count == 1
682 }).count() != 1
683 {
684 return Err(GfError::Validation(
685 "staged revert participant inventory differs from the validated complete snapshot"
686 .into(),
687 ));
688 }
689 Ok(())
690 },
691 |parent, _| {
692 if parent.generation_uuid() != expected_parent_uuid {
693 return Err(GfError::Validation(
694 "staged revert parent changed after composite validation".into(),
695 ));
696 }
697 Ok(())
698 },
699 )?
700 .publish()?
701 }
702 };
703 let resolved = resolve_verified_generation(
704 &root,
705 receipt.generation_uuid,
706 receipt.generation_manifest_sha256,
707 )?;
708 Ok((
709 CheckpointReceipt {
710 operation: "revert_to_checkpoint",
711 operation_uuid: request.operation_uuid,
712 checkpoint_uuid: checkpoint.checkpoint_uuid,
713 name: requested_name,
714 source_generation_uuid: source.generation_uuid(),
715 prior_current_generation_uuid: Some(original_prior_uuid),
716 result_generation_uuid: Some(receipt.generation_uuid),
717 registry_revision,
718 committed_at: restored_at,
719 },
720 resolved,
721 ))
722}
723
724fn release_revert_replay_locks(checkpoint: &File, writer: &File) -> Result<(), GfError> {
725 let checkpoint_unlock = FileExt::unlock(checkpoint);
726 let writer_unlock = FileExt::unlock(writer);
727 finish_revert_replay_lock_handoff(checkpoint_unlock, writer_unlock)
728}
729
730fn finish_revert_replay_lock_handoff(
731 checkpoint_unlock: std::io::Result<()>,
732 writer_unlock: std::io::Result<()>,
733) -> Result<(), GfError> {
734 checkpoint_unlock.map_err(|error| {
735 GfError::Storage(format!(
736 "checkpoint revert replay lock handoff failed at checkpoints.lock: {error}"
737 ))
738 })?;
739 writer_unlock.map_err(|error| {
740 GfError::Storage(format!(
741 "checkpoint revert replay lock handoff failed at writer.lock: {error}"
742 ))
743 })
744}
745
746fn validate_revert_replay_request(
747 request: &CheckpointRevertRequest,
748 name: &str,
749 reason: &str,
750 extension: &RevertJournalExtension,
751) -> Result<(), GfError> {
752 if extension.operation_uuid != request.operation_uuid.to_string()
753 || extension.checkpoint_name != name
754 || extension.reason != reason
755 || extension.actor_uuid != request.actor_uuid.map(|value| value.to_string())
756 {
757 return Err(project_error(
758 ProjectErrorCode::TransactionConflict,
759 "revert operation UUID was reused with different canonical request bytes",
760 ));
761 }
762 Ok(())
763}
764
765fn revert_receipt(
766 request: &CheckpointRevertRequest,
767 name: &str,
768 extension: &RevertJournalExtension,
769 result_generation_uuid: Uuid,
770) -> Result<CheckpointReceipt, GfError> {
771 Ok(CheckpointReceipt {
772 operation: "revert_to_checkpoint",
773 operation_uuid: request.operation_uuid,
774 checkpoint_uuid: parse_uuid(&extension.checkpoint_uuid)?,
775 name: name.to_owned(),
776 source_generation_uuid: parse_uuid(&extension.source_generation_uuid)?,
777 prior_current_generation_uuid: Some(parse_uuid(&extension.prior_current_generation_uuid)?),
778 result_generation_uuid: Some(result_generation_uuid),
779 registry_revision: extension.registry_revision,
780 committed_at: extension.restored_at,
781 })
782}
783
784pub(crate) struct CheckpointRetentionRoots {
785 _checkpoint_lock: File,
786 pub(crate) roots: Vec<(Uuid, [u8; 32])>,
787}
788
789pub(crate) fn checkpoint_retention_roots_after_writer_lock(
790 root: &Path,
791) -> Result<CheckpointRetentionRoots, GfError> {
792 let lock_root = ensure_machine_directory(root, Path::new(LOCKS_DIR))?;
793 let checkpoint_lock = open_regular_lock(&lock_root.join(CHECKPOINT_LOCK_FILE))?;
794 if !FileExt::try_lock_exclusive(&checkpoint_lock).map_err(storage_io)? {
795 return Err(project_error(
796 ProjectErrorCode::WriterBusy,
797 "recovery could not acquire checkpoints.lock after writer.lock",
798 ));
799 }
800 let checkpoint_root = checkpoint_root(root)?;
801 recover_pair(&checkpoint_root)?;
802 let registry = read_registry(&checkpoint_root)?;
803 let roots = registry
804 .active
805 .into_iter()
806 .map(|row| {
807 let digest = decode_digest(&row.generation_manifest_sha256)?;
808 Ok((row.generation_uuid, digest))
809 })
810 .collect::<Result<Vec<_>, GfError>>()?;
811 Ok(CheckpointRetentionRoots {
812 _checkpoint_lock: checkpoint_lock,
813 roots,
814 })
815}
816
817fn create_receipt(row: &CheckpointRecord) -> CheckpointReceipt {
818 CheckpointReceipt {
819 operation: "checkpoint",
820 operation_uuid: row.create_operation_uuid,
821 checkpoint_uuid: row.checkpoint_uuid,
822 name: row.name.clone(),
823 source_generation_uuid: row.generation_uuid,
824 prior_current_generation_uuid: None,
825 result_generation_uuid: None,
826 registry_revision: row.created_revision,
827 committed_at: row.created_at,
828 }
829}
830
831fn create_tombstone_receipt(row: &CheckpointTombstone) -> CheckpointReceipt {
832 CheckpointReceipt {
833 operation: "checkpoint",
834 operation_uuid: row.create_operation_uuid,
835 checkpoint_uuid: row.checkpoint_uuid,
836 name: row.name.clone(),
837 source_generation_uuid: row.generation_uuid,
838 prior_current_generation_uuid: None,
839 result_generation_uuid: None,
840 registry_revision: row.created_revision,
841 committed_at: row.created_at,
842 }
843}
844
845fn delete_receipt(row: &CheckpointTombstone) -> CheckpointReceipt {
846 CheckpointReceipt {
847 operation: "delete_checkpoint",
848 operation_uuid: row.delete_operation_uuid,
849 checkpoint_uuid: row.checkpoint_uuid,
850 name: row.name.clone(),
851 source_generation_uuid: row.generation_uuid,
852 prior_current_generation_uuid: None,
853 result_generation_uuid: None,
854 registry_revision: row.deleted_revision,
855 committed_at: row.deleted_at,
856 }
857}
858
859fn canonical_project_root(path: &Path) -> Result<PathBuf, GfError> {
860 let metadata = fs::symlink_metadata(path).map_err(storage_io)?;
861 if metadata.file_type().is_symlink() || !metadata.is_dir() {
862 return Err(project_error(
863 ProjectErrorCode::UnsupportedProjectFormat,
864 "project root must be a real local directory, not a link",
865 ));
866 }
867 std::fs::canonicalize(path).map_err(storage_io)
868}
869
870fn checkpoint_root(root: &Path) -> Result<PathBuf, GfError> {
871 ensure_machine_directory(root, Path::new(CHECKPOINTS_DIR))
872}
873
874fn acquire_mutation_locks(root: &Path) -> Result<MutationLocks, GfError> {
875 let lock_root = ensure_machine_directory(root, Path::new(LOCKS_DIR))?;
876 sync_directory(root)?;
877 let writer = open_regular_lock(&lock_root.join(WRITER_LOCK_FILE))?;
878 if !FileExt::try_lock_exclusive(&writer).map_err(storage_io)? {
879 return Err(project_error(
880 ProjectErrorCode::WriterBusy,
881 "checkpoint mutation could not acquire writer.lock",
882 ));
883 }
884 let checkpoint = open_regular_lock(&lock_root.join(CHECKPOINT_LOCK_FILE))?;
885 if !FileExt::try_lock_exclusive(&checkpoint).map_err(storage_io)? {
886 return Err(project_error(
887 ProjectErrorCode::WriterBusy,
888 "checkpoint mutation could not acquire checkpoints.lock",
889 ));
890 }
891 Ok(MutationLocks {
892 writer: Some(writer),
893 checkpoint: Some(checkpoint),
894 })
895}
896
897fn acquire_checkpoint_read_lock(root: &Path) -> Result<File, GfError> {
898 let lock_root = ensure_machine_directory(root, Path::new(LOCKS_DIR))?;
899 let checkpoint = open_regular_lock(&lock_root.join(CHECKPOINT_LOCK_FILE))?;
900 if !FileExt::try_lock_shared(&checkpoint).map_err(storage_io)? {
901 return Err(project_error(
902 ProjectErrorCode::WriterBusy,
903 "checkpoint read could not acquire checkpoints.lock",
904 ));
905 }
906 Ok(checkpoint)
907}
908
909fn read_registry_for_read(
910 root: &Path,
911 checkpoint_root: &Path,
912) -> Result<(File, Registry), GfError> {
913 let checkpoint = acquire_checkpoint_read_lock(root)?;
914 if !checkpoint_root.join(INTENT_FILE).exists() {
915 return read_registry(checkpoint_root).map(|registry| (checkpoint, registry));
916 }
917 drop(checkpoint);
918 {
919 let _locks = acquire_mutation_locks(root)?;
920 recover_pair(checkpoint_root)?;
921 }
922 let checkpoint = acquire_checkpoint_read_lock(root)?;
923 let registry = read_registry(checkpoint_root)?;
924 Ok((checkpoint, registry))
925}
926
927fn read_registry(root: &Path) -> Result<Registry, GfError> {
928 let registry_path = root.join(REGISTRY_FILE);
929 let checksum_path = root.join(CHECKSUM_FILE);
930 if !registry_path.exists() && !checksum_path.exists() {
931 return Ok(Registry::empty());
932 }
933 let bytes = read_regular_bounded(®istry_path, MAX_REGISTRY_BYTES)?;
934 let checksum = read_regular_bounded(&checksum_path, 128)?;
935 let expected = format!("{}\n", hex(&Sha256::digest(&bytes).into()));
936 if checksum != expected.as_bytes() {
937 return Err(registry_corrupt(
938 "registry checksum does not match exact bytes",
939 ));
940 }
941 let registry: Registry = serde_json::from_slice(&bytes)
942 .map_err(|_| registry_corrupt("registry JSON is malformed"))?;
943 if registry.canonical_bytes()? != bytes {
944 return Err(registry_corrupt("registry JSON is noncanonical"));
945 }
946 Ok(registry)
947}
948
949fn commit_registry(
950 root: &Path,
951 registry: &Registry,
952 transaction_uuid: Uuid,
953) -> Result<(), GfError> {
954 let next = registry.canonical_bytes()?;
955 let next_digest = hex(&Sha256::digest(&next).into());
956 let previous = read_valid_pair(root)?;
957 let registry_temp = format!(".registry.{transaction_uuid}.json.next");
958 let checksum_temp = format!(".registry.{transaction_uuid}.sha256.next");
959 prepare_temp_path(&root.join(®istry_temp))?;
960 prepare_temp_path(&root.join(&checksum_temp))?;
961 write_new_synced(&root.join(®istry_temp), &next)?;
962 write_new_synced(
963 &root.join(&checksum_temp),
964 format!("{next_digest}\n").as_bytes(),
965 )?;
966 sync_directory(root)?;
967 project_failpoint::hit(
968 "checkpoint.registry.after_file_fsync",
969 Some(transaction_uuid),
970 None,
971 "REGISTRY_STAGED",
972 false,
973 )?;
974 let intent = RegistryIntent {
975 transaction_uuid,
976 previous_revision: previous.as_ref().map(|(registry, _)| registry.revision),
977 previous_sha256: previous.as_ref().map(|(_, digest)| digest.clone()),
978 next_revision: registry.revision,
979 next_sha256: next_digest,
980 registry_temp: registry_temp.clone(),
981 checksum_temp: checksum_temp.clone(),
982 };
983 write_intent(root, &intent)?;
984 project_failpoint::hit(
985 "checkpoint.registry.before_replace",
986 Some(transaction_uuid),
987 None,
988 "REGISTRY_INTENT_DURABLE",
989 false,
990 )?;
991 fs::rename(root.join(®istry_temp), root.join(REGISTRY_FILE)).map_err(storage_io)?;
992 project_failpoint::hit(
993 "checkpoint.registry.after_replace",
994 Some(transaction_uuid),
995 None,
996 "REGISTRY_REPLACED",
997 true,
998 )?;
999 fs::rename(root.join(&checksum_temp), root.join(CHECKSUM_FILE)).map_err(storage_io)?;
1000 sync_directory(root)?;
1001 project_failpoint::hit(
1002 "checkpoint.registry.after_dir_fsync",
1003 Some(transaction_uuid),
1004 None,
1005 "REGISTRY_DURABLE",
1006 true,
1007 )?;
1008 fs::remove_file(root.join(INTENT_FILE)).map_err(storage_io)?;
1009 sync_directory(root)
1010}
1011
1012fn recover_pair(root: &Path) -> Result<(), GfError> {
1013 let intent_path = root.join(INTENT_FILE);
1014 if !intent_path.exists() {
1015 read_registry(root)?;
1016 return Ok(());
1017 }
1018 let bytes = read_regular_bounded(&intent_path, 16 * 1024)?;
1019 let intent: RegistryIntent = serde_json::from_slice(&bytes)
1020 .map_err(|_| registry_corrupt("registry intent is malformed"))?;
1021 let mut canonical = serde_json::to_vec(&intent).map_err(registry_serde)?;
1022 canonical.push(b'\n');
1023 if canonical != bytes
1024 || !valid_private_name(&intent.registry_temp, intent.transaction_uuid, "json")
1025 || !valid_private_name(&intent.checksum_temp, intent.transaction_uuid, "sha256")
1026 {
1027 return Err(registry_corrupt(
1028 "registry intent is noncanonical or names unsafe files",
1029 ));
1030 }
1031 if let Ok(Some((current, digest))) = read_valid_pair(root) {
1032 if current.revision == intent.next_revision && digest == intent.next_sha256 {
1033 cleanup_intent(root, &intent)?;
1034 return Ok(());
1035 }
1036 if Some(current.revision) == intent.previous_revision
1037 && Some(digest) == intent.previous_sha256
1038 {
1039 validate_staged_pair(root, &intent)?;
1040 cleanup_intent(root, &intent)?;
1041 return Ok(());
1042 }
1043 }
1044 if intent.previous_revision.is_none()
1045 && !root.join(REGISTRY_FILE).exists()
1046 && !root.join(CHECKSUM_FILE).exists()
1047 {
1048 validate_staged_pair(root, &intent)?;
1049 cleanup_intent(root, &intent)?;
1050 return Ok(());
1051 }
1052 let registry_bytes = read_regular_bounded(&root.join(REGISTRY_FILE), MAX_REGISTRY_BYTES)?;
1053 if hex(&Sha256::digest(®istry_bytes).into()) == intent.next_sha256 {
1054 let checksum_bytes = read_regular_bounded(&root.join(&intent.checksum_temp), 128)?;
1055 if checksum_bytes == format!("{}\n", intent.next_sha256).as_bytes() {
1056 fs::rename(root.join(&intent.checksum_temp), root.join(CHECKSUM_FILE))
1057 .map_err(storage_io)?;
1058 sync_directory(root)?;
1059 cleanup_intent(root, &intent)?;
1060 read_registry(root)?;
1061 return Ok(());
1062 }
1063 }
1064 Err(registry_corrupt(
1065 "registry transaction is not a validated previous or staged next state",
1066 ))
1067}
1068
1069fn validate_staged_pair(root: &Path, intent: &RegistryIntent) -> Result<(), GfError> {
1070 let registry_bytes =
1071 read_regular_bounded(&root.join(&intent.registry_temp), MAX_REGISTRY_BYTES)?;
1072 let checksum_bytes = read_regular_bounded(&root.join(&intent.checksum_temp), 128)?;
1073 let digest = hex(&Sha256::digest(®istry_bytes).into());
1074 if digest != intent.next_sha256
1075 || checksum_bytes != format!("{}\n", intent.next_sha256).as_bytes()
1076 {
1077 return Err(registry_corrupt(
1078 "registry staged pair does not match its durable intent",
1079 ));
1080 }
1081 let registry: Registry = serde_json::from_slice(®istry_bytes)
1082 .map_err(|_| registry_corrupt("staged checkpoint registry is malformed"))?;
1083 if registry.canonical_bytes()? != registry_bytes || registry.revision != intent.next_revision {
1084 return Err(registry_corrupt(
1085 "registry staged pair is noncanonical or has the wrong revision",
1086 ));
1087 }
1088 Ok(())
1089}
1090
1091fn read_valid_pair(root: &Path) -> Result<Option<(Registry, String)>, GfError> {
1092 if !root.join(REGISTRY_FILE).exists() && !root.join(CHECKSUM_FILE).exists() {
1093 return Ok(None);
1094 }
1095 let registry = read_registry(root)?;
1096 let digest = hex(&Sha256::digest(registry.canonical_bytes()?).into());
1097 Ok(Some((registry, digest)))
1098}
1099
1100fn cleanup_intent(root: &Path, intent: &RegistryIntent) -> Result<(), GfError> {
1101 for name in [&intent.registry_temp, &intent.checksum_temp] {
1102 let path = root.join(name);
1103 if path.exists() {
1104 validate_single_link_regular(&path, "registry transaction temporary file")?;
1105 }
1106 match fs::remove_file(path) {
1107 Ok(()) => {}
1108 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
1109 Err(error) => return Err(storage_io(error)),
1110 }
1111 }
1112 fs::remove_file(root.join(INTENT_FILE)).map_err(storage_io)?;
1113 sync_directory(root)
1114}
1115
1116fn write_intent(root: &Path, intent: &RegistryIntent) -> Result<(), GfError> {
1117 let temp = root.join(format!(".registry.{}.txn.next", intent.transaction_uuid));
1118 let mut bytes = serde_json::to_vec(intent).map_err(registry_serde)?;
1119 bytes.push(b'\n');
1120 prepare_temp_path(&temp)?;
1121 write_new_synced(&temp, &bytes)?;
1122 project_failpoint::hit(
1123 "checkpoint.registry.after_intent_file_fsync",
1124 Some(intent.transaction_uuid),
1125 None,
1126 "REGISTRY_INTENT_STAGED",
1127 false,
1128 )?;
1129 fs::rename(temp, root.join(INTENT_FILE)).map_err(storage_io)?;
1130 sync_directory(root)
1131}
1132
1133fn write_new_synced(path: &Path, bytes: &[u8]) -> Result<(), GfError> {
1134 let mut options = OpenOptions::new();
1135 options.write(true).create_new(true);
1136 #[cfg(unix)]
1137 {
1138 use std::os::unix::fs::OpenOptionsExt;
1139 options.mode(0o600);
1140 }
1141 let mut file = options.open(path).map_err(storage_io)?;
1142 file.write_all(bytes).map_err(storage_io)?;
1143 file.sync_all().map_err(storage_io)
1144}
1145
1146fn prepare_temp_path(path: &Path) -> Result<(), GfError> {
1147 if !path.exists() {
1148 return Ok(());
1149 }
1150 let metadata = fs::symlink_metadata(path).map_err(storage_io)?;
1151 if !metadata.file_type().is_file() {
1152 return Err(registry_corrupt(
1153 "checkpoint registry temporary path is linked or special",
1154 ));
1155 }
1156 #[cfg(unix)]
1157 {
1158 use std::os::unix::fs::MetadataExt;
1159 if metadata.nlink() != 1 {
1160 return Err(registry_corrupt(
1161 "checkpoint registry temporary path is hard-linked",
1162 ));
1163 }
1164 }
1165 fs::remove_file(path).map_err(storage_io)
1166}
1167
1168fn read_regular_bounded(path: &Path, max: u64) -> Result<Vec<u8>, GfError> {
1169 let metadata = fs::symlink_metadata(path)
1170 .map_err(|_| registry_corrupt("checkpoint registry file is missing"))?;
1171 if !metadata.file_type().is_file() || metadata.len() > max {
1172 return Err(registry_corrupt(
1173 "checkpoint registry file is linked, special, or oversized",
1174 ));
1175 }
1176 #[cfg(unix)]
1177 {
1178 use std::os::unix::fs::MetadataExt;
1179 if metadata.nlink() != 1 {
1180 return Err(registry_corrupt("checkpoint registry file is hard-linked"));
1181 }
1182 }
1183 let mut options = OpenOptions::new();
1184 options.read(true);
1185 #[cfg(unix)]
1186 {
1187 use std::os::unix::fs::OpenOptionsExt;
1188 options.custom_flags(libc::O_NOFOLLOW);
1189 }
1190 let file = options.open(path).map_err(|_| {
1191 registry_corrupt("checkpoint registry file could not be opened without following links")
1192 })?;
1193 let opened = file.metadata().map_err(storage_io)?;
1194 if !opened.is_file() || opened.len() != metadata.len() {
1195 return Err(registry_corrupt(
1196 "checkpoint registry file identity changed while opening",
1197 ));
1198 }
1199 #[cfg(unix)]
1200 {
1201 use std::os::unix::fs::MetadataExt;
1202 if opened.dev() != metadata.dev() || opened.ino() != metadata.ino() || opened.nlink() != 1 {
1203 return Err(registry_corrupt(
1204 "checkpoint registry file identity changed while opening",
1205 ));
1206 }
1207 }
1208 let capacity = usize::try_from(metadata.len())
1209 .map_err(|_| registry_corrupt("checkpoint registry file length exceeds address space"))?;
1210 let mut bytes = Vec::with_capacity(capacity);
1211 file.take(max + 1)
1212 .read_to_end(&mut bytes)
1213 .map_err(storage_io)?;
1214 if bytes.len() as u64 > max {
1215 return Err(registry_corrupt(
1216 "checkpoint registry file exceeds its read bound",
1217 ));
1218 }
1219 Ok(bytes)
1220}
1221
1222fn validate_registry(registry: &Registry) -> Result<(), GfError> {
1223 if registry.format != "graphforge-checkpoints"
1224 || registry.format_version != 1
1225 || registry.active.len() > MAX_ACTIVE
1226 || registry.tombstones.len() > MAX_TOMBSTONES
1227 {
1228 return Err(registry_corrupt(
1229 "checkpoint registry header or bounds are invalid",
1230 ));
1231 }
1232 if !registry.active.windows(2).all(|pair| {
1233 (&pair[0].name, pair[0].checkpoint_uuid) < (&pair[1].name, pair[1].checkpoint_uuid)
1234 }) {
1235 return Err(registry_corrupt(
1236 "active checkpoints are not strictly sorted",
1237 ));
1238 }
1239 if !registry.tombstones.windows(2).all(|pair| {
1240 (pair[0].deleted_revision, pair[0].checkpoint_uuid)
1241 < (pair[1].deleted_revision, pair[1].checkpoint_uuid)
1242 }) {
1243 return Err(registry_corrupt(
1244 "checkpoint tombstones are not strictly sorted",
1245 ));
1246 }
1247 let mut names = BTreeSet::new();
1248 let mut checkpoint_uuids = BTreeSet::new();
1249 let mut create_operations = BTreeSet::new();
1250 let mut delete_operations = BTreeSet::new();
1251 for row in ®istry.active {
1252 validate_checkpoint_content(&CheckpointContentRef {
1253 label: "active checkpoint",
1254 checkpoint_uuid: row.checkpoint_uuid,
1255 create_operation_uuid: row.create_operation_uuid,
1256 name: &row.name,
1257 description: row.description.as_deref(),
1258 created_by: row.created_by,
1259 generation_manifest_sha256: &row.generation_manifest_sha256,
1260 create_request_sha256: &row.create_request_sha256,
1261 })?;
1262 if row.created_revision == 0
1263 || row.created_revision > registry.revision
1264 || !names.insert(row.name.as_str())
1265 || !checkpoint_uuids.insert(row.checkpoint_uuid)
1266 || !create_operations.insert(row.create_operation_uuid)
1267 {
1268 return Err(registry_corrupt(
1269 "active checkpoint identities or revision are inconsistent",
1270 ));
1271 }
1272 }
1273 for row in ®istry.tombstones {
1274 validate_name(&row.name)
1277 .map_err(|_| registry_corrupt("checkpoint tombstone name is invalid"))?;
1278 validate_description(row.description.as_deref())
1279 .map_err(|_| registry_corrupt("checkpoint tombstone description is invalid"))?;
1280 validate_digest(&row.generation_manifest_sha256)?;
1281 validate_digest(&row.create_request_sha256)?;
1282 validate_digest(&row.delete_request_sha256)?;
1283 validate_record_identity(
1284 row.checkpoint_uuid,
1285 row.create_operation_uuid,
1286 &row.name,
1287 row.description.as_deref(),
1288 row.created_by,
1289 &row.create_request_sha256,
1290 )?;
1291 let expected_delete =
1292 delete_request_digest_values(row.delete_operation_uuid, &row.name, row.deleted_by);
1293 if row.created_revision == 0
1294 || row.created_revision >= row.deleted_revision
1295 || row.deleted_revision > registry.revision
1296 || !checkpoint_uuids.insert(row.checkpoint_uuid)
1297 || !create_operations.insert(row.create_operation_uuid)
1298 || delete_operations.contains(&row.create_operation_uuid)
1299 || !delete_operations.insert(row.delete_operation_uuid)
1300 || create_operations.contains(&row.delete_operation_uuid)
1301 || row.delete_request_sha256 != hex(&expected_delete)
1302 {
1303 return Err(registry_corrupt(
1304 "checkpoint tombstone identities or revisions are inconsistent",
1305 ));
1306 }
1307 }
1308 Ok(())
1309}
1310
1311struct CheckpointContentRef<'a> {
1312 label: &'static str,
1313 checkpoint_uuid: Uuid,
1314 create_operation_uuid: Uuid,
1315 name: &'a str,
1316 description: Option<&'a str>,
1317 created_by: Option<Uuid>,
1318 generation_manifest_sha256: &'a str,
1319 create_request_sha256: &'a str,
1320}
1321
1322fn validate_checkpoint_content(content: &CheckpointContentRef<'_>) -> Result<(), GfError> {
1323 validate_name(content.name)
1324 .map_err(|_| registry_corrupt(format!("{} name is invalid", content.label)))?;
1325 validate_description(content.description)
1326 .map_err(|_| registry_corrupt(format!("{} description is invalid", content.label)))?;
1327 validate_digest(content.generation_manifest_sha256)?;
1328 validate_digest(content.create_request_sha256)?;
1329 validate_record_identity(
1330 content.checkpoint_uuid,
1331 content.create_operation_uuid,
1332 content.name,
1333 content.description,
1334 content.created_by,
1335 content.create_request_sha256,
1336 )
1337}
1338
1339fn validate_single_link_regular(path: &Path, label: &str) -> Result<(), GfError> {
1340 let metadata = fs::symlink_metadata(path).map_err(storage_io)?;
1341 if !metadata.file_type().is_file() {
1342 return Err(registry_corrupt(format!("{label} is linked or special")));
1343 }
1344 #[cfg(unix)]
1345 {
1346 use std::os::unix::fs::MetadataExt;
1347 if metadata.nlink() != 1 {
1348 return Err(registry_corrupt(format!("{label} is hard-linked")));
1349 }
1350 }
1351 Ok(())
1352}
1353
1354fn validate_name(value: &str) -> Result<String, GfError> {
1355 let normalized: String = value.nfc().collect();
1356 if normalized != value
1357 || value.is_empty()
1358 || value.len() > MAX_NAME_BYTES
1359 || value.trim() != value
1360 || value == "."
1361 || value == ".."
1362 || value.contains(" ")
1363 || !value
1364 .chars()
1365 .all(|ch| ch.is_alphanumeric() || matches!(ch, ' ' | '_' | '-' | '.'))
1366 {
1367 return Err(GfError::Validation(
1368 "checkpoint name is not canonical NFC content or violates the 1-128 byte grammar"
1369 .into(),
1370 ));
1371 }
1372 Ok(normalized)
1373}
1374
1375fn validate_description(value: Option<&str>) -> Result<(), GfError> {
1376 if value.is_some_and(|value| {
1377 value.len() > MAX_DESCRIPTION_BYTES || value.chars().any(char::is_control)
1378 }) {
1379 return Err(GfError::Validation(
1380 "checkpoint description exceeds 1024 UTF-8 bytes or contains controls".into(),
1381 ));
1382 }
1383 Ok(())
1384}
1385
1386fn validate_digest(value: &str) -> Result<(), GfError> {
1387 if value.len() != 64
1388 || !value
1389 .bytes()
1390 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1391 {
1392 return Err(registry_corrupt("checkpoint digest is noncanonical"));
1393 }
1394 Ok(())
1395}
1396
1397fn decode_digest(value: &str) -> Result<[u8; 32], GfError> {
1398 validate_digest(value)?;
1399 let mut digest = [0_u8; 32];
1400 for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
1401 let text = std::str::from_utf8(pair)
1402 .map_err(|_| registry_corrupt("checkpoint digest is not UTF-8"))?;
1403 digest[index] = u8::from_str_radix(text, 16)
1404 .map_err(|_| registry_corrupt("checkpoint digest is not lowercase hex"))?;
1405 }
1406 Ok(digest)
1407}
1408
1409fn create_request_digest(request: &CheckpointCreateRequest, name: &str) -> [u8; 32] {
1410 create_request_digest_values(
1411 request.operation_uuid,
1412 name,
1413 request.description.as_deref(),
1414 request.actor_uuid,
1415 )
1416}
1417
1418fn create_request_digest_values(
1419 operation_uuid: Uuid,
1420 name: &str,
1421 description: Option<&str>,
1422 actor_uuid: Option<Uuid>,
1423) -> [u8; 32] {
1424 let mut hasher = Sha256::new();
1425 hasher.update(b"graphforge-checkpoint-create-request/1");
1426 hasher.update(operation_uuid.as_bytes());
1427 append_bytes(&mut hasher, name.as_bytes());
1428 match description {
1429 Some(value) => {
1430 hasher.update([1]);
1431 append_bytes(&mut hasher, value.as_bytes());
1432 }
1433 None => hasher.update([0]),
1434 }
1435 append_actor(&mut hasher, actor_uuid);
1436 hasher.finalize().into()
1437}
1438
1439fn delete_request_digest(request: &CheckpointDeleteRequest, name: &str) -> [u8; 32] {
1440 delete_request_digest_values(request.operation_uuid, name, request.actor_uuid)
1441}
1442
1443fn delete_request_digest_values(
1444 operation_uuid: Uuid,
1445 name: &str,
1446 actor_uuid: Option<Uuid>,
1447) -> [u8; 32] {
1448 let mut hasher = Sha256::new();
1449 hasher.update(b"graphforge-checkpoint-delete-request/1");
1450 hasher.update(operation_uuid.as_bytes());
1451 append_bytes(&mut hasher, name.as_bytes());
1452 append_actor(&mut hasher, actor_uuid);
1453 hasher.finalize().into()
1454}
1455
1456fn validate_record_identity(
1457 checkpoint: Uuid,
1458 operation: Uuid,
1459 name: &str,
1460 description: Option<&str>,
1461 actor: Option<Uuid>,
1462 request_hex: &str,
1463) -> Result<(), GfError> {
1464 let request = create_request_digest_values(operation, name, description, actor);
1465 if request_hex != hex(&request) || checkpoint != checkpoint_uuid(operation, request) {
1466 return Err(registry_corrupt(
1467 "checkpoint deterministic identity or create request digest is inconsistent",
1468 ));
1469 }
1470 Ok(())
1471}
1472
1473fn checkpoint_uuid(operation_uuid: Uuid, request_digest: [u8; 32]) -> Uuid {
1474 let mut hasher = Sha256::new();
1475 hasher.update(b"graphforge-checkpoint-uuid/1");
1476 hasher.update(operation_uuid.as_bytes());
1477 hasher.update(request_digest);
1478 graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1479}
1480
1481fn append_bytes(hasher: &mut Sha256, bytes: &[u8]) {
1482 hasher.update(
1483 u32::try_from(bytes.len())
1484 .expect("validated checkpoint strings fit u32")
1485 .to_be_bytes(),
1486 );
1487 hasher.update(bytes);
1488}
1489fn append_actor(hasher: &mut Sha256, actor: Option<Uuid>) {
1490 match actor {
1491 Some(value) => {
1492 hasher.update([1]);
1493 hasher.update(value.as_bytes());
1494 }
1495 None => hasher.update([0]),
1496 }
1497}
1498fn valid_private_name(name: &str, uuid: Uuid, kind: &str) -> bool {
1499 name == format!(".registry.{uuid}.{kind}.next")
1500}
1501fn hex(bytes: &[u8; 32]) -> String {
1502 let mut output = String::with_capacity(64);
1503 for byte in bytes {
1504 write!(&mut output, "{byte:02x}").expect("writing hexadecimal to String cannot fail");
1505 }
1506 output
1507}
1508
1509fn parse_uuid(value: &str) -> Result<Uuid, GfError> {
1510 Uuid::parse_str(value).map_err(|_| registry_corrupt("revert journal UUID is invalid"))
1511}
1512
1513fn validate_reason(value: &str) -> Result<String, GfError> {
1514 let trimmed = value.trim();
1515 if trimmed.is_empty() || trimmed.len() > MAX_REASON_BYTES {
1516 return Err(GfError::Validation(
1517 "checkpoint revert reason must contain 1..=1024 UTF-8 bytes after trimming".into(),
1518 ));
1519 }
1520 Ok(trimmed.to_owned())
1521}
1522
1523fn revert_request_digest(
1524 operation_uuid: Uuid,
1525 name: &str,
1526 checkpoint_uuid: Uuid,
1527 source_generation_uuid: Uuid,
1528 source_manifest_sha256: [u8; 32],
1529 reason: &str,
1530 actor_uuid: Option<Uuid>,
1531) -> [u8; 32] {
1532 let mut hasher = Sha256::new();
1533 hasher.update(b"graphforge-checkpoint-revert-request/1");
1534 hasher.update(operation_uuid.as_bytes());
1535 append_bytes(&mut hasher, name.as_bytes());
1536 hasher.update(checkpoint_uuid.as_bytes());
1537 hasher.update(source_generation_uuid.as_bytes());
1538 hasher.update(source_manifest_sha256);
1539 append_bytes(&mut hasher, reason.as_bytes());
1540 append_actor(&mut hasher, actor_uuid);
1541 hasher.finalize().into()
1542}
1543
1544fn revert_transaction_uuid(operation_uuid: Uuid) -> Uuid {
1545 let mut hasher = Sha256::new();
1546 hasher.update(b"graphforge-checkpoint-revert-transaction/1");
1547 hasher.update(operation_uuid.as_bytes());
1548 graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1549}
1550
1551fn restoration_uuid(operation_uuid: Uuid, request_digest: [u8; 32]) -> Uuid {
1552 let mut hasher = Sha256::new();
1553 hasher.update(b"graphforge-restoration-transition-uuid/1");
1554 hasher.update(operation_uuid.as_bytes());
1555 hasher.update(request_digest);
1556 graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1557}
1558
1559fn restored_generation_uuid(
1560 transaction_uuid: Uuid,
1561 checkpoint_uuid: Uuid,
1562 source_generation_uuid: Uuid,
1563 source_manifest_sha256: [u8; 32],
1564 prior_current_generation_uuid: Uuid,
1565 restored_at: i64,
1566 request_digest: [u8; 32],
1567) -> Uuid {
1568 let mut hasher = Sha256::new();
1569 hasher.update(b"graphforge-checkpoint-restored-generation/1");
1570 hasher.update(transaction_uuid.as_bytes());
1571 hasher.update(checkpoint_uuid.as_bytes());
1572 hasher.update(source_generation_uuid.as_bytes());
1573 hasher.update(source_manifest_sha256);
1574 hasher.update(prior_current_generation_uuid.as_bytes());
1575 hasher.update(restored_at.to_be_bytes());
1576 hasher.update(request_digest);
1577 graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1578}
1579
1580fn snapshot_to_participant(
1581 snapshot: crate::ProjectParticipantSnapshot,
1582) -> Result<ProjectParticipant, GfError> {
1583 let encoding = match snapshot.encoding.as_str() {
1584 "parquet" => ProjectParticipantEncoding::Parquet,
1585 "arrow" => ProjectParticipantEncoding::Arrow,
1586 "json" => ProjectParticipantEncoding::Json,
1587 _ => {
1588 return Err(registry_corrupt(
1589 "checkpoint participant encoding is unsupported",
1590 ));
1591 }
1592 };
1593 Ok(ProjectParticipant {
1594 capability_id: snapshot.capability_id,
1595 capability_version: snapshot.capability_version,
1596 record_family_id: snapshot.record_family_id,
1597 record_version: snapshot.record_version,
1598 encoding,
1599 schema_fingerprint: snapshot.schema_fingerprint,
1600 row_count: snapshot.row_count,
1601 bytes: snapshot.bytes,
1602 })
1603}
1604
1605#[allow(clippy::too_many_arguments)]
1606fn restoration_participant(
1607 restoration_uuid: Uuid,
1608 checkpoint_uuid: Uuid,
1609 source_generation_uuid: Uuid,
1610 source_manifest_sha256: [u8; 32],
1611 prior_current_generation_uuid: Uuid,
1612 restored_generation_uuid: Uuid,
1613 operation_uuid: Uuid,
1614 actor_uuid: Option<Uuid>,
1615 reason: &str,
1616 restored_at: i64,
1617) -> Result<ProjectParticipant, GfError> {
1618 let schema = Arc::new(Schema::new(vec![
1619 uuid_field("restoration_uuid", false),
1620 uuid_field("checkpoint_uuid", false),
1621 uuid_field("source_generation_uuid", false),
1622 Field::new(
1623 "source_manifest_sha256",
1624 DataType::FixedSizeBinary(32),
1625 false,
1626 ),
1627 uuid_field("prior_current_generation_uuid", false),
1628 uuid_field("restored_generation_uuid", false),
1629 uuid_field("operation_uuid", false),
1630 uuid_field("actor_uuid", true),
1631 Field::new("reason", DataType::Utf8, false),
1632 Field::new(
1633 "restored_at",
1634 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
1635 false,
1636 ),
1637 Field::new("contract_version", DataType::UInt32, false),
1638 ]));
1639 let mut columns = Vec::<ArrayRef>::new();
1640 for value in [
1641 Some(restoration_uuid),
1642 Some(checkpoint_uuid),
1643 Some(source_generation_uuid),
1644 Some(prior_current_generation_uuid),
1645 Some(restored_generation_uuid),
1646 Some(operation_uuid),
1647 actor_uuid,
1648 ] {
1649 let mut builder = FixedSizeBinaryBuilder::with_capacity(1, 16);
1650 match value {
1651 Some(uuid) => builder.append_value(uuid.as_bytes()).map_err(arrow_error)?,
1652 None => builder.append_null(),
1653 }
1654 columns.push(Arc::new(builder.finish()));
1655 }
1656 let mut source_digest = FixedSizeBinaryBuilder::with_capacity(1, 32);
1657 source_digest
1658 .append_value(source_manifest_sha256)
1659 .map_err(arrow_error)?;
1660 columns.insert(3, Arc::new(source_digest.finish()));
1661 columns.push(Arc::new(StringArray::from(vec![reason])));
1662 columns.push(Arc::new(
1663 TimestampMicrosecondArray::from(vec![restored_at]).with_timezone("UTC"),
1664 ));
1665 columns.push(Arc::new(UInt32Array::from(vec![
1666 RESTORATION_CONTRACT_VERSION,
1667 ])));
1668 let batch = RecordBatch::try_new(Arc::clone(&schema), columns).map_err(arrow_error)?;
1669 let properties = WriterProperties::builder()
1670 .set_created_by("graphforge-restoration-transition/1".into())
1671 .build();
1672 let mut writer =
1673 ArrowWriter::try_new(Vec::new(), schema, Some(properties)).map_err(parquet_error)?;
1674 writer.write(&batch).map_err(parquet_error)?;
1675 let bytes = writer.into_inner().map_err(parquet_error)?;
1676 let schema_fingerprint = fingerprint(
1677 CanonicalDomain::Schema,
1678 CANONICAL_CONTRACT_VERSION,
1679 b"restoration_transition/1|restoration_uuid:fixed[16]:required|checkpoint_uuid:fixed[16]:required|source_generation_uuid:fixed[16]:required|source_manifest_sha256:fixed[32]:required|prior_current_generation_uuid:fixed[16]:required|restored_generation_uuid:fixed[16]:required|operation_uuid:fixed[16]:required|actor_uuid:fixed[16]:optional|reason:utf8:required|restored_at:timestamp_us_utc:required|contract_version:u32:required",
1680 )
1681 .map_err(|error| GfError::Validation(error.to_string()))?;
1682 Ok(ProjectParticipant {
1683 capability_id: crate::WORKSPACE_CAPABILITY_ID.into(),
1684 capability_version: crate::WORKSPACE_CAPABILITY_VERSION,
1685 record_family_id: RESTORATION_FAMILY.into(),
1686 record_version: RESTORATION_CONTRACT_VERSION,
1687 encoding: ProjectParticipantEncoding::Parquet,
1688 schema_fingerprint,
1689 row_count: 1,
1690 bytes,
1691 })
1692}
1693
1694fn uuid_field(name: &str, nullable: bool) -> Field {
1695 Field::new(name, DataType::FixedSizeBinary(16), nullable)
1696}
1697
1698fn arrow_error(error: arrow::error::ArrowError) -> GfError {
1699 let message = format!("restoration Arrow encoding failed: {error}");
1700 drop(error);
1701 GfError::Storage(message)
1702}
1703
1704fn parquet_error(error: parquet::errors::ParquetError) -> GfError {
1705 let message = format!("restoration Parquet encoding failed: {error}");
1706 drop(error);
1707 GfError::Storage(message)
1708}
1709fn utc_micros() -> Result<i64, GfError> {
1710 let value = SystemTime::now()
1711 .duration_since(UNIX_EPOCH)
1712 .map_err(|_| GfError::Storage("system clock is before Unix epoch".into()))?
1713 .as_micros();
1714 i64::try_from(value).map_err(|_| GfError::Storage("UTC microsecond timestamp overflow".into()))
1715}
1716fn registry_serde(error: impl std::fmt::Display) -> GfError {
1717 GfError::Storage(format!("checkpoint registry encoding failed: {error}"))
1718}
1719fn registry_corrupt(message: impl Into<String>) -> GfError {
1720 project_error(ProjectErrorCode::CheckpointRegistryCorrupt, message)
1721}
1722fn project_error(code: ProjectErrorCode, message: impl Into<String>) -> GfError {
1723 GfError::Project {
1724 code,
1725 message: message.into(),
1726 }
1727}
1728fn storage_io(error: impl std::fmt::Display) -> GfError {
1729 GfError::Storage(format!("checkpoint registry I/O failed: {error}"))
1730}
1731
1732#[cfg(test)]
1733mod tests {
1734 use super::*;
1735 use std::collections::BTreeMap;
1736 use std::io::{BufRead, BufReader};
1737 use std::panic::{AssertUnwindSafe, catch_unwind, resume_unwind};
1738 use std::process::{Command, Stdio};
1739 use std::sync::mpsc;
1740 use std::time::Duration;
1741 use tempfile::tempdir;
1742 use wait_timeout::ChildExt;
1743
1744 const TEST_DEADLINE: Duration = Duration::from_secs(1);
1745 const CHILD_DEADLINE: Duration = Duration::from_secs(10);
1746
1747 struct WriterLockHolder {
1748 release: Option<mpsc::SyncSender<()>>,
1749 worker: Option<std::thread::JoinHandle<()>>,
1750 }
1751
1752 impl WriterLockHolder {
1753 fn finish(mut self) -> Result<(), String> {
1754 let release = self
1755 .release
1756 .take()
1757 .ok_or_else(|| "phase=main release sender missing".to_owned())?;
1758 let release_result = release
1759 .send(())
1760 .map_err(|error| format!("phase=main release holder error={error}"));
1761 let join_result = self
1762 .worker
1763 .take()
1764 .ok_or_else(|| "phase=main holder worker missing".to_owned())?
1765 .join()
1766 .map_err(|_| "phase=main holder worker panicked".to_owned());
1767 release_result.and(join_result)
1768 }
1769 }
1770
1771 impl Drop for WriterLockHolder {
1772 fn drop(&mut self) {
1773 if let Some(release) = self.release.take() {
1774 let _ = release.send(());
1775 }
1776 if let Some(worker) = self.worker.take() {
1777 let _ = worker.join();
1778 }
1779 }
1780 }
1781
1782 fn while_writer_lock_is_held<T>(root: &Path, action: impl FnOnce() -> T) -> T {
1783 let writer_path = root.join(LOCKS_DIR).join(WRITER_LOCK_FILE);
1784 let worker_path = writer_path.clone();
1785 let (ready_sender, ready_receiver) = mpsc::sync_channel(0);
1786 let (release_sender, release_receiver) = mpsc::sync_channel(0);
1787 let worker = std::thread::Builder::new()
1788 .name("checkpoint-writer-lock-holder".into())
1789 .spawn(move || {
1790 let writer =
1791 open_regular_lock(&worker_path).expect("phase=holder open writer.lock");
1792 assert!(
1793 FileExt::try_lock_exclusive(&writer).expect("phase=holder acquire writer.lock"),
1794 "phase=holder writer.lock unexpectedly busy"
1795 );
1796 ready_sender.send(()).expect("phase=holder publish ready");
1797 release_receiver.recv().expect("phase=holder await release");
1798 FileExt::unlock(&writer).expect("phase=holder release writer.lock");
1799 })
1800 .expect("phase=holder spawn");
1801 let holder = WriterLockHolder {
1802 release: Some(release_sender),
1803 worker: Some(worker),
1804 };
1805 if let Err(error) = ready_receiver.recv_timeout(TEST_DEADLINE) {
1806 drop(ready_receiver);
1807 let cleanup = holder.finish();
1808 panic!("phase=main await held writer.lock error={error}; cleanup={cleanup:?}");
1809 }
1810 let result = catch_unwind(AssertUnwindSafe(action));
1811 let cleanup = holder.finish();
1812 match result {
1813 Ok(value) => {
1814 cleanup.unwrap_or_else(|error| panic!("phase=main holder cleanup error={error}"));
1815 value
1816 }
1817 Err(original) => {
1818 let _ = cleanup;
1819 resume_unwind(original);
1820 }
1821 }
1822 }
1823
1824 struct BoundedChild {
1825 child: std::process::Child,
1826 reaped: bool,
1827 }
1828
1829 impl BoundedChild {
1830 fn wait(mut self, phase: &str) -> std::process::ExitStatus {
1831 let mut failures = Vec::new();
1832 match self.child.wait_timeout(CHILD_DEADLINE) {
1833 Ok(Some(status)) => {
1834 self.reaped = true;
1835 return status;
1836 }
1837 Ok(None) => failures.push(format!("wait timeout={CHILD_DEADLINE:?}")),
1838 Err(error) => failures.push(format!("wait error={error}")),
1839 }
1840 if let Err(error) = self.child.kill() {
1841 failures.push(format!("kill error={error}"));
1842 }
1843 match self.child.wait_timeout(TEST_DEADLINE) {
1844 Ok(Some(status)) => {
1845 self.reaped = true;
1846 failures.push(format!("killed_status={status}"));
1847 }
1848 Ok(None) => failures.push(format!("reap timeout={TEST_DEADLINE:?}")),
1849 Err(error) => failures.push(format!("reap error={error}")),
1850 }
1851 panic!("phase={phase} child cleanup failures={failures:?}");
1852 }
1853 }
1854
1855 impl Drop for BoundedChild {
1856 fn drop(&mut self) {
1857 if !self.reaped {
1858 let mut failures = Vec::new();
1859 if let Err(error) = self.child.kill() {
1860 failures.push(format!("kill error={error}"));
1861 }
1862 match self.child.wait_timeout(TEST_DEADLINE) {
1863 Ok(Some(_)) => self.reaped = true,
1864 Ok(None) => failures.push(format!("reap timeout={TEST_DEADLINE:?}")),
1865 Err(error) => failures.push(format!("reap error={error}")),
1866 }
1867 if !failures.is_empty() {
1868 eprintln!("phase=drop child cleanup failures={failures:?}");
1869 }
1870 }
1871 }
1872 }
1873
1874 fn recover_checkpoint_pair_after_lock_handoff(root: &Path, phase: &str) {
1875 checkpoint_lock_handoff(root, phase, true);
1876 }
1877
1878 fn preserve_checkpoint_intent_after_lock_handoff(root: &Path, phase: &str) {
1879 checkpoint_lock_handoff(root, phase, false);
1880 }
1881
1882 fn checkpoint_lock_handoff(root: &Path, phase: &str, recover_durable_intent: bool) {
1883 let lock_root = root.join(LOCKS_DIR);
1884 let writer_path = lock_root.join(WRITER_LOCK_FILE);
1885 let checkpoint_path = lock_root.join(CHECKPOINT_LOCK_FILE);
1886 let checkpoint_root = root.join(CHECKPOINTS_DIR);
1887 let worker_writer_path = writer_path.clone();
1888 let worker_checkpoint_path = checkpoint_path.clone();
1889 let (sender, receiver) = mpsc::sync_channel(0);
1890 std::thread::Builder::new()
1891 .name("checkpoint-lock-handoff-recovery".into())
1892 .spawn(move || {
1893 let result = (|| {
1894 let writer = open_regular_lock(&worker_writer_path)
1895 .map_err(|error| format!("open writer.lock failed: {error}"))?;
1896 FileExt::lock_exclusive(&writer)
1897 .map_err(|error| format!("acquire writer.lock failed: {error}"))?;
1898
1899 let checkpoint = match open_regular_lock(&worker_checkpoint_path) {
1900 Ok(checkpoint) => checkpoint,
1901 Err(error) => {
1902 let writer_unlock = FileExt::unlock(&writer);
1903 return Err(format!(
1904 "open checkpoints.lock failed: {error}; \
1905 writer_unlock={writer_unlock:?}"
1906 ));
1907 }
1908 };
1909 if let Err(error) = FileExt::lock_exclusive(&checkpoint) {
1910 let writer_unlock = FileExt::unlock(&writer);
1911 return Err(format!(
1912 "acquire checkpoints.lock failed: {error}; writer_unlock={writer_unlock:?}"
1913 ));
1914 }
1915
1916 let recovery = if recover_durable_intent
1917 && checkpoint_root.join(INTENT_FILE).exists()
1918 {
1919 recover_pair(&checkpoint_root)
1920 .map_err(|error| format!("recover durable checkpoint intent failed: {error}"))
1921 } else {
1922 Ok(())
1923 };
1924 let checkpoint_unlock = FileExt::unlock(&checkpoint)
1925 .map_err(|error| format!("unlock checkpoints.lock failed: {error}"));
1926 let writer_unlock = FileExt::unlock(&writer)
1927 .map_err(|error| format!("unlock writer.lock failed: {error}"));
1928
1929 recovery?;
1930 checkpoint_unlock?;
1931 writer_unlock
1932 })();
1933 let _ = sender.send(result);
1934 })
1935 .unwrap();
1936 match receiver.recv_timeout(Duration::from_secs(1)) {
1937 Ok(Ok(())) => {}
1938 Ok(Err(error)) => panic!(
1939 "checkpoint lock handoff/recovery failed at {phase}; writer_path={}; \
1940 checkpoint_path={}: {error}",
1941 writer_path.display(),
1942 checkpoint_path.display()
1943 ),
1944 Err(error) => panic!(
1945 "checkpoint lock handoff/recovery timed out at {phase}; writer_path={}; \
1946 checkpoint_path={}; timeout=1s; channel={error}",
1947 writer_path.display(),
1948 checkpoint_path.display()
1949 ),
1950 }
1951 }
1952
1953 fn publish_clone(root: &Path) -> Uuid {
1954 let selected = crate::resolve_project_generation(root).unwrap();
1955 let capabilities = selected
1956 .capabilities()
1957 .into_iter()
1958 .map(|entry| crate::ProjectCapability {
1959 capability_id: entry.capability_id,
1960 capability_version: entry.capability_version,
1961 })
1962 .collect();
1963 let participants = selected
1964 .participant_snapshots()
1965 .unwrap()
1966 .into_iter()
1967 .map(|entry| crate::ProjectParticipant {
1968 capability_id: entry.capability_id,
1969 capability_version: entry.capability_version,
1970 record_family_id: entry.record_family_id,
1971 record_version: entry.record_version,
1972 encoding: match entry.encoding.as_str() {
1973 "arrow" => crate::ProjectParticipantEncoding::Arrow,
1974 "json" => crate::ProjectParticipantEncoding::Json,
1975 "parquet" => crate::ProjectParticipantEncoding::Parquet,
1976 other => panic!("unexpected participant encoding {other}"),
1977 },
1978 schema_fingerprint: entry.schema_fingerprint,
1979 row_count: entry.row_count,
1980 bytes: entry.bytes,
1981 })
1982 .collect();
1983 let generation_uuid = Uuid::now_v7();
1984 let request = crate::ProjectGenerationRequest {
1985 transaction_uuid: Uuid::now_v7(),
1986 generation_uuid,
1987 capabilities,
1988 participants,
1989 };
1990 let crate::ProjectStageOutcome::Staged(staged) =
1991 crate::stage_project_generation(root, &request).unwrap()
1992 else {
1993 panic!("fresh publication unexpectedly replayed");
1994 };
1995 staged
1996 .validate(|_| Ok(()), |_, _| Ok(()))
1997 .unwrap()
1998 .publish()
1999 .unwrap();
2000 generation_uuid
2001 }
2002
2003 fn create_request(operation_uuid: Uuid, name: &str) -> CheckpointCreateRequest {
2004 CheckpointCreateRequest {
2005 operation_uuid,
2006 name: name.into(),
2007 description: Some("release candidate".into()),
2008 actor_uuid: Some(Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000aa").unwrap()),
2009 }
2010 }
2011
2012 fn write_raw_registry(root: &Path, registry: &Registry) {
2013 let checkpoint_root = root.join(CHECKPOINTS_DIR);
2014 let mut bytes = serde_json::to_vec(registry).unwrap();
2015 bytes.push(b'\n');
2016 fs::write(checkpoint_root.join(REGISTRY_FILE), &bytes).unwrap();
2017 fs::write(
2018 checkpoint_root.join(CHECKSUM_FILE),
2019 format!("{}\n", hex(&Sha256::digest(&bytes).into())),
2020 )
2021 .unwrap();
2022 }
2023
2024 fn install_registry_intent(
2025 checkpoint_root: &Path,
2026 previous: Option<&Registry>,
2027 next: &Registry,
2028 ) -> RegistryIntent {
2029 fs::create_dir_all(checkpoint_root).unwrap();
2030 let transaction_uuid = Uuid::now_v7();
2031 let next_bytes = next.canonical_bytes().unwrap();
2032 let next_sha256 = hex(&Sha256::digest(&next_bytes).into());
2033 let registry_temp = format!(".registry.{transaction_uuid}.json.next");
2034 let checksum_temp = format!(".registry.{transaction_uuid}.sha256.next");
2035 fs::write(checkpoint_root.join(®istry_temp), &next_bytes).unwrap();
2036 fs::write(
2037 checkpoint_root.join(&checksum_temp),
2038 format!("{next_sha256}\n"),
2039 )
2040 .unwrap();
2041 let intent = RegistryIntent {
2042 transaction_uuid,
2043 previous_revision: previous.map(|registry| registry.revision),
2044 previous_sha256: previous
2045 .map(|registry| hex(&Sha256::digest(registry.canonical_bytes().unwrap()).into())),
2046 next_revision: next.revision,
2047 next_sha256,
2048 registry_temp,
2049 checksum_temp,
2050 };
2051 let mut intent_bytes = serde_json::to_vec(&intent).unwrap();
2052 intent_bytes.push(b'\n');
2053 fs::write(checkpoint_root.join(INTENT_FILE), intent_bytes).unwrap();
2054 intent
2055 }
2056
2057 #[test]
2058 fn revert_identity_matches_frozen_golden_vector() {
2059 let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000003").unwrap();
2060 let checkpoint = Uuid::parse_str("4084179c-38db-8b6b-9b6e-c0b0a855e002").unwrap();
2061 let source = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000b0").unwrap();
2062 let prior = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000d0").unwrap();
2063 let actor = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000aa").unwrap();
2064 let source_digest = [0x11; 32];
2065 let request_digest = revert_request_digest(
2066 operation,
2067 "Release 1.0",
2068 checkpoint,
2069 source,
2070 source_digest,
2071 "restore release candidate",
2072 Some(actor),
2073 );
2074 assert_eq!(
2075 hex(&request_digest),
2076 "dff3755629942d1189060b117cb70dc864428fcc3b28d5c2d22d3924c3690e93"
2077 );
2078 let transaction = revert_transaction_uuid(operation);
2079 assert_eq!(
2080 transaction.to_string(),
2081 "908d637b-d6e6-8508-919e-d4e708e037b2"
2082 );
2083 assert_eq!(
2084 restoration_uuid(operation, request_digest).to_string(),
2085 "9e1f160c-badb-80c9-beaa-ae580910bf8a"
2086 );
2087 assert_eq!(
2088 restored_generation_uuid(
2089 transaction,
2090 checkpoint,
2091 source,
2092 source_digest,
2093 prior,
2094 1_720_000_000_123_456,
2095 request_digest,
2096 )
2097 .to_string(),
2098 "5dc02888-2064-8892-a0e1-2c00968ba0cc"
2099 );
2100 }
2101
2102 #[test]
2103 fn revert_publishes_child_preserves_registry_and_replays_after_delete() {
2104 let directory = tempdir().unwrap();
2105 crate::open_or_initialize_project(directory.path()).unwrap();
2106 let source = crate::resolve_project_generation(directory.path()).unwrap();
2107 let created = create_checkpoint(
2108 directory.path(),
2109 &create_request(Uuid::from_u128(40), "Before"),
2110 )
2111 .unwrap();
2112 let prior_current = publish_clone(directory.path());
2113 let request = CheckpointRevertRequest {
2114 operation_uuid: Uuid::from_u128(41),
2115 name: "Before".into(),
2116 reason: " restore known state ".into(),
2117 actor_uuid: None,
2118 };
2119 let (receipt, restored) = revert_checkpoint(
2120 directory.path(),
2121 &request,
2122 || Ok(1_720_000_000_123_456),
2123 |_| Ok(()),
2124 )
2125 .unwrap();
2126 assert_eq!(restored.parent_generation_uuid(), Some(prior_current));
2127 assert_eq!(receipt.source_generation_uuid, source.generation_uuid());
2128 assert_eq!(receipt.prior_current_generation_uuid, Some(prior_current));
2129 assert_eq!(receipt.registry_revision, created.registry_revision);
2130 assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2131 let restoration_count = restored
2132 .participant_descriptors()
2133 .unwrap()
2134 .iter()
2135 .filter(|row| row.record_family_id == RESTORATION_FAMILY)
2136 .count();
2137 assert_eq!(restoration_count, 1);
2138
2139 delete_checkpoint(
2140 directory.path(),
2141 &CheckpointDeleteRequest {
2142 operation_uuid: Uuid::from_u128(42),
2143 name: "Before".into(),
2144 actor_uuid: None,
2145 },
2146 )
2147 .unwrap();
2148 let (replay, replayed_generation) = revert_checkpoint(
2149 directory.path(),
2150 &request,
2151 || panic!("published replay sampled clock"),
2152 |_| Ok(()),
2153 )
2154 .unwrap();
2155 assert_eq!(replay, receipt);
2156 assert_eq!(replay.prior_current_generation_uuid, Some(prior_current));
2157 assert_eq!(
2158 replayed_generation.generation_uuid(),
2159 restored.generation_uuid()
2160 );
2161 checkpoint_lock_handoff(
2162 directory.path(),
2163 "action=revert published-replay return",
2164 false,
2165 );
2166
2167 let mut conflict = request;
2168 conflict.reason = "different".into();
2169 let conflict_error =
2170 revert_checkpoint(directory.path(), &conflict, || Ok(0), |_| Ok(())).unwrap_err();
2171 assert_eq!(conflict_error.code(), "GF_IDEMPOTENCY_CONFLICT");
2172 checkpoint_lock_handoff(
2173 directory.path(),
2174 "action=revert published-replay conflict return",
2175 false,
2176 );
2177 }
2178
2179 #[test]
2180 fn revert_replay_lock_handoff_fails_closed_with_stable_storage_errors() {
2181 let checkpoint_error = finish_revert_replay_lock_handoff(
2182 Err(std::io::Error::other("checkpoint unlock failed")),
2183 Ok(()),
2184 )
2185 .unwrap_err();
2186 assert_eq!(checkpoint_error.code(), "GF_IO");
2187 assert_eq!(
2188 checkpoint_error.to_string(),
2189 "storage error: checkpoint revert replay lock handoff failed at checkpoints.lock: checkpoint unlock failed"
2190 );
2191
2192 let writer_error = finish_revert_replay_lock_handoff(
2193 Ok(()),
2194 Err(std::io::Error::other("writer unlock failed")),
2195 )
2196 .unwrap_err();
2197 assert_eq!(writer_error.code(), "GF_IO");
2198 assert_eq!(
2199 writer_error.to_string(),
2200 "storage error: checkpoint revert replay lock handoff failed at writer.lock: writer unlock failed"
2201 );
2202 }
2203
2204 #[test]
2205 fn revert_validation_failure_preserves_prior_current() {
2206 let directory = tempdir().unwrap();
2207 crate::open_or_initialize_project(directory.path()).unwrap();
2208 create_checkpoint(
2209 directory.path(),
2210 &create_request(Uuid::from_u128(50), "Before"),
2211 )
2212 .unwrap();
2213 let prior = publish_clone(directory.path());
2214 let error = revert_checkpoint(
2215 directory.path(),
2216 &CheckpointRevertRequest {
2217 operation_uuid: Uuid::from_u128(51),
2218 name: "Before".into(),
2219 reason: "must fail closed".into(),
2220 actor_uuid: None,
2221 },
2222 || Ok(1_720_000_000_123_456),
2223 |_| Err(GfError::Validation("injected composite failure".into())),
2224 )
2225 .unwrap_err();
2226 assert_eq!(error.code(), "GF_VALIDATION");
2227 checkpoint_lock_handoff(
2228 directory.path(),
2229 "action=revert validation-error return",
2230 false,
2231 );
2232 assert_eq!(
2233 crate::resolve_project_generation(directory.path())
2234 .unwrap()
2235 .generation_uuid(),
2236 prior
2237 );
2238 assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2239 }
2240
2241 #[cfg(unix)]
2242 #[test]
2243 fn mutation_lock_guard_unlocks_checkpoint_with_retained_duplicate_open() {
2244 let directory = tempdir().unwrap();
2245 crate::open_or_initialize_project(directory.path()).unwrap();
2246 let locks = acquire_mutation_locks(directory.path()).unwrap();
2247 let retained = locks.checkpoint.as_ref().unwrap().try_clone().unwrap();
2248 drop(locks);
2249
2250 let checkpoint =
2251 open_regular_lock(&directory.path().join(LOCKS_DIR).join(CHECKPOINT_LOCK_FILE))
2252 .unwrap();
2253 assert!(FileExt::try_lock_exclusive(&checkpoint).unwrap());
2254 FileExt::unlock(&checkpoint).unwrap();
2255 drop(retained);
2256 }
2257
2258 #[test]
2259 fn create_list_delete_and_replays_are_deterministic() {
2260 let directory = tempdir().unwrap();
2261 crate::open_or_initialize_project(directory.path()).unwrap();
2262 let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000001").unwrap();
2263 let request = create_request(operation, "Release 1.0");
2264 assert_eq!(
2265 hex(&create_request_digest(&request, "Release 1.0")),
2266 "01c7bf2f2c443d85d31ff80fef4a36484e31402213e8d145371867bdb2addbe8"
2267 );
2268 let created = create_checkpoint(directory.path(), &request).unwrap();
2269 assert_eq!(
2270 created.checkpoint_uuid,
2271 Uuid::parse_str("4084179c-38db-8b6b-9b6e-c0b0a855e002").unwrap()
2272 );
2273 let replayed = create_checkpoint(directory.path(), &request).unwrap();
2274 assert_eq!(created, replayed);
2275 assert_eq!(created.registry_revision, 1);
2276
2277 let rows = list_checkpoints(directory.path()).unwrap();
2278 assert_eq!(rows.len(), 1);
2279 assert_eq!(rows[0].checkpoint_uuid, created.checkpoint_uuid);
2280 assert_eq!(rows[0].generation_uuid, created.source_generation_uuid);
2281
2282 let delete = CheckpointDeleteRequest {
2283 operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000002").unwrap(),
2284 name: "Release 1.0".into(),
2285 actor_uuid: request.actor_uuid,
2286 };
2287 assert_eq!(
2288 hex(&delete_request_digest(&delete, "Release 1.0")),
2289 "9e6e15801f66ea4f7f58755c505fc1964e48ab87625459f388e2c23659135bb3"
2290 );
2291 let deleted = delete_checkpoint(directory.path(), &delete).unwrap();
2292 assert_eq!(
2293 deleted,
2294 delete_checkpoint(directory.path(), &delete).unwrap()
2295 );
2296 assert_eq!(deleted.registry_revision, 2);
2297 assert!(list_checkpoints(directory.path()).unwrap().is_empty());
2298 assert_eq!(
2299 created,
2300 create_checkpoint(directory.path(), &request).unwrap()
2301 );
2302 let changed_replay = create_request(operation, "Release 1.1");
2303 assert_eq!(
2304 create_checkpoint(directory.path(), &changed_replay)
2305 .unwrap_err()
2306 .code(),
2307 "GF_IDEMPOTENCY_CONFLICT"
2308 );
2309 assert!(list_checkpoints(directory.path()).unwrap().is_empty());
2310 }
2311
2312 #[test]
2313 fn identity_is_stable_across_independent_projects() {
2314 let first = tempdir().unwrap();
2315 let second = tempdir().unwrap();
2316 crate::open_or_initialize_project(first.path()).unwrap();
2317 crate::open_or_initialize_project(second.path()).unwrap();
2318 let request = create_request(
2319 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000010").unwrap(),
2320 "Stable",
2321 );
2322 let left = create_checkpoint(first.path(), &request).unwrap();
2323 let right = create_checkpoint(second.path(), &request).unwrap();
2324 assert_eq!(left.checkpoint_uuid, right.checkpoint_uuid);
2325 }
2326
2327 #[test]
2328 fn opened_checkpoint_generation_remains_pinned_after_delete() {
2329 let directory = tempdir().unwrap();
2330 crate::open_or_initialize_project(directory.path()).unwrap();
2331 let created = create_checkpoint(
2332 directory.path(),
2333 &create_request(Uuid::now_v7(), "Pinned View"),
2334 )
2335 .unwrap();
2336 let (row, opened) = open_checkpoint_generation(directory.path(), "Pinned View").unwrap();
2337 assert_eq!(row.checkpoint_uuid, created.checkpoint_uuid);
2338 assert_eq!(opened.generation_uuid(), created.source_generation_uuid);
2339 delete_checkpoint(
2340 directory.path(),
2341 &CheckpointDeleteRequest {
2342 operation_uuid: Uuid::now_v7(),
2343 name: "Pinned View".into(),
2344 actor_uuid: None,
2345 },
2346 )
2347 .unwrap();
2348 assert_eq!(
2349 open_checkpoint_generation(directory.path(), "Pinned View")
2350 .unwrap_err()
2351 .code(),
2352 "GF_CHECKPOINT_NOT_FOUND"
2353 );
2354 assert_eq!(opened.generation_uuid(), created.source_generation_uuid);
2355 assert!(opened.participant_snapshots().is_ok());
2356 }
2357
2358 #[test]
2359 fn conflicts_names_and_corruption_fail_closed() {
2360 let directory = tempdir().unwrap();
2361 crate::open_or_initialize_project(directory.path()).unwrap();
2362 let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000020").unwrap();
2363 create_checkpoint(directory.path(), &create_request(operation, "Safe.Name")).unwrap();
2364
2365 let conflict =
2366 create_checkpoint(directory.path(), &create_request(operation, "Other")).unwrap_err();
2367 assert_eq!(conflict.code(), "GF_IDEMPOTENCY_CONFLICT");
2368 let exists = create_checkpoint(
2369 directory.path(),
2370 &create_request(
2371 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000021").unwrap(),
2372 "Safe.Name",
2373 ),
2374 )
2375 .unwrap_err();
2376 assert_eq!(exists.code(), "GF_CHECKPOINT_EXISTS");
2377 for invalid in ["", "../escape", "two spaces", " e", "e ", ".", ".."] {
2378 let error =
2379 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), invalid))
2380 .unwrap_err();
2381 assert_eq!(error.code(), "GF_VALIDATION", "name={invalid:?}");
2382 }
2383
2384 fs::write(
2385 directory.path().join(CHECKPOINTS_DIR).join(CHECKSUM_FILE),
2386 b"0000000000000000000000000000000000000000000000000000000000000000\n",
2387 )
2388 .unwrap();
2389 let error = list_checkpoints(directory.path()).unwrap_err();
2390 assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2391 }
2392
2393 #[test]
2394 fn checksummed_but_impossible_registry_identities_fail_closed() {
2395 let directory = tempdir().unwrap();
2396 crate::open_or_initialize_project(directory.path()).unwrap();
2397 create_checkpoint(
2398 directory.path(),
2399 &create_request(Uuid::now_v7(), "Tampered"),
2400 )
2401 .unwrap();
2402 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2403 let mut registry = read_registry(&checkpoint_root).unwrap();
2404 registry.active[0].checkpoint_uuid = Uuid::now_v7();
2405 write_raw_registry(directory.path(), ®istry);
2406 assert_eq!(
2407 list_checkpoints(directory.path()).unwrap_err().code(),
2408 "GF_CHECKPOINT_REGISTRY_CORRUPT"
2409 );
2410
2411 let directory = tempdir().unwrap();
2412 crate::open_or_initialize_project(directory.path()).unwrap();
2413 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Overlap")).unwrap();
2414 delete_checkpoint(
2415 directory.path(),
2416 &CheckpointDeleteRequest {
2417 operation_uuid: Uuid::now_v7(),
2418 name: "Overlap".into(),
2419 actor_uuid: None,
2420 },
2421 )
2422 .unwrap();
2423 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2424 let mut registry = read_registry(&checkpoint_root).unwrap();
2425 let row = &mut registry.tombstones[0];
2426 row.delete_operation_uuid = row.create_operation_uuid;
2427 row.delete_request_sha256 = hex(&delete_request_digest_values(
2428 row.delete_operation_uuid,
2429 &row.name,
2430 row.deleted_by,
2431 ));
2432 write_raw_registry(directory.path(), ®istry);
2433 assert_eq!(
2434 list_checkpoints(directory.path()).unwrap_err().code(),
2435 "GF_CHECKPOINT_REGISTRY_CORRUPT"
2436 );
2437 }
2438
2439 #[test]
2440 fn exact_input_bounds_and_writer_lock_are_enforced() {
2441 assert!(validate_name(&"a".repeat(MAX_NAME_BYTES)).is_ok());
2442 assert_eq!(
2443 validate_name(&"a".repeat(MAX_NAME_BYTES + 1))
2444 .unwrap_err()
2445 .code(),
2446 "GF_VALIDATION"
2447 );
2448 assert!(validate_description(Some(&"d".repeat(MAX_DESCRIPTION_BYTES))).is_ok());
2449 assert_eq!(
2450 validate_description(Some(&"d".repeat(MAX_DESCRIPTION_BYTES + 1)))
2451 .unwrap_err()
2452 .code(),
2453 "GF_VALIDATION"
2454 );
2455
2456 let directory = tempdir().unwrap();
2457 let selected = crate::open_or_initialize_project(directory.path()).unwrap();
2458 let root = selected.container_root().to_owned();
2459 let locks = acquire_mutation_locks(&root).unwrap();
2460 let error = create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Busy"))
2461 .unwrap_err();
2462 assert_eq!(error.code(), "GF_WRITER_BUSY");
2463 drop(locks);
2464 }
2465
2466 #[cfg(unix)]
2467 #[test]
2468 fn linked_registry_surfaces_fail_closed_without_following_targets() {
2469 use std::os::unix::fs::symlink;
2470
2471 for hard in [false, true] {
2472 let directory = tempdir().unwrap();
2473 crate::open_or_initialize_project(directory.path()).unwrap();
2474 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Linked")).unwrap();
2475 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2476 let checksum = checkpoint_root.join(CHECKSUM_FILE);
2477 let external = directory.path().join("external-checksum");
2478 fs::rename(&checksum, &external).unwrap();
2479 if hard {
2480 fs::hard_link(&external, &checksum).unwrap();
2481 } else {
2482 symlink(&external, &checksum).unwrap();
2483 }
2484 let external_before = fs::read(&external).unwrap();
2485 let error = while_writer_lock_is_held(directory.path(), || {
2486 list_checkpoints(directory.path()).unwrap_err()
2487 });
2488 assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2489 assert_eq!(fs::read(&external).unwrap(), external_before);
2490 }
2491 }
2492
2493 #[test]
2494 fn no_intent_registry_corruption_wins_over_writer_contention() {
2495 let directory = tempdir().unwrap();
2496 crate::open_or_initialize_project(directory.path()).unwrap();
2497 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Corrupt")).unwrap();
2498 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2499 fs::write(
2500 checkpoint_root.join(CHECKSUM_FILE),
2501 b"not-the-registry-digest\n",
2502 )
2503 .unwrap();
2504
2505 let expected = read_registry(&checkpoint_root).unwrap_err();
2506 let error = while_writer_lock_is_held(directory.path(), || {
2507 list_checkpoints(directory.path()).unwrap_err()
2508 });
2509 assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2510 assert_eq!(error.to_string(), expected.to_string());
2511 }
2512
2513 #[test]
2514 fn intent_recovery_contention_preserves_writer_busy_and_intent() {
2515 let directory = tempdir().unwrap();
2516 crate::open_or_initialize_project(directory.path()).unwrap();
2517 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2518 let child = Command::new(std::env::current_exe().unwrap())
2519 .args([
2520 "--exact",
2521 "project_checkpoints::tests::checkpoint_failpoint_helper",
2522 "--ignored",
2523 ])
2524 .env(
2525 "GRAPHFORGE_PROJECT_FAILPOINTS",
2526 "graphforge-internal-subprocess-v1",
2527 )
2528 .env(
2529 "GRAPHFORGE_PROJECT_FAILPOINT",
2530 "checkpoint.registry.before_replace",
2531 )
2532 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2533 .spawn()
2534 .unwrap();
2535 let status = BoundedChild {
2536 child,
2537 reaped: false,
2538 }
2539 .wait("intent-recovery-contention failpoint=checkpoint.registry.before_replace");
2540 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2541 let intent_path = checkpoint_root.join(INTENT_FILE);
2542 let intent = fs::read(&intent_path).unwrap();
2543 let staged = fs::read_dir(&checkpoint_root)
2544 .unwrap()
2545 .filter_map(Result::ok)
2546 .filter(|entry| {
2547 entry
2548 .file_name()
2549 .to_string_lossy()
2550 .starts_with(".registry.")
2551 })
2552 .map(|entry| (entry.file_name(), fs::read(entry.path()).unwrap()))
2553 .collect::<BTreeMap<_, _>>();
2554 assert_eq!(staged.len(), 2);
2555
2556 let error = while_writer_lock_is_held(directory.path(), || {
2557 list_checkpoints(directory.path()).unwrap_err()
2558 });
2559 assert_eq!(error.code(), "GF_WRITER_BUSY");
2560 assert_eq!(fs::read(intent_path).unwrap(), intent);
2561 let staged_after = fs::read_dir(&checkpoint_root)
2562 .unwrap()
2563 .filter_map(Result::ok)
2564 .filter(|entry| {
2565 entry
2566 .file_name()
2567 .to_string_lossy()
2568 .starts_with(".registry.")
2569 })
2570 .map(|entry| (entry.file_name(), fs::read(entry.path()).unwrap()))
2571 .collect::<BTreeMap<_, _>>();
2572 assert_eq!(staged_after, staged);
2573 }
2574
2575 #[cfg(unix)]
2576 #[test]
2577 fn linked_project_root_is_rejected_before_checkpoint_access() {
2578 use std::os::unix::fs::symlink;
2579
2580 let directory = tempdir().unwrap();
2581 let project = directory.path().join("project");
2582 fs::create_dir(&project).unwrap();
2583 crate::open_or_initialize_project(&project).unwrap();
2584 let linked = directory.path().join("linked-project");
2585 symlink(&project, &linked).unwrap();
2586 assert_eq!(
2587 list_checkpoints(&linked).unwrap_err().code(),
2588 "GF_UNSUPPORTED_PROJECT_FORMAT"
2589 );
2590 }
2591
2592 #[test]
2593 fn checkpoint_pin_and_open_lease_control_recovery_cleanup() {
2594 let directory = tempdir().unwrap();
2595 crate::open_or_initialize_project(directory.path()).unwrap();
2596 let pinned_generation = publish_clone(directory.path());
2597 let created =
2598 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Pinned")).unwrap();
2599 assert_eq!(created.source_generation_uuid, pinned_generation);
2600 for _ in 0..4 {
2601 publish_clone(directory.path());
2602 }
2603 crate::recover_project_transactions(directory.path()).unwrap();
2604 recover_checkpoint_pair_after_lock_handoff(
2605 directory.path(),
2606 "action=delete-pinned-checkpoint recovery-complete",
2607 );
2608 let generation_path = directory
2609 .path()
2610 .join(crate::project_publication::GENERATIONS_DIR)
2611 .join(pinned_generation.hyphenated().to_string());
2612 assert!(
2613 generation_path.exists(),
2614 "active checkpoint lost its generation"
2615 );
2616
2617 let lease =
2618 crate::project_publication::open_regular_lock(&generation_path.join("lease.lock"))
2619 .unwrap();
2620 FileExt::lock_shared(&lease).unwrap();
2621 delete_checkpoint(
2622 directory.path(),
2623 &CheckpointDeleteRequest {
2624 operation_uuid: Uuid::now_v7(),
2625 name: "Pinned".into(),
2626 actor_uuid: None,
2627 },
2628 )
2629 .unwrap();
2630 crate::recover_project_transactions(directory.path()).unwrap();
2631 assert!(generation_path.exists(), "an open lease was invalidated");
2632 FileExt::unlock(&lease).unwrap();
2633 drop(lease);
2634 crate::recover_project_transactions(directory.path()).unwrap();
2635 assert!(
2636 !generation_path.exists(),
2637 "deleted pin did not permit later GC"
2638 );
2639 }
2640
2641 #[test]
2642 fn checkpoint_pin_survives_process_restart() {
2643 let directory = tempdir().unwrap();
2644 crate::open_or_initialize_project(directory.path()).unwrap();
2645 let pinned_generation = publish_clone(directory.path());
2646 create_checkpoint(
2647 directory.path(),
2648 &create_request(Uuid::from_u128(70), "Restart Pin"),
2649 )
2650 .unwrap();
2651 for _ in 0..4 {
2652 publish_clone(directory.path());
2653 }
2654
2655 let mut child = Command::new(std::env::current_exe().unwrap())
2656 .args([
2657 "--exact",
2658 "project_checkpoints::tests::checkpoint_failpoint_helper",
2659 "--ignored",
2660 "--nocapture",
2661 ])
2662 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2663 .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", "hold-open")
2664 .stdin(Stdio::piped())
2665 .stdout(Stdio::piped())
2666 .spawn()
2667 .unwrap();
2668 let mut output = BufReader::new(child.stdout.take().unwrap());
2669 let mut ready = String::new();
2670 while ready != "ready\n" {
2671 ready.clear();
2672 assert_ne!(
2673 output.read_line(&mut ready).unwrap(),
2674 0,
2675 "child exited before ready"
2676 );
2677 }
2678
2679 delete_checkpoint(
2680 directory.path(),
2681 &CheckpointDeleteRequest {
2682 operation_uuid: Uuid::from_u128(71),
2683 name: "Restart Pin".into(),
2684 actor_uuid: None,
2685 },
2686 )
2687 .unwrap();
2688 crate::recover_project_transactions(directory.path()).unwrap();
2689 let generation_path = directory
2690 .path()
2691 .join(crate::project_publication::GENERATIONS_DIR)
2692 .join(pinned_generation.hyphenated().to_string());
2693 assert!(
2694 generation_path.exists(),
2695 "subprocess lease was not retained"
2696 );
2697
2698 child.stdin.take().unwrap().write_all(b"release\n").unwrap();
2699 assert!(child.wait().unwrap().success());
2700 crate::recover_project_transactions(directory.path()).unwrap();
2701 assert!(
2702 !generation_path.exists(),
2703 "generation survived after the restarted reader released its lease"
2704 );
2705 }
2706
2707 #[test]
2708 fn checkpoint_cleanup_removes_all_transient_resources() {
2709 let directory = tempdir().unwrap();
2710 crate::open_or_initialize_project(directory.path()).unwrap();
2711 create_checkpoint(
2712 directory.path(),
2713 &create_request(Uuid::from_u128(80), "Cleanup"),
2714 )
2715 .unwrap();
2716 publish_clone(directory.path());
2717 delete_checkpoint(
2718 directory.path(),
2719 &CheckpointDeleteRequest {
2720 operation_uuid: Uuid::from_u128(81),
2721 name: "Cleanup".into(),
2722 actor_uuid: None,
2723 },
2724 )
2725 .unwrap();
2726 crate::recover_project_transactions(directory.path()).unwrap();
2727
2728 let selected = crate::resolve_project_generation(directory.path()).unwrap();
2729 let root = selected.container_root().to_owned();
2730 drop(selected);
2731 let checkpoint_root = root.join(CHECKPOINTS_DIR);
2732 let checkpoint_entries = fs::read_dir(&checkpoint_root)
2733 .unwrap()
2734 .map(|entry| entry.unwrap().file_name().into_string().unwrap())
2735 .collect::<BTreeSet<_>>();
2736 assert_eq!(
2737 checkpoint_entries,
2738 BTreeSet::from([REGISTRY_FILE.into(), CHECKSUM_FILE.into()]),
2739 "checkpoint transaction staging leaked"
2740 );
2741 let trash = root.join("trash");
2742 assert!(
2743 !trash.exists() || fs::read_dir(&trash).unwrap().next().is_none(),
2744 "recovery trash was not emptied"
2745 );
2746 assert!(
2747 !root.join("cache").exists(),
2748 "checkpoint lifecycle leaked process-cache state to disk"
2749 );
2750 for entry in fs::read_dir(root.join(crate::project_publication::GENERATIONS_DIR)).unwrap() {
2751 let path = entry.unwrap().path();
2752 let name = path.file_name().unwrap().to_str().unwrap();
2753 Uuid::parse_str(name).expect("generation staging entry leaked");
2754 let lease = open_regular_lock(&path.join("lease.lock")).unwrap();
2755 assert!(FileExt::try_lock_exclusive(&lease).unwrap());
2756 FileExt::unlock(&lease).unwrap();
2757 }
2758 let lock_root = root.join(LOCKS_DIR);
2759 for name in [WRITER_LOCK_FILE, CHECKPOINT_LOCK_FILE] {
2760 let lock = open_regular_lock(&lock_root.join(name)).unwrap();
2761 assert!(FileExt::try_lock_exclusive(&lock).unwrap(), "{name} leaked");
2762 FileExt::unlock(&lock).unwrap();
2763 }
2764 }
2765
2766 #[test]
2767 #[ignore = "subprocess failpoint helper"]
2768 fn checkpoint_failpoint_helper() {
2769 let root = std::env::var("GRAPHFORGE_CHECKPOINT_TEST_ROOT").unwrap();
2770 let action = std::env::var("GRAPHFORGE_CHECKPOINT_TEST_ACTION");
2771 if action.as_deref() == Ok("hold-open") {
2772 let (_, opened) = open_checkpoint_generation(root, "Restart Pin").unwrap();
2773 println!("ready");
2774 std::io::stdout().flush().unwrap();
2775 let mut release = String::new();
2776 std::io::stdin().read_line(&mut release).unwrap();
2777 assert_eq!(release, "release\n");
2778 assert!(opened.participant_snapshots().is_ok());
2779 } else if action.as_deref() == Ok("revert") {
2780 revert_checkpoint(
2781 root,
2782 &CheckpointRevertRequest {
2783 operation_uuid: Uuid::from_u128(61),
2784 name: "Base".into(),
2785 reason: "crash recovery".into(),
2786 actor_uuid: None,
2787 },
2788 || Ok(1_720_000_000_123_456),
2789 |_| Ok(()),
2790 )
2791 .unwrap();
2792 } else if action.as_deref() == Ok("delete") {
2793 delete_checkpoint(
2794 root,
2795 &CheckpointDeleteRequest {
2796 operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000031")
2797 .unwrap(),
2798 name: "Base".into(),
2799 actor_uuid: None,
2800 },
2801 )
2802 .unwrap();
2803 } else {
2804 let request = create_request(
2805 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
2806 "Crash",
2807 );
2808 create_checkpoint(root, &request).unwrap();
2809 }
2810 }
2811
2812 #[test]
2813 fn revert_publication_failpoint_matrix() {
2814 let failpoints = [
2815 ("project.after_journal_preparing", false),
2816 ("project.after_participant_dir_fsync", false),
2817 ("project.after_journal_staged", false),
2818 ("project.after_domain_validation", false),
2819 ("project.after_composite_validation", false),
2820 ("project.after_journal_validated", false),
2821 ("project.after_manifest_write", false),
2822 ("project.after_manifest_fsync", false),
2823 ("project.after_generation_dir_fsync", false),
2824 ("project.after_journal_durable", false),
2825 ("project.after_current_temp_write", false),
2826 ("project.after_current_temp_fsync", false),
2827 ("project.before_current_replace", false),
2828 ("project.after_current_replace", true),
2829 ("project.after_root_fsync", true),
2830 ("project.after_journal_published", true),
2831 ];
2832 for (failpoint, committed) in failpoints {
2833 let directory = tempdir().unwrap();
2834 crate::open_or_initialize_project(directory.path()).unwrap();
2835 create_checkpoint(
2836 directory.path(),
2837 &create_request(Uuid::from_u128(60), "Base"),
2838 )
2839 .unwrap();
2840 let prior = publish_clone(directory.path());
2841 let status = Command::new(std::env::current_exe().unwrap())
2842 .args([
2843 "--exact",
2844 "project_checkpoints::tests::checkpoint_failpoint_helper",
2845 "--ignored",
2846 ])
2847 .env(
2848 "GRAPHFORGE_PROJECT_FAILPOINTS",
2849 "graphforge-internal-subprocess-v1",
2850 )
2851 .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
2852 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2853 .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", "revert")
2854 .status()
2855 .unwrap();
2856 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2857 recover_checkpoint_pair_after_lock_handoff(
2858 directory.path(),
2859 &format!("action=revert failpoint={failpoint} committed={committed}"),
2860 );
2861 crate::recover_project_transactions(directory.path()).unwrap();
2862 let recovered = crate::resolve_project_generation(directory.path()).unwrap();
2863 assert_eq!(
2864 recovered.generation_uuid() != prior,
2865 committed,
2866 "{failpoint}"
2867 );
2868 recover_checkpoint_pair_after_lock_handoff(
2869 directory.path(),
2870 &format!(
2871 "action=revert parent-recovery-complete failpoint={failpoint} \
2872 committed={committed}"
2873 ),
2874 );
2875
2876 let (receipt, replayed) = revert_checkpoint(
2877 directory.path(),
2878 &CheckpointRevertRequest {
2879 operation_uuid: Uuid::from_u128(61),
2880 name: "Base".into(),
2881 reason: "crash recovery".into(),
2882 actor_uuid: None,
2883 },
2884 || Ok(1_720_000_000_123_456),
2885 |_| Ok(()),
2886 )
2887 .unwrap();
2888 assert_eq!(
2889 receipt.result_generation_uuid,
2890 Some(replayed.generation_uuid())
2891 );
2892 recover_checkpoint_pair_after_lock_handoff(
2893 directory.path(),
2894 &format!(
2895 "action=revert replay-complete failpoint={failpoint} committed={committed}"
2896 ),
2897 );
2898 assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2899 }
2900 }
2901
2902 #[test]
2903 fn registry_failpoints_recover_exact_previous_or_next_revision() {
2904 for (failpoint, committed) in [
2905 ("checkpoint.registry.after_intent_file_fsync", false),
2906 ("checkpoint.registry.after_file_fsync", false),
2907 ("checkpoint.registry.before_replace", false),
2908 ("checkpoint.registry.after_replace", true),
2909 ("checkpoint.registry.after_dir_fsync", true),
2910 ] {
2911 let directory = tempdir().unwrap();
2912 crate::open_or_initialize_project(directory.path()).unwrap();
2913 let status = Command::new(std::env::current_exe().unwrap())
2914 .args([
2915 "--exact",
2916 "project_checkpoints::tests::checkpoint_failpoint_helper",
2917 "--ignored",
2918 ])
2919 .env(
2920 "GRAPHFORGE_PROJECT_FAILPOINTS",
2921 "graphforge-internal-subprocess-v1",
2922 )
2923 .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
2924 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2925 .status()
2926 .unwrap();
2927 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2928 recover_checkpoint_pair_after_lock_handoff(
2929 directory.path(),
2930 &format!("action=create-unseeded failpoint={failpoint} committed={committed}"),
2931 );
2932 let rows = list_checkpoints(directory.path()).unwrap();
2933 assert_eq!(rows.len(), usize::from(committed), "{failpoint}");
2934 if !committed {
2935 recover_checkpoint_pair_after_lock_handoff(
2936 directory.path(),
2937 &format!(
2938 "action=create-unseeded parent-read-complete failpoint={failpoint} \
2939 committed={committed}"
2940 ),
2941 );
2942 let request = create_request(
2943 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
2944 "Crash",
2945 );
2946 create_checkpoint(directory.path(), &request).unwrap();
2947 recover_checkpoint_pair_after_lock_handoff(
2948 directory.path(),
2949 &format!(
2950 "action=create-unseeded replay-complete failpoint={failpoint} \
2951 committed={committed}"
2952 ),
2953 );
2954 assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2955 }
2956 }
2957 }
2958
2959 #[test]
2960 fn seeded_create_and_delete_failpoints_recover_exact_previous_or_next_revision() {
2961 let failpoints = [
2962 ("checkpoint.registry.after_intent_file_fsync", false),
2963 ("checkpoint.registry.after_file_fsync", false),
2964 ("checkpoint.registry.before_replace", false),
2965 ("checkpoint.registry.after_replace", true),
2966 ("checkpoint.registry.after_dir_fsync", true),
2967 ];
2968 for action in ["create", "delete"] {
2969 for (failpoint, committed) in failpoints {
2970 let directory = tempdir().unwrap();
2971 crate::open_or_initialize_project(directory.path()).unwrap();
2972 create_checkpoint(
2973 directory.path(),
2974 &create_request(
2975 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000029").unwrap(),
2976 "Base",
2977 ),
2978 )
2979 .unwrap();
2980 let status = Command::new(std::env::current_exe().unwrap())
2981 .args([
2982 "--exact",
2983 "project_checkpoints::tests::checkpoint_failpoint_helper",
2984 "--ignored",
2985 ])
2986 .env(
2987 "GRAPHFORGE_PROJECT_FAILPOINTS",
2988 "graphforge-internal-subprocess-v1",
2989 )
2990 .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
2991 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2992 .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", action)
2993 .status()
2994 .unwrap();
2995 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2996 recover_checkpoint_pair_after_lock_handoff(
2997 directory.path(),
2998 &format!("action={action} failpoint={failpoint} committed={committed}"),
2999 );
3000 let rows = list_checkpoints(directory.path()).unwrap();
3001 let expected = match (action, committed) {
3002 ("create", true) => vec!["Base", "Crash"],
3003 ("create", false) | ("delete", false) => vec!["Base"],
3004 ("delete", true) => vec![],
3005 _ => unreachable!(),
3006 };
3007 assert_eq!(
3008 rows.iter().map(|row| row.name.as_str()).collect::<Vec<_>>(),
3009 expected,
3010 "{action} {failpoint}"
3011 );
3012 recover_checkpoint_pair_after_lock_handoff(
3013 directory.path(),
3014 &format!(
3015 "action={action} parent-read-complete failpoint={failpoint} \
3016 committed={committed}"
3017 ),
3018 );
3019
3020 if action == "create" {
3021 let replay = create_checkpoint(
3022 directory.path(),
3023 &create_request(
3024 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
3025 "Crash",
3026 ),
3027 )
3028 .unwrap();
3029 assert_eq!(replay.registry_revision, 2);
3030 } else {
3031 let replay = delete_checkpoint(
3032 directory.path(),
3033 &CheckpointDeleteRequest {
3034 operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000031")
3035 .unwrap(),
3036 name: "Base".into(),
3037 actor_uuid: None,
3038 },
3039 )
3040 .unwrap();
3041 assert_eq!(replay.registry_revision, 2);
3042 }
3043 recover_checkpoint_pair_after_lock_handoff(
3044 directory.path(),
3045 &format!(
3046 "action={action} replay-complete failpoint={failpoint} \
3047 committed={committed}"
3048 ),
3049 );
3050 let final_rows = list_checkpoints(directory.path()).unwrap();
3051 let final_names = final_rows
3052 .iter()
3053 .map(|row| row.name.as_str())
3054 .collect::<Vec<_>>();
3055 if action == "create" {
3056 assert_eq!(final_names, vec!["Base", "Crash"], "{failpoint}");
3057 } else {
3058 assert!(final_names.is_empty(), "{failpoint}");
3059 }
3060 }
3061 }
3062 }
3063
3064 #[test]
3065 fn wave9_durable_registry_intent_recovers_every_atomic_pair_boundary() {
3066 for boundary in [
3067 "first-staged",
3068 "previous-staged",
3069 "registry-replaced",
3070 "next-complete",
3071 ] {
3072 let directory = tempdir().unwrap();
3073 crate::open_or_initialize_project(directory.path()).unwrap();
3074 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
3075 let previous = Registry::empty();
3076 let mut next = Registry::empty();
3077 next.revision = 1;
3078
3079 let has_previous = boundary != "first-staged";
3080 fs::create_dir_all(&checkpoint_root).unwrap();
3081 if has_previous {
3082 write_raw_registry(directory.path(), &previous);
3083 }
3084 let intent =
3085 install_registry_intent(&checkpoint_root, has_previous.then_some(&previous), &next);
3086 match boundary {
3087 "first-staged" | "previous-staged" => {}
3088 "registry-replaced" => {
3089 fs::rename(
3090 checkpoint_root.join(&intent.registry_temp),
3091 checkpoint_root.join(REGISTRY_FILE),
3092 )
3093 .unwrap();
3094 }
3095 "next-complete" => {
3096 fs::rename(
3097 checkpoint_root.join(&intent.registry_temp),
3098 checkpoint_root.join(REGISTRY_FILE),
3099 )
3100 .unwrap();
3101 fs::rename(
3102 checkpoint_root.join(&intent.checksum_temp),
3103 checkpoint_root.join(CHECKSUM_FILE),
3104 )
3105 .unwrap();
3106 }
3107 _ => unreachable!(),
3108 }
3109
3110 recover_pair(&checkpoint_root).unwrap();
3111
3112 assert!(!checkpoint_root.join(INTENT_FILE).exists(), "{boundary}");
3113 assert!(
3114 !checkpoint_root.join(&intent.registry_temp).exists(),
3115 "{boundary}"
3116 );
3117 assert!(
3118 !checkpoint_root.join(&intent.checksum_temp).exists(),
3119 "{boundary}"
3120 );
3121 let recovered = read_registry(&checkpoint_root).unwrap();
3122 let expected_revision =
3123 usize::from(matches!(boundary, "registry-replaced" | "next-complete"));
3124 assert_eq!(recovered.revision, expected_revision as u64, "{boundary}");
3125 }
3126 }
3127
3128 #[test]
3129 fn wave9_registry_intent_rejects_unsafe_names_and_wrong_staged_revision() {
3130 let directory = tempdir().unwrap();
3131 crate::open_or_initialize_project(directory.path()).unwrap();
3132 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
3133 let mut next = Registry::empty();
3134 next.revision = 1;
3135 let mut intent = install_registry_intent(&checkpoint_root, None, &next);
3136 intent.registry_temp = "../registry.json".into();
3137 let mut bytes = serde_json::to_vec(&intent).unwrap();
3138 bytes.push(b'\n');
3139 fs::write(checkpoint_root.join(INTENT_FILE), bytes).unwrap();
3140 assert_eq!(
3141 recover_pair(&checkpoint_root).unwrap_err().code(),
3142 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3143 );
3144
3145 fs::remove_dir_all(&checkpoint_root).unwrap();
3146 fs::create_dir(&checkpoint_root).unwrap();
3147 let mut intent = install_registry_intent(&checkpoint_root, None, &next);
3148 intent.next_revision = 2;
3149 let mut bytes = serde_json::to_vec(&intent).unwrap();
3150 bytes.push(b'\n');
3151 fs::write(checkpoint_root.join(INTENT_FILE), bytes).unwrap();
3152 assert_eq!(
3153 recover_pair(&checkpoint_root).unwrap_err().code(),
3154 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3155 );
3156 }
3157
3158 #[test]
3159 fn recovery_rejects_missing_or_tampered_staged_pair() {
3160 for tamper_checksum in [false, true] {
3161 let directory = tempdir().unwrap();
3162 crate::open_or_initialize_project(directory.path()).unwrap();
3163 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Base")).unwrap();
3164 let status = Command::new(std::env::current_exe().unwrap())
3165 .args([
3166 "--exact",
3167 "project_checkpoints::tests::checkpoint_failpoint_helper",
3168 "--ignored",
3169 ])
3170 .env(
3171 "GRAPHFORGE_PROJECT_FAILPOINTS",
3172 "graphforge-internal-subprocess-v1",
3173 )
3174 .env(
3175 "GRAPHFORGE_PROJECT_FAILPOINT",
3176 "checkpoint.registry.before_replace",
3177 )
3178 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
3179 .status()
3180 .unwrap();
3181 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
3182 preserve_checkpoint_intent_after_lock_handoff(
3183 directory.path(),
3184 &format!(
3185 "action=create-tamper failpoint=checkpoint.registry.before_replace \
3186 committed=false tamper_checksum={tamper_checksum}"
3187 ),
3188 );
3189 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
3190 let intent_bytes = fs::read(checkpoint_root.join(INTENT_FILE)).unwrap();
3191 let intent: RegistryIntent = serde_json::from_slice(&intent_bytes).unwrap();
3192 if tamper_checksum {
3193 fs::write(
3194 checkpoint_root.join(intent.checksum_temp),
3195 b"0000000000000000000000000000000000000000000000000000000000000000\n",
3196 )
3197 .unwrap();
3198 } else {
3199 fs::remove_file(checkpoint_root.join(intent.registry_temp)).unwrap();
3200 }
3201 assert_eq!(
3202 list_checkpoints(directory.path()).unwrap_err().code(),
3203 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3204 );
3205 }
3206 }
3207
3208 #[test]
3209 fn checkpoint_text_and_identity_boundaries_are_canonical() {
3210 for invalid in [
3211 "",
3212 " leading",
3213 "trailing ",
3214 ".",
3215 "..",
3216 "two spaces",
3217 "bad/name",
3218 ] {
3219 assert_eq!(validate_name(invalid).unwrap_err().code(), "GF_VALIDATION");
3220 }
3221 assert_eq!(validate_name("Résumé.v1").unwrap(), "Résumé.v1");
3222 let combining_acute = char::from_u32(0x301).unwrap();
3223 let decomposed = format!("Re{combining_acute}sume{combining_acute}");
3224 assert_eq!(
3225 validate_name(&decomposed).unwrap_err().code(),
3226 "GF_VALIDATION"
3227 );
3228 assert!(validate_description(None).is_ok());
3229 assert!(validate_description(Some("bounded description")).is_ok());
3230 assert_eq!(
3231 validate_description(Some("contains\ncontrol"))
3232 .unwrap_err()
3233 .code(),
3234 "GF_VALIDATION"
3235 );
3236 assert_eq!(
3237 validate_description(Some(&"x".repeat(MAX_DESCRIPTION_BYTES + 1)))
3238 .unwrap_err()
3239 .code(),
3240 "GF_VALIDATION"
3241 );
3242 assert_eq!(
3243 validate_reason(" restored after audit ").unwrap(),
3244 "restored after audit"
3245 );
3246 assert_eq!(validate_reason(" ").unwrap_err().code(), "GF_VALIDATION");
3247
3248 let operation = Uuid::now_v7();
3249 let actor = Uuid::now_v7();
3250 let digest = create_request_digest_values(operation, "baseline", Some("desc"), Some(actor));
3251 let checkpoint = checkpoint_uuid(operation, digest);
3252 let encoded = hex(&digest);
3253 assert_eq!(decode_digest(&encoded).unwrap(), digest);
3254 assert!(
3255 validate_record_identity(
3256 checkpoint,
3257 operation,
3258 "baseline",
3259 Some("desc"),
3260 Some(actor),
3261 &encoded,
3262 )
3263 .is_ok()
3264 );
3265 assert_eq!(
3266 validate_record_identity(
3267 Uuid::nil(),
3268 operation,
3269 "baseline",
3270 Some("desc"),
3271 Some(actor),
3272 &encoded,
3273 )
3274 .unwrap_err()
3275 .code(),
3276 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3277 );
3278 for malformed in ["0", &"A".repeat(64), &"g".repeat(64)] {
3279 assert_eq!(
3280 decode_digest(malformed).unwrap_err().code(),
3281 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3282 );
3283 }
3284 assert!(valid_private_name(
3285 &format!(".registry.{operation}.json.next"),
3286 operation,
3287 "json"
3288 ));
3289 assert!(!valid_private_name(
3290 ".registry.other.json.next",
3291 operation,
3292 "json"
3293 ));
3294 assert_eq!(
3295 parse_uuid("not-a-uuid").unwrap_err().code(),
3296 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3297 );
3298 }
3299
3300 #[test]
3301 fn empty_checkpoint_registry_has_stable_canonical_bytes() {
3302 let registry = Registry::empty();
3303 let first = registry.canonical_bytes().unwrap();
3304 let second = registry.canonical_bytes().unwrap();
3305 assert_eq!(first, second);
3306 assert!(first.ends_with(b"\n"));
3307 let decoded: Registry = serde_json::from_slice(&first).unwrap();
3308 assert_eq!(decoded, registry);
3309 }
3310
3311 #[test]
3312 fn checkpoint_operation_identities_remain_disjoint_across_tombstones() {
3313 let root = tempdir().unwrap();
3314 crate::open_or_initialize_project(root.path()).unwrap();
3315 let create_operation = Uuid::now_v7();
3316 let create = create_request(create_operation, "Baseline");
3317 let created = create_checkpoint(root.path(), &create).unwrap();
3318 let exact_create_replay = create_checkpoint(root.path(), &create).unwrap();
3319 assert_eq!(exact_create_replay, created);
3320
3321 let changed_create = CheckpointCreateRequest {
3322 name: "Changed".into(),
3323 ..create.clone()
3324 };
3325 assert_eq!(
3326 create_checkpoint(root.path(), &changed_create)
3327 .unwrap_err()
3328 .code(),
3329 "GF_IDEMPOTENCY_CONFLICT"
3330 );
3331 assert_eq!(list_checkpoints(root.path()).unwrap().len(), 1);
3332
3333 let delete_operation = Uuid::now_v7();
3334 let delete = CheckpointDeleteRequest {
3335 operation_uuid: delete_operation,
3336 name: "Baseline".into(),
3337 actor_uuid: create.actor_uuid,
3338 };
3339 let deleted = delete_checkpoint(root.path(), &delete).unwrap();
3340 let exact_delete_replay = delete_checkpoint(root.path(), &delete).unwrap();
3341 assert_eq!(exact_delete_replay, deleted);
3342 assert!(list_checkpoints(root.path()).unwrap().is_empty());
3343
3344 let tombstone_create_replay = create_checkpoint(root.path(), &create).unwrap();
3345 assert_eq!(
3346 tombstone_create_replay.checkpoint_uuid,
3347 created.checkpoint_uuid
3348 );
3349 assert_eq!(tombstone_create_replay, created);
3350 assert!(list_checkpoints(root.path()).unwrap().is_empty());
3351
3352 let changed_delete = CheckpointDeleteRequest {
3353 name: "Other".into(),
3354 ..delete.clone()
3355 };
3356 assert_eq!(
3357 delete_checkpoint(root.path(), &changed_delete)
3358 .unwrap_err()
3359 .code(),
3360 "GF_IDEMPOTENCY_CONFLICT"
3361 );
3362 assert_eq!(
3363 create_checkpoint(root.path(), &create_request(delete_operation, "Other"))
3364 .unwrap_err()
3365 .code(),
3366 "GF_IDEMPOTENCY_CONFLICT"
3367 );
3368 assert_eq!(
3369 delete_checkpoint(
3370 root.path(),
3371 &CheckpointDeleteRequest {
3372 operation_uuid: create_operation,
3373 name: "Missing".into(),
3374 actor_uuid: None,
3375 },
3376 )
3377 .unwrap_err()
3378 .code(),
3379 "GF_IDEMPOTENCY_CONFLICT"
3380 );
3381 assert_eq!(
3382 delete_checkpoint(
3383 root.path(),
3384 &CheckpointDeleteRequest {
3385 operation_uuid: Uuid::now_v7(),
3386 name: "Missing".into(),
3387 actor_uuid: None,
3388 },
3389 )
3390 .unwrap_err()
3391 .code(),
3392 "GF_CHECKPOINT_NOT_FOUND"
3393 );
3394 assert!(list_checkpoints(root.path()).unwrap().is_empty());
3395 }
3396
3397 #[test]
3398 fn registry_header_sort_and_revision_validation_matrix_uses_durable_records() {
3399 let directory = tempdir().unwrap();
3400 crate::open_or_initialize_project(directory.path()).unwrap();
3401 create_checkpoint(
3402 directory.path(),
3403 &create_request(Uuid::from_u128(701), "Alpha"),
3404 )
3405 .unwrap();
3406 create_checkpoint(
3407 directory.path(),
3408 &create_request(Uuid::from_u128(702), "Beta"),
3409 )
3410 .unwrap();
3411 let root = directory.path().join(CHECKPOINTS_DIR);
3412 let stable = read_registry(&root).unwrap();
3413 assert!(validate_registry(&stable).is_ok());
3414
3415 let mutations: Vec<Box<dyn Fn(&mut Registry)>> = vec![
3416 Box::new(|registry| registry.format = "future".into()),
3417 Box::new(|registry| registry.format_version = 2),
3418 Box::new(|registry| registry.active.reverse()),
3419 Box::new(|registry| registry.active[0].created_revision = 0),
3420 Box::new(|registry| registry.active[1].name = registry.active[0].name.clone()),
3421 Box::new(|registry| {
3422 registry.active[1].checkpoint_uuid = registry.active[0].checkpoint_uuid
3423 }),
3424 Box::new(|registry| {
3425 registry.active[1].create_operation_uuid = registry.active[0].create_operation_uuid
3426 }),
3427 ];
3428 for mutate in mutations {
3429 let mut candidate = stable.clone();
3430 mutate(&mut candidate);
3431 assert_eq!(
3432 validate_registry(&candidate).unwrap_err().code(),
3433 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3434 );
3435 }
3436 assert_eq!(read_registry(&root).unwrap(), stable);
3437 }
3438
3439 #[test]
3440 fn tombstone_identity_revision_and_operation_disjointness_matrix_is_total() {
3441 let directory = tempdir().unwrap();
3442 crate::open_or_initialize_project(directory.path()).unwrap();
3443 create_checkpoint(
3444 directory.path(),
3445 &create_request(Uuid::from_u128(801), "Deleted"),
3446 )
3447 .unwrap();
3448 delete_checkpoint(
3449 directory.path(),
3450 &CheckpointDeleteRequest {
3451 operation_uuid: Uuid::from_u128(802),
3452 name: "Deleted".into(),
3453 actor_uuid: None,
3454 },
3455 )
3456 .unwrap();
3457 let root = directory.path().join(CHECKPOINTS_DIR);
3458 let stable = read_registry(&root).unwrap();
3459 assert_eq!(stable.active.len(), 0);
3460 assert_eq!(stable.tombstones.len(), 1);
3461 assert!(validate_registry(&stable).is_ok());
3462
3463 let mutations: Vec<Box<dyn Fn(&mut Registry)>> = vec![
3464 Box::new(|registry| registry.tombstones[0].name = " bad".into()),
3465 Box::new(|registry| registry.tombstones[0].description = Some("bad\nvalue".into())),
3466 Box::new(|registry| registry.tombstones[0].generation_manifest_sha256 = "bad".into()),
3467 Box::new(|registry| registry.tombstones[0].create_request_sha256 = "bad".into()),
3468 Box::new(|registry| registry.tombstones[0].delete_request_sha256 = "bad".into()),
3469 Box::new(|registry| registry.tombstones[0].checkpoint_uuid = Uuid::nil()),
3470 Box::new(|registry| registry.tombstones[0].created_revision = 0),
3471 Box::new(|registry| {
3472 registry.tombstones[0].deleted_revision = registry.tombstones[0].created_revision
3473 }),
3474 Box::new(|registry| registry.tombstones[0].deleted_revision = registry.revision + 1),
3475 Box::new(|registry| {
3476 registry.tombstones[0].delete_operation_uuid =
3477 registry.tombstones[0].create_operation_uuid
3478 }),
3479 Box::new(|registry| registry.revision = 0),
3480 ];
3481 for mutate in mutations {
3482 let mut candidate = stable.clone();
3483 mutate(&mut candidate);
3484 assert_eq!(
3485 validate_registry(&candidate).unwrap_err().code(),
3486 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3487 );
3488 }
3489
3490 let mut duplicate = stable.clone();
3491 let mut second = duplicate.tombstones[0].clone();
3492 second.deleted_revision += 1;
3493 duplicate.revision = second.deleted_revision;
3494 duplicate.tombstones.push(second);
3495 assert_eq!(
3496 validate_registry(&duplicate).unwrap_err().code(),
3497 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3498 );
3499 assert_eq!(read_registry(&root).unwrap(), stable);
3500 }
3501
3502 #[test]
3503 fn public_checkpoint_operations_reject_cross_kind_uuid_reuse_after_reopen() {
3504 let root = tempdir().unwrap();
3505 crate::open_or_initialize_project(root.path()).unwrap();
3506 let operation = Uuid::now_v7();
3507 create_checkpoint(root.path(), &create_request(operation, "release")).unwrap();
3508
3509 let delete = CheckpointDeleteRequest {
3510 operation_uuid: operation,
3511 name: "release".into(),
3512 actor_uuid: None,
3513 };
3514 assert_eq!(
3515 delete_checkpoint(root.path(), &delete).unwrap_err().code(),
3516 "GF_IDEMPOTENCY_CONFLICT"
3517 );
3518
3519 let delete_operation = Uuid::now_v7();
3520 delete_checkpoint(
3521 root.path(),
3522 &CheckpointDeleteRequest {
3523 operation_uuid: delete_operation,
3524 name: "release".into(),
3525 actor_uuid: None,
3526 },
3527 )
3528 .unwrap();
3529 assert_eq!(
3530 create_checkpoint(
3531 root.path(),
3532 &create_request(delete_operation, "replacement")
3533 )
3534 .unwrap_err()
3535 .code(),
3536 "GF_IDEMPOTENCY_CONFLICT"
3537 );
3538 assert!(list_checkpoints(root.path()).unwrap().is_empty());
3539 }
3540}