Skip to main content

graphforge_storage/
project_checkpoints.rs

1//! Durable named references to verified immutable project generations.
2
3use 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/// Input for an idempotent checkpoint creation.
51#[derive(Debug, Clone)]
52pub struct CheckpointCreateRequest {
53    /// Canonical operation UUID.
54    pub operation_uuid: Uuid,
55    /// Human-facing checkpoint name (content, never a path).
56    pub name: String,
57    /// Optional bounded description.
58    pub description: Option<String>,
59    /// Optional actor identity.
60    pub actor_uuid: Option<Uuid>,
61}
62
63/// Input for an idempotent checkpoint deletion.
64#[derive(Debug, Clone)]
65pub struct CheckpointDeleteRequest {
66    /// Canonical operation UUID.
67    pub operation_uuid: Uuid,
68    /// Exact normalized checkpoint name.
69    pub name: String,
70    /// Optional actor identity.
71    pub actor_uuid: Option<Uuid>,
72}
73
74/// Internal complete-workspace revert request after the API selects its clock.
75#[derive(Debug, Clone)]
76pub struct CheckpointRevertRequest {
77    /// Caller-controlled idempotency UUID.
78    pub operation_uuid: Uuid,
79    /// Canonical checkpoint name.
80    pub name: String,
81    /// Bounded human restoration reason.
82    pub reason: String,
83    /// Optional actor identity.
84    pub actor_uuid: Option<Uuid>,
85}
86
87/// One active durable checkpoint.
88#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
89#[serde(deny_unknown_fields)]
90pub struct CheckpointRecord {
91    /// Stable deterministic identity.
92    pub checkpoint_uuid: Uuid,
93    /// Normalized display name.
94    pub name: String,
95    /// Exact pinned generation.
96    pub generation_uuid: Uuid,
97    /// Digest of that generation's canonical manifest.
98    pub generation_manifest_sha256: String,
99    /// Optional description.
100    pub description: Option<String>,
101    /// Engine-supplied UTC microseconds.
102    pub created_at: i64,
103    /// Optional actor identity.
104    pub created_by: Option<Uuid>,
105    /// Idempotency operation.
106    pub create_operation_uuid: Uuid,
107    /// Canonical request digest.
108    pub create_request_sha256: String,
109    /// Registry revision that originally committed this checkpoint.
110    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/// Stable mutation receipt.
134#[derive(Debug, Clone, PartialEq, Eq)]
135pub struct CheckpointReceipt {
136    /// Operation name (`checkpoint`, `delete_checkpoint`, or `revert_to_checkpoint`).
137    pub operation: &'static str,
138    /// Idempotency UUID.
139    pub operation_uuid: Uuid,
140    /// Stable checkpoint UUID.
141    pub checkpoint_uuid: Uuid,
142    /// Checkpoint name.
143    pub name: String,
144    /// Pinned generation UUID.
145    pub source_generation_uuid: Uuid,
146    /// Generation that was current immediately before a revert; absent for registry-only operations.
147    pub prior_current_generation_uuid: Option<Uuid>,
148    /// Newly published generation for revert; absent for registry-only operations.
149    pub result_generation_uuid: Option<Uuid>,
150    /// Resulting registry revision.
151    pub registry_revision: u64,
152    /// Original commit time in UTC microseconds.
153    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
236/// Create a checkpoint pinned to the post-lock validated `CURRENT` generation.
237pub 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, &registry, request.operation_uuid)?;
327    Ok(create_receipt(&row))
328}
329
330/// Delete one active checkpoint while preserving any already-open generation lease.
331pub 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, &registry, request.operation_uuid)?;
415    Ok(delete_receipt(&tombstone))
416}
417
418/// Return active checkpoints in canonical `(name, checkpoint_uuid)` order.
419pub 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
428/// Resolve and lifetime-pin the exact generation named by an active checkpoint.
429pub 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/// Publish a complete-workspace restoration as a new child generation.
466#[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(&registry_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(&registry_temp))?;
960    prepare_temp_path(&root.join(&checksum_temp))?;
961    write_new_synced(&root.join(&registry_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(&registry_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(&registry_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(&registry_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(&registry_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 &registry.active {
1252        validate_checkpoint_content(&CheckpointContentRef {
1253            label: "active checkpoint",
1254            checkpoint_uuid: row.checkpoint_uuid,
1255            create_operation_uuid: row.create_operation_uuid,
1256            name: &row.name,
1257            description: row.description.as_deref(),
1258            created_by: row.created_by,
1259            generation_manifest_sha256: &row.generation_manifest_sha256,
1260            create_request_sha256: &row.create_request_sha256,
1261        })?;
1262        if row.created_revision == 0
1263            || row.created_revision > registry.revision
1264            || !names.insert(row.name.as_str())
1265            || !checkpoint_uuids.insert(row.checkpoint_uuid)
1266            || !create_operations.insert(row.create_operation_uuid)
1267        {
1268            return Err(registry_corrupt(
1269                "active checkpoint identities or revision are inconsistent",
1270            ));
1271        }
1272    }
1273    for row in &registry.tombstones {
1274        // Preserve pre-consolidation order: content fields, then delete digest,
1275        // then deterministic create-request identity (error precedence).
1276        validate_name(&row.name)
1277            .map_err(|_| registry_corrupt("checkpoint tombstone name is invalid"))?;
1278        validate_description(row.description.as_deref())
1279            .map_err(|_| registry_corrupt("checkpoint tombstone description is invalid"))?;
1280        validate_digest(&row.generation_manifest_sha256)?;
1281        validate_digest(&row.create_request_sha256)?;
1282        validate_digest(&row.delete_request_sha256)?;
1283        validate_record_identity(
1284            row.checkpoint_uuid,
1285            row.create_operation_uuid,
1286            &row.name,
1287            row.description.as_deref(),
1288            row.created_by,
1289            &row.create_request_sha256,
1290        )?;
1291        let expected_delete =
1292            delete_request_digest_values(row.delete_operation_uuid, &row.name, row.deleted_by);
1293        if row.created_revision == 0
1294            || row.created_revision >= row.deleted_revision
1295            || row.deleted_revision > registry.revision
1296            || !checkpoint_uuids.insert(row.checkpoint_uuid)
1297            || !create_operations.insert(row.create_operation_uuid)
1298            || delete_operations.contains(&row.create_operation_uuid)
1299            || !delete_operations.insert(row.delete_operation_uuid)
1300            || create_operations.contains(&row.delete_operation_uuid)
1301            || row.delete_request_sha256 != hex(&expected_delete)
1302        {
1303            return Err(registry_corrupt(
1304                "checkpoint tombstone identities or revisions are inconsistent",
1305            ));
1306        }
1307    }
1308    Ok(())
1309}
1310
1311struct CheckpointContentRef<'a> {
1312    label: &'static str,
1313    checkpoint_uuid: Uuid,
1314    create_operation_uuid: Uuid,
1315    name: &'a str,
1316    description: Option<&'a str>,
1317    created_by: Option<Uuid>,
1318    generation_manifest_sha256: &'a str,
1319    create_request_sha256: &'a str,
1320}
1321
1322fn validate_checkpoint_content(content: &CheckpointContentRef<'_>) -> Result<(), GfError> {
1323    validate_name(content.name)
1324        .map_err(|_| registry_corrupt(format!("{} name is invalid", content.label)))?;
1325    validate_description(content.description)
1326        .map_err(|_| registry_corrupt(format!("{} description is invalid", content.label)))?;
1327    validate_digest(content.generation_manifest_sha256)?;
1328    validate_digest(content.create_request_sha256)?;
1329    validate_record_identity(
1330        content.checkpoint_uuid,
1331        content.create_operation_uuid,
1332        content.name,
1333        content.description,
1334        content.created_by,
1335        content.create_request_sha256,
1336    )
1337}
1338
1339fn validate_single_link_regular(path: &Path, label: &str) -> Result<(), GfError> {
1340    let metadata = fs::symlink_metadata(path).map_err(storage_io)?;
1341    if !metadata.file_type().is_file() {
1342        return Err(registry_corrupt(format!("{label} is linked or special")));
1343    }
1344    #[cfg(unix)]
1345    {
1346        use std::os::unix::fs::MetadataExt;
1347        if metadata.nlink() != 1 {
1348            return Err(registry_corrupt(format!("{label} is hard-linked")));
1349        }
1350    }
1351    Ok(())
1352}
1353
1354fn validate_name(value: &str) -> Result<String, GfError> {
1355    let normalized: String = value.nfc().collect();
1356    if normalized != value
1357        || value.is_empty()
1358        || value.len() > MAX_NAME_BYTES
1359        || value.trim() != value
1360        || value == "."
1361        || value == ".."
1362        || value.contains("  ")
1363        || !value
1364            .chars()
1365            .all(|ch| ch.is_alphanumeric() || matches!(ch, ' ' | '_' | '-' | '.'))
1366    {
1367        return Err(GfError::Validation(
1368            "checkpoint name is not canonical NFC content or violates the 1-128 byte grammar"
1369                .into(),
1370        ));
1371    }
1372    Ok(normalized)
1373}
1374
1375fn validate_description(value: Option<&str>) -> Result<(), GfError> {
1376    if value.is_some_and(|value| {
1377        value.len() > MAX_DESCRIPTION_BYTES || value.chars().any(char::is_control)
1378    }) {
1379        return Err(GfError::Validation(
1380            "checkpoint description exceeds 1024 UTF-8 bytes or contains controls".into(),
1381        ));
1382    }
1383    Ok(())
1384}
1385
1386fn validate_digest(value: &str) -> Result<(), GfError> {
1387    if value.len() != 64
1388        || !value
1389            .bytes()
1390            .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1391    {
1392        return Err(registry_corrupt("checkpoint digest is noncanonical"));
1393    }
1394    Ok(())
1395}
1396
1397fn decode_digest(value: &str) -> Result<[u8; 32], GfError> {
1398    validate_digest(value)?;
1399    let mut digest = [0_u8; 32];
1400    for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
1401        let text = std::str::from_utf8(pair)
1402            .map_err(|_| registry_corrupt("checkpoint digest is not UTF-8"))?;
1403        digest[index] = u8::from_str_radix(text, 16)
1404            .map_err(|_| registry_corrupt("checkpoint digest is not lowercase hex"))?;
1405    }
1406    Ok(digest)
1407}
1408
1409fn create_request_digest(request: &CheckpointCreateRequest, name: &str) -> [u8; 32] {
1410    create_request_digest_values(
1411        request.operation_uuid,
1412        name,
1413        request.description.as_deref(),
1414        request.actor_uuid,
1415    )
1416}
1417
1418fn create_request_digest_values(
1419    operation_uuid: Uuid,
1420    name: &str,
1421    description: Option<&str>,
1422    actor_uuid: Option<Uuid>,
1423) -> [u8; 32] {
1424    let mut hasher = Sha256::new();
1425    hasher.update(b"graphforge-checkpoint-create-request/1");
1426    hasher.update(operation_uuid.as_bytes());
1427    append_bytes(&mut hasher, name.as_bytes());
1428    match description {
1429        Some(value) => {
1430            hasher.update([1]);
1431            append_bytes(&mut hasher, value.as_bytes());
1432        }
1433        None => hasher.update([0]),
1434    }
1435    append_actor(&mut hasher, actor_uuid);
1436    hasher.finalize().into()
1437}
1438
1439fn delete_request_digest(request: &CheckpointDeleteRequest, name: &str) -> [u8; 32] {
1440    delete_request_digest_values(request.operation_uuid, name, request.actor_uuid)
1441}
1442
1443fn delete_request_digest_values(
1444    operation_uuid: Uuid,
1445    name: &str,
1446    actor_uuid: Option<Uuid>,
1447) -> [u8; 32] {
1448    let mut hasher = Sha256::new();
1449    hasher.update(b"graphforge-checkpoint-delete-request/1");
1450    hasher.update(operation_uuid.as_bytes());
1451    append_bytes(&mut hasher, name.as_bytes());
1452    append_actor(&mut hasher, actor_uuid);
1453    hasher.finalize().into()
1454}
1455
1456fn validate_record_identity(
1457    checkpoint: Uuid,
1458    operation: Uuid,
1459    name: &str,
1460    description: Option<&str>,
1461    actor: Option<Uuid>,
1462    request_hex: &str,
1463) -> Result<(), GfError> {
1464    let request = create_request_digest_values(operation, name, description, actor);
1465    if request_hex != hex(&request) || checkpoint != checkpoint_uuid(operation, request) {
1466        return Err(registry_corrupt(
1467            "checkpoint deterministic identity or create request digest is inconsistent",
1468        ));
1469    }
1470    Ok(())
1471}
1472
1473fn checkpoint_uuid(operation_uuid: Uuid, request_digest: [u8; 32]) -> Uuid {
1474    let mut hasher = Sha256::new();
1475    hasher.update(b"graphforge-checkpoint-uuid/1");
1476    hasher.update(operation_uuid.as_bytes());
1477    hasher.update(request_digest);
1478    graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1479}
1480
1481fn append_bytes(hasher: &mut Sha256, bytes: &[u8]) {
1482    hasher.update(
1483        u32::try_from(bytes.len())
1484            .expect("validated checkpoint strings fit u32")
1485            .to_be_bytes(),
1486    );
1487    hasher.update(bytes);
1488}
1489fn append_actor(hasher: &mut Sha256, actor: Option<Uuid>) {
1490    match actor {
1491        Some(value) => {
1492            hasher.update([1]);
1493            hasher.update(value.as_bytes());
1494        }
1495        None => hasher.update([0]),
1496    }
1497}
1498fn valid_private_name(name: &str, uuid: Uuid, kind: &str) -> bool {
1499    name == format!(".registry.{uuid}.{kind}.next")
1500}
1501fn hex(bytes: &[u8; 32]) -> String {
1502    let mut output = String::with_capacity(64);
1503    for byte in bytes {
1504        write!(&mut output, "{byte:02x}").expect("writing hexadecimal to String cannot fail");
1505    }
1506    output
1507}
1508
1509fn parse_uuid(value: &str) -> Result<Uuid, GfError> {
1510    Uuid::parse_str(value).map_err(|_| registry_corrupt("revert journal UUID is invalid"))
1511}
1512
1513fn validate_reason(value: &str) -> Result<String, GfError> {
1514    let trimmed = value.trim();
1515    if trimmed.is_empty() || trimmed.len() > MAX_REASON_BYTES {
1516        return Err(GfError::Validation(
1517            "checkpoint revert reason must contain 1..=1024 UTF-8 bytes after trimming".into(),
1518        ));
1519    }
1520    Ok(trimmed.to_owned())
1521}
1522
1523fn revert_request_digest(
1524    operation_uuid: Uuid,
1525    name: &str,
1526    checkpoint_uuid: Uuid,
1527    source_generation_uuid: Uuid,
1528    source_manifest_sha256: [u8; 32],
1529    reason: &str,
1530    actor_uuid: Option<Uuid>,
1531) -> [u8; 32] {
1532    let mut hasher = Sha256::new();
1533    hasher.update(b"graphforge-checkpoint-revert-request/1");
1534    hasher.update(operation_uuid.as_bytes());
1535    append_bytes(&mut hasher, name.as_bytes());
1536    hasher.update(checkpoint_uuid.as_bytes());
1537    hasher.update(source_generation_uuid.as_bytes());
1538    hasher.update(source_manifest_sha256);
1539    append_bytes(&mut hasher, reason.as_bytes());
1540    append_actor(&mut hasher, actor_uuid);
1541    hasher.finalize().into()
1542}
1543
1544fn revert_transaction_uuid(operation_uuid: Uuid) -> Uuid {
1545    let mut hasher = Sha256::new();
1546    hasher.update(b"graphforge-checkpoint-revert-transaction/1");
1547    hasher.update(operation_uuid.as_bytes());
1548    graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1549}
1550
1551fn restoration_uuid(operation_uuid: Uuid, request_digest: [u8; 32]) -> Uuid {
1552    let mut hasher = Sha256::new();
1553    hasher.update(b"graphforge-restoration-transition-uuid/1");
1554    hasher.update(operation_uuid.as_bytes());
1555    hasher.update(request_digest);
1556    graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1557}
1558
1559fn restored_generation_uuid(
1560    transaction_uuid: Uuid,
1561    checkpoint_uuid: Uuid,
1562    source_generation_uuid: Uuid,
1563    source_manifest_sha256: [u8; 32],
1564    prior_current_generation_uuid: Uuid,
1565    restored_at: i64,
1566    request_digest: [u8; 32],
1567) -> Uuid {
1568    let mut hasher = Sha256::new();
1569    hasher.update(b"graphforge-checkpoint-restored-generation/1");
1570    hasher.update(transaction_uuid.as_bytes());
1571    hasher.update(checkpoint_uuid.as_bytes());
1572    hasher.update(source_generation_uuid.as_bytes());
1573    hasher.update(source_manifest_sha256);
1574    hasher.update(prior_current_generation_uuid.as_bytes());
1575    hasher.update(restored_at.to_be_bytes());
1576    hasher.update(request_digest);
1577    graphforge_core::canonical::uuid_v8(hasher.finalize().into())
1578}
1579
1580fn snapshot_to_participant(
1581    snapshot: crate::ProjectParticipantSnapshot,
1582) -> Result<ProjectParticipant, GfError> {
1583    let encoding = match snapshot.encoding.as_str() {
1584        "parquet" => ProjectParticipantEncoding::Parquet,
1585        "arrow" => ProjectParticipantEncoding::Arrow,
1586        "json" => ProjectParticipantEncoding::Json,
1587        _ => {
1588            return Err(registry_corrupt(
1589                "checkpoint participant encoding is unsupported",
1590            ));
1591        }
1592    };
1593    Ok(ProjectParticipant {
1594        capability_id: snapshot.capability_id,
1595        capability_version: snapshot.capability_version,
1596        record_family_id: snapshot.record_family_id,
1597        record_version: snapshot.record_version,
1598        encoding,
1599        schema_fingerprint: snapshot.schema_fingerprint,
1600        row_count: snapshot.row_count,
1601        bytes: snapshot.bytes,
1602    })
1603}
1604
1605#[allow(clippy::too_many_arguments)]
1606fn restoration_participant(
1607    restoration_uuid: Uuid,
1608    checkpoint_uuid: Uuid,
1609    source_generation_uuid: Uuid,
1610    source_manifest_sha256: [u8; 32],
1611    prior_current_generation_uuid: Uuid,
1612    restored_generation_uuid: Uuid,
1613    operation_uuid: Uuid,
1614    actor_uuid: Option<Uuid>,
1615    reason: &str,
1616    restored_at: i64,
1617) -> Result<ProjectParticipant, GfError> {
1618    let schema = Arc::new(Schema::new(vec![
1619        uuid_field("restoration_uuid", false),
1620        uuid_field("checkpoint_uuid", false),
1621        uuid_field("source_generation_uuid", false),
1622        Field::new(
1623            "source_manifest_sha256",
1624            DataType::FixedSizeBinary(32),
1625            false,
1626        ),
1627        uuid_field("prior_current_generation_uuid", false),
1628        uuid_field("restored_generation_uuid", false),
1629        uuid_field("operation_uuid", false),
1630        uuid_field("actor_uuid", true),
1631        Field::new("reason", DataType::Utf8, false),
1632        Field::new(
1633            "restored_at",
1634            DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
1635            false,
1636        ),
1637        Field::new("contract_version", DataType::UInt32, false),
1638    ]));
1639    let mut columns = Vec::<ArrayRef>::new();
1640    for value in [
1641        Some(restoration_uuid),
1642        Some(checkpoint_uuid),
1643        Some(source_generation_uuid),
1644        Some(prior_current_generation_uuid),
1645        Some(restored_generation_uuid),
1646        Some(operation_uuid),
1647        actor_uuid,
1648    ] {
1649        let mut builder = FixedSizeBinaryBuilder::with_capacity(1, 16);
1650        match value {
1651            Some(uuid) => builder.append_value(uuid.as_bytes()).map_err(arrow_error)?,
1652            None => builder.append_null(),
1653        }
1654        columns.push(Arc::new(builder.finish()));
1655    }
1656    let mut source_digest = FixedSizeBinaryBuilder::with_capacity(1, 32);
1657    source_digest
1658        .append_value(source_manifest_sha256)
1659        .map_err(arrow_error)?;
1660    columns.insert(3, Arc::new(source_digest.finish()));
1661    columns.push(Arc::new(StringArray::from(vec![reason])));
1662    columns.push(Arc::new(
1663        TimestampMicrosecondArray::from(vec![restored_at]).with_timezone("UTC"),
1664    ));
1665    columns.push(Arc::new(UInt32Array::from(vec![
1666        RESTORATION_CONTRACT_VERSION,
1667    ])));
1668    let batch = RecordBatch::try_new(Arc::clone(&schema), columns).map_err(arrow_error)?;
1669    let properties = WriterProperties::builder()
1670        .set_created_by("graphforge-restoration-transition/1".into())
1671        .build();
1672    let mut writer =
1673        ArrowWriter::try_new(Vec::new(), schema, Some(properties)).map_err(parquet_error)?;
1674    writer.write(&batch).map_err(parquet_error)?;
1675    let bytes = writer.into_inner().map_err(parquet_error)?;
1676    let schema_fingerprint = fingerprint(
1677        CanonicalDomain::Schema,
1678        CANONICAL_CONTRACT_VERSION,
1679        b"restoration_transition/1|restoration_uuid:fixed[16]:required|checkpoint_uuid:fixed[16]:required|source_generation_uuid:fixed[16]:required|source_manifest_sha256:fixed[32]:required|prior_current_generation_uuid:fixed[16]:required|restored_generation_uuid:fixed[16]:required|operation_uuid:fixed[16]:required|actor_uuid:fixed[16]:optional|reason:utf8:required|restored_at:timestamp_us_utc:required|contract_version:u32:required",
1680    )
1681    .map_err(|error| GfError::Validation(error.to_string()))?;
1682    Ok(ProjectParticipant {
1683        capability_id: crate::WORKSPACE_CAPABILITY_ID.into(),
1684        capability_version: crate::WORKSPACE_CAPABILITY_VERSION,
1685        record_family_id: RESTORATION_FAMILY.into(),
1686        record_version: RESTORATION_CONTRACT_VERSION,
1687        encoding: ProjectParticipantEncoding::Parquet,
1688        schema_fingerprint,
1689        row_count: 1,
1690        bytes,
1691    })
1692}
1693
1694fn uuid_field(name: &str, nullable: bool) -> Field {
1695    Field::new(name, DataType::FixedSizeBinary(16), nullable)
1696}
1697
1698fn arrow_error(error: arrow::error::ArrowError) -> GfError {
1699    let message = format!("restoration Arrow encoding failed: {error}");
1700    drop(error);
1701    GfError::Storage(message)
1702}
1703
1704fn parquet_error(error: parquet::errors::ParquetError) -> GfError {
1705    let message = format!("restoration Parquet encoding failed: {error}");
1706    drop(error);
1707    GfError::Storage(message)
1708}
1709fn utc_micros() -> Result<i64, GfError> {
1710    let value = SystemTime::now()
1711        .duration_since(UNIX_EPOCH)
1712        .map_err(|_| GfError::Storage("system clock is before Unix epoch".into()))?
1713        .as_micros();
1714    i64::try_from(value).map_err(|_| GfError::Storage("UTC microsecond timestamp overflow".into()))
1715}
1716fn registry_serde(error: impl std::fmt::Display) -> GfError {
1717    GfError::Storage(format!("checkpoint registry encoding failed: {error}"))
1718}
1719fn registry_corrupt(message: impl Into<String>) -> GfError {
1720    project_error(ProjectErrorCode::CheckpointRegistryCorrupt, message)
1721}
1722fn project_error(code: ProjectErrorCode, message: impl Into<String>) -> GfError {
1723    GfError::Project {
1724        code,
1725        message: message.into(),
1726    }
1727}
1728fn storage_io(error: impl std::fmt::Display) -> GfError {
1729    GfError::Storage(format!("checkpoint registry I/O failed: {error}"))
1730}
1731
1732#[cfg(test)]
1733mod tests {
1734    use super::*;
1735    use std::collections::BTreeMap;
1736    use std::io::{BufRead, BufReader};
1737    use std::panic::{AssertUnwindSafe, catch_unwind, resume_unwind};
1738    use std::process::{Command, Stdio};
1739    use std::sync::mpsc;
1740    use std::time::Duration;
1741    use tempfile::tempdir;
1742    use wait_timeout::ChildExt;
1743
1744    const TEST_DEADLINE: Duration = Duration::from_secs(1);
1745    const CHILD_DEADLINE: Duration = Duration::from_secs(10);
1746
1747    struct WriterLockHolder {
1748        release: Option<mpsc::SyncSender<()>>,
1749        worker: Option<std::thread::JoinHandle<()>>,
1750    }
1751
1752    impl WriterLockHolder {
1753        fn finish(mut self) -> Result<(), String> {
1754            let release = self
1755                .release
1756                .take()
1757                .ok_or_else(|| "phase=main release sender missing".to_owned())?;
1758            let release_result = release
1759                .send(())
1760                .map_err(|error| format!("phase=main release holder error={error}"));
1761            let join_result = self
1762                .worker
1763                .take()
1764                .ok_or_else(|| "phase=main holder worker missing".to_owned())?
1765                .join()
1766                .map_err(|_| "phase=main holder worker panicked".to_owned());
1767            release_result.and(join_result)
1768        }
1769    }
1770
1771    impl Drop for WriterLockHolder {
1772        fn drop(&mut self) {
1773            if let Some(release) = self.release.take() {
1774                let _ = release.send(());
1775            }
1776            if let Some(worker) = self.worker.take() {
1777                let _ = worker.join();
1778            }
1779        }
1780    }
1781
1782    fn while_writer_lock_is_held<T>(root: &Path, action: impl FnOnce() -> T) -> T {
1783        let writer_path = root.join(LOCKS_DIR).join(WRITER_LOCK_FILE);
1784        let worker_path = writer_path.clone();
1785        let (ready_sender, ready_receiver) = mpsc::sync_channel(0);
1786        let (release_sender, release_receiver) = mpsc::sync_channel(0);
1787        let worker = std::thread::Builder::new()
1788            .name("checkpoint-writer-lock-holder".into())
1789            .spawn(move || {
1790                let writer =
1791                    open_regular_lock(&worker_path).expect("phase=holder open writer.lock");
1792                assert!(
1793                    FileExt::try_lock_exclusive(&writer).expect("phase=holder acquire writer.lock"),
1794                    "phase=holder writer.lock unexpectedly busy"
1795                );
1796                ready_sender.send(()).expect("phase=holder publish ready");
1797                release_receiver.recv().expect("phase=holder await release");
1798                FileExt::unlock(&writer).expect("phase=holder release writer.lock");
1799            })
1800            .expect("phase=holder spawn");
1801        let holder = WriterLockHolder {
1802            release: Some(release_sender),
1803            worker: Some(worker),
1804        };
1805        if let Err(error) = ready_receiver.recv_timeout(TEST_DEADLINE) {
1806            drop(ready_receiver);
1807            let cleanup = holder.finish();
1808            panic!("phase=main await held writer.lock error={error}; cleanup={cleanup:?}");
1809        }
1810        let result = catch_unwind(AssertUnwindSafe(action));
1811        let cleanup = holder.finish();
1812        match result {
1813            Ok(value) => {
1814                cleanup.unwrap_or_else(|error| panic!("phase=main holder cleanup error={error}"));
1815                value
1816            }
1817            Err(original) => {
1818                let _ = cleanup;
1819                resume_unwind(original);
1820            }
1821        }
1822    }
1823
1824    struct BoundedChild {
1825        child: std::process::Child,
1826        reaped: bool,
1827    }
1828
1829    impl BoundedChild {
1830        fn wait(mut self, phase: &str) -> std::process::ExitStatus {
1831            let mut failures = Vec::new();
1832            match self.child.wait_timeout(CHILD_DEADLINE) {
1833                Ok(Some(status)) => {
1834                    self.reaped = true;
1835                    return status;
1836                }
1837                Ok(None) => failures.push(format!("wait timeout={CHILD_DEADLINE:?}")),
1838                Err(error) => failures.push(format!("wait error={error}")),
1839            }
1840            if let Err(error) = self.child.kill() {
1841                failures.push(format!("kill error={error}"));
1842            }
1843            match self.child.wait_timeout(TEST_DEADLINE) {
1844                Ok(Some(status)) => {
1845                    self.reaped = true;
1846                    failures.push(format!("killed_status={status}"));
1847                }
1848                Ok(None) => failures.push(format!("reap timeout={TEST_DEADLINE:?}")),
1849                Err(error) => failures.push(format!("reap error={error}")),
1850            }
1851            panic!("phase={phase} child cleanup failures={failures:?}");
1852        }
1853    }
1854
1855    impl Drop for BoundedChild {
1856        fn drop(&mut self) {
1857            if !self.reaped {
1858                let mut failures = Vec::new();
1859                if let Err(error) = self.child.kill() {
1860                    failures.push(format!("kill error={error}"));
1861                }
1862                match self.child.wait_timeout(TEST_DEADLINE) {
1863                    Ok(Some(_)) => self.reaped = true,
1864                    Ok(None) => failures.push(format!("reap timeout={TEST_DEADLINE:?}")),
1865                    Err(error) => failures.push(format!("reap error={error}")),
1866                }
1867                if !failures.is_empty() {
1868                    eprintln!("phase=drop child cleanup failures={failures:?}");
1869                }
1870            }
1871        }
1872    }
1873
1874    fn recover_checkpoint_pair_after_lock_handoff(root: &Path, phase: &str) {
1875        checkpoint_lock_handoff(root, phase, true);
1876    }
1877
1878    fn preserve_checkpoint_intent_after_lock_handoff(root: &Path, phase: &str) {
1879        checkpoint_lock_handoff(root, phase, false);
1880    }
1881
1882    fn checkpoint_lock_handoff(root: &Path, phase: &str, recover_durable_intent: bool) {
1883        let lock_root = root.join(LOCKS_DIR);
1884        let writer_path = lock_root.join(WRITER_LOCK_FILE);
1885        let checkpoint_path = lock_root.join(CHECKPOINT_LOCK_FILE);
1886        let checkpoint_root = root.join(CHECKPOINTS_DIR);
1887        let worker_writer_path = writer_path.clone();
1888        let worker_checkpoint_path = checkpoint_path.clone();
1889        let (sender, receiver) = mpsc::sync_channel(0);
1890        std::thread::Builder::new()
1891            .name("checkpoint-lock-handoff-recovery".into())
1892            .spawn(move || {
1893                let result = (|| {
1894                    let writer = open_regular_lock(&worker_writer_path)
1895                        .map_err(|error| format!("open writer.lock failed: {error}"))?;
1896                    FileExt::lock_exclusive(&writer)
1897                        .map_err(|error| format!("acquire writer.lock failed: {error}"))?;
1898
1899                    let checkpoint = match open_regular_lock(&worker_checkpoint_path) {
1900                        Ok(checkpoint) => checkpoint,
1901                        Err(error) => {
1902                            let writer_unlock = FileExt::unlock(&writer);
1903                            return Err(format!(
1904                                "open checkpoints.lock failed: {error}; \
1905                                 writer_unlock={writer_unlock:?}"
1906                            ));
1907                        }
1908                    };
1909                    if let Err(error) = FileExt::lock_exclusive(&checkpoint) {
1910                        let writer_unlock = FileExt::unlock(&writer);
1911                        return Err(format!(
1912                            "acquire checkpoints.lock failed: {error}; writer_unlock={writer_unlock:?}"
1913                        ));
1914                    }
1915
1916                    let recovery = if recover_durable_intent
1917                        && checkpoint_root.join(INTENT_FILE).exists()
1918                    {
1919                        recover_pair(&checkpoint_root)
1920                            .map_err(|error| format!("recover durable checkpoint intent failed: {error}"))
1921                    } else {
1922                        Ok(())
1923                    };
1924                    let checkpoint_unlock = FileExt::unlock(&checkpoint)
1925                        .map_err(|error| format!("unlock checkpoints.lock failed: {error}"));
1926                    let writer_unlock = FileExt::unlock(&writer)
1927                        .map_err(|error| format!("unlock writer.lock failed: {error}"));
1928
1929                    recovery?;
1930                    checkpoint_unlock?;
1931                    writer_unlock
1932                })();
1933                let _ = sender.send(result);
1934            })
1935            .unwrap();
1936        match receiver.recv_timeout(Duration::from_secs(1)) {
1937            Ok(Ok(())) => {}
1938            Ok(Err(error)) => panic!(
1939                "checkpoint lock handoff/recovery failed at {phase}; writer_path={}; \
1940                 checkpoint_path={}: {error}",
1941                writer_path.display(),
1942                checkpoint_path.display()
1943            ),
1944            Err(error) => panic!(
1945                "checkpoint lock handoff/recovery timed out at {phase}; writer_path={}; \
1946                 checkpoint_path={}; timeout=1s; channel={error}",
1947                writer_path.display(),
1948                checkpoint_path.display()
1949            ),
1950        }
1951    }
1952
1953    fn publish_clone(root: &Path) -> Uuid {
1954        let selected = crate::resolve_project_generation(root).unwrap();
1955        let capabilities = selected
1956            .capabilities()
1957            .into_iter()
1958            .map(|entry| crate::ProjectCapability {
1959                capability_id: entry.capability_id,
1960                capability_version: entry.capability_version,
1961            })
1962            .collect();
1963        let participants = selected
1964            .participant_snapshots()
1965            .unwrap()
1966            .into_iter()
1967            .map(|entry| crate::ProjectParticipant {
1968                capability_id: entry.capability_id,
1969                capability_version: entry.capability_version,
1970                record_family_id: entry.record_family_id,
1971                record_version: entry.record_version,
1972                encoding: match entry.encoding.as_str() {
1973                    "arrow" => crate::ProjectParticipantEncoding::Arrow,
1974                    "json" => crate::ProjectParticipantEncoding::Json,
1975                    "parquet" => crate::ProjectParticipantEncoding::Parquet,
1976                    other => panic!("unexpected participant encoding {other}"),
1977                },
1978                schema_fingerprint: entry.schema_fingerprint,
1979                row_count: entry.row_count,
1980                bytes: entry.bytes,
1981            })
1982            .collect();
1983        let generation_uuid = Uuid::now_v7();
1984        let request = crate::ProjectGenerationRequest {
1985            transaction_uuid: Uuid::now_v7(),
1986            generation_uuid,
1987            capabilities,
1988            participants,
1989        };
1990        let crate::ProjectStageOutcome::Staged(staged) =
1991            crate::stage_project_generation(root, &request).unwrap()
1992        else {
1993            panic!("fresh publication unexpectedly replayed");
1994        };
1995        staged
1996            .validate(|_| Ok(()), |_, _| Ok(()))
1997            .unwrap()
1998            .publish()
1999            .unwrap();
2000        generation_uuid
2001    }
2002
2003    fn create_request(operation_uuid: Uuid, name: &str) -> CheckpointCreateRequest {
2004        CheckpointCreateRequest {
2005            operation_uuid,
2006            name: name.into(),
2007            description: Some("release candidate".into()),
2008            actor_uuid: Some(Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000aa").unwrap()),
2009        }
2010    }
2011
2012    fn write_raw_registry(root: &Path, registry: &Registry) {
2013        let checkpoint_root = root.join(CHECKPOINTS_DIR);
2014        let mut bytes = serde_json::to_vec(registry).unwrap();
2015        bytes.push(b'\n');
2016        fs::write(checkpoint_root.join(REGISTRY_FILE), &bytes).unwrap();
2017        fs::write(
2018            checkpoint_root.join(CHECKSUM_FILE),
2019            format!("{}\n", hex(&Sha256::digest(&bytes).into())),
2020        )
2021        .unwrap();
2022    }
2023
2024    fn install_registry_intent(
2025        checkpoint_root: &Path,
2026        previous: Option<&Registry>,
2027        next: &Registry,
2028    ) -> RegistryIntent {
2029        fs::create_dir_all(checkpoint_root).unwrap();
2030        let transaction_uuid = Uuid::now_v7();
2031        let next_bytes = next.canonical_bytes().unwrap();
2032        let next_sha256 = hex(&Sha256::digest(&next_bytes).into());
2033        let registry_temp = format!(".registry.{transaction_uuid}.json.next");
2034        let checksum_temp = format!(".registry.{transaction_uuid}.sha256.next");
2035        fs::write(checkpoint_root.join(&registry_temp), &next_bytes).unwrap();
2036        fs::write(
2037            checkpoint_root.join(&checksum_temp),
2038            format!("{next_sha256}\n"),
2039        )
2040        .unwrap();
2041        let intent = RegistryIntent {
2042            transaction_uuid,
2043            previous_revision: previous.map(|registry| registry.revision),
2044            previous_sha256: previous
2045                .map(|registry| hex(&Sha256::digest(registry.canonical_bytes().unwrap()).into())),
2046            next_revision: next.revision,
2047            next_sha256,
2048            registry_temp,
2049            checksum_temp,
2050        };
2051        let mut intent_bytes = serde_json::to_vec(&intent).unwrap();
2052        intent_bytes.push(b'\n');
2053        fs::write(checkpoint_root.join(INTENT_FILE), intent_bytes).unwrap();
2054        intent
2055    }
2056
2057    #[test]
2058    fn revert_identity_matches_frozen_golden_vector() {
2059        let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000003").unwrap();
2060        let checkpoint = Uuid::parse_str("4084179c-38db-8b6b-9b6e-c0b0a855e002").unwrap();
2061        let source = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000b0").unwrap();
2062        let prior = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000d0").unwrap();
2063        let actor = Uuid::parse_str("018f0f4e-7b8c-7000-8000-0000000000aa").unwrap();
2064        let source_digest = [0x11; 32];
2065        let request_digest = revert_request_digest(
2066            operation,
2067            "Release 1.0",
2068            checkpoint,
2069            source,
2070            source_digest,
2071            "restore release candidate",
2072            Some(actor),
2073        );
2074        assert_eq!(
2075            hex(&request_digest),
2076            "dff3755629942d1189060b117cb70dc864428fcc3b28d5c2d22d3924c3690e93"
2077        );
2078        let transaction = revert_transaction_uuid(operation);
2079        assert_eq!(
2080            transaction.to_string(),
2081            "908d637b-d6e6-8508-919e-d4e708e037b2"
2082        );
2083        assert_eq!(
2084            restoration_uuid(operation, request_digest).to_string(),
2085            "9e1f160c-badb-80c9-beaa-ae580910bf8a"
2086        );
2087        assert_eq!(
2088            restored_generation_uuid(
2089                transaction,
2090                checkpoint,
2091                source,
2092                source_digest,
2093                prior,
2094                1_720_000_000_123_456,
2095                request_digest,
2096            )
2097            .to_string(),
2098            "5dc02888-2064-8892-a0e1-2c00968ba0cc"
2099        );
2100    }
2101
2102    #[test]
2103    fn revert_publishes_child_preserves_registry_and_replays_after_delete() {
2104        let directory = tempdir().unwrap();
2105        crate::open_or_initialize_project(directory.path()).unwrap();
2106        let source = crate::resolve_project_generation(directory.path()).unwrap();
2107        let created = create_checkpoint(
2108            directory.path(),
2109            &create_request(Uuid::from_u128(40), "Before"),
2110        )
2111        .unwrap();
2112        let prior_current = publish_clone(directory.path());
2113        let request = CheckpointRevertRequest {
2114            operation_uuid: Uuid::from_u128(41),
2115            name: "Before".into(),
2116            reason: " restore known state ".into(),
2117            actor_uuid: None,
2118        };
2119        let (receipt, restored) = revert_checkpoint(
2120            directory.path(),
2121            &request,
2122            || Ok(1_720_000_000_123_456),
2123            |_| Ok(()),
2124        )
2125        .unwrap();
2126        assert_eq!(restored.parent_generation_uuid(), Some(prior_current));
2127        assert_eq!(receipt.source_generation_uuid, source.generation_uuid());
2128        assert_eq!(receipt.prior_current_generation_uuid, Some(prior_current));
2129        assert_eq!(receipt.registry_revision, created.registry_revision);
2130        assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2131        let restoration_count = restored
2132            .participant_descriptors()
2133            .unwrap()
2134            .iter()
2135            .filter(|row| row.record_family_id == RESTORATION_FAMILY)
2136            .count();
2137        assert_eq!(restoration_count, 1);
2138
2139        delete_checkpoint(
2140            directory.path(),
2141            &CheckpointDeleteRequest {
2142                operation_uuid: Uuid::from_u128(42),
2143                name: "Before".into(),
2144                actor_uuid: None,
2145            },
2146        )
2147        .unwrap();
2148        let (replay, replayed_generation) = revert_checkpoint(
2149            directory.path(),
2150            &request,
2151            || panic!("published replay sampled clock"),
2152            |_| Ok(()),
2153        )
2154        .unwrap();
2155        assert_eq!(replay, receipt);
2156        assert_eq!(replay.prior_current_generation_uuid, Some(prior_current));
2157        assert_eq!(
2158            replayed_generation.generation_uuid(),
2159            restored.generation_uuid()
2160        );
2161        checkpoint_lock_handoff(
2162            directory.path(),
2163            "action=revert published-replay return",
2164            false,
2165        );
2166
2167        let mut conflict = request;
2168        conflict.reason = "different".into();
2169        let conflict_error =
2170            revert_checkpoint(directory.path(), &conflict, || Ok(0), |_| Ok(())).unwrap_err();
2171        assert_eq!(conflict_error.code(), "GF_IDEMPOTENCY_CONFLICT");
2172        checkpoint_lock_handoff(
2173            directory.path(),
2174            "action=revert published-replay conflict return",
2175            false,
2176        );
2177    }
2178
2179    #[test]
2180    fn revert_replay_lock_handoff_fails_closed_with_stable_storage_errors() {
2181        let checkpoint_error = finish_revert_replay_lock_handoff(
2182            Err(std::io::Error::other("checkpoint unlock failed")),
2183            Ok(()),
2184        )
2185        .unwrap_err();
2186        assert_eq!(checkpoint_error.code(), "GF_IO");
2187        assert_eq!(
2188            checkpoint_error.to_string(),
2189            "storage error: checkpoint revert replay lock handoff failed at checkpoints.lock: checkpoint unlock failed"
2190        );
2191
2192        let writer_error = finish_revert_replay_lock_handoff(
2193            Ok(()),
2194            Err(std::io::Error::other("writer unlock failed")),
2195        )
2196        .unwrap_err();
2197        assert_eq!(writer_error.code(), "GF_IO");
2198        assert_eq!(
2199            writer_error.to_string(),
2200            "storage error: checkpoint revert replay lock handoff failed at writer.lock: writer unlock failed"
2201        );
2202    }
2203
2204    #[test]
2205    fn revert_validation_failure_preserves_prior_current() {
2206        let directory = tempdir().unwrap();
2207        crate::open_or_initialize_project(directory.path()).unwrap();
2208        create_checkpoint(
2209            directory.path(),
2210            &create_request(Uuid::from_u128(50), "Before"),
2211        )
2212        .unwrap();
2213        let prior = publish_clone(directory.path());
2214        let error = revert_checkpoint(
2215            directory.path(),
2216            &CheckpointRevertRequest {
2217                operation_uuid: Uuid::from_u128(51),
2218                name: "Before".into(),
2219                reason: "must fail closed".into(),
2220                actor_uuid: None,
2221            },
2222            || Ok(1_720_000_000_123_456),
2223            |_| Err(GfError::Validation("injected composite failure".into())),
2224        )
2225        .unwrap_err();
2226        assert_eq!(error.code(), "GF_VALIDATION");
2227        checkpoint_lock_handoff(
2228            directory.path(),
2229            "action=revert validation-error return",
2230            false,
2231        );
2232        assert_eq!(
2233            crate::resolve_project_generation(directory.path())
2234                .unwrap()
2235                .generation_uuid(),
2236            prior
2237        );
2238        assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2239    }
2240
2241    #[cfg(unix)]
2242    #[test]
2243    fn mutation_lock_guard_unlocks_checkpoint_with_retained_duplicate_open() {
2244        let directory = tempdir().unwrap();
2245        crate::open_or_initialize_project(directory.path()).unwrap();
2246        let locks = acquire_mutation_locks(directory.path()).unwrap();
2247        let retained = locks.checkpoint.as_ref().unwrap().try_clone().unwrap();
2248        drop(locks);
2249
2250        let checkpoint =
2251            open_regular_lock(&directory.path().join(LOCKS_DIR).join(CHECKPOINT_LOCK_FILE))
2252                .unwrap();
2253        assert!(FileExt::try_lock_exclusive(&checkpoint).unwrap());
2254        FileExt::unlock(&checkpoint).unwrap();
2255        drop(retained);
2256    }
2257
2258    #[test]
2259    fn create_list_delete_and_replays_are_deterministic() {
2260        let directory = tempdir().unwrap();
2261        crate::open_or_initialize_project(directory.path()).unwrap();
2262        let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000001").unwrap();
2263        let request = create_request(operation, "Release 1.0");
2264        assert_eq!(
2265            hex(&create_request_digest(&request, "Release 1.0")),
2266            "01c7bf2f2c443d85d31ff80fef4a36484e31402213e8d145371867bdb2addbe8"
2267        );
2268        let created = create_checkpoint(directory.path(), &request).unwrap();
2269        assert_eq!(
2270            created.checkpoint_uuid,
2271            Uuid::parse_str("4084179c-38db-8b6b-9b6e-c0b0a855e002").unwrap()
2272        );
2273        let replayed = create_checkpoint(directory.path(), &request).unwrap();
2274        assert_eq!(created, replayed);
2275        assert_eq!(created.registry_revision, 1);
2276
2277        let rows = list_checkpoints(directory.path()).unwrap();
2278        assert_eq!(rows.len(), 1);
2279        assert_eq!(rows[0].checkpoint_uuid, created.checkpoint_uuid);
2280        assert_eq!(rows[0].generation_uuid, created.source_generation_uuid);
2281
2282        let delete = CheckpointDeleteRequest {
2283            operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000002").unwrap(),
2284            name: "Release 1.0".into(),
2285            actor_uuid: request.actor_uuid,
2286        };
2287        assert_eq!(
2288            hex(&delete_request_digest(&delete, "Release 1.0")),
2289            "9e6e15801f66ea4f7f58755c505fc1964e48ab87625459f388e2c23659135bb3"
2290        );
2291        let deleted = delete_checkpoint(directory.path(), &delete).unwrap();
2292        assert_eq!(
2293            deleted,
2294            delete_checkpoint(directory.path(), &delete).unwrap()
2295        );
2296        assert_eq!(deleted.registry_revision, 2);
2297        assert!(list_checkpoints(directory.path()).unwrap().is_empty());
2298        assert_eq!(
2299            created,
2300            create_checkpoint(directory.path(), &request).unwrap()
2301        );
2302        let changed_replay = create_request(operation, "Release 1.1");
2303        assert_eq!(
2304            create_checkpoint(directory.path(), &changed_replay)
2305                .unwrap_err()
2306                .code(),
2307            "GF_IDEMPOTENCY_CONFLICT"
2308        );
2309        assert!(list_checkpoints(directory.path()).unwrap().is_empty());
2310    }
2311
2312    #[test]
2313    fn identity_is_stable_across_independent_projects() {
2314        let first = tempdir().unwrap();
2315        let second = tempdir().unwrap();
2316        crate::open_or_initialize_project(first.path()).unwrap();
2317        crate::open_or_initialize_project(second.path()).unwrap();
2318        let request = create_request(
2319            Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000010").unwrap(),
2320            "Stable",
2321        );
2322        let left = create_checkpoint(first.path(), &request).unwrap();
2323        let right = create_checkpoint(second.path(), &request).unwrap();
2324        assert_eq!(left.checkpoint_uuid, right.checkpoint_uuid);
2325    }
2326
2327    #[test]
2328    fn opened_checkpoint_generation_remains_pinned_after_delete() {
2329        let directory = tempdir().unwrap();
2330        crate::open_or_initialize_project(directory.path()).unwrap();
2331        let created = create_checkpoint(
2332            directory.path(),
2333            &create_request(Uuid::now_v7(), "Pinned View"),
2334        )
2335        .unwrap();
2336        let (row, opened) = open_checkpoint_generation(directory.path(), "Pinned View").unwrap();
2337        assert_eq!(row.checkpoint_uuid, created.checkpoint_uuid);
2338        assert_eq!(opened.generation_uuid(), created.source_generation_uuid);
2339        delete_checkpoint(
2340            directory.path(),
2341            &CheckpointDeleteRequest {
2342                operation_uuid: Uuid::now_v7(),
2343                name: "Pinned View".into(),
2344                actor_uuid: None,
2345            },
2346        )
2347        .unwrap();
2348        assert_eq!(
2349            open_checkpoint_generation(directory.path(), "Pinned View")
2350                .unwrap_err()
2351                .code(),
2352            "GF_CHECKPOINT_NOT_FOUND"
2353        );
2354        assert_eq!(opened.generation_uuid(), created.source_generation_uuid);
2355        assert!(opened.participant_snapshots().is_ok());
2356    }
2357
2358    #[test]
2359    fn conflicts_names_and_corruption_fail_closed() {
2360        let directory = tempdir().unwrap();
2361        crate::open_or_initialize_project(directory.path()).unwrap();
2362        let operation = Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000020").unwrap();
2363        create_checkpoint(directory.path(), &create_request(operation, "Safe.Name")).unwrap();
2364
2365        let conflict =
2366            create_checkpoint(directory.path(), &create_request(operation, "Other")).unwrap_err();
2367        assert_eq!(conflict.code(), "GF_IDEMPOTENCY_CONFLICT");
2368        let exists = create_checkpoint(
2369            directory.path(),
2370            &create_request(
2371                Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000021").unwrap(),
2372                "Safe.Name",
2373            ),
2374        )
2375        .unwrap_err();
2376        assert_eq!(exists.code(), "GF_CHECKPOINT_EXISTS");
2377        for invalid in ["", "../escape", "two  spaces", " e", "e ", ".", ".."] {
2378            let error =
2379                create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), invalid))
2380                    .unwrap_err();
2381            assert_eq!(error.code(), "GF_VALIDATION", "name={invalid:?}");
2382        }
2383
2384        fs::write(
2385            directory.path().join(CHECKPOINTS_DIR).join(CHECKSUM_FILE),
2386            b"0000000000000000000000000000000000000000000000000000000000000000\n",
2387        )
2388        .unwrap();
2389        let error = list_checkpoints(directory.path()).unwrap_err();
2390        assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2391    }
2392
2393    #[test]
2394    fn checksummed_but_impossible_registry_identities_fail_closed() {
2395        let directory = tempdir().unwrap();
2396        crate::open_or_initialize_project(directory.path()).unwrap();
2397        create_checkpoint(
2398            directory.path(),
2399            &create_request(Uuid::now_v7(), "Tampered"),
2400        )
2401        .unwrap();
2402        let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2403        let mut registry = read_registry(&checkpoint_root).unwrap();
2404        registry.active[0].checkpoint_uuid = Uuid::now_v7();
2405        write_raw_registry(directory.path(), &registry);
2406        assert_eq!(
2407            list_checkpoints(directory.path()).unwrap_err().code(),
2408            "GF_CHECKPOINT_REGISTRY_CORRUPT"
2409        );
2410
2411        let directory = tempdir().unwrap();
2412        crate::open_or_initialize_project(directory.path()).unwrap();
2413        create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Overlap")).unwrap();
2414        delete_checkpoint(
2415            directory.path(),
2416            &CheckpointDeleteRequest {
2417                operation_uuid: Uuid::now_v7(),
2418                name: "Overlap".into(),
2419                actor_uuid: None,
2420            },
2421        )
2422        .unwrap();
2423        let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2424        let mut registry = read_registry(&checkpoint_root).unwrap();
2425        let row = &mut registry.tombstones[0];
2426        row.delete_operation_uuid = row.create_operation_uuid;
2427        row.delete_request_sha256 = hex(&delete_request_digest_values(
2428            row.delete_operation_uuid,
2429            &row.name,
2430            row.deleted_by,
2431        ));
2432        write_raw_registry(directory.path(), &registry);
2433        assert_eq!(
2434            list_checkpoints(directory.path()).unwrap_err().code(),
2435            "GF_CHECKPOINT_REGISTRY_CORRUPT"
2436        );
2437    }
2438
2439    #[test]
2440    fn exact_input_bounds_and_writer_lock_are_enforced() {
2441        assert!(validate_name(&"a".repeat(MAX_NAME_BYTES)).is_ok());
2442        assert_eq!(
2443            validate_name(&"a".repeat(MAX_NAME_BYTES + 1))
2444                .unwrap_err()
2445                .code(),
2446            "GF_VALIDATION"
2447        );
2448        assert!(validate_description(Some(&"d".repeat(MAX_DESCRIPTION_BYTES))).is_ok());
2449        assert_eq!(
2450            validate_description(Some(&"d".repeat(MAX_DESCRIPTION_BYTES + 1)))
2451                .unwrap_err()
2452                .code(),
2453            "GF_VALIDATION"
2454        );
2455
2456        let directory = tempdir().unwrap();
2457        let selected = crate::open_or_initialize_project(directory.path()).unwrap();
2458        let root = selected.container_root().to_owned();
2459        let locks = acquire_mutation_locks(&root).unwrap();
2460        let error = create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Busy"))
2461            .unwrap_err();
2462        assert_eq!(error.code(), "GF_WRITER_BUSY");
2463        drop(locks);
2464    }
2465
2466    #[cfg(unix)]
2467    #[test]
2468    fn linked_registry_surfaces_fail_closed_without_following_targets() {
2469        use std::os::unix::fs::symlink;
2470
2471        for hard in [false, true] {
2472            let directory = tempdir().unwrap();
2473            crate::open_or_initialize_project(directory.path()).unwrap();
2474            create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Linked")).unwrap();
2475            let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2476            let checksum = checkpoint_root.join(CHECKSUM_FILE);
2477            let external = directory.path().join("external-checksum");
2478            fs::rename(&checksum, &external).unwrap();
2479            if hard {
2480                fs::hard_link(&external, &checksum).unwrap();
2481            } else {
2482                symlink(&external, &checksum).unwrap();
2483            }
2484            let external_before = fs::read(&external).unwrap();
2485            let error = while_writer_lock_is_held(directory.path(), || {
2486                list_checkpoints(directory.path()).unwrap_err()
2487            });
2488            assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2489            assert_eq!(fs::read(&external).unwrap(), external_before);
2490        }
2491    }
2492
2493    #[test]
2494    fn no_intent_registry_corruption_wins_over_writer_contention() {
2495        let directory = tempdir().unwrap();
2496        crate::open_or_initialize_project(directory.path()).unwrap();
2497        create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Corrupt")).unwrap();
2498        let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2499        fs::write(
2500            checkpoint_root.join(CHECKSUM_FILE),
2501            b"not-the-registry-digest\n",
2502        )
2503        .unwrap();
2504
2505        let expected = read_registry(&checkpoint_root).unwrap_err();
2506        let error = while_writer_lock_is_held(directory.path(), || {
2507            list_checkpoints(directory.path()).unwrap_err()
2508        });
2509        assert_eq!(error.code(), "GF_CHECKPOINT_REGISTRY_CORRUPT");
2510        assert_eq!(error.to_string(), expected.to_string());
2511    }
2512
2513    #[test]
2514    fn intent_recovery_contention_preserves_writer_busy_and_intent() {
2515        let directory = tempdir().unwrap();
2516        crate::open_or_initialize_project(directory.path()).unwrap();
2517        let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
2518        let child = Command::new(std::env::current_exe().unwrap())
2519            .args([
2520                "--exact",
2521                "project_checkpoints::tests::checkpoint_failpoint_helper",
2522                "--ignored",
2523            ])
2524            .env(
2525                "GRAPHFORGE_PROJECT_FAILPOINTS",
2526                "graphforge-internal-subprocess-v1",
2527            )
2528            .env(
2529                "GRAPHFORGE_PROJECT_FAILPOINT",
2530                "checkpoint.registry.before_replace",
2531            )
2532            .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2533            .spawn()
2534            .unwrap();
2535        let status = BoundedChild {
2536            child,
2537            reaped: false,
2538        }
2539        .wait("intent-recovery-contention failpoint=checkpoint.registry.before_replace");
2540        assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2541        let intent_path = checkpoint_root.join(INTENT_FILE);
2542        let intent = fs::read(&intent_path).unwrap();
2543        let staged = fs::read_dir(&checkpoint_root)
2544            .unwrap()
2545            .filter_map(Result::ok)
2546            .filter(|entry| {
2547                entry
2548                    .file_name()
2549                    .to_string_lossy()
2550                    .starts_with(".registry.")
2551            })
2552            .map(|entry| (entry.file_name(), fs::read(entry.path()).unwrap()))
2553            .collect::<BTreeMap<_, _>>();
2554        assert_eq!(staged.len(), 2);
2555
2556        let error = while_writer_lock_is_held(directory.path(), || {
2557            list_checkpoints(directory.path()).unwrap_err()
2558        });
2559        assert_eq!(error.code(), "GF_WRITER_BUSY");
2560        assert_eq!(fs::read(intent_path).unwrap(), intent);
2561        let staged_after = fs::read_dir(&checkpoint_root)
2562            .unwrap()
2563            .filter_map(Result::ok)
2564            .filter(|entry| {
2565                entry
2566                    .file_name()
2567                    .to_string_lossy()
2568                    .starts_with(".registry.")
2569            })
2570            .map(|entry| (entry.file_name(), fs::read(entry.path()).unwrap()))
2571            .collect::<BTreeMap<_, _>>();
2572        assert_eq!(staged_after, staged);
2573    }
2574
2575    #[cfg(unix)]
2576    #[test]
2577    fn linked_project_root_is_rejected_before_checkpoint_access() {
2578        use std::os::unix::fs::symlink;
2579
2580        let directory = tempdir().unwrap();
2581        let project = directory.path().join("project");
2582        fs::create_dir(&project).unwrap();
2583        crate::open_or_initialize_project(&project).unwrap();
2584        let linked = directory.path().join("linked-project");
2585        symlink(&project, &linked).unwrap();
2586        assert_eq!(
2587            list_checkpoints(&linked).unwrap_err().code(),
2588            "GF_UNSUPPORTED_PROJECT_FORMAT"
2589        );
2590    }
2591
2592    #[test]
2593    fn checkpoint_pin_and_open_lease_control_recovery_cleanup() {
2594        let directory = tempdir().unwrap();
2595        crate::open_or_initialize_project(directory.path()).unwrap();
2596        let pinned_generation = publish_clone(directory.path());
2597        let created =
2598            create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Pinned")).unwrap();
2599        assert_eq!(created.source_generation_uuid, pinned_generation);
2600        for _ in 0..4 {
2601            publish_clone(directory.path());
2602        }
2603        crate::recover_project_transactions(directory.path()).unwrap();
2604        recover_checkpoint_pair_after_lock_handoff(
2605            directory.path(),
2606            "action=delete-pinned-checkpoint recovery-complete",
2607        );
2608        let generation_path = directory
2609            .path()
2610            .join(crate::project_publication::GENERATIONS_DIR)
2611            .join(pinned_generation.hyphenated().to_string());
2612        assert!(
2613            generation_path.exists(),
2614            "active checkpoint lost its generation"
2615        );
2616
2617        let lease =
2618            crate::project_publication::open_regular_lock(&generation_path.join("lease.lock"))
2619                .unwrap();
2620        FileExt::lock_shared(&lease).unwrap();
2621        delete_checkpoint(
2622            directory.path(),
2623            &CheckpointDeleteRequest {
2624                operation_uuid: Uuid::now_v7(),
2625                name: "Pinned".into(),
2626                actor_uuid: None,
2627            },
2628        )
2629        .unwrap();
2630        crate::recover_project_transactions(directory.path()).unwrap();
2631        assert!(generation_path.exists(), "an open lease was invalidated");
2632        FileExt::unlock(&lease).unwrap();
2633        drop(lease);
2634        crate::recover_project_transactions(directory.path()).unwrap();
2635        assert!(
2636            !generation_path.exists(),
2637            "deleted pin did not permit later GC"
2638        );
2639    }
2640
2641    #[test]
2642    fn checkpoint_pin_survives_process_restart() {
2643        let directory = tempdir().unwrap();
2644        crate::open_or_initialize_project(directory.path()).unwrap();
2645        let pinned_generation = publish_clone(directory.path());
2646        create_checkpoint(
2647            directory.path(),
2648            &create_request(Uuid::from_u128(70), "Restart Pin"),
2649        )
2650        .unwrap();
2651        for _ in 0..4 {
2652            publish_clone(directory.path());
2653        }
2654
2655        let mut child = Command::new(std::env::current_exe().unwrap())
2656            .args([
2657                "--exact",
2658                "project_checkpoints::tests::checkpoint_failpoint_helper",
2659                "--ignored",
2660                "--nocapture",
2661            ])
2662            .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2663            .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", "hold-open")
2664            .stdin(Stdio::piped())
2665            .stdout(Stdio::piped())
2666            .spawn()
2667            .unwrap();
2668        let mut output = BufReader::new(child.stdout.take().unwrap());
2669        let mut ready = String::new();
2670        while ready != "ready\n" {
2671            ready.clear();
2672            assert_ne!(
2673                output.read_line(&mut ready).unwrap(),
2674                0,
2675                "child exited before ready"
2676            );
2677        }
2678
2679        delete_checkpoint(
2680            directory.path(),
2681            &CheckpointDeleteRequest {
2682                operation_uuid: Uuid::from_u128(71),
2683                name: "Restart Pin".into(),
2684                actor_uuid: None,
2685            },
2686        )
2687        .unwrap();
2688        crate::recover_project_transactions(directory.path()).unwrap();
2689        let generation_path = directory
2690            .path()
2691            .join(crate::project_publication::GENERATIONS_DIR)
2692            .join(pinned_generation.hyphenated().to_string());
2693        assert!(
2694            generation_path.exists(),
2695            "subprocess lease was not retained"
2696        );
2697
2698        child.stdin.take().unwrap().write_all(b"release\n").unwrap();
2699        assert!(child.wait().unwrap().success());
2700        crate::recover_project_transactions(directory.path()).unwrap();
2701        assert!(
2702            !generation_path.exists(),
2703            "generation survived after the restarted reader released its lease"
2704        );
2705    }
2706
2707    #[test]
2708    fn checkpoint_cleanup_removes_all_transient_resources() {
2709        let directory = tempdir().unwrap();
2710        crate::open_or_initialize_project(directory.path()).unwrap();
2711        create_checkpoint(
2712            directory.path(),
2713            &create_request(Uuid::from_u128(80), "Cleanup"),
2714        )
2715        .unwrap();
2716        publish_clone(directory.path());
2717        delete_checkpoint(
2718            directory.path(),
2719            &CheckpointDeleteRequest {
2720                operation_uuid: Uuid::from_u128(81),
2721                name: "Cleanup".into(),
2722                actor_uuid: None,
2723            },
2724        )
2725        .unwrap();
2726        crate::recover_project_transactions(directory.path()).unwrap();
2727
2728        let selected = crate::resolve_project_generation(directory.path()).unwrap();
2729        let root = selected.container_root().to_owned();
2730        drop(selected);
2731        let checkpoint_root = root.join(CHECKPOINTS_DIR);
2732        let checkpoint_entries = fs::read_dir(&checkpoint_root)
2733            .unwrap()
2734            .map(|entry| entry.unwrap().file_name().into_string().unwrap())
2735            .collect::<BTreeSet<_>>();
2736        assert_eq!(
2737            checkpoint_entries,
2738            BTreeSet::from([REGISTRY_FILE.into(), CHECKSUM_FILE.into()]),
2739            "checkpoint transaction staging leaked"
2740        );
2741        let trash = root.join("trash");
2742        assert!(
2743            !trash.exists() || fs::read_dir(&trash).unwrap().next().is_none(),
2744            "recovery trash was not emptied"
2745        );
2746        assert!(
2747            !root.join("cache").exists(),
2748            "checkpoint lifecycle leaked process-cache state to disk"
2749        );
2750        for entry in fs::read_dir(root.join(crate::project_publication::GENERATIONS_DIR)).unwrap() {
2751            let path = entry.unwrap().path();
2752            let name = path.file_name().unwrap().to_str().unwrap();
2753            Uuid::parse_str(name).expect("generation staging entry leaked");
2754            let lease = open_regular_lock(&path.join("lease.lock")).unwrap();
2755            assert!(FileExt::try_lock_exclusive(&lease).unwrap());
2756            FileExt::unlock(&lease).unwrap();
2757        }
2758        let lock_root = root.join(LOCKS_DIR);
2759        for name in [WRITER_LOCK_FILE, CHECKPOINT_LOCK_FILE] {
2760            let lock = open_regular_lock(&lock_root.join(name)).unwrap();
2761            assert!(FileExt::try_lock_exclusive(&lock).unwrap(), "{name} leaked");
2762            FileExt::unlock(&lock).unwrap();
2763        }
2764    }
2765
2766    #[test]
2767    #[ignore = "subprocess failpoint helper"]
2768    fn checkpoint_failpoint_helper() {
2769        let root = std::env::var("GRAPHFORGE_CHECKPOINT_TEST_ROOT").unwrap();
2770        let action = std::env::var("GRAPHFORGE_CHECKPOINT_TEST_ACTION");
2771        if action.as_deref() == Ok("hold-open") {
2772            let (_, opened) = open_checkpoint_generation(root, "Restart Pin").unwrap();
2773            println!("ready");
2774            std::io::stdout().flush().unwrap();
2775            let mut release = String::new();
2776            std::io::stdin().read_line(&mut release).unwrap();
2777            assert_eq!(release, "release\n");
2778            assert!(opened.participant_snapshots().is_ok());
2779        } else if action.as_deref() == Ok("revert") {
2780            revert_checkpoint(
2781                root,
2782                &CheckpointRevertRequest {
2783                    operation_uuid: Uuid::from_u128(61),
2784                    name: "Base".into(),
2785                    reason: "crash recovery".into(),
2786                    actor_uuid: None,
2787                },
2788                || Ok(1_720_000_000_123_456),
2789                |_| Ok(()),
2790            )
2791            .unwrap();
2792        } else if action.as_deref() == Ok("delete") {
2793            delete_checkpoint(
2794                root,
2795                &CheckpointDeleteRequest {
2796                    operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000031")
2797                        .unwrap(),
2798                    name: "Base".into(),
2799                    actor_uuid: None,
2800                },
2801            )
2802            .unwrap();
2803        } else {
2804            let request = create_request(
2805                Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
2806                "Crash",
2807            );
2808            create_checkpoint(root, &request).unwrap();
2809        }
2810    }
2811
2812    #[test]
2813    fn revert_publication_failpoint_matrix() {
2814        let failpoints = [
2815            ("project.after_journal_preparing", false),
2816            ("project.after_participant_dir_fsync", false),
2817            ("project.after_journal_staged", false),
2818            ("project.after_domain_validation", false),
2819            ("project.after_composite_validation", false),
2820            ("project.after_journal_validated", false),
2821            ("project.after_manifest_write", false),
2822            ("project.after_manifest_fsync", false),
2823            ("project.after_generation_dir_fsync", false),
2824            ("project.after_journal_durable", false),
2825            ("project.after_current_temp_write", false),
2826            ("project.after_current_temp_fsync", false),
2827            ("project.before_current_replace", false),
2828            ("project.after_current_replace", true),
2829            ("project.after_root_fsync", true),
2830            ("project.after_journal_published", true),
2831        ];
2832        for (failpoint, committed) in failpoints {
2833            let directory = tempdir().unwrap();
2834            crate::open_or_initialize_project(directory.path()).unwrap();
2835            create_checkpoint(
2836                directory.path(),
2837                &create_request(Uuid::from_u128(60), "Base"),
2838            )
2839            .unwrap();
2840            let prior = publish_clone(directory.path());
2841            let status = Command::new(std::env::current_exe().unwrap())
2842                .args([
2843                    "--exact",
2844                    "project_checkpoints::tests::checkpoint_failpoint_helper",
2845                    "--ignored",
2846                ])
2847                .env(
2848                    "GRAPHFORGE_PROJECT_FAILPOINTS",
2849                    "graphforge-internal-subprocess-v1",
2850                )
2851                .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
2852                .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2853                .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", "revert")
2854                .status()
2855                .unwrap();
2856            assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2857            recover_checkpoint_pair_after_lock_handoff(
2858                directory.path(),
2859                &format!("action=revert failpoint={failpoint} committed={committed}"),
2860            );
2861            crate::recover_project_transactions(directory.path()).unwrap();
2862            let recovered = crate::resolve_project_generation(directory.path()).unwrap();
2863            assert_eq!(
2864                recovered.generation_uuid() != prior,
2865                committed,
2866                "{failpoint}"
2867            );
2868            recover_checkpoint_pair_after_lock_handoff(
2869                directory.path(),
2870                &format!(
2871                    "action=revert parent-recovery-complete failpoint={failpoint} \
2872                     committed={committed}"
2873                ),
2874            );
2875
2876            let (receipt, replayed) = revert_checkpoint(
2877                directory.path(),
2878                &CheckpointRevertRequest {
2879                    operation_uuid: Uuid::from_u128(61),
2880                    name: "Base".into(),
2881                    reason: "crash recovery".into(),
2882                    actor_uuid: None,
2883                },
2884                || Ok(1_720_000_000_123_456),
2885                |_| Ok(()),
2886            )
2887            .unwrap();
2888            assert_eq!(
2889                receipt.result_generation_uuid,
2890                Some(replayed.generation_uuid())
2891            );
2892            recover_checkpoint_pair_after_lock_handoff(
2893                directory.path(),
2894                &format!(
2895                    "action=revert replay-complete failpoint={failpoint} committed={committed}"
2896                ),
2897            );
2898            assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2899        }
2900    }
2901
2902    #[test]
2903    fn registry_failpoints_recover_exact_previous_or_next_revision() {
2904        for (failpoint, committed) in [
2905            ("checkpoint.registry.after_intent_file_fsync", false),
2906            ("checkpoint.registry.after_file_fsync", false),
2907            ("checkpoint.registry.before_replace", false),
2908            ("checkpoint.registry.after_replace", true),
2909            ("checkpoint.registry.after_dir_fsync", true),
2910        ] {
2911            let directory = tempdir().unwrap();
2912            crate::open_or_initialize_project(directory.path()).unwrap();
2913            let status = Command::new(std::env::current_exe().unwrap())
2914                .args([
2915                    "--exact",
2916                    "project_checkpoints::tests::checkpoint_failpoint_helper",
2917                    "--ignored",
2918                ])
2919                .env(
2920                    "GRAPHFORGE_PROJECT_FAILPOINTS",
2921                    "graphforge-internal-subprocess-v1",
2922                )
2923                .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
2924                .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2925                .status()
2926                .unwrap();
2927            assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2928            recover_checkpoint_pair_after_lock_handoff(
2929                directory.path(),
2930                &format!("action=create-unseeded failpoint={failpoint} committed={committed}"),
2931            );
2932            let rows = list_checkpoints(directory.path()).unwrap();
2933            assert_eq!(rows.len(), usize::from(committed), "{failpoint}");
2934            if !committed {
2935                recover_checkpoint_pair_after_lock_handoff(
2936                    directory.path(),
2937                    &format!(
2938                        "action=create-unseeded parent-read-complete failpoint={failpoint} \
2939                         committed={committed}"
2940                    ),
2941                );
2942                let request = create_request(
2943                    Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
2944                    "Crash",
2945                );
2946                create_checkpoint(directory.path(), &request).unwrap();
2947                recover_checkpoint_pair_after_lock_handoff(
2948                    directory.path(),
2949                    &format!(
2950                        "action=create-unseeded replay-complete failpoint={failpoint} \
2951                         committed={committed}"
2952                    ),
2953                );
2954                assert_eq!(list_checkpoints(directory.path()).unwrap().len(), 1);
2955            }
2956        }
2957    }
2958
2959    #[test]
2960    fn seeded_create_and_delete_failpoints_recover_exact_previous_or_next_revision() {
2961        let failpoints = [
2962            ("checkpoint.registry.after_intent_file_fsync", false),
2963            ("checkpoint.registry.after_file_fsync", false),
2964            ("checkpoint.registry.before_replace", false),
2965            ("checkpoint.registry.after_replace", true),
2966            ("checkpoint.registry.after_dir_fsync", true),
2967        ];
2968        for action in ["create", "delete"] {
2969            for (failpoint, committed) in failpoints {
2970                let directory = tempdir().unwrap();
2971                crate::open_or_initialize_project(directory.path()).unwrap();
2972                create_checkpoint(
2973                    directory.path(),
2974                    &create_request(
2975                        Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000029").unwrap(),
2976                        "Base",
2977                    ),
2978                )
2979                .unwrap();
2980                let status = Command::new(std::env::current_exe().unwrap())
2981                    .args([
2982                        "--exact",
2983                        "project_checkpoints::tests::checkpoint_failpoint_helper",
2984                        "--ignored",
2985                    ])
2986                    .env(
2987                        "GRAPHFORGE_PROJECT_FAILPOINTS",
2988                        "graphforge-internal-subprocess-v1",
2989                    )
2990                    .env("GRAPHFORGE_PROJECT_FAILPOINT", failpoint)
2991                    .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
2992                    .env("GRAPHFORGE_CHECKPOINT_TEST_ACTION", action)
2993                    .status()
2994                    .unwrap();
2995                assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
2996                recover_checkpoint_pair_after_lock_handoff(
2997                    directory.path(),
2998                    &format!("action={action} failpoint={failpoint} committed={committed}"),
2999                );
3000                let rows = list_checkpoints(directory.path()).unwrap();
3001                let expected = match (action, committed) {
3002                    ("create", true) => vec!["Base", "Crash"],
3003                    ("create", false) | ("delete", false) => vec!["Base"],
3004                    ("delete", true) => vec![],
3005                    _ => unreachable!(),
3006                };
3007                assert_eq!(
3008                    rows.iter().map(|row| row.name.as_str()).collect::<Vec<_>>(),
3009                    expected,
3010                    "{action} {failpoint}"
3011                );
3012                recover_checkpoint_pair_after_lock_handoff(
3013                    directory.path(),
3014                    &format!(
3015                        "action={action} parent-read-complete failpoint={failpoint} \
3016                         committed={committed}"
3017                    ),
3018                );
3019
3020                if action == "create" {
3021                    let replay = create_checkpoint(
3022                        directory.path(),
3023                        &create_request(
3024                            Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000030").unwrap(),
3025                            "Crash",
3026                        ),
3027                    )
3028                    .unwrap();
3029                    assert_eq!(replay.registry_revision, 2);
3030                } else {
3031                    let replay = delete_checkpoint(
3032                        directory.path(),
3033                        &CheckpointDeleteRequest {
3034                            operation_uuid: Uuid::parse_str("018f0f4e-7b8c-7000-8000-000000000031")
3035                                .unwrap(),
3036                            name: "Base".into(),
3037                            actor_uuid: None,
3038                        },
3039                    )
3040                    .unwrap();
3041                    assert_eq!(replay.registry_revision, 2);
3042                }
3043                recover_checkpoint_pair_after_lock_handoff(
3044                    directory.path(),
3045                    &format!(
3046                        "action={action} replay-complete failpoint={failpoint} \
3047                         committed={committed}"
3048                    ),
3049                );
3050                let final_rows = list_checkpoints(directory.path()).unwrap();
3051                let final_names = final_rows
3052                    .iter()
3053                    .map(|row| row.name.as_str())
3054                    .collect::<Vec<_>>();
3055                if action == "create" {
3056                    assert_eq!(final_names, vec!["Base", "Crash"], "{failpoint}");
3057                } else {
3058                    assert!(final_names.is_empty(), "{failpoint}");
3059                }
3060            }
3061        }
3062    }
3063
3064    #[test]
3065    fn wave9_durable_registry_intent_recovers_every_atomic_pair_boundary() {
3066        for boundary in [
3067            "first-staged",
3068            "previous-staged",
3069            "registry-replaced",
3070            "next-complete",
3071        ] {
3072            let directory = tempdir().unwrap();
3073            crate::open_or_initialize_project(directory.path()).unwrap();
3074            let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
3075            let previous = Registry::empty();
3076            let mut next = Registry::empty();
3077            next.revision = 1;
3078
3079            let has_previous = boundary != "first-staged";
3080            fs::create_dir_all(&checkpoint_root).unwrap();
3081            if has_previous {
3082                write_raw_registry(directory.path(), &previous);
3083            }
3084            let intent =
3085                install_registry_intent(&checkpoint_root, has_previous.then_some(&previous), &next);
3086            match boundary {
3087                "first-staged" | "previous-staged" => {}
3088                "registry-replaced" => {
3089                    fs::rename(
3090                        checkpoint_root.join(&intent.registry_temp),
3091                        checkpoint_root.join(REGISTRY_FILE),
3092                    )
3093                    .unwrap();
3094                }
3095                "next-complete" => {
3096                    fs::rename(
3097                        checkpoint_root.join(&intent.registry_temp),
3098                        checkpoint_root.join(REGISTRY_FILE),
3099                    )
3100                    .unwrap();
3101                    fs::rename(
3102                        checkpoint_root.join(&intent.checksum_temp),
3103                        checkpoint_root.join(CHECKSUM_FILE),
3104                    )
3105                    .unwrap();
3106                }
3107                _ => unreachable!(),
3108            }
3109
3110            recover_pair(&checkpoint_root).unwrap();
3111
3112            assert!(!checkpoint_root.join(INTENT_FILE).exists(), "{boundary}");
3113            assert!(
3114                !checkpoint_root.join(&intent.registry_temp).exists(),
3115                "{boundary}"
3116            );
3117            assert!(
3118                !checkpoint_root.join(&intent.checksum_temp).exists(),
3119                "{boundary}"
3120            );
3121            let recovered = read_registry(&checkpoint_root).unwrap();
3122            let expected_revision =
3123                usize::from(matches!(boundary, "registry-replaced" | "next-complete"));
3124            assert_eq!(recovered.revision, expected_revision as u64, "{boundary}");
3125        }
3126    }
3127
3128    #[test]
3129    fn wave9_registry_intent_rejects_unsafe_names_and_wrong_staged_revision() {
3130        let directory = tempdir().unwrap();
3131        crate::open_or_initialize_project(directory.path()).unwrap();
3132        let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
3133        let mut next = Registry::empty();
3134        next.revision = 1;
3135        let mut intent = install_registry_intent(&checkpoint_root, None, &next);
3136        intent.registry_temp = "../registry.json".into();
3137        let mut bytes = serde_json::to_vec(&intent).unwrap();
3138        bytes.push(b'\n');
3139        fs::write(checkpoint_root.join(INTENT_FILE), bytes).unwrap();
3140        assert_eq!(
3141            recover_pair(&checkpoint_root).unwrap_err().code(),
3142            "GF_CHECKPOINT_REGISTRY_CORRUPT"
3143        );
3144
3145        fs::remove_dir_all(&checkpoint_root).unwrap();
3146        fs::create_dir(&checkpoint_root).unwrap();
3147        let mut intent = install_registry_intent(&checkpoint_root, None, &next);
3148        intent.next_revision = 2;
3149        let mut bytes = serde_json::to_vec(&intent).unwrap();
3150        bytes.push(b'\n');
3151        fs::write(checkpoint_root.join(INTENT_FILE), bytes).unwrap();
3152        assert_eq!(
3153            recover_pair(&checkpoint_root).unwrap_err().code(),
3154            "GF_CHECKPOINT_REGISTRY_CORRUPT"
3155        );
3156    }
3157
3158    #[test]
3159    fn recovery_rejects_missing_or_tampered_staged_pair() {
3160        for tamper_checksum in [false, true] {
3161            let directory = tempdir().unwrap();
3162            crate::open_or_initialize_project(directory.path()).unwrap();
3163            create_checkpoint(directory.path(), &create_request(Uuid::now_v7(), "Base")).unwrap();
3164            let status = Command::new(std::env::current_exe().unwrap())
3165                .args([
3166                    "--exact",
3167                    "project_checkpoints::tests::checkpoint_failpoint_helper",
3168                    "--ignored",
3169                ])
3170                .env(
3171                    "GRAPHFORGE_PROJECT_FAILPOINTS",
3172                    "graphforge-internal-subprocess-v1",
3173                )
3174                .env(
3175                    "GRAPHFORGE_PROJECT_FAILPOINT",
3176                    "checkpoint.registry.before_replace",
3177                )
3178                .env("GRAPHFORGE_CHECKPOINT_TEST_ROOT", directory.path())
3179                .status()
3180                .unwrap();
3181            assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
3182            preserve_checkpoint_intent_after_lock_handoff(
3183                directory.path(),
3184                &format!(
3185                    "action=create-tamper failpoint=checkpoint.registry.before_replace \
3186                     committed=false tamper_checksum={tamper_checksum}"
3187                ),
3188            );
3189            let checkpoint_root = directory.path().join(CHECKPOINTS_DIR);
3190            let intent_bytes = fs::read(checkpoint_root.join(INTENT_FILE)).unwrap();
3191            let intent: RegistryIntent = serde_json::from_slice(&intent_bytes).unwrap();
3192            if tamper_checksum {
3193                fs::write(
3194                    checkpoint_root.join(intent.checksum_temp),
3195                    b"0000000000000000000000000000000000000000000000000000000000000000\n",
3196                )
3197                .unwrap();
3198            } else {
3199                fs::remove_file(checkpoint_root.join(intent.registry_temp)).unwrap();
3200            }
3201            assert_eq!(
3202                list_checkpoints(directory.path()).unwrap_err().code(),
3203                "GF_CHECKPOINT_REGISTRY_CORRUPT"
3204            );
3205        }
3206    }
3207
3208    #[test]
3209    fn checkpoint_text_and_identity_boundaries_are_canonical() {
3210        for invalid in [
3211            "",
3212            " leading",
3213            "trailing ",
3214            ".",
3215            "..",
3216            "two  spaces",
3217            "bad/name",
3218        ] {
3219            assert_eq!(validate_name(invalid).unwrap_err().code(), "GF_VALIDATION");
3220        }
3221        assert_eq!(validate_name("Résumé.v1").unwrap(), "Résumé.v1");
3222        let combining_acute = char::from_u32(0x301).unwrap();
3223        let decomposed = format!("Re{combining_acute}sume{combining_acute}");
3224        assert_eq!(
3225            validate_name(&decomposed).unwrap_err().code(),
3226            "GF_VALIDATION"
3227        );
3228        assert!(validate_description(None).is_ok());
3229        assert!(validate_description(Some("bounded description")).is_ok());
3230        assert_eq!(
3231            validate_description(Some("contains\ncontrol"))
3232                .unwrap_err()
3233                .code(),
3234            "GF_VALIDATION"
3235        );
3236        assert_eq!(
3237            validate_description(Some(&"x".repeat(MAX_DESCRIPTION_BYTES + 1)))
3238                .unwrap_err()
3239                .code(),
3240            "GF_VALIDATION"
3241        );
3242        assert_eq!(
3243            validate_reason("  restored after audit  ").unwrap(),
3244            "restored after audit"
3245        );
3246        assert_eq!(validate_reason("   ").unwrap_err().code(), "GF_VALIDATION");
3247
3248        let operation = Uuid::now_v7();
3249        let actor = Uuid::now_v7();
3250        let digest = create_request_digest_values(operation, "baseline", Some("desc"), Some(actor));
3251        let checkpoint = checkpoint_uuid(operation, digest);
3252        let encoded = hex(&digest);
3253        assert_eq!(decode_digest(&encoded).unwrap(), digest);
3254        assert!(
3255            validate_record_identity(
3256                checkpoint,
3257                operation,
3258                "baseline",
3259                Some("desc"),
3260                Some(actor),
3261                &encoded,
3262            )
3263            .is_ok()
3264        );
3265        assert_eq!(
3266            validate_record_identity(
3267                Uuid::nil(),
3268                operation,
3269                "baseline",
3270                Some("desc"),
3271                Some(actor),
3272                &encoded,
3273            )
3274            .unwrap_err()
3275            .code(),
3276            "GF_CHECKPOINT_REGISTRY_CORRUPT"
3277        );
3278        for malformed in ["0", &"A".repeat(64), &"g".repeat(64)] {
3279            assert_eq!(
3280                decode_digest(malformed).unwrap_err().code(),
3281                "GF_CHECKPOINT_REGISTRY_CORRUPT"
3282            );
3283        }
3284        assert!(valid_private_name(
3285            &format!(".registry.{operation}.json.next"),
3286            operation,
3287            "json"
3288        ));
3289        assert!(!valid_private_name(
3290            ".registry.other.json.next",
3291            operation,
3292            "json"
3293        ));
3294        assert_eq!(
3295            parse_uuid("not-a-uuid").unwrap_err().code(),
3296            "GF_CHECKPOINT_REGISTRY_CORRUPT"
3297        );
3298    }
3299
3300    #[test]
3301    fn empty_checkpoint_registry_has_stable_canonical_bytes() {
3302        let registry = Registry::empty();
3303        let first = registry.canonical_bytes().unwrap();
3304        let second = registry.canonical_bytes().unwrap();
3305        assert_eq!(first, second);
3306        assert!(first.ends_with(b"\n"));
3307        let decoded: Registry = serde_json::from_slice(&first).unwrap();
3308        assert_eq!(decoded, registry);
3309    }
3310
3311    #[test]
3312    fn checkpoint_operation_identities_remain_disjoint_across_tombstones() {
3313        let root = tempdir().unwrap();
3314        crate::open_or_initialize_project(root.path()).unwrap();
3315        let create_operation = Uuid::now_v7();
3316        let create = create_request(create_operation, "Baseline");
3317        let created = create_checkpoint(root.path(), &create).unwrap();
3318        let exact_create_replay = create_checkpoint(root.path(), &create).unwrap();
3319        assert_eq!(exact_create_replay, created);
3320
3321        let changed_create = CheckpointCreateRequest {
3322            name: "Changed".into(),
3323            ..create.clone()
3324        };
3325        assert_eq!(
3326            create_checkpoint(root.path(), &changed_create)
3327                .unwrap_err()
3328                .code(),
3329            "GF_IDEMPOTENCY_CONFLICT"
3330        );
3331        assert_eq!(list_checkpoints(root.path()).unwrap().len(), 1);
3332
3333        let delete_operation = Uuid::now_v7();
3334        let delete = CheckpointDeleteRequest {
3335            operation_uuid: delete_operation,
3336            name: "Baseline".into(),
3337            actor_uuid: create.actor_uuid,
3338        };
3339        let deleted = delete_checkpoint(root.path(), &delete).unwrap();
3340        let exact_delete_replay = delete_checkpoint(root.path(), &delete).unwrap();
3341        assert_eq!(exact_delete_replay, deleted);
3342        assert!(list_checkpoints(root.path()).unwrap().is_empty());
3343
3344        let tombstone_create_replay = create_checkpoint(root.path(), &create).unwrap();
3345        assert_eq!(
3346            tombstone_create_replay.checkpoint_uuid,
3347            created.checkpoint_uuid
3348        );
3349        assert_eq!(tombstone_create_replay, created);
3350        assert!(list_checkpoints(root.path()).unwrap().is_empty());
3351
3352        let changed_delete = CheckpointDeleteRequest {
3353            name: "Other".into(),
3354            ..delete.clone()
3355        };
3356        assert_eq!(
3357            delete_checkpoint(root.path(), &changed_delete)
3358                .unwrap_err()
3359                .code(),
3360            "GF_IDEMPOTENCY_CONFLICT"
3361        );
3362        assert_eq!(
3363            create_checkpoint(root.path(), &create_request(delete_operation, "Other"))
3364                .unwrap_err()
3365                .code(),
3366            "GF_IDEMPOTENCY_CONFLICT"
3367        );
3368        assert_eq!(
3369            delete_checkpoint(
3370                root.path(),
3371                &CheckpointDeleteRequest {
3372                    operation_uuid: create_operation,
3373                    name: "Missing".into(),
3374                    actor_uuid: None,
3375                },
3376            )
3377            .unwrap_err()
3378            .code(),
3379            "GF_IDEMPOTENCY_CONFLICT"
3380        );
3381        assert_eq!(
3382            delete_checkpoint(
3383                root.path(),
3384                &CheckpointDeleteRequest {
3385                    operation_uuid: Uuid::now_v7(),
3386                    name: "Missing".into(),
3387                    actor_uuid: None,
3388                },
3389            )
3390            .unwrap_err()
3391            .code(),
3392            "GF_CHECKPOINT_NOT_FOUND"
3393        );
3394        assert!(list_checkpoints(root.path()).unwrap().is_empty());
3395    }
3396
3397    #[test]
3398    fn registry_header_sort_and_revision_validation_matrix_uses_durable_records() {
3399        let directory = tempdir().unwrap();
3400        crate::open_or_initialize_project(directory.path()).unwrap();
3401        create_checkpoint(
3402            directory.path(),
3403            &create_request(Uuid::from_u128(701), "Alpha"),
3404        )
3405        .unwrap();
3406        create_checkpoint(
3407            directory.path(),
3408            &create_request(Uuid::from_u128(702), "Beta"),
3409        )
3410        .unwrap();
3411        let root = directory.path().join(CHECKPOINTS_DIR);
3412        let stable = read_registry(&root).unwrap();
3413        assert!(validate_registry(&stable).is_ok());
3414
3415        let mutations: Vec<Box<dyn Fn(&mut Registry)>> = vec![
3416            Box::new(|registry| registry.format = "future".into()),
3417            Box::new(|registry| registry.format_version = 2),
3418            Box::new(|registry| registry.active.reverse()),
3419            Box::new(|registry| registry.active[0].created_revision = 0),
3420            Box::new(|registry| registry.active[1].name = registry.active[0].name.clone()),
3421            Box::new(|registry| {
3422                registry.active[1].checkpoint_uuid = registry.active[0].checkpoint_uuid
3423            }),
3424            Box::new(|registry| {
3425                registry.active[1].create_operation_uuid = registry.active[0].create_operation_uuid
3426            }),
3427        ];
3428        for mutate in mutations {
3429            let mut candidate = stable.clone();
3430            mutate(&mut candidate);
3431            assert_eq!(
3432                validate_registry(&candidate).unwrap_err().code(),
3433                "GF_CHECKPOINT_REGISTRY_CORRUPT"
3434            );
3435        }
3436        assert_eq!(read_registry(&root).unwrap(), stable);
3437    }
3438
3439    #[test]
3440    fn tombstone_identity_revision_and_operation_disjointness_matrix_is_total() {
3441        let directory = tempdir().unwrap();
3442        crate::open_or_initialize_project(directory.path()).unwrap();
3443        create_checkpoint(
3444            directory.path(),
3445            &create_request(Uuid::from_u128(801), "Deleted"),
3446        )
3447        .unwrap();
3448        delete_checkpoint(
3449            directory.path(),
3450            &CheckpointDeleteRequest {
3451                operation_uuid: Uuid::from_u128(802),
3452                name: "Deleted".into(),
3453                actor_uuid: None,
3454            },
3455        )
3456        .unwrap();
3457        let root = directory.path().join(CHECKPOINTS_DIR);
3458        let stable = read_registry(&root).unwrap();
3459        assert_eq!(stable.active.len(), 0);
3460        assert_eq!(stable.tombstones.len(), 1);
3461        assert!(validate_registry(&stable).is_ok());
3462
3463        let mutations: Vec<Box<dyn Fn(&mut Registry)>> = vec![
3464            Box::new(|registry| registry.tombstones[0].name = " bad".into()),
3465            Box::new(|registry| registry.tombstones[0].description = Some("bad\nvalue".into())),
3466            Box::new(|registry| registry.tombstones[0].generation_manifest_sha256 = "bad".into()),
3467            Box::new(|registry| registry.tombstones[0].create_request_sha256 = "bad".into()),
3468            Box::new(|registry| registry.tombstones[0].delete_request_sha256 = "bad".into()),
3469            Box::new(|registry| registry.tombstones[0].checkpoint_uuid = Uuid::nil()),
3470            Box::new(|registry| registry.tombstones[0].created_revision = 0),
3471            Box::new(|registry| {
3472                registry.tombstones[0].deleted_revision = registry.tombstones[0].created_revision
3473            }),
3474            Box::new(|registry| registry.tombstones[0].deleted_revision = registry.revision + 1),
3475            Box::new(|registry| {
3476                registry.tombstones[0].delete_operation_uuid =
3477                    registry.tombstones[0].create_operation_uuid
3478            }),
3479            Box::new(|registry| registry.revision = 0),
3480        ];
3481        for mutate in mutations {
3482            let mut candidate = stable.clone();
3483            mutate(&mut candidate);
3484            assert_eq!(
3485                validate_registry(&candidate).unwrap_err().code(),
3486                "GF_CHECKPOINT_REGISTRY_CORRUPT"
3487            );
3488        }
3489
3490        let mut duplicate = stable.clone();
3491        let mut second = duplicate.tombstones[0].clone();
3492        second.deleted_revision += 1;
3493        duplicate.revision = second.deleted_revision;
3494        duplicate.tombstones.push(second);
3495        assert_eq!(
3496            validate_registry(&duplicate).unwrap_err().code(),
3497            "GF_CHECKPOINT_REGISTRY_CORRUPT"
3498        );
3499        assert_eq!(read_registry(&root).unwrap(), stable);
3500    }
3501
3502    #[test]
3503    fn public_checkpoint_operations_reject_cross_kind_uuid_reuse_after_reopen() {
3504        let root = tempdir().unwrap();
3505        crate::open_or_initialize_project(root.path()).unwrap();
3506        let operation = Uuid::now_v7();
3507        create_checkpoint(root.path(), &create_request(operation, "release")).unwrap();
3508
3509        let delete = CheckpointDeleteRequest {
3510            operation_uuid: operation,
3511            name: "release".into(),
3512            actor_uuid: None,
3513        };
3514        assert_eq!(
3515            delete_checkpoint(root.path(), &delete).unwrap_err().code(),
3516            "GF_IDEMPOTENCY_CONFLICT"
3517        );
3518
3519        let delete_operation = Uuid::now_v7();
3520        delete_checkpoint(
3521            root.path(),
3522            &CheckpointDeleteRequest {
3523                operation_uuid: delete_operation,
3524                name: "release".into(),
3525                actor_uuid: None,
3526            },
3527        )
3528        .unwrap();
3529        assert_eq!(
3530            create_checkpoint(
3531                root.path(),
3532                &create_request(delete_operation, "replacement")
3533            )
3534            .unwrap_err()
3535            .code(),
3536            "GF_IDEMPOTENCY_CONFLICT"
3537        );
3538        assert!(list_checkpoints(root.path()).unwrap().is_empty());
3539    }
3540}