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_name(&row.name)
1253 .map_err(|_| registry_corrupt("active checkpoint name is invalid"))?;
1254 validate_description(row.description.as_deref())
1255 .map_err(|_| registry_corrupt("active checkpoint description is invalid"))?;
1256 validate_digest(&row.generation_manifest_sha256)?;
1257 validate_digest(&row.create_request_sha256)?;
1258 validate_record_identity(
1259 row.checkpoint_uuid,
1260 row.create_operation_uuid,
1261 &row.name,
1262 row.description.as_deref(),
1263 row.created_by,
1264 &row.create_request_sha256,
1265 )?;
1266 if row.created_revision == 0
1267 || row.created_revision > registry.revision
1268 || !names.insert(row.name.as_str())
1269 || !checkpoint_uuids.insert(row.checkpoint_uuid)
1270 || !create_operations.insert(row.create_operation_uuid)
1271 {
1272 return Err(registry_corrupt(
1273 "active checkpoint identities or revision are inconsistent",
1274 ));
1275 }
1276 }
1277 for row in ®istry.tombstones {
1278 validate_name(&row.name)
1279 .map_err(|_| registry_corrupt("checkpoint tombstone name is invalid"))?;
1280 validate_description(row.description.as_deref())
1281 .map_err(|_| registry_corrupt("checkpoint tombstone description is invalid"))?;
1282 validate_digest(&row.generation_manifest_sha256)?;
1283 validate_digest(&row.create_request_sha256)?;
1284 validate_digest(&row.delete_request_sha256)?;
1285 validate_record_identity(
1286 row.checkpoint_uuid,
1287 row.create_operation_uuid,
1288 &row.name,
1289 row.description.as_deref(),
1290 row.created_by,
1291 &row.create_request_sha256,
1292 )?;
1293 let expected_delete =
1294 delete_request_digest_values(row.delete_operation_uuid, &row.name, row.deleted_by);
1295 if row.created_revision == 0
1296 || row.created_revision >= row.deleted_revision
1297 || row.deleted_revision > registry.revision
1298 || !checkpoint_uuids.insert(row.checkpoint_uuid)
1299 || !create_operations.insert(row.create_operation_uuid)
1300 || delete_operations.contains(&row.create_operation_uuid)
1301 || !delete_operations.insert(row.delete_operation_uuid)
1302 || create_operations.contains(&row.delete_operation_uuid)
1303 || row.delete_request_sha256 != hex(&expected_delete)
1304 {
1305 return Err(registry_corrupt(
1306 "checkpoint tombstone identities or revisions are inconsistent",
1307 ));
1308 }
1309 }
1310 Ok(())
1311}
1312
1313fn validate_single_link_regular(path: &Path, label: &str) -> Result<(), GfError> {
1314 let metadata = fs::symlink_metadata(path).map_err(storage_io)?;
1315 if !metadata.file_type().is_file() {
1316 return Err(registry_corrupt(format!("{label} is linked or special")));
1317 }
1318 #[cfg(unix)]
1319 {
1320 use std::os::unix::fs::MetadataExt;
1321 if metadata.nlink() != 1 {
1322 return Err(registry_corrupt(format!("{label} is hard-linked")));
1323 }
1324 }
1325 Ok(())
1326}
1327
1328fn validate_name(value: &str) -> Result<String, GfError> {
1329 let normalized: String = value.nfc().collect();
1330 if normalized != value
1331 || value.is_empty()
1332 || value.len() > MAX_NAME_BYTES
1333 || value.trim() != value
1334 || value == "."
1335 || value == ".."
1336 || value.contains(" ")
1337 || !value
1338 .chars()
1339 .all(|ch| ch.is_alphanumeric() || matches!(ch, ' ' | '_' | '-' | '.'))
1340 {
1341 return Err(GfError::Validation(
1342 "checkpoint name is not canonical NFC content or violates the 1-128 byte grammar"
1343 .into(),
1344 ));
1345 }
1346 Ok(normalized)
1347}
1348
1349fn validate_description(value: Option<&str>) -> Result<(), GfError> {
1350 if value.is_some_and(|value| {
1351 value.len() > MAX_DESCRIPTION_BYTES || value.chars().any(char::is_control)
1352 }) {
1353 return Err(GfError::Validation(
1354 "checkpoint description exceeds 1024 UTF-8 bytes or contains controls".into(),
1355 ));
1356 }
1357 Ok(())
1358}
1359
1360fn validate_digest(value: &str) -> Result<(), GfError> {
1361 if value.len() != 64
1362 || !value
1363 .bytes()
1364 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1365 {
1366 return Err(registry_corrupt("checkpoint digest is noncanonical"));
1367 }
1368 Ok(())
1369}
1370
1371fn decode_digest(value: &str) -> Result<[u8; 32], GfError> {
1372 validate_digest(value)?;
1373 let mut digest = [0_u8; 32];
1374 for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
1375 let text = std::str::from_utf8(pair)
1376 .map_err(|_| registry_corrupt("checkpoint digest is not UTF-8"))?;
1377 digest[index] = u8::from_str_radix(text, 16)
1378 .map_err(|_| registry_corrupt("checkpoint digest is not lowercase hex"))?;
1379 }
1380 Ok(digest)
1381}
1382
1383fn create_request_digest(request: &CheckpointCreateRequest, name: &str) -> [u8; 32] {
1384 create_request_digest_values(
1385 request.operation_uuid,
1386 name,
1387 request.description.as_deref(),
1388 request.actor_uuid,
1389 )
1390}
1391
1392fn create_request_digest_values(
1393 operation_uuid: Uuid,
1394 name: &str,
1395 description: Option<&str>,
1396 actor_uuid: Option<Uuid>,
1397) -> [u8; 32] {
1398 let mut hasher = Sha256::new();
1399 hasher.update(b"graphforge-checkpoint-create-request/1");
1400 hasher.update(operation_uuid.as_bytes());
1401 append_bytes(&mut hasher, name.as_bytes());
1402 match description {
1403 Some(value) => {
1404 hasher.update([1]);
1405 append_bytes(&mut hasher, value.as_bytes());
1406 }
1407 None => hasher.update([0]),
1408 }
1409 append_actor(&mut hasher, actor_uuid);
1410 hasher.finalize().into()
1411}
1412
1413fn delete_request_digest(request: &CheckpointDeleteRequest, name: &str) -> [u8; 32] {
1414 delete_request_digest_values(request.operation_uuid, name, request.actor_uuid)
1415}
1416
1417fn delete_request_digest_values(
1418 operation_uuid: Uuid,
1419 name: &str,
1420 actor_uuid: Option<Uuid>,
1421) -> [u8; 32] {
1422 let mut hasher = Sha256::new();
1423 hasher.update(b"graphforge-checkpoint-delete-request/1");
1424 hasher.update(operation_uuid.as_bytes());
1425 append_bytes(&mut hasher, name.as_bytes());
1426 append_actor(&mut hasher, actor_uuid);
1427 hasher.finalize().into()
1428}
1429
1430fn validate_record_identity(
1431 checkpoint: Uuid,
1432 operation: Uuid,
1433 name: &str,
1434 description: Option<&str>,
1435 actor: Option<Uuid>,
1436 request_hex: &str,
1437) -> Result<(), GfError> {
1438 let request = create_request_digest_values(operation, name, description, actor);
1439 if request_hex != hex(&request) || checkpoint != checkpoint_uuid(operation, request) {
1440 return Err(registry_corrupt(
1441 "checkpoint deterministic identity or create request digest is inconsistent",
1442 ));
1443 }
1444 Ok(())
1445}
1446
1447fn checkpoint_uuid(operation_uuid: Uuid, request_digest: [u8; 32]) -> Uuid {
1448 let mut hasher = Sha256::new();
1449 hasher.update(b"graphforge-checkpoint-uuid/1");
1450 hasher.update(operation_uuid.as_bytes());
1451 hasher.update(request_digest);
1452 graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1453}
1454
1455fn append_bytes(hasher: &mut Sha256, bytes: &[u8]) {
1456 hasher.update(
1457 u32::try_from(bytes.len())
1458 .expect("validated checkpoint strings fit u32")
1459 .to_be_bytes(),
1460 );
1461 hasher.update(bytes);
1462}
1463fn append_actor(hasher: &mut Sha256, actor: Option<Uuid>) {
1464 match actor {
1465 Some(value) => {
1466 hasher.update([1]);
1467 hasher.update(value.as_bytes());
1468 }
1469 None => hasher.update([0]),
1470 }
1471}
1472fn valid_private_name(name: &str, uuid: Uuid, kind: &str) -> bool {
1473 name == format!(".registry.{uuid}.{kind}.next")
1474}
1475fn hex(bytes: &[u8; 32]) -> String {
1476 let mut output = String::with_capacity(64);
1477 for byte in bytes {
1478 write!(&mut output, "{byte:02x}").expect("writing hexadecimal to String cannot fail");
1479 }
1480 output
1481}
1482
1483fn parse_uuid(value: &str) -> Result<Uuid, GfError> {
1484 Uuid::parse_str(value).map_err(|_| registry_corrupt("revert journal UUID is invalid"))
1485}
1486
1487fn validate_reason(value: &str) -> Result<String, GfError> {
1488 let trimmed = value.trim();
1489 if trimmed.is_empty() || trimmed.len() > MAX_REASON_BYTES {
1490 return Err(GfError::Validation(
1491 "checkpoint revert reason must contain 1..=1024 UTF-8 bytes after trimming".into(),
1492 ));
1493 }
1494 Ok(trimmed.to_owned())
1495}
1496
1497fn revert_request_digest(
1498 operation_uuid: Uuid,
1499 name: &str,
1500 checkpoint_uuid: Uuid,
1501 source_generation_uuid: Uuid,
1502 source_manifest_sha256: [u8; 32],
1503 reason: &str,
1504 actor_uuid: Option<Uuid>,
1505) -> [u8; 32] {
1506 let mut hasher = Sha256::new();
1507 hasher.update(b"graphforge-checkpoint-revert-request/1");
1508 hasher.update(operation_uuid.as_bytes());
1509 append_bytes(&mut hasher, name.as_bytes());
1510 hasher.update(checkpoint_uuid.as_bytes());
1511 hasher.update(source_generation_uuid.as_bytes());
1512 hasher.update(source_manifest_sha256);
1513 append_bytes(&mut hasher, reason.as_bytes());
1514 append_actor(&mut hasher, actor_uuid);
1515 hasher.finalize().into()
1516}
1517
1518fn revert_transaction_uuid(operation_uuid: Uuid) -> Uuid {
1519 let mut hasher = Sha256::new();
1520 hasher.update(b"graphforge-checkpoint-revert-transaction/1");
1521 hasher.update(operation_uuid.as_bytes());
1522 graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1523}
1524
1525fn restoration_uuid(operation_uuid: Uuid, request_digest: [u8; 32]) -> Uuid {
1526 let mut hasher = Sha256::new();
1527 hasher.update(b"graphforge-restoration-transition-uuid/1");
1528 hasher.update(operation_uuid.as_bytes());
1529 hasher.update(request_digest);
1530 graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1531}
1532
1533fn restored_generation_uuid(
1534 transaction_uuid: Uuid,
1535 checkpoint_uuid: Uuid,
1536 source_generation_uuid: Uuid,
1537 source_manifest_sha256: [u8; 32],
1538 prior_current_generation_uuid: Uuid,
1539 restored_at: i64,
1540 request_digest: [u8; 32],
1541) -> Uuid {
1542 let mut hasher = Sha256::new();
1543 hasher.update(b"graphforge-checkpoint-restored-generation/1");
1544 hasher.update(transaction_uuid.as_bytes());
1545 hasher.update(checkpoint_uuid.as_bytes());
1546 hasher.update(source_generation_uuid.as_bytes());
1547 hasher.update(source_manifest_sha256);
1548 hasher.update(prior_current_generation_uuid.as_bytes());
1549 hasher.update(restored_at.to_be_bytes());
1550 hasher.update(request_digest);
1551 graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1552}
1553
1554fn snapshot_to_participant(
1555 snapshot: crate::ProjectParticipantSnapshot,
1556) -> Result<ProjectParticipant, GfError> {
1557 let encoding = match snapshot.encoding.as_str() {
1558 "parquet" => ProjectParticipantEncoding::Parquet,
1559 "arrow" => ProjectParticipantEncoding::Arrow,
1560 "json" => ProjectParticipantEncoding::Json,
1561 _ => {
1562 return Err(registry_corrupt(
1563 "checkpoint participant encoding is unsupported",
1564 ));
1565 }
1566 };
1567 Ok(ProjectParticipant {
1568 capability_id: snapshot.capability_id,
1569 capability_version: snapshot.capability_version,
1570 record_family_id: snapshot.record_family_id,
1571 record_version: snapshot.record_version,
1572 encoding,
1573 schema_fingerprint: snapshot.schema_fingerprint,
1574 row_count: snapshot.row_count,
1575 bytes: snapshot.bytes,
1576 })
1577}
1578
1579#[allow(clippy::too_many_arguments)]
1580fn restoration_participant(
1581 restoration_uuid: Uuid,
1582 checkpoint_uuid: Uuid,
1583 source_generation_uuid: Uuid,
1584 source_manifest_sha256: [u8; 32],
1585 prior_current_generation_uuid: Uuid,
1586 restored_generation_uuid: Uuid,
1587 operation_uuid: Uuid,
1588 actor_uuid: Option<Uuid>,
1589 reason: &str,
1590 restored_at: i64,
1591) -> Result<ProjectParticipant, GfError> {
1592 let schema = Arc::new(Schema::new(vec![
1593 uuid_field("restoration_uuid", false),
1594 uuid_field("checkpoint_uuid", false),
1595 uuid_field("source_generation_uuid", false),
1596 Field::new(
1597 "source_manifest_sha256",
1598 DataType::FixedSizeBinary(32),
1599 false,
1600 ),
1601 uuid_field("prior_current_generation_uuid", false),
1602 uuid_field("restored_generation_uuid", false),
1603 uuid_field("operation_uuid", false),
1604 uuid_field("actor_uuid", true),
1605 Field::new("reason", DataType::Utf8, false),
1606 Field::new(
1607 "restored_at",
1608 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
1609 false,
1610 ),
1611 Field::new("contract_version", DataType::UInt32, false),
1612 ]));
1613 let mut columns = Vec::<ArrayRef>::new();
1614 for value in [
1615 Some(restoration_uuid),
1616 Some(checkpoint_uuid),
1617 Some(source_generation_uuid),
1618 Some(prior_current_generation_uuid),
1619 Some(restored_generation_uuid),
1620 Some(operation_uuid),
1621 actor_uuid,
1622 ] {
1623 let mut builder = FixedSizeBinaryBuilder::with_capacity(1, 16);
1624 match value {
1625 Some(uuid) => builder.append_value(uuid.as_bytes()).map_err(arrow_error)?,
1626 None => builder.append_null(),
1627 }
1628 columns.push(Arc::new(builder.finish()));
1629 }
1630 let mut source_digest = FixedSizeBinaryBuilder::with_capacity(1, 32);
1631 source_digest
1632 .append_value(source_manifest_sha256)
1633 .map_err(arrow_error)?;
1634 columns.insert(3, Arc::new(source_digest.finish()));
1635 columns.push(Arc::new(StringArray::from(vec![reason])));
1636 columns.push(Arc::new(
1637 TimestampMicrosecondArray::from(vec![restored_at]).with_timezone("UTC"),
1638 ));
1639 columns.push(Arc::new(UInt32Array::from(vec![
1640 RESTORATION_CONTRACT_VERSION,
1641 ])));
1642 let batch = RecordBatch::try_new(Arc::clone(&schema), columns).map_err(arrow_error)?;
1643 let properties = WriterProperties::builder()
1644 .set_created_by("graphforge-restoration-transition/1".into())
1645 .build();
1646 let mut writer =
1647 ArrowWriter::try_new(Vec::new(), schema, Some(properties)).map_err(parquet_error)?;
1648 writer.write(&batch).map_err(parquet_error)?;
1649 let bytes = writer.into_inner().map_err(parquet_error)?;
1650 let schema_fingerprint = fingerprint(
1651 CanonicalDomain::Schema,
1652 CANONICAL_CONTRACT_VERSION,
1653 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",
1654 )
1655 .map_err(|error| GfError::Validation(error.to_string()))?;
1656 Ok(ProjectParticipant {
1657 capability_id: crate::WORKSPACE_CAPABILITY_ID.into(),
1658 capability_version: crate::WORKSPACE_CAPABILITY_VERSION,
1659 record_family_id: RESTORATION_FAMILY.into(),
1660 record_version: RESTORATION_CONTRACT_VERSION,
1661 encoding: ProjectParticipantEncoding::Parquet,
1662 schema_fingerprint,
1663 row_count: 1,
1664 bytes,
1665 })
1666}
1667
1668fn uuid_field(name: &str, nullable: bool) -> Field {
1669 Field::new(name, DataType::FixedSizeBinary(16), nullable)
1670}
1671
1672fn arrow_error(error: arrow::error::ArrowError) -> GfError {
1673 let message = format!("restoration Arrow encoding failed: {error}");
1674 drop(error);
1675 GfError::Storage(message)
1676}
1677
1678fn parquet_error(error: parquet::errors::ParquetError) -> GfError {
1679 let message = format!("restoration Parquet encoding failed: {error}");
1680 drop(error);
1681 GfError::Storage(message)
1682}
1683fn utc_micros() -> Result<i64, GfError> {
1684 let value = SystemTime::now()
1685 .duration_since(UNIX_EPOCH)
1686 .map_err(|_| GfError::Storage("system clock is before Unix epoch".into()))?
1687 .as_micros();
1688 i64::try_from(value).map_err(|_| GfError::Storage("UTC microsecond timestamp overflow".into()))
1689}
1690fn registry_serde(error: impl std::fmt::Display) -> GfError {
1691 GfError::Storage(format!("checkpoint registry encoding failed: {error}"))
1692}
1693fn registry_corrupt(message: impl Into<String>) -> GfError {
1694 project_error(ProjectErrorCode::CheckpointRegistryCorrupt, message)
1695}
1696fn project_error(code: ProjectErrorCode, message: impl Into<String>) -> GfError {
1697 GfError::Project {
1698 code,
1699 message: message.into(),
1700 }
1701}
1702fn storage_io(error: impl std::fmt::Display) -> GfError {
1703 GfError::Storage(format!("checkpoint registry I/O failed: {error}"))
1704}
1705
1706#[cfg(test)]
1707mod tests {
1708 use super::*;
1709 use std::collections::BTreeMap;
1710 use std::io::{BufRead, BufReader};
1711 use std::panic::{AssertUnwindSafe, catch_unwind, resume_unwind};
1712 use std::process::{Command, Stdio};
1713 use std::sync::mpsc;
1714 use std::time::Duration;
1715 use tempfile::tempdir;
1716 use wait_timeout::ChildExt;
1717
1718 const TEST_DEADLINE: Duration = Duration::from_secs(1);
1719 const CHILD_DEADLINE: Duration = Duration::from_secs(10);
1720
1721 struct WriterLockHolder {
1722 release: Option<mpsc::SyncSender<()>>,
1723 worker: Option<std::thread::JoinHandle<()>>,
1724 }
1725
1726 impl WriterLockHolder {
1727 fn finish(mut self) -> Result<(), String> {
1728 let release = self
1729 .release
1730 .take()
1731 .ok_or_else(|| "phase=main release sender missing".to_owned())?;
1732 let release_result = release
1733 .send(())
1734 .map_err(|error| format!("phase=main release holder error={error}"));
1735 let join_result = self
1736 .worker
1737 .take()
1738 .ok_or_else(|| "phase=main holder worker missing".to_owned())?
1739 .join()
1740 .map_err(|_| "phase=main holder worker panicked".to_owned());
1741 release_result.and(join_result)
1742 }
1743 }
1744
1745 impl Drop for WriterLockHolder {
1746 fn drop(&mut self) {
1747 if let Some(release) = self.release.take() {
1748 let _ = release.send(());
1749 }
1750 if let Some(worker) = self.worker.take() {
1751 let _ = worker.join();
1752 }
1753 }
1754 }
1755
1756 fn while_writer_lock_is_held<T>(root: &Path, action: impl FnOnce() -> T) -> T {
1757 let writer_path = root.join(LOCKS_DIR).join(WRITER_LOCK_FILE);
1758 let worker_path = writer_path.clone();
1759 let (ready_sender, ready_receiver) = mpsc::sync_channel(0);
1760 let (release_sender, release_receiver) = mpsc::sync_channel(0);
1761 let worker = std::thread::Builder::new()
1762 .name("checkpoint-writer-lock-holder".into())
1763 .spawn(move || {
1764 let writer =
1765 open_regular_lock(&worker_path).expect("phase=holder open writer.lock");
1766 assert!(
1767 FileExt::try_lock_exclusive(&writer).expect("phase=holder acquire writer.lock"),
1768 "phase=holder writer.lock unexpectedly busy"
1769 );
1770 ready_sender.send(()).expect("phase=holder publish ready");
1771 release_receiver.recv().expect("phase=holder await release");
1772 FileExt::unlock(&writer).expect("phase=holder release writer.lock");
1773 })
1774 .expect("phase=holder spawn");
1775 let holder = WriterLockHolder {
1776 release: Some(release_sender),
1777 worker: Some(worker),
1778 };
1779 if let Err(error) = ready_receiver.recv_timeout(TEST_DEADLINE) {
1780 drop(ready_receiver);
1781 let cleanup = holder.finish();
1782 panic!("phase=main await held writer.lock error={error}; cleanup={cleanup:?}");
1783 }
1784 let result = catch_unwind(AssertUnwindSafe(action));
1785 let cleanup = holder.finish();
1786 match result {
1787 Ok(value) => {
1788 cleanup.unwrap_or_else(|error| panic!("phase=main holder cleanup error={error}"));
1789 value
1790 }
1791 Err(original) => {
1792 let _ = cleanup;
1793 resume_unwind(original);
1794 }
1795 }
1796 }
1797
1798 struct BoundedChild {
1799 child: std::process::Child,
1800 reaped: bool,
1801 }
1802
1803 impl BoundedChild {
1804 fn wait(mut self, phase: &str) -> std::process::ExitStatus {
1805 let mut failures = Vec::new();
1806 match self.child.wait_timeout(CHILD_DEADLINE) {
1807 Ok(Some(status)) => {
1808 self.reaped = true;
1809 return status;
1810 }
1811 Ok(None) => failures.push(format!("wait timeout={CHILD_DEADLINE:?}")),
1812 Err(error) => failures.push(format!("wait error={error}")),
1813 }
1814 if let Err(error) = self.child.kill() {
1815 failures.push(format!("kill error={error}"));
1816 }
1817 match self.child.wait_timeout(TEST_DEADLINE) {
1818 Ok(Some(status)) => {
1819 self.reaped = true;
1820 failures.push(format!("killed_status={status}"));
1821 }
1822 Ok(None) => failures.push(format!("reap timeout={TEST_DEADLINE:?}")),
1823 Err(error) => failures.push(format!("reap error={error}")),
1824 }
1825 panic!("phase={phase} child cleanup failures={failures:?}");
1826 }
1827 }
1828
1829 impl Drop for BoundedChild {
1830 fn drop(&mut self) {
1831 if !self.reaped {
1832 let mut failures = Vec::new();
1833 if let Err(error) = self.child.kill() {
1834 failures.push(format!("kill error={error}"));
1835 }
1836 match self.child.wait_timeout(TEST_DEADLINE) {
1837 Ok(Some(_)) => self.reaped = true,
1838 Ok(None) => failures.push(format!("reap timeout={TEST_DEADLINE:?}")),
1839 Err(error) => failures.push(format!("reap error={error}")),
1840 }
1841 if !failures.is_empty() {
1842 eprintln!("phase=drop child cleanup failures={failures:?}");
1843 }
1844 }
1845 }
1846 }
1847
1848 fn recover_checkpoint_pair_after_lock_handoff(root: &Path, phase: &str) {
1849 checkpoint_lock_handoff(root, phase, true);
1850 }
1851
1852 fn preserve_checkpoint_intent_after_lock_handoff(root: &Path, phase: &str) {
1853 checkpoint_lock_handoff(root, phase, false);
1854 }
1855
1856 fn checkpoint_lock_handoff(root: &Path, phase: &str, recover_durable_intent: bool) {
1857 let lock_root = root.join(LOCKS_DIR);
1858 let writer_path = lock_root.join(WRITER_LOCK_FILE);
1859 let checkpoint_path = lock_root.join(CHECKPOINT_LOCK_FILE);
1860 let checkpoint_root = root.join(CHECKPOINTS_DIR);
1861 let worker_writer_path = writer_path.clone();
1862 let worker_checkpoint_path = checkpoint_path.clone();
1863 let (sender, receiver) = mpsc::sync_channel(0);
1864 std::thread::Builder::new()
1865 .name("checkpoint-lock-handoff-recovery".into())
1866 .spawn(move || {
1867 let result = (|| {
1868 let writer = open_regular_lock(&worker_writer_path)
1869 .map_err(|error| format!("open writer.lock failed: {error}"))?;
1870 FileExt::lock_exclusive(&writer)
1871 .map_err(|error| format!("acquire writer.lock failed: {error}"))?;
1872
1873 let checkpoint = match open_regular_lock(&worker_checkpoint_path) {
1874 Ok(checkpoint) => checkpoint,
1875 Err(error) => {
1876 let writer_unlock = FileExt::unlock(&writer);
1877 return Err(format!(
1878 "open checkpoints.lock failed: {error}; \
1879 writer_unlock={writer_unlock:?}"
1880 ));
1881 }
1882 };
1883 if let Err(error) = FileExt::lock_exclusive(&checkpoint) {
1884 let writer_unlock = FileExt::unlock(&writer);
1885 return Err(format!(
1886 "acquire checkpoints.lock failed: {error}; writer_unlock={writer_unlock:?}"
1887 ));
1888 }
1889
1890 let recovery = if recover_durable_intent
1891 && checkpoint_root.join(INTENT_FILE).exists()
1892 {
1893 recover_pair(&checkpoint_root)
1894 .map_err(|error| format!("recover durable checkpoint intent failed: {error}"))
1895 } else {
1896 Ok(())
1897 };
1898 let checkpoint_unlock = FileExt::unlock(&checkpoint)
1899 .map_err(|error| format!("unlock checkpoints.lock failed: {error}"));
1900 let writer_unlock = FileExt::unlock(&writer)
1901 .map_err(|error| format!("unlock writer.lock failed: {error}"));
1902
1903 recovery?;
1904 checkpoint_unlock?;
1905 writer_unlock
1906 })();
1907 let _ = sender.send(result);
1908 })
1909 .unwrap();
1910 match receiver.recv_timeout(Duration::from_secs(1)) {
1911 Ok(Ok(())) => {}
1912 Ok(Err(error)) => panic!(
1913 "checkpoint lock handoff/recovery failed at {phase}; writer_path={}; \
1914 checkpoint_path={}: {error}",
1915 writer_path.display(),
1916 checkpoint_path.display()
1917 ),
1918 Err(error) => panic!(
1919 "checkpoint lock handoff/recovery timed out at {phase}; writer_path={}; \
1920 checkpoint_path={}; timeout=1s; channel={error}",
1921 writer_path.display(),
1922 checkpoint_path.display()
1923 ),
1924 }
1925 }
1926
1927 fn publish_clone(root: &Path) -> Uuid {
1928 let selected = crate::resolve_project_generation(root).unwrap();
1929 let capabilities = selected
1930 .capabilities()
1931 .into_iter()
1932 .map(|entry| crate::ProjectCapability {
1933 capability_id: entry.capability_id,
1934 capability_version: entry.capability_version,
1935 })
1936 .collect();
1937 let participants = selected
1938 .participant_snapshots()
1939 .unwrap()
1940 .into_iter()
1941 .map(|entry| crate::ProjectParticipant {
1942 capability_id: entry.capability_id,
1943 capability_version: entry.capability_version,
1944 record_family_id: entry.record_family_id,
1945 record_version: entry.record_version,
1946 encoding: match entry.encoding.as_str() {
1947 "arrow" => crate::ProjectParticipantEncoding::Arrow,
1948 "json" => crate::ProjectParticipantEncoding::Json,
1949 "parquet" => crate::ProjectParticipantEncoding::Parquet,
1950 other => panic!("unexpected participant encoding {other}"),
1951 },
1952 schema_fingerprint: entry.schema_fingerprint,
1953 row_count: entry.row_count,
1954 bytes: entry.bytes,
1955 })
1956 .collect();
1957 let generation_uuid = Uuid::now_v7();
1958 let request = crate::ProjectGenerationRequest {
1959 transaction_uuid: Uuid::now_v7(),
1960 generation_uuid,
1961 capabilities,
1962 participants,
1963 };
1964 let crate::ProjectStageOutcome::Staged(staged) =
1965 crate::stage_project_generation(root, &request).unwrap()
1966 else {
1967 panic!("fresh publication unexpectedly replayed");
1968 };
1969 staged
1970 .validate(|_| Ok(()), |_, _| Ok(()))
1971 .unwrap()
1972 .publish()
1973 .unwrap();
1974 generation_uuid
1975 }
1976
1977 fn create_request(operation_uuid: Uuid, name: &str) -> CheckpointCreateRequest {
1978 CheckpointCreateRequest {
1979 operation_uuid,
1980 name: name.into(),
1981 description: Some("release candidate".into()),
1982 actor_uuid: Some(Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000aa").unwrap()),
1983 }
1984 }
1985
1986 fn write_raw_registry(root: &Path, registry: &Registry) {
1987 let checkpoint_root = root.join(CHECKPOINTS_DIR);
1988 let mut bytes = serde_json::to_vec(registry).unwrap();
1989 bytes.push(b'\n');
1990 fs::write(checkpoint_root.join(REGISTRY_FILE), &bytes).unwrap();
1991 fs::write(
1992 checkpoint_root.join(CHECKSUM_FILE),
1993 format!("{}\n", hex(&Sha256::digest(&bytes).into())),
1994 )
1995 .unwrap();
1996 }
1997
1998 #[test]
1999 fn revert_identity_matches_frozen_golden_vector() {
2000 let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000003").unwrap();
2001 let checkpoint = Uuid::parse_str("4084179c-38db-8b6b-9b6e-c0b0a855e002").unwrap();
2002 let source = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000b0").unwrap();
2003 let prior = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000d0").unwrap();
2004 let actor = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000aa").unwrap();
2005 let source_digest = [0x11; 32];
2006 let request_digest = revert_request_digest(
2007 operation,
2008 "Release 1.0",
2009 checkpoint,
2010 source,
2011 source_digest,
2012 "restore release candidate",
2013 Some(actor),
2014 );
2015 assert_eq!(
2016 hex(&request_digest),
2017 "dff3755629942d1189060b117cb70dc864428fcc3b28d5c2d22d3924c3690e93"
2018 );
2019 let transaction = revert_transaction_uuid(operation);
2020 assert_eq!(
2021 transaction.to_string(),
2022 "908d637b-d6e6-8508-919e-d4e708e037b2"
2023 );
2024 assert_eq!(
2025 restoration_uuid(operation, request_digest).to_string(),
2026 "9e1f160c-badb-80c9-beaa-ae580910bf8a"
2027 );
2028 assert_eq!(
2029 restored_generation_uuid(
2030 transaction,
2031 checkpoint,
2032 source,
2033 source_digest,
2034 prior,
2035 1_720_000_000_123_456,
2036 request_digest,
2037 )
2038 .to_string(),
2039 "5dc02888-2064-8892-a0e1-2c00968ba0cc"
2040 );
2041 }
2042
2043 #[test]
2044 fn revert_publishes_child_preserves_registry_and_replays_after_delete() {
2045 let directory = tempdir().unwrap();
2046 crate::open_or_initialize_project(directory.path()).unwrap();
2047 let source = crate::resolve_project_generation(directory.path()).unwrap();
2048 let created = create_checkpoint(
2049 directory.path(),
2050 &create_request(Uuid::from_u128(40), "Before"),
2051 )
2052 .unwrap();
2053 let prior_current = publish_clone(directory.path());
2054 let request = CheckpointRevertRequest {
2055 operation_uuid: Uuid::from_u128(41),
2056 name: "Before".into(),
2057 reason: " restore known state ".into(),
2058 actor_uuid: None,
2059 };
2060 let (receipt, restored) = revert_checkpoint(
2061 directory.path(),
2062 &request,
2063 || Ok(1_720_000_000_123_456),
2064 |_| Ok(()),
2065 )
2066 .unwrap();
2067 assert_eq!(restored.parent_generation_uuid(), Some(prior_current));
2068 assert_eq!(receipt.source_generation_uuid, source.generation_uuid());
2069 assert_eq!(receipt.prior_current_generation_uuid, Some(prior_current));
2070 assert_eq!(receipt.registry_revision, created.registry_revision);
2071 assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2072 let restoration_count = restored
2073 .participant_descriptors()
2074 .unwrap()
2075 .iter()
2076 .filter(|row| row.record_family_id == RESTORATION_FAMILY)
2077 .count();
2078 assert_eq!(restoration_count, 1);
2079
2080 delete_checkpoint(
2081 directory.path(),
2082 &CheckpointDeleteRequest {
2083 operation_uuid: Uuid::from_u128(42),
2084 name: "Before".into(),
2085 actor_uuid: None,
2086 },
2087 )
2088 .unwrap();
2089 let (replay, replayed_generation) = revert_checkpoint(
2090 directory.path(),
2091 &request,
2092 || panic!("published replay sampled clock"),
2093 |_| Ok(()),
2094 )
2095 .unwrap();
2096 assert_eq!(replay, receipt);
2097 assert_eq!(replay.prior_current_generation_uuid, Some(prior_current));
2098 assert_eq!(
2099 replayed_generation.generation_uuid(),
2100 restored.generation_uuid()
2101 );
2102 checkpoint_lock_handoff(
2103 directory.path(),
2104 "action=revert published-replay return",
2105 false,
2106 );
2107
2108 let mut conflict = request;
2109 conflict.reason = "different".into();
2110 let conflict_error =
2111 revert_checkpoint(directory.path(), &conflict, || Ok(0), |_| Ok(())).unwrap_err();
2112 assert_eq!(conflict_error.code(), "GF_IDEMPOTENCY_CONFLICT");
2113 checkpoint_lock_handoff(
2114 directory.path(),
2115 "action=revert published-replay conflict return",
2116 false,
2117 );
2118 }
2119
2120 #[test]
2121 fn revert_replay_lock_handoff_fails_closed_with_stable_storage_errors() {
2122 let checkpoint_error = finish_revert_replay_lock_handoff(
2123 Err(std::io::Error::other("checkpoint unlock failed")),
2124 Ok(()),
2125 )
2126 .unwrap_err();
2127 assert_eq!(checkpoint_error.code(), "GF_IO");
2128 assert_eq!(
2129 checkpoint_error.to_string(),
2130 "storage error: checkpoint revert replay lock handoff failed at checkpoints.lock: checkpoint unlock failed"
2131 );
2132
2133 let writer_error = finish_revert_replay_lock_handoff(
2134 Ok(()),
2135 Err(std::io::Error::other("writer unlock failed")),
2136 )
2137 .unwrap_err();
2138 assert_eq!(writer_error.code(), "GF_IO");
2139 assert_eq!(
2140 writer_error.to_string(),
2141 "storage error: checkpoint revert replay lock handoff failed at writer.lock: writer unlock failed"
2142 );
2143 }
2144
2145 #[test]
2146 fn revert_validation_failure_preserves_prior_current() {
2147 let directory = tempdir().unwrap();
2148 crate::open_or_initialize_project(directory.path()).unwrap();
2149 create_checkpoint(
2150 directory.path(),
2151 &create_request(Uuid::from_u128(50), "Before"),
2152 )
2153 .unwrap();
2154 let prior = publish_clone(directory.path());
2155 let error = revert_checkpoint(
2156 directory.path(),
2157 &CheckpointRevertRequest {
2158 operation_uuid: Uuid::from_u128(51),
2159 name: "Before".into(),
2160 reason: "must fail closed".into(),
2161 actor_uuid: None,
2162 },
2163 || Ok(1_720_000_000_123_456),
2164 |_| Err(GfError::Validation("injected composite failure".into())),
2165 )
2166 .unwrap_err();
2167 assert_eq!(error.code(), "GF_VALIDATION");
2168 checkpoint_lock_handoff(
2169 directory.path(),
2170 "action=revert validation-error return",
2171 false,
2172 );
2173 assert_eq!(
2174 crate::resolve_project_generation(directory.path())
2175 .unwrap()
2176 .generation_uuid(),
2177 prior
2178 );
2179 assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2180 }
2181
2182 #[cfg(unix)]
2183 #[test]
2184 fn mutation_lock_guard_unlocks_checkpoint_with_retained_duplicate_open() {
2185 let directory = tempdir().unwrap();
2186 crate::open_or_initialize_project(directory.path()).unwrap();
2187 let locks = acquire_mutation_locks(directory.path()).unwrap();
2188 let retained = locks.checkpoint.as_ref().unwrap().try_clone().unwrap();
2189 drop(locks);
2190
2191 let checkpoint =
2192 open_regular_lock(&directory.path().join(LOCKS_DIR).join(CHECKPOINT_LOCK_FILE))
2193 .unwrap();
2194 assert!(FileExt::try_lock_exclusive(&checkpoint).unwrap());
2195 FileExt::unlock(&checkpoint).unwrap();
2196 drop(retained);
2197 }
2198
2199 #[test]
2200 fn create_list_delete_and_replays_are_deterministic() {
2201 let directory = tempdir().unwrap();
2202 crate::open_or_initialize_project(directory.path()).unwrap();
2203 let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000001").unwrap();
2204 let request = create_request(operation, "Release 1.0");
2205 assert_eq!(
2206 hex(&create_request_digest(&request, "Release 1.0")),
2207 "01c7bf2f2c443d85d31ff80fef4a36484e31402213e8d145371867bdb2addbe8"
2208 );
2209 let created = create_checkpoint(directory.path(), &request).unwrap();
2210 assert_eq!(
2211 created.checkpoint_uuid,
2212 Uuid::parse_str("4084179c-38db-8b6b-9b6e-c0b0a855e002").unwrap()
2213 );
2214 let replayed = create_checkpoint(directory.path(), &request).unwrap();
2215 assert_eq!(created, replayed);
2216 assert_eq!(created.registry_revision, 1);
2217
2218 let rows = list_checkpoints(directory.path()).unwrap();
2219 assert_eq!(rows.len(), 1);
2220 assert_eq!(rows[0].checkpoint_uuid, created.checkpoint_uuid);
2221 assert_eq!(rows[0].generation_uuid, created.source_generation_uuid);
2222
2223 let delete = CheckpointDeleteRequest {
2224 operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000002").unwrap(),
2225 name: "Release 1.0".into(),
2226 actor_uuid: request.actor_uuid,
2227 };
2228 assert_eq!(
2229 hex(&delete_request_digest(&delete, "Release 1.0")),
2230 "9e6e15801f66ea4f7f58755c505fc1964e48ab87625459f388e2c23659135bb3"
2231 );
2232 let deleted = delete_checkpoint(directory.path(), &delete).unwrap();
2233 assert_eq!(
2234 deleted,
2235 delete_checkpoint(directory.path(), &delete).unwrap()
2236 );
2237 assert_eq!(deleted.registry_revision, 2);
2238 assert!(list_checkpoints(directory.path()).unwrap().is_empty());
2239 assert_eq!(
2240 created,
2241 create_checkpoint(directory.path(), &request).unwrap()
2242 );
2243 }
2244
2245 #[test]
2246 fn identity_is_stable_across_independent_projects() {
2247 let first = tempdir().unwrap();
2248 let second = tempdir().unwrap();
2249 crate::open_or_initialize_project(first.path()).unwrap();
2250 crate::open_or_initialize_project(second.path()).unwrap();
2251 let request = create_request(
2252 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000010").unwrap(),
2253 "Stable",
2254 );
2255 let left = create_checkpoint(first.path(), &request).unwrap();
2256 let right = create_checkpoint(second.path(), &request).unwrap();
2257 assert_eq!(left.checkpoint_uuid, right.checkpoint_uuid);
2258 }
2259
2260 #[test]
2261 fn opened_checkpoint_generation_remains_pinned_after_delete() {
2262 let directory = tempdir().unwrap();
2263 crate::open_or_initialize_project(directory.path()).unwrap();
2264 let created = create_checkpoint(
2265 directory.path(),
2266 &create_request(Uuid::now_v7(), "Pinned View"),
2267 )
2268 .unwrap();
2269 let (row, opened) = open_checkpoint_generation(directory.path(), "Pinned View").unwrap();
2270 assert_eq!(row.checkpoint_uuid, created.checkpoint_uuid);
2271 assert_eq!(opened.generation_uuid(), created.source_generation_uuid);
2272 delete_checkpoint(
2273 directory.path(),
2274 &CheckpointDeleteRequest {
2275 operation_uuid: Uuid::now_v7(),
2276 name: "Pinned View".into(),
2277 actor_uuid: None,
2278 },
2279 )
2280 .unwrap();
2281 assert_eq!(
2282 open_checkpoint_generation(directory.path(), "Pinned View")
2283 .unwrap_err()
2284 .code(),
2285 "GF_CHECKPOINT_NOT_FOUND"
2286 );
2287 assert_eq!(opened.generation_uuid(), created.source_generation_uuid);
2288 assert!(opened.participant_snapshots().is_ok());
2289 }
2290
2291 #[test]
2292 fn conflicts_names_and_corruption_fail_closed() {
2293 let directory = tempdir().unwrap();
2294 crate::open_or_initialize_project(directory.path()).unwrap();
2295 let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000020").unwrap();
2296 create_checkpoint(directory.path(), &create_request(operation, "Safe.Name")).unwrap();
2297
2298 let conflict =
2299 create_checkpoint(directory.path(), &create_request(operation, "Other")).unwrap_err();
2300 assert_eq!(conflict.code(), "GF_IDEMPOTENCY_CONFLICT");
2301 let exists = create_checkpoint(
2302 directory.path(),
2303 &create_request(
2304 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000021").unwrap(),
2305 "Safe.Name",
2306 ),
2307 )
2308 .unwrap_err();
2309 assert_eq!(exists.code(), "GF_CHECKPOINT_EXISTS");
2310 for invalid in ["", "../escape", "two spaces", " e", "e ", ".", ".."] {
2311 let error =
2312 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), invalid))
2313 .unwrap_err();
2314 assert_eq!(error.code(), "GF_VALIDATION", "name={invalid:?}");
2315 }
2316
2317 fs::write(
2318 directory.path().join(CHECKPOINTS_DIR).join(CHECKSUM_FILE),
2319 b"0000000000000000000000000000000000000000000000000000000000000000\n",
2320 )
2321 .unwrap();
2322 let error = list_checkpoints(directory.path()).unwrap_err();
2323 assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2324 }
2325
2326 #[test]
2327 fn checksummed_but_impossible_registry_identities_fail_closed() {
2328 let directory = tempdir().unwrap();
2329 crate::open_or_initialize_project(directory.path()).unwrap();
2330 create_checkpoint(
2331 directory.path(),
2332 &create_request(Uuid::now_v7(), "Tampered"),
2333 )
2334 .unwrap();
2335 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2336 let mut registry = read_registry(&checkpoint_root).unwrap();
2337 registry.active[0].checkpoint_uuid = Uuid::now_v7();
2338 write_raw_registry(directory.path(), ®istry);
2339 assert_eq!(
2340 list_checkpoints(directory.path()).unwrap_err().code(),
2341 "GF_CHECKPOINT_REGISTRY_CORRUPT"
2342 );
2343
2344 let directory = tempdir().unwrap();
2345 crate::open_or_initialize_project(directory.path()).unwrap();
2346 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Overlap")).unwrap();
2347 delete_checkpoint(
2348 directory.path(),
2349 &CheckpointDeleteRequest {
2350 operation_uuid: Uuid::now_v7(),
2351 name: "Overlap".into(),
2352 actor_uuid: None,
2353 },
2354 )
2355 .unwrap();
2356 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2357 let mut registry = read_registry(&checkpoint_root).unwrap();
2358 let row = &mut registry.tombstones[0];
2359 row.delete_operation_uuid = row.create_operation_uuid;
2360 row.delete_request_sha256 = hex(&delete_request_digest_values(
2361 row.delete_operation_uuid,
2362 &row.name,
2363 row.deleted_by,
2364 ));
2365 write_raw_registry(directory.path(), ®istry);
2366 assert_eq!(
2367 list_checkpoints(directory.path()).unwrap_err().code(),
2368 "GF_CHECKPOINT_REGISTRY_CORRUPT"
2369 );
2370 }
2371
2372 #[test]
2373 fn exact_input_bounds_and_writer_lock_are_enforced() {
2374 assert!(validate_name(&"a".repeat(MAX_NAME_BYTES)).is_ok());
2375 assert_eq!(
2376 validate_name(&"a".repeat(MAX_NAME_BYTES + 1))
2377 .unwrap_err()
2378 .code(),
2379 "GF_VALIDATION"
2380 );
2381 assert!(validate_description(Some(&"d".repeat(MAX_DESCRIPTION_BYTES))).is_ok());
2382 assert_eq!(
2383 validate_description(Some(&"d".repeat(MAX_DESCRIPTION_BYTES + 1)))
2384 .unwrap_err()
2385 .code(),
2386 "GF_VALIDATION"
2387 );
2388
2389 let directory = tempdir().unwrap();
2390 let selected = crate::open_or_initialize_project(directory.path()).unwrap();
2391 let root = selected.container_root().to_owned();
2392 let locks = acquire_mutation_locks(&root).unwrap();
2393 let error = create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Busy"))
2394 .unwrap_err();
2395 assert_eq!(error.code(), "GF_WRITER_BUSY");
2396 drop(locks);
2397 }
2398
2399 #[cfg(unix)]
2400 #[test]
2401 fn linked_registry_surfaces_fail_closed_without_following_targets() {
2402 use std::os::unix::fs::symlink;
2403
2404 for hard in [false, true] {
2405 let directory = tempdir().unwrap();
2406 crate::open_or_initialize_project(directory.path()).unwrap();
2407 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Linked")).unwrap();
2408 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2409 let checksum = checkpoint_root.join(CHECKSUM_FILE);
2410 let external = directory.path().join("external-checksum");
2411 fs::rename(&checksum, &external).unwrap();
2412 if hard {
2413 fs::hard_link(&external, &checksum).unwrap();
2414 } else {
2415 symlink(&external, &checksum).unwrap();
2416 }
2417 let external_before = fs::read(&external).unwrap();
2418 let error = while_writer_lock_is_held(directory.path(), || {
2419 list_checkpoints(directory.path()).unwrap_err()
2420 });
2421 assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2422 assert_eq!(fs::read(&external).unwrap(), external_before);
2423 }
2424 }
2425
2426 #[test]
2427 fn no_intent_registry_corruption_wins_over_writer_contention() {
2428 let directory = tempdir().unwrap();
2429 crate::open_or_initialize_project(directory.path()).unwrap();
2430 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Corrupt")).unwrap();
2431 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2432 fs::write(
2433 checkpoint_root.join(CHECKSUM_FILE),
2434 b"not-the-registry-digest\n",
2435 )
2436 .unwrap();
2437
2438 let expected = read_registry(&checkpoint_root).unwrap_err();
2439 let error = while_writer_lock_is_held(directory.path(), || {
2440 list_checkpoints(directory.path()).unwrap_err()
2441 });
2442 assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2443 assert_eq!(error.to_string(), expected.to_string());
2444 }
2445
2446 #[test]
2447 fn intent_recovery_contention_preserves_writer_busy_and_intent() {
2448 let directory = tempdir().unwrap();
2449 crate::open_or_initialize_project(directory.path()).unwrap();
2450 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2451 let child = Command::new(std::env::current_exe().unwrap())
2452 .args([
2453 "--exact",
2454 "project_checkpoints::tests::checkpoint_failpoint_helper",
2455 "--ignored",
2456 ])
2457 .env(
2458 "GRAPHFORGE_PROJECT_FAILPOINTS",
2459 "graphforge-internal-subprocess-v1",
2460 )
2461 .env(
2462 "GRAPHFORGE_PROJECT_FAILPOINT",
2463 "checkpoint.registry.before_replace",
2464 )
2465 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2466 .spawn()
2467 .unwrap();
2468 let status = BoundedChild {
2469 child,
2470 reaped: false,
2471 }
2472 .wait("intent-recovery-contention failpoint=checkpoint.registry.before_replace");
2473 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2474 let intent_path = checkpoint_root.join(INTENT_FILE);
2475 let intent = fs::read(&intent_path).unwrap();
2476 let staged = fs::read_dir(&checkpoint_root)
2477 .unwrap()
2478 .filter_map(Result::ok)
2479 .filter(|entry| {
2480 entry
2481 .file_name()
2482 .to_string_lossy()
2483 .starts_with(".registry.")
2484 })
2485 .map(|entry| (entry.file_name(), fs::read(entry.path()).unwrap()))
2486 .collect::<BTreeMap<_, _>>();
2487 assert_eq!(staged.len(), 2);
2488
2489 let error = while_writer_lock_is_held(directory.path(), || {
2490 list_checkpoints(directory.path()).unwrap_err()
2491 });
2492 assert_eq!(error.code(), "GF_WRITER_BUSY");
2493 assert_eq!(fs::read(intent_path).unwrap(), intent);
2494 let staged_after = fs::read_dir(&checkpoint_root)
2495 .unwrap()
2496 .filter_map(Result::ok)
2497 .filter(|entry| {
2498 entry
2499 .file_name()
2500 .to_string_lossy()
2501 .starts_with(".registry.")
2502 })
2503 .map(|entry| (entry.file_name(), fs::read(entry.path()).unwrap()))
2504 .collect::<BTreeMap<_, _>>();
2505 assert_eq!(staged_after, staged);
2506 }
2507
2508 #[cfg(unix)]
2509 #[test]
2510 fn linked_project_root_is_rejected_before_checkpoint_access() {
2511 use std::os::unix::fs::symlink;
2512
2513 let directory = tempdir().unwrap();
2514 let project = directory.path().join("project");
2515 fs::create_dir(&project).unwrap();
2516 crate::open_or_initialize_project(&project).unwrap();
2517 let linked = directory.path().join("linked-project");
2518 symlink(&project, &linked).unwrap();
2519 assert_eq!(
2520 list_checkpoints(&linked).unwrap_err().code(),
2521 "GF_UNSUPPORTED_PROJECT_FORMAT"
2522 );
2523 }
2524
2525 #[test]
2526 fn checkpoint_pin_and_open_lease_control_recovery_cleanup() {
2527 let directory = tempdir().unwrap();
2528 crate::open_or_initialize_project(directory.path()).unwrap();
2529 let pinned_generation = publish_clone(directory.path());
2530 let created =
2531 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Pinned")).unwrap();
2532 assert_eq!(created.source_generation_uuid, pinned_generation);
2533 for _ in 0..4 {
2534 publish_clone(directory.path());
2535 }
2536 crate::recover_project_transactions(directory.path()).unwrap();
2537 recover_checkpoint_pair_after_lock_handoff(
2538 directory.path(),
2539 "action=delete-pinned-checkpoint recovery-complete",
2540 );
2541 let generation_path = directory
2542 .path()
2543 .join(crate::project_publication::GENERATIONS_DIR)
2544 .join(pinned_generation.hyphenated().to_string());
2545 assert!(
2546 generation_path.exists(),
2547 "active checkpoint lost its generation"
2548 );
2549
2550 let lease =
2551 crate::project_publication::open_regular_lock(&generation_path.join("lease.lock"))
2552 .unwrap();
2553 FileExt::lock_shared(&lease).unwrap();
2554 delete_checkpoint(
2555 directory.path(),
2556 &CheckpointDeleteRequest {
2557 operation_uuid: Uuid::now_v7(),
2558 name: "Pinned".into(),
2559 actor_uuid: None,
2560 },
2561 )
2562 .unwrap();
2563 crate::recover_project_transactions(directory.path()).unwrap();
2564 assert!(generation_path.exists(), "an open lease was invalidated");
2565 FileExt::unlock(&lease).unwrap();
2566 drop(lease);
2567 crate::recover_project_transactions(directory.path()).unwrap();
2568 assert!(
2569 !generation_path.exists(),
2570 "deleted pin did not permit later GC"
2571 );
2572 }
2573
2574 #[test]
2575 fn checkpoint_pin_survives_process_restart() {
2576 let directory = tempdir().unwrap();
2577 crate::open_or_initialize_project(directory.path()).unwrap();
2578 let pinned_generation = publish_clone(directory.path());
2579 create_checkpoint(
2580 directory.path(),
2581 &create_request(Uuid::from_u128(70), "Restart Pin"),
2582 )
2583 .unwrap();
2584 for _ in 0..4 {
2585 publish_clone(directory.path());
2586 }
2587
2588 let mut child = Command::new(std::env::current_exe().unwrap())
2589 .args([
2590 "--exact",
2591 "project_checkpoints::tests::checkpoint_failpoint_helper",
2592 "--ignored",
2593 "--nocapture",
2594 ])
2595 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2596 .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", "hold-open")
2597 .stdin(Stdio::piped())
2598 .stdout(Stdio::piped())
2599 .spawn()
2600 .unwrap();
2601 let mut output = BufReader::new(child.stdout.take().unwrap());
2602 let mut ready = String::new();
2603 while ready != "ready\n" {
2604 ready.clear();
2605 assert_ne!(
2606 output.read_line(&mut ready).unwrap(),
2607 0,
2608 "child exited before ready"
2609 );
2610 }
2611
2612 delete_checkpoint(
2613 directory.path(),
2614 &CheckpointDeleteRequest {
2615 operation_uuid: Uuid::from_u128(71),
2616 name: "Restart Pin".into(),
2617 actor_uuid: None,
2618 },
2619 )
2620 .unwrap();
2621 crate::recover_project_transactions(directory.path()).unwrap();
2622 let generation_path = directory
2623 .path()
2624 .join(crate::project_publication::GENERATIONS_DIR)
2625 .join(pinned_generation.hyphenated().to_string());
2626 assert!(
2627 generation_path.exists(),
2628 "subprocess lease was not retained"
2629 );
2630
2631 child.stdin.take().unwrap().write_all(b"release\n").unwrap();
2632 assert!(child.wait().unwrap().success());
2633 crate::recover_project_transactions(directory.path()).unwrap();
2634 assert!(
2635 !generation_path.exists(),
2636 "generation survived after the restarted reader released its lease"
2637 );
2638 }
2639
2640 #[test]
2641 fn checkpoint_cleanup_removes_all_transient_resources() {
2642 let directory = tempdir().unwrap();
2643 crate::open_or_initialize_project(directory.path()).unwrap();
2644 create_checkpoint(
2645 directory.path(),
2646 &create_request(Uuid::from_u128(80), "Cleanup"),
2647 )
2648 .unwrap();
2649 publish_clone(directory.path());
2650 delete_checkpoint(
2651 directory.path(),
2652 &CheckpointDeleteRequest {
2653 operation_uuid: Uuid::from_u128(81),
2654 name: "Cleanup".into(),
2655 actor_uuid: None,
2656 },
2657 )
2658 .unwrap();
2659 crate::recover_project_transactions(directory.path()).unwrap();
2660
2661 let selected = crate::resolve_project_generation(directory.path()).unwrap();
2662 let root = selected.container_root().to_owned();
2663 drop(selected);
2664 let checkpoint_root = root.join(CHECKPOINTS_DIR);
2665 let checkpoint_entries = fs::read_dir(&checkpoint_root)
2666 .unwrap()
2667 .map(|entry| entry.unwrap().file_name().into_string().unwrap())
2668 .collect::<BTreeSet<_>>();
2669 assert_eq!(
2670 checkpoint_entries,
2671 BTreeSet::from([REGISTRY_FILE.into(), CHECKSUM_FILE.into()]),
2672 "checkpoint transaction staging leaked"
2673 );
2674 let trash = root.join("trash");
2675 assert!(
2676 !trash.exists() || fs::read_dir(&trash).unwrap().next().is_none(),
2677 "recovery trash was not emptied"
2678 );
2679 assert!(
2680 !root.join("cache").exists(),
2681 "checkpoint lifecycle leaked process-cache state to disk"
2682 );
2683 for entry in fs::read_dir(root.join(crate::project_publication::GENERATIONS_DIR)).unwrap() {
2684 let path = entry.unwrap().path();
2685 let name = path.file_name().unwrap().to_str().unwrap();
2686 Uuid::parse_str(name).expect("generation staging entry leaked");
2687 let lease = open_regular_lock(&path.join("lease.lock")).unwrap();
2688 assert!(FileExt::try_lock_exclusive(&lease).unwrap());
2689 FileExt::unlock(&lease).unwrap();
2690 }
2691 let lock_root = root.join(LOCKS_DIR);
2692 for name in [WRITER_LOCK_FILE, CHECKPOINT_LOCK_FILE] {
2693 let lock = open_regular_lock(&lock_root.join(name)).unwrap();
2694 assert!(FileExt::try_lock_exclusive(&lock).unwrap(), "{name} leaked");
2695 FileExt::unlock(&lock).unwrap();
2696 }
2697 }
2698
2699 #[test]
2700 #[ignore = "subprocess failpoint helper"]
2701 fn checkpoint_failpoint_helper() {
2702 let root = std::env::var("GRAPHFORGE_CHECKPOINT_TEST_ROOT").unwrap();
2703 let action = std::env::var("GRAPHFORGE_CHECKPOINT_TEST_ACTION");
2704 if action.as_deref() == Ok("hold-open") {
2705 let (_, opened) = open_checkpoint_generation(root, "Restart Pin").unwrap();
2706 println!("ready");
2707 std::io::stdout().flush().unwrap();
2708 let mut release = String::new();
2709 std::io::stdin().read_line(&mut release).unwrap();
2710 assert_eq!(release, "release\n");
2711 assert!(opened.participant_snapshots().is_ok());
2712 } else if action.as_deref() == Ok("revert") {
2713 revert_checkpoint(
2714 root,
2715 &CheckpointRevertRequest {
2716 operation_uuid: Uuid::from_u128(61),
2717 name: "Base".into(),
2718 reason: "crash recovery".into(),
2719 actor_uuid: None,
2720 },
2721 || Ok(1_720_000_000_123_456),
2722 |_| Ok(()),
2723 )
2724 .unwrap();
2725 } else if action.as_deref() == Ok("delete") {
2726 delete_checkpoint(
2727 root,
2728 &CheckpointDeleteRequest {
2729 operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000031")
2730 .unwrap(),
2731 name: "Base".into(),
2732 actor_uuid: None,
2733 },
2734 )
2735 .unwrap();
2736 } else {
2737 let request = create_request(
2738 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
2739 "Crash",
2740 );
2741 create_checkpoint(root, &request).unwrap();
2742 }
2743 }
2744
2745 #[test]
2746 fn revert_publication_failpoint_matrix() {
2747 let failpoints = [
2748 ("project.after_journal_preparing", false),
2749 ("project.after_participant_dir_fsync", false),
2750 ("project.after_journal_staged", false),
2751 ("project.after_domain_validation", false),
2752 ("project.after_composite_validation", false),
2753 ("project.after_journal_validated", false),
2754 ("project.after_manifest_write", false),
2755 ("project.after_manifest_fsync", false),
2756 ("project.after_generation_dir_fsync", false),
2757 ("project.after_journal_durable", false),
2758 ("project.after_current_temp_write", false),
2759 ("project.after_current_temp_fsync", false),
2760 ("project.before_current_replace", false),
2761 ("project.after_current_replace", true),
2762 ("project.after_root_fsync", true),
2763 ("project.after_journal_published", true),
2764 ];
2765 for (failpoint, committed) in failpoints {
2766 let directory = tempdir().unwrap();
2767 crate::open_or_initialize_project(directory.path()).unwrap();
2768 create_checkpoint(
2769 directory.path(),
2770 &create_request(Uuid::from_u128(60), "Base"),
2771 )
2772 .unwrap();
2773 let prior = publish_clone(directory.path());
2774 let status = Command::new(std::env::current_exe().unwrap())
2775 .args([
2776 "--exact",
2777 "project_checkpoints::tests::checkpoint_failpoint_helper",
2778 "--ignored",
2779 ])
2780 .env(
2781 "GRAPHFORGE_PROJECT_FAILPOINTS",
2782 "graphforge-internal-subprocess-v1",
2783 )
2784 .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
2785 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2786 .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", "revert")
2787 .status()
2788 .unwrap();
2789 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2790 recover_checkpoint_pair_after_lock_handoff(
2791 directory.path(),
2792 &format!("action=revert failpoint={failpoint} committed={committed}"),
2793 );
2794 crate::recover_project_transactions(directory.path()).unwrap();
2795 let recovered = crate::resolve_project_generation(directory.path()).unwrap();
2796 assert_eq!(
2797 recovered.generation_uuid() != prior,
2798 committed,
2799 "{failpoint}"
2800 );
2801 recover_checkpoint_pair_after_lock_handoff(
2802 directory.path(),
2803 &format!(
2804 "action=revert parent-recovery-complete failpoint={failpoint} \
2805 committed={committed}"
2806 ),
2807 );
2808
2809 let (receipt, replayed) = revert_checkpoint(
2810 directory.path(),
2811 &CheckpointRevertRequest {
2812 operation_uuid: Uuid::from_u128(61),
2813 name: "Base".into(),
2814 reason: "crash recovery".into(),
2815 actor_uuid: None,
2816 },
2817 || Ok(1_720_000_000_123_456),
2818 |_| Ok(()),
2819 )
2820 .unwrap();
2821 assert_eq!(
2822 receipt.result_generation_uuid,
2823 Some(replayed.generation_uuid())
2824 );
2825 recover_checkpoint_pair_after_lock_handoff(
2826 directory.path(),
2827 &format!(
2828 "action=revert replay-complete failpoint={failpoint} committed={committed}"
2829 ),
2830 );
2831 assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2832 }
2833 }
2834
2835 #[test]
2836 fn registry_failpoints_recover_exact_previous_or_next_revision() {
2837 for (failpoint, committed) in [
2838 ("checkpoint.registry.after_intent_file_fsync", false),
2839 ("checkpoint.registry.after_file_fsync", false),
2840 ("checkpoint.registry.before_replace", false),
2841 ("checkpoint.registry.after_replace", true),
2842 ("checkpoint.registry.after_dir_fsync", true),
2843 ] {
2844 let directory = tempdir().unwrap();
2845 crate::open_or_initialize_project(directory.path()).unwrap();
2846 let status = Command::new(std::env::current_exe().unwrap())
2847 .args([
2848 "--exact",
2849 "project_checkpoints::tests::checkpoint_failpoint_helper",
2850 "--ignored",
2851 ])
2852 .env(
2853 "GRAPHFORGE_PROJECT_FAILPOINTS",
2854 "graphforge-internal-subprocess-v1",
2855 )
2856 .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
2857 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2858 .status()
2859 .unwrap();
2860 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2861 recover_checkpoint_pair_after_lock_handoff(
2862 directory.path(),
2863 &format!("action=create-unseeded failpoint={failpoint} committed={committed}"),
2864 );
2865 let rows = list_checkpoints(directory.path()).unwrap();
2866 assert_eq!(rows.len(), usize::from(committed), "{failpoint}");
2867 if !committed {
2868 recover_checkpoint_pair_after_lock_handoff(
2869 directory.path(),
2870 &format!(
2871 "action=create-unseeded parent-read-complete failpoint={failpoint} \
2872 committed={committed}"
2873 ),
2874 );
2875 let request = create_request(
2876 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
2877 "Crash",
2878 );
2879 create_checkpoint(directory.path(), &request).unwrap();
2880 recover_checkpoint_pair_after_lock_handoff(
2881 directory.path(),
2882 &format!(
2883 "action=create-unseeded replay-complete failpoint={failpoint} \
2884 committed={committed}"
2885 ),
2886 );
2887 assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2888 }
2889 }
2890 }
2891
2892 #[test]
2893 fn seeded_create_and_delete_failpoints_recover_exact_previous_or_next_revision() {
2894 let failpoints = [
2895 ("checkpoint.registry.after_intent_file_fsync", false),
2896 ("checkpoint.registry.after_file_fsync", false),
2897 ("checkpoint.registry.before_replace", false),
2898 ("checkpoint.registry.after_replace", true),
2899 ("checkpoint.registry.after_dir_fsync", true),
2900 ];
2901 for action in ["create", "delete"] {
2902 for (failpoint, committed) in failpoints {
2903 let directory = tempdir().unwrap();
2904 crate::open_or_initialize_project(directory.path()).unwrap();
2905 create_checkpoint(
2906 directory.path(),
2907 &create_request(
2908 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000029").unwrap(),
2909 "Base",
2910 ),
2911 )
2912 .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 .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", action)
2926 .status()
2927 .unwrap();
2928 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2929 recover_checkpoint_pair_after_lock_handoff(
2930 directory.path(),
2931 &format!("action={action} failpoint={failpoint} committed={committed}"),
2932 );
2933 let rows = list_checkpoints(directory.path()).unwrap();
2934 let expected = match (action, committed) {
2935 ("create", true) => vec!["Base", "Crash"],
2936 ("create", false) | ("delete", false) => vec!["Base"],
2937 ("delete", true) => vec![],
2938 _ => unreachable!(),
2939 };
2940 assert_eq!(
2941 rows.iter().map(|row| row.name.as_str()).collect::<Vec<_>>(),
2942 expected,
2943 "{action} {failpoint}"
2944 );
2945 recover_checkpoint_pair_after_lock_handoff(
2946 directory.path(),
2947 &format!(
2948 "action={action} parent-read-complete failpoint={failpoint} \
2949 committed={committed}"
2950 ),
2951 );
2952
2953 if action == "create" {
2954 let replay = create_checkpoint(
2955 directory.path(),
2956 &create_request(
2957 Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
2958 "Crash",
2959 ),
2960 )
2961 .unwrap();
2962 assert_eq!(replay.registry_revision, 2);
2963 } else {
2964 let replay = delete_checkpoint(
2965 directory.path(),
2966 &CheckpointDeleteRequest {
2967 operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000031")
2968 .unwrap(),
2969 name: "Base".into(),
2970 actor_uuid: None,
2971 },
2972 )
2973 .unwrap();
2974 assert_eq!(replay.registry_revision, 2);
2975 }
2976 recover_checkpoint_pair_after_lock_handoff(
2977 directory.path(),
2978 &format!(
2979 "action={action} replay-complete failpoint={failpoint} \
2980 committed={committed}"
2981 ),
2982 );
2983 let final_rows = list_checkpoints(directory.path()).unwrap();
2984 let final_names = final_rows
2985 .iter()
2986 .map(|row| row.name.as_str())
2987 .collect::<Vec<_>>();
2988 if action == "create" {
2989 assert_eq!(final_names, vec!["Base", "Crash"], "{failpoint}");
2990 } else {
2991 assert!(final_names.is_empty(), "{failpoint}");
2992 }
2993 }
2994 }
2995 }
2996
2997 #[test]
2998 fn recovery_rejects_missing_or_tampered_staged_pair() {
2999 for tamper_checksum in [false, true] {
3000 let directory = tempdir().unwrap();
3001 crate::open_or_initialize_project(directory.path()).unwrap();
3002 create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Base")).unwrap();
3003 let status = Command::new(std::env::current_exe().unwrap())
3004 .args([
3005 "--exact",
3006 "project_checkpoints::tests::checkpoint_failpoint_helper",
3007 "--ignored",
3008 ])
3009 .env(
3010 "GRAPHFORGE_PROJECT_FAILPOINTS",
3011 "graphforge-internal-subprocess-v1",
3012 )
3013 .env(
3014 "GRAPHFORGE_PROJECT_FAILPOINT",
3015 "checkpoint.registry.before_replace",
3016 )
3017 .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
3018 .status()
3019 .unwrap();
3020 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
3021 preserve_checkpoint_intent_after_lock_handoff(
3022 directory.path(),
3023 &format!(
3024 "action=create-tamper failpoint=checkpoint.registry.before_replace \
3025 committed=false tamper_checksum={tamper_checksum}"
3026 ),
3027 );
3028 let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
3029 let intent_bytes = fs::read(checkpoint_root.join(INTENT_FILE)).unwrap();
3030 let intent: RegistryIntent = serde_json::from_slice(&intent_bytes).unwrap();
3031 if tamper_checksum {
3032 fs::write(
3033 checkpoint_root.join(intent.checksum_temp),
3034 b"0000000000000000000000000000000000000000000000000000000000000000\n",
3035 )
3036 .unwrap();
3037 } else {
3038 fs::remove_file(checkpoint_root.join(intent.registry_temp)).unwrap();
3039 }
3040 assert_eq!(
3041 list_checkpoints(directory.path()).unwrap_err().code(),
3042 "GF_CHECKPOINT_REGISTRY_CORRUPT"
3043 );
3044 }
3045 }
3046}