1use std::fs::{File, OpenOptions};
8use std::io::{Read, Write};
9use std::path::{Path, PathBuf};
10
11use atomicwrites::{AllowOverwrite, AtomicFile};
12use fs4::fs_std::FileExt;
13use graphforge_core::{GfError, ProjectErrorCode};
14use serde::{Deserialize, Serialize};
15use sha2::{Digest, Sha256};
16use uuid::Uuid;
17
18use crate::project_failpoint;
19use crate::project_generation::{
20 CURRENT_FILE, ResolvedProjectGeneration, resolve_project_generation,
21};
22
23pub(crate) const LOCKS_DIR: &str = "locks";
24pub(crate) const WRITER_LOCK_FILE: &str = "writer.lock";
25pub(crate) const TRANSACTION_LOCKS_DIR: &str = "transactions";
26pub(crate) const TRANSACTIONS_DIR: &str = "transactions";
27pub(crate) const ATTEMPTS_DIR: &str = "attempts";
28pub(crate) const GENERATIONS_DIR: &str = "generations";
29const PARTICIPANTS_DIR: &str = "participants";
30const LEASE_FILE: &str = "lease.lock";
31const MANIFEST_FILE: &str = "manifest.json";
32const MAX_JOURNAL_BYTES: u64 = 1024 * 1024;
33
34#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
36pub enum ProjectParticipantEncoding {
37 Parquet,
39 Arrow,
41 Json,
43}
44
45impl ProjectParticipantEncoding {
46 const fn extension(self) -> &'static str {
47 match self {
48 Self::Parquet => "parquet",
49 Self::Arrow => "arrow",
50 Self::Json => "json",
51 }
52 }
53}
54
55#[derive(Debug, Clone)]
57pub struct ProjectParticipant {
58 pub capability_id: String,
60 pub capability_version: u32,
62 pub record_family_id: String,
64 pub record_version: u32,
66 pub encoding: ProjectParticipantEncoding,
68 pub schema_fingerprint: [u8; 32],
70 pub row_count: u64,
72 pub bytes: Vec<u8>,
74}
75
76#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
78pub struct ProjectCapability {
79 pub capability_id: String,
81 pub capability_version: u32,
83}
84
85#[derive(Debug, Clone)]
92pub struct ProjectGenerationRequest {
93 pub transaction_uuid: Uuid,
95 pub generation_uuid: Uuid,
97 pub capabilities: Vec<ProjectCapability>,
99 pub participants: Vec<ProjectParticipant>,
101}
102
103#[derive(Debug, Clone, PartialEq, Eq)]
105pub struct StagedParticipant {
106 pub capability_id: String,
108 pub capability_version: u32,
110 pub record_family_id: String,
112 pub record_version: u32,
114 pub relative_path: String,
116 pub encoding: String,
118 pub byte_length: u64,
120 pub row_count: u64,
122 pub schema_fingerprint: String,
124 pub content_sha256: String,
126}
127
128#[derive(Debug, Clone, PartialEq, Eq)]
130pub struct ProjectPublicationReceipt {
131 pub transaction_uuid: Uuid,
133 pub generation_uuid: Uuid,
135 pub generation_manifest_sha256: [u8; 32],
137 pub idempotent_replay: bool,
139}
140
141pub enum ProjectStageOutcome {
143 Staged(Box<StagedProjectGeneration>),
145 AlreadyPublished(ProjectPublicationReceipt),
147}
148
149pub struct StagedProjectGeneration {
151 root: PathBuf,
152 publication_lock: PublicationLock,
153 parent: ResolvedProjectGeneration,
154 transaction_uuid: Uuid,
155 generation_uuid: Uuid,
156 generation_root: PathBuf,
157 requires_promotion: bool,
158 request_fingerprint: String,
159 operation_fingerprint: String,
160 capabilities: Vec<ProjectCapability>,
161 participants: Vec<StagedParticipant>,
162 revert: Option<RevertJournalExtension>,
163}
164
165enum PublicationLock {
166 Exclusive(File),
167 Optimistic(File),
168}
169
170#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
172#[serde(deny_unknown_fields)]
173pub(crate) struct RevertJournalExtension {
174 pub(crate) operation_uuid: String,
175 pub(crate) request_sha256: String,
176 pub(crate) checkpoint_uuid: String,
177 pub(crate) checkpoint_name: String,
178 pub(crate) source_generation_uuid: String,
179 pub(crate) source_manifest_sha256: String,
180 pub(crate) prior_current_generation_uuid: String,
181 pub(crate) restored_at: i64,
182 pub(crate) reason: String,
183 pub(crate) actor_uuid: Option<String>,
184 pub(crate) restoration_uuid: String,
185 pub(crate) registry_revision: u64,
186}
187
188pub struct ValidatedProjectGeneration(StagedProjectGeneration);
190
191pub fn stage_project_generation(
201 container_root: impl AsRef<Path>,
202 request: &ProjectGenerationRequest,
203) -> Result<ProjectStageOutcome, GfError> {
204 stage_project_generation_inner(container_root.as_ref(), request)
205 .map_err(|error| map_stage_error(request, error))
206}
207
208pub fn stage_project_generation_optimistic(
220 container_root: impl AsRef<Path>,
221 request: &ProjectGenerationRequest,
222 operation_fingerprint: [u8; 32],
223) -> Result<ProjectStageOutcome, GfError> {
224 stage_project_generation_optimistic_inner(
225 container_root.as_ref(),
226 request,
227 operation_fingerprint,
228 )
229 .map_err(|error| map_stage_error(request, error))
230}
231
232fn map_stage_error(request: &ProjectGenerationRequest, error: GfError) -> GfError {
233 match error {
234 GfError::Storage(message) => publication_error(request, "STAGE", false, &message),
235 other => other,
236 }
237}
238
239fn stage_project_generation_inner(
240 container_root: &Path,
241 request: &ProjectGenerationRequest,
242) -> Result<ProjectStageOutcome, GfError> {
243 let root = canonical_supported_root(container_root)?;
244 let writer_lock = acquire_writer_lock(&root, request)?;
245 project_failpoint::hit(
246 "project.after_writer_lock",
247 Some(request.transaction_uuid),
248 Some(request.generation_uuid),
249 "WRITER_LOCK",
250 false,
251 )?;
252 let parent = resolve_project_generation(&root)?;
253 stage_project_generation_with_lock(root, writer_lock, parent, request, None)
254}
255
256fn stage_project_generation_optimistic_inner(
257 container_root: &Path,
258 request: &ProjectGenerationRequest,
259 operation_fingerprint: [u8; 32],
260) -> Result<ProjectStageOutcome, GfError> {
261 let root = canonical_supported_root(container_root)?;
262 let transaction_lock = acquire_transaction_lock(&root, request)?;
263 let parent = resolve_project_generation(&root)?;
264 stage_project_generation_inner_with_locks(
265 root,
266 PublicationLock::Optimistic(transaction_lock),
267 parent,
268 request,
269 None,
270 Some(operation_fingerprint),
271 )
272}
273
274pub(crate) fn stage_project_generation_with_lock(
277 root: PathBuf,
278 writer_lock: File,
279 parent: ResolvedProjectGeneration,
280 request: &ProjectGenerationRequest,
281 revert: Option<RevertJournalExtension>,
282) -> Result<ProjectStageOutcome, GfError> {
283 stage_project_generation_inner_with_locks(
284 root,
285 PublicationLock::Exclusive(writer_lock),
286 parent,
287 request,
288 revert,
289 None,
290 )
291}
292
293fn stage_project_generation_inner_with_locks(
294 root: PathBuf,
295 publication_lock: PublicationLock,
296 parent: ResolvedProjectGeneration,
297 request: &ProjectGenerationRequest,
298 revert: Option<RevertJournalExtension>,
299 operation_fingerprint: Option<[u8; 32]>,
300) -> Result<ProjectStageOutcome, GfError> {
301 validate_request(request)?;
302 let (capabilities, participants, request_fingerprint) = request_metadata(request)?;
303 let operation_fingerprint =
304 operation_fingerprint.map_or_else(|| request_fingerprint.clone(), hex_digest);
305 let transactions_dir = ensure_machine_directory(&root, Path::new(TRANSACTIONS_DIR))?;
306 sync_directory(&root)?;
307 let journal_path =
308 transactions_dir.join(format!("{}.json", request.transaction_uuid.hyphenated()));
309 if journal_path.exists()
310 && let Some(outcome) = handle_existing_journal(
311 &root,
312 request,
313 &request_fingerprint,
314 &operation_fingerprint,
315 revert.as_ref(),
316 &journal_path,
317 )?
318 {
319 return Ok(outcome);
320 }
321
322 let requires_promotion = matches!(publication_lock, PublicationLock::Optimistic(_));
323 let generation_root =
324 prepare_generation_directory(&root, request, &request_fingerprint, requires_promotion)?;
325 write_journal(
326 &journal_path,
327 &JournalRecord::new(
328 request,
329 Some(parent.generation_uuid()),
330 JournalPhase::Preparing,
331 (request_fingerprint.clone(), operation_fingerprint.clone()),
332 &participants,
333 None,
334 revert.clone(),
335 ),
336 )?;
337 project_failpoint::hit(
338 "project.after_journal_preparing",
339 Some(request.transaction_uuid),
340 Some(request.generation_uuid),
341 "PREPARING",
342 false,
343 )?;
344
345 stage_participant_files(request, &generation_root, &participants)?;
346 sync_participant_directories(&generation_root.join(PARTICIPANTS_DIR), &participants)?;
347 project_failpoint::hit(
348 "project.after_participant_dir_fsync",
349 Some(request.transaction_uuid),
350 Some(request.generation_uuid),
351 "STAGED",
352 false,
353 )?;
354 write_journal(
355 &journal_path,
356 &JournalRecord::new(
357 request,
358 Some(parent.generation_uuid()),
359 JournalPhase::Staged,
360 (request_fingerprint.clone(), operation_fingerprint.clone()),
361 &participants,
362 None,
363 revert.clone(),
364 ),
365 )?;
366 project_failpoint::hit(
367 "project.after_journal_staged",
368 Some(request.transaction_uuid),
369 Some(request.generation_uuid),
370 "STAGED",
371 false,
372 )?;
373
374 Ok(ProjectStageOutcome::Staged(Box::new(
375 StagedProjectGeneration {
376 root,
377 publication_lock,
378 parent,
379 transaction_uuid: request.transaction_uuid,
380 generation_uuid: request.generation_uuid,
381 generation_root,
382 requires_promotion,
383 request_fingerprint,
384 operation_fingerprint,
385 capabilities,
386 participants,
387 revert,
388 },
389 )))
390}
391
392fn acquire_writer_lock(root: &Path, request: &ProjectGenerationRequest) -> Result<File, GfError> {
393 acquire_writer_lock_for_parts(root, request.transaction_uuid, request.generation_uuid)
394}
395
396fn acquire_writer_lock_for_parts(
397 root: &Path,
398 transaction_uuid: Uuid,
399 generation_uuid: Uuid,
400) -> Result<File, GfError> {
401 let lock_dir = ensure_machine_directory(root, Path::new(LOCKS_DIR))?;
402 sync_directory(root)?;
403 let writer_lock = open_regular_lock(&lock_dir.join(WRITER_LOCK_FILE))?;
404 if !FileExt::try_lock_exclusive(&writer_lock).map_err(publication_io)? {
405 return Err(project_error(
406 ProjectErrorCode::WriterBusy,
407 format!(
408 "transaction_uuid={} generation_uuid={} phase=WRITER_LOCK committed=false cause=busy",
409 transaction_uuid.hyphenated(),
410 generation_uuid.hyphenated()
411 ),
412 ));
413 }
414 Ok(writer_lock)
415}
416
417fn wait_for_writer_lock(root: &Path) -> Result<File, GfError> {
418 let lock_dir = ensure_machine_directory(root, Path::new(LOCKS_DIR))?;
419 sync_directory(root)?;
420 let writer_lock = open_regular_lock(&lock_dir.join(WRITER_LOCK_FILE))?;
421 FileExt::lock_exclusive(&writer_lock).map_err(publication_io)?;
422 Ok(writer_lock)
423}
424
425fn acquire_transaction_lock(
426 root: &Path,
427 request: &ProjectGenerationRequest,
428) -> Result<File, GfError> {
429 let lock = open_transaction_lock(root, request.transaction_uuid)?;
430 if !FileExt::try_lock_exclusive(&lock).map_err(publication_io)? {
431 return Err(project_error(
432 ProjectErrorCode::WriterBusy,
433 format!(
434 "transaction_uuid={} generation_uuid={} phase=TRANSACTION_LOCK committed=false cause=busy",
435 request.transaction_uuid.hyphenated(),
436 request.generation_uuid.hyphenated()
437 ),
438 ));
439 }
440 Ok(lock)
441}
442
443pub(crate) fn open_transaction_lock(root: &Path, transaction_uuid: Uuid) -> Result<File, GfError> {
444 let lock_dir =
445 ensure_machine_directory(root, &Path::new(LOCKS_DIR).join(TRANSACTION_LOCKS_DIR))?;
446 open_regular_lock(&lock_dir.join(format!("{}.lock", transaction_uuid.hyphenated())))
447}
448
449fn handle_existing_journal(
450 root: &Path,
451 request: &ProjectGenerationRequest,
452 request_fingerprint: &str,
453 operation_fingerprint: &str,
454 expected_revert: Option<&RevertJournalExtension>,
455 journal_path: &Path,
456) -> Result<Option<ProjectStageOutcome>, GfError> {
457 let journal = read_journal(journal_path)?;
458 if journal.operation_fingerprint() != operation_fingerprint
459 || journal.generation_uuid != request.generation_uuid.hyphenated().to_string()
460 || journal.revert.as_ref() != expected_revert
461 {
462 return Err(transaction_conflict(request));
463 }
464 if journal.phase == JournalPhase::Aborted {
465 let generation_name = request.generation_uuid.hyphenated().to_string();
466 if root.join(GENERATIONS_DIR).join(&generation_name).exists()
467 || root.join("trash").join(&generation_name).exists()
468 {
469 return Err(publication_error(
470 request,
471 "ABORTED",
472 false,
473 "aborted transaction cleanup is incomplete; run recovery again",
474 ));
475 }
476 cleanup_aborted_attempts(root, request.transaction_uuid)?;
477 return Ok(None);
478 }
479 if journal.request_fingerprint != request_fingerprint
480 && journal.phase != JournalPhase::Published
481 {
482 return Err(transaction_conflict(request));
483 }
484 if journal.phase != JournalPhase::Published {
485 return Err(publication_error(
486 request,
487 "PREPARING",
488 false,
489 "an interrupted transaction requires recovery",
490 ));
491 }
492 let digest = journal
493 .generation_manifest_sha256
494 .as_deref()
495 .and_then(parse_digest)
496 .ok_or_else(|| {
497 project_error(
498 ProjectErrorCode::ProjectCorrupt,
499 "published transaction journal has no valid manifest digest",
500 )
501 })?;
502 let manifest_path = root
503 .join(GENERATIONS_DIR)
504 .join(request.generation_uuid.hyphenated().to_string())
505 .join(MANIFEST_FILE);
506 let manifest_bytes = std::fs::read(&manifest_path).map_err(publication_io)?;
507 let actual: [u8; 32] = Sha256::digest(&manifest_bytes).into();
508 if actual != digest {
509 return Err(project_error(
510 ProjectErrorCode::ProjectCorrupt,
511 "published transaction manifest does not match its journal",
512 ));
513 }
514 Ok(Some(ProjectStageOutcome::AlreadyPublished(
515 ProjectPublicationReceipt {
516 transaction_uuid: request.transaction_uuid,
517 generation_uuid: request.generation_uuid,
518 generation_manifest_sha256: digest,
519 idempotent_replay: true,
520 },
521 )))
522}
523
524fn cleanup_aborted_attempts(root: &Path, transaction_uuid: Uuid) -> Result<(), GfError> {
525 let attempts_root = root.join(ATTEMPTS_DIR);
526 let transaction_root = attempts_root.join(transaction_uuid.hyphenated().to_string());
527 let metadata = match std::fs::symlink_metadata(&transaction_root) {
528 Ok(metadata) => metadata,
529 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
530 Err(error) => return Err(publication_io(error)),
531 };
532 if !metadata.is_dir() || metadata.file_type().is_symlink() {
533 return Err(project_error(
534 ProjectErrorCode::ProjectCorrupt,
535 "aborted transaction attempt path is linked or not a directory",
536 ));
537 }
538 std::fs::remove_dir_all(&transaction_root).map_err(publication_io)?;
539 sync_directory(&attempts_root)
540}
541
542pub(crate) fn load_revert_journal_extension(
544 root: &Path,
545 transaction_uuid: Uuid,
546) -> Result<Option<RevertJournalExtension>, GfError> {
547 let path = root
548 .join(TRANSACTIONS_DIR)
549 .join(format!("{}.json", transaction_uuid.hyphenated()));
550 if !path.exists() {
551 return Ok(None);
552 }
553 Ok(read_journal(&path)?.revert)
554}
555
556pub(crate) fn load_published_revert(
558 root: &Path,
559 transaction_uuid: Uuid,
560) -> Result<Option<(RevertJournalExtension, ProjectPublicationReceipt)>, GfError> {
561 let path = root
562 .join(TRANSACTIONS_DIR)
563 .join(format!("{}.json", transaction_uuid.hyphenated()));
564 if !path.exists() {
565 return Ok(None);
566 }
567 let journal = read_journal(&path)?;
568 let Some(revert) = journal.revert else {
569 return Ok(None);
570 };
571 if journal.phase != JournalPhase::Published {
572 return Ok(None);
573 }
574 let generation_uuid = Uuid::parse_str(&journal.generation_uuid).map_err(|_| {
575 project_error(
576 ProjectErrorCode::ProjectCorrupt,
577 "published journal has an invalid generation UUID",
578 )
579 })?;
580 let digest = journal
581 .generation_manifest_sha256
582 .as_deref()
583 .and_then(parse_digest)
584 .ok_or_else(|| {
585 project_error(
586 ProjectErrorCode::ProjectCorrupt,
587 "published journal has an invalid manifest digest",
588 )
589 })?;
590 Ok(Some((
591 revert,
592 ProjectPublicationReceipt {
593 transaction_uuid,
594 generation_uuid,
595 generation_manifest_sha256: digest,
596 idempotent_replay: true,
597 },
598 )))
599}
600
601pub fn published_project_transaction(
606 root: &Path,
607 transaction_uuid: Uuid,
608) -> Result<Option<ProjectPublicationReceipt>, GfError> {
609 let path = root
610 .join(TRANSACTIONS_DIR)
611 .join(format!("{}.json", transaction_uuid.hyphenated()));
612 if !path.exists() {
613 return Ok(None);
614 }
615 let journal = read_journal(&path)?;
616 if journal.phase != JournalPhase::Published {
617 return Ok(None);
618 }
619 let generation_uuid = Uuid::parse_str(&journal.generation_uuid).map_err(|_| {
620 project_error(
621 ProjectErrorCode::ProjectCorrupt,
622 "published journal has an invalid generation UUID",
623 )
624 })?;
625 let digest = journal
626 .generation_manifest_sha256
627 .as_deref()
628 .and_then(parse_digest)
629 .ok_or_else(|| {
630 project_error(
631 ProjectErrorCode::ProjectCorrupt,
632 "published journal has an invalid manifest digest",
633 )
634 })?;
635 let manifest_path = root
636 .join(GENERATIONS_DIR)
637 .join(generation_uuid.hyphenated().to_string())
638 .join(MANIFEST_FILE);
639 let actual: [u8; 32] =
640 Sha256::digest(std::fs::read(manifest_path).map_err(publication_io)?).into();
641 if actual != digest {
642 return Err(project_error(
643 ProjectErrorCode::ProjectCorrupt,
644 "published transaction manifest does not match its journal",
645 ));
646 }
647 Ok(Some(ProjectPublicationReceipt {
648 transaction_uuid,
649 generation_uuid,
650 generation_manifest_sha256: digest,
651 idempotent_replay: true,
652 }))
653}
654
655fn prepare_generation_directory(
656 root: &Path,
657 request: &ProjectGenerationRequest,
658 request_fingerprint: &str,
659 requires_promotion: bool,
660) -> Result<PathBuf, GfError> {
661 let relative_root = if requires_promotion {
662 Path::new(ATTEMPTS_DIR)
663 .join(request.transaction_uuid.hyphenated().to_string())
664 .join(request_fingerprint)
665 } else {
666 Path::new(GENERATIONS_DIR).join(request.generation_uuid.hyphenated().to_string())
667 };
668 let generation_root = root.join(&relative_root);
669 if generation_root.exists() {
670 return Err(transaction_conflict(request));
671 }
672 ensure_machine_directory(root, &relative_root.join(PARTICIPANTS_DIR))?;
673 Ok(generation_root)
674}
675
676fn stage_participant_files(
677 request: &ProjectGenerationRequest,
678 generation_root: &Path,
679 participants: &[StagedParticipant],
680) -> Result<(), GfError> {
681 for metadata in participants {
682 let input = request
683 .participants
684 .iter()
685 .find(|candidate| {
686 candidate.capability_id == metadata.capability_id
687 && candidate.record_family_id == metadata.record_family_id
688 })
689 .expect("validated canonical metadata has one source participant");
690 let destination = generation_root
691 .join(PARTICIPANTS_DIR)
692 .join(&metadata.relative_path);
693 let parent_dir = destination
694 .parent()
695 .expect("machine-derived participant path has a parent");
696 let relative_parent = parent_dir
697 .strip_prefix(generation_root)
698 .expect("machine-derived participant parent is contained");
699 ensure_machine_directory(generation_root, relative_parent)?;
700 let mut file = OpenOptions::new()
701 .write(true)
702 .create_new(true)
703 .open(&destination)
704 .map_err(publication_io)?;
705 file.write_all(&input.bytes).map_err(publication_io)?;
706 project_failpoint::hit(
707 "project.after_participant_write",
708 Some(request.transaction_uuid),
709 Some(request.generation_uuid),
710 "STAGED",
711 false,
712 )?;
713 file.sync_all().map_err(publication_io)?;
714 project_failpoint::hit(
715 "project.after_participant_fsync",
716 Some(request.transaction_uuid),
717 Some(request.generation_uuid),
718 "STAGED",
719 false,
720 )?;
721 verify_participant_file(&destination, metadata)?;
722 }
723 Ok(())
724}
725
726impl StagedProjectGeneration {
727 #[must_use]
729 pub fn participants(&self) -> &[StagedParticipant] {
730 &self.participants
731 }
732
733 pub fn validate<D, C>(
738 self,
739 domain_validation: D,
740 composite_validation: C,
741 ) -> Result<ValidatedProjectGeneration, GfError>
742 where
743 D: FnOnce(&[StagedParticipant]) -> Result<(), GfError>,
744 C: FnOnce(&ResolvedProjectGeneration, &[StagedParticipant]) -> Result<(), GfError>,
745 {
746 for participant in &self.participants {
747 verify_participant_file(
748 &self
749 .generation_root
750 .join(PARTICIPANTS_DIR)
751 .join(&participant.relative_path),
752 participant,
753 )?;
754 }
755 domain_validation(&self.participants)?;
756 project_failpoint::hit(
757 "project.after_domain_validation",
758 Some(self.transaction_uuid),
759 Some(self.generation_uuid),
760 "VALIDATED",
761 false,
762 )?;
763 if let Err(error) = composite_validation(&self.parent, &self.participants) {
764 if self.requires_promotion && error.code() == "GF_WRITE_CONFLICT" {
771 abort_stale_generation(&self)?;
772 }
773 return Err(error);
774 }
775 project_failpoint::hit(
776 "project.after_composite_validation",
777 Some(self.transaction_uuid),
778 Some(self.generation_uuid),
779 "VALIDATED",
780 false,
781 )?;
782 write_journal(
783 &self.journal_path(),
784 &self.journal(JournalPhase::Validated, None),
785 )?;
786 project_failpoint::hit(
787 "project.after_journal_validated",
788 Some(self.transaction_uuid),
789 Some(self.generation_uuid),
790 "VALIDATED",
791 false,
792 )?;
793 Ok(ValidatedProjectGeneration(self))
794 }
795
796 fn journal_path(&self) -> PathBuf {
797 self.root
798 .join(TRANSACTIONS_DIR)
799 .join(format!("{}.json", self.transaction_uuid.hyphenated()))
800 }
801
802 fn journal(&self, phase: JournalPhase, manifest_sha256: Option<String>) -> JournalRecord {
803 JournalRecord {
804 format: "graphforge-transaction".into(),
805 format_version: 1,
806 transaction_uuid: self.transaction_uuid.hyphenated().to_string(),
807 generation_uuid: self.generation_uuid.hyphenated().to_string(),
808 parent_generation_uuid: Some(self.parent.generation_uuid().hyphenated().to_string()),
809 phase,
810 request_fingerprint: self.request_fingerprint.clone(),
811 operation_fingerprint: Some(self.operation_fingerprint.clone()),
812 participant_paths: self
813 .participants
814 .iter()
815 .map(|participant| participant.relative_path.clone())
816 .collect(),
817 generation_manifest_sha256: manifest_sha256,
818 revert: self.revert.clone(),
819 }
820 }
821}
822
823impl ValidatedProjectGeneration {
824 pub fn publish(self) -> Result<ProjectPublicationReceipt, GfError> {
830 let commit_lock = self.prepare_commit_lock()?;
831 let result = self.publish_inner().map_err(|error| {
832 if matches!(error, GfError::Project { .. }) {
833 error
834 } else {
835 publication_error_from_parts(
836 self.0.transaction_uuid,
837 self.0.generation_uuid,
838 "DURABLE",
839 false,
840 &error.to_string(),
841 )
842 }
843 });
844 drop(commit_lock);
845 result
846 }
847
848 fn prepare_commit_lock(&self) -> Result<Option<File>, GfError> {
849 if matches!(self.0.publication_lock, PublicationLock::Exclusive(_)) {
850 return Ok(None);
851 }
852 let staged = &self.0;
853 let writer_lock = wait_for_writer_lock(&staged.root)?;
854 project_failpoint::hit(
855 "project.after_optimistic_commit_lock",
856 Some(staged.transaction_uuid),
857 Some(staged.generation_uuid),
858 "COMMIT_LOCK",
859 false,
860 )?;
861 let current = resolve_project_generation(&staged.root)?;
862 if current.generation_uuid() != staged.parent.generation_uuid() {
863 abort_stale_generation(staged)?;
864 return Err(project_error(
865 ProjectErrorCode::WriteConflict,
866 format!(
867 "transaction_uuid={} generation_uuid={} phase=COMMIT_LOCK committed=false cause=stale_parent expected_parent={} actual_parent={}",
868 staged.transaction_uuid.hyphenated(),
869 staged.generation_uuid.hyphenated(),
870 staged.parent.generation_uuid().hyphenated(),
871 current.generation_uuid().hyphenated()
872 ),
873 ));
874 }
875 Ok(Some(writer_lock))
876 }
877
878 fn publish_inner(&self) -> Result<ProjectPublicationReceipt, GfError> {
879 let staged = &self.0;
880 let manifest_sha256 = make_generation_durable(staged)?;
881 replace_current(staged, manifest_sha256)?;
882 finish_published_generation(staged, manifest_sha256)?;
883 Ok(ProjectPublicationReceipt {
884 transaction_uuid: staged.transaction_uuid,
885 generation_uuid: staged.generation_uuid,
886 generation_manifest_sha256: manifest_sha256,
887 idempotent_replay: false,
888 })
889 }
890}
891
892fn abort_stale_generation(staged: &StagedProjectGeneration) -> Result<(), GfError> {
893 write_journal(
894 &staged.journal_path(),
895 &staged.journal(JournalPhase::Aborted, None),
896 )?;
897 if staged.generation_root.exists() {
898 std::fs::remove_dir_all(&staged.generation_root).map_err(publication_io)?;
899 sync_directory(
900 staged
901 .generation_root
902 .parent()
903 .expect("machine attempt path has a parent"),
904 )?;
905 }
906 Ok(())
907}
908
909fn make_generation_durable(staged: &StagedProjectGeneration) -> Result<[u8; 32], GfError> {
910 let lease_path = staged.generation_root.join(LEASE_FILE);
911 let lease = OpenOptions::new()
912 .write(true)
913 .create_new(true)
914 .open(&lease_path)
915 .map_err(publication_io)?;
916 lease.sync_all().map_err(publication_io)?;
917 drop(lease);
921
922 let manifest = GenerationManifestRecord {
923 format: "graphforge-generation".into(),
924 format_version: 1,
925 generation_uuid: staged.generation_uuid.hyphenated().to_string(),
926 parent_generation_uuid: Some(staged.parent.generation_uuid().hyphenated().to_string()),
927 transaction_uuid: staged.transaction_uuid.hyphenated().to_string(),
928 capabilities: staged
929 .capabilities
930 .iter()
931 .map(|capability| CapabilityRecord {
932 capability_id: capability.capability_id.clone(),
933 capability_version: capability.capability_version,
934 })
935 .collect(),
936 participants: staged.participants.clone(),
937 };
938 let manifest_bytes = canonical_line(&manifest)?;
939 let manifest_path = staged.generation_root.join(MANIFEST_FILE);
940 let manifest_file = write_new(&manifest_path, &manifest_bytes)?;
941 project_failpoint::hit(
942 "project.after_manifest_write",
943 Some(staged.transaction_uuid),
944 Some(staged.generation_uuid),
945 "DURABLE",
946 false,
947 )?;
948 manifest_file.sync_all().map_err(publication_io)?;
949 drop(manifest_file);
952 project_failpoint::hit(
953 "project.after_manifest_fsync",
954 Some(staged.transaction_uuid),
955 Some(staged.generation_uuid),
956 "DURABLE",
957 false,
958 )?;
959 let manifest_sha256: [u8; 32] = Sha256::digest(&manifest_bytes).into();
960 for participant in &staged.participants {
961 verify_participant_file(
962 &staged
963 .generation_root
964 .join(PARTICIPANTS_DIR)
965 .join(&participant.relative_path),
966 participant,
967 )?;
968 }
969 verify_exact_file(&manifest_path, &manifest_bytes)?;
970 sync_participant_directories(
971 &staged.generation_root.join(PARTICIPANTS_DIR),
972 &staged.participants,
973 )?;
974 sync_directory(&staged.generation_root)?;
975 project_failpoint::hit(
976 "project.after_generation_dir_fsync",
977 Some(staged.transaction_uuid),
978 Some(staged.generation_uuid),
979 "DURABLE",
980 false,
981 )?;
982 if staged.requires_promotion {
983 promote_optimistic_generation(staged)?;
984 } else {
985 sync_directory(&staged.root.join(GENERATIONS_DIR))?;
986 }
987 write_journal(
988 &staged.journal_path(),
989 &staged.journal(JournalPhase::Durable, Some(hex_digest(manifest_sha256))),
990 )?;
991 project_failpoint::hit(
992 "project.after_journal_durable",
993 Some(staged.transaction_uuid),
994 Some(staged.generation_uuid),
995 "DURABLE",
996 false,
997 )?;
998 Ok(manifest_sha256)
999}
1000
1001fn promote_optimistic_generation(staged: &StagedProjectGeneration) -> Result<(), GfError> {
1002 let generations_root = staged.root.join(GENERATIONS_DIR);
1003 let destination = generations_root.join(staged.generation_uuid.hyphenated().to_string());
1004 if destination.exists() {
1005 return Err(project_error(
1006 ProjectErrorCode::TransactionConflict,
1007 format!(
1008 "transaction_uuid={} generation_uuid={} phase=PROMOTE committed=false cause=generation_exists",
1009 staged.transaction_uuid.hyphenated(),
1010 staged.generation_uuid.hyphenated()
1011 ),
1012 ));
1013 }
1014 std::fs::rename(&staged.generation_root, &destination).map_err(publication_io)?;
1015 let transaction_attempt_root = staged
1016 .generation_root
1017 .parent()
1018 .expect("machine attempt path has a parent");
1019 sync_directory(transaction_attempt_root)?;
1020 std::fs::remove_dir(transaction_attempt_root).map_err(publication_io)?;
1021 sync_directory(&staged.root.join(ATTEMPTS_DIR))?;
1022 sync_directory(&generations_root)?;
1023 project_failpoint::hit(
1024 "project.after_optimistic_promotion",
1025 Some(staged.transaction_uuid),
1026 Some(staged.generation_uuid),
1027 "DURABLE",
1028 false,
1029 )
1030}
1031
1032fn replace_current(
1033 staged: &StagedProjectGeneration,
1034 manifest_sha256: [u8; 32],
1035) -> Result<(), GfError> {
1036 let current = CurrentRecord {
1037 format: "graphforge-project".into(),
1038 format_version: 1,
1039 generation_uuid: staged.generation_uuid.hyphenated().to_string(),
1040 generation_manifest_sha256: hex_digest(manifest_sha256),
1041 };
1042 let current_bytes = canonical_line(¤t)?;
1043 let current_path = staged.root.join(CURRENT_FILE);
1044 AtomicFile::new(¤t_path, AllowOverwrite)
1045 .write(|file| {
1046 file.write_all(¤t_bytes)?;
1047 failpoint_as_io(
1048 "project.after_current_temp_write",
1049 staged.transaction_uuid,
1050 staged.generation_uuid,
1051 "CURRENT",
1052 false,
1053 )?;
1054 file.sync_all()?;
1055 failpoint_as_io(
1056 "project.after_current_temp_fsync",
1057 staged.transaction_uuid,
1058 staged.generation_uuid,
1059 "CURRENT",
1060 false,
1061 )?;
1062 failpoint_as_io(
1063 "project.before_current_replace",
1064 staged.transaction_uuid,
1065 staged.generation_uuid,
1066 "CURRENT",
1067 false,
1068 )
1069 })
1070 .map_err(|error| {
1071 publication_error_from_parts(
1072 staged.transaction_uuid,
1073 staged.generation_uuid,
1074 "CURRENT",
1075 false,
1076 &error.to_string(),
1077 )
1078 })?;
1079 project_failpoint::hit(
1080 "project.after_current_replace",
1081 Some(staged.transaction_uuid),
1082 Some(staged.generation_uuid),
1083 "CURRENT",
1084 true,
1085 )
1086}
1087
1088fn finish_published_generation(
1089 staged: &StagedProjectGeneration,
1090 manifest_sha256: [u8; 32],
1091) -> Result<(), GfError> {
1092 sync_directory(&staged.root).map_err(|error| {
1095 publication_error_from_parts(
1096 staged.transaction_uuid,
1097 staged.generation_uuid,
1098 "CURRENT",
1099 true,
1100 &error.to_string(),
1101 )
1102 })?;
1103 project_failpoint::hit(
1104 "project.after_root_fsync",
1105 Some(staged.transaction_uuid),
1106 Some(staged.generation_uuid),
1107 "PUBLISHED",
1108 true,
1109 )?;
1110 write_journal(
1111 &staged.journal_path(),
1112 &staged.journal(JournalPhase::Published, Some(hex_digest(manifest_sha256))),
1113 )
1114 .map_err(|error| {
1115 publication_error_from_parts(
1116 staged.transaction_uuid,
1117 staged.generation_uuid,
1118 "PUBLISHED",
1119 true,
1120 &error.to_string(),
1121 )
1122 })?;
1123 project_failpoint::hit(
1124 "project.after_journal_published",
1125 Some(staged.transaction_uuid),
1126 Some(staged.generation_uuid),
1127 "PUBLISHED",
1128 true,
1129 )?;
1130
1131 let resolved = resolve_project_generation(&staged.root).map_err(|error| {
1132 publication_error_from_parts(
1133 staged.transaction_uuid,
1134 staged.generation_uuid,
1135 "PUBLISHED",
1136 true,
1137 &error.to_string(),
1138 )
1139 })?;
1140 if resolved.generation_uuid() != staged.generation_uuid
1141 || resolved.manifest_sha256() != manifest_sha256
1142 {
1143 return Err(publication_error_from_parts(
1144 staged.transaction_uuid,
1145 staged.generation_uuid,
1146 "PUBLISHED",
1147 true,
1148 "published CURRENT did not resolve to exact generation bytes",
1149 ));
1150 }
1151 Ok(())
1152}
1153
1154impl Drop for StagedProjectGeneration {
1155 fn drop(&mut self) {
1156 match &self.publication_lock {
1157 PublicationLock::Exclusive(lock) | PublicationLock::Optimistic(lock) => {
1158 let _ = FileExt::unlock(lock);
1159 }
1160 }
1161 }
1162}
1163
1164#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1165#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
1166pub(crate) enum JournalPhase {
1167 Preparing,
1168 Staged,
1169 Validated,
1170 Durable,
1171 Published,
1172 Aborted,
1173}
1174
1175#[derive(Debug, Serialize, Deserialize)]
1176#[serde(deny_unknown_fields)]
1177pub(crate) struct JournalRecord {
1178 pub(crate) format: String,
1179 pub(crate) format_version: u32,
1180 pub(crate) transaction_uuid: String,
1181 pub(crate) generation_uuid: String,
1182 pub(crate) parent_generation_uuid: Option<String>,
1183 pub(crate) phase: JournalPhase,
1184 pub(crate) request_fingerprint: String,
1185 #[serde(default, skip_serializing_if = "Option::is_none")]
1186 pub(crate) operation_fingerprint: Option<String>,
1187 pub(crate) participant_paths: Vec<String>,
1188 pub(crate) generation_manifest_sha256: Option<String>,
1189 #[serde(default, skip_serializing_if = "Option::is_none")]
1190 pub(crate) revert: Option<RevertJournalExtension>,
1191}
1192
1193impl JournalRecord {
1194 fn new(
1195 request: &ProjectGenerationRequest,
1196 parent: Option<Uuid>,
1197 phase: JournalPhase,
1198 fingerprints: (String, String),
1199 participants: &[StagedParticipant],
1200 generation_manifest_sha256: Option<String>,
1201 revert: Option<RevertJournalExtension>,
1202 ) -> Self {
1203 let (request_fingerprint, operation_fingerprint) = fingerprints;
1204 Self {
1205 format: "graphforge-transaction".into(),
1206 format_version: 1,
1207 transaction_uuid: request.transaction_uuid.hyphenated().to_string(),
1208 generation_uuid: request.generation_uuid.hyphenated().to_string(),
1209 parent_generation_uuid: parent.map(|uuid| uuid.hyphenated().to_string()),
1210 phase,
1211 request_fingerprint,
1212 operation_fingerprint: Some(operation_fingerprint),
1213 participant_paths: participants
1214 .iter()
1215 .map(|participant| participant.relative_path.clone())
1216 .collect(),
1217 generation_manifest_sha256,
1218 revert,
1219 }
1220 }
1221
1222 pub(crate) fn operation_fingerprint(&self) -> &str {
1223 self.operation_fingerprint
1224 .as_deref()
1225 .unwrap_or(&self.request_fingerprint)
1226 }
1227}
1228
1229#[derive(Debug, Serialize)]
1230struct RequestFingerprint<'a> {
1231 format: &'static str,
1232 format_version: u32,
1233 transaction_uuid: String,
1234 generation_uuid: String,
1235 capabilities: &'a [ProjectCapability],
1236 participants: &'a [StagedParticipant],
1237}
1238
1239#[derive(Debug, Serialize)]
1240struct CurrentRecord {
1241 format: String,
1242 format_version: u32,
1243 generation_uuid: String,
1244 generation_manifest_sha256: String,
1245}
1246
1247#[derive(Debug, Serialize)]
1248struct GenerationManifestRecord {
1249 format: String,
1250 format_version: u32,
1251 generation_uuid: String,
1252 parent_generation_uuid: Option<String>,
1253 transaction_uuid: String,
1254 capabilities: Vec<CapabilityRecord>,
1255 participants: Vec<StagedParticipant>,
1256}
1257
1258#[derive(Debug, Serialize)]
1259struct CapabilityRecord {
1260 capability_id: String,
1261 capability_version: u32,
1262}
1263
1264impl Serialize for StagedParticipant {
1265 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1266 where
1267 S: serde::Serializer,
1268 {
1269 #[derive(Serialize)]
1270 struct Ordered<'a> {
1271 capability_id: &'a str,
1272 capability_version: u32,
1273 record_family_id: &'a str,
1274 record_version: u32,
1275 relative_path: &'a str,
1276 encoding: &'a str,
1277 byte_length: u64,
1278 row_count: u64,
1279 schema_fingerprint: &'a str,
1280 content_sha256: &'a str,
1281 }
1282 Ordered {
1283 capability_id: &self.capability_id,
1284 capability_version: self.capability_version,
1285 record_family_id: &self.record_family_id,
1286 record_version: self.record_version,
1287 relative_path: &self.relative_path,
1288 encoding: &self.encoding,
1289 byte_length: self.byte_length,
1290 row_count: self.row_count,
1291 schema_fingerprint: &self.schema_fingerprint,
1292 content_sha256: &self.content_sha256,
1293 }
1294 .serialize(serializer)
1295 }
1296}
1297
1298fn canonical_supported_root(root: &Path) -> Result<PathBuf, GfError> {
1299 resolve_project_generation(root).map(|resolved| resolved.container_root().to_owned())
1301}
1302
1303fn validate_request(request: &ProjectGenerationRequest) -> Result<(), GfError> {
1304 if request.capabilities.is_empty() {
1305 return Err(project_error(
1306 ProjectErrorCode::PublicationFailed,
1307 "a generation must declare at least one capability",
1308 ));
1309 }
1310 for capability in &request.capabilities {
1311 validate_machine_id(&capability.capability_id)?;
1312 if capability.capability_version == 0 {
1313 return Err(project_error(
1314 ProjectErrorCode::PublicationFailed,
1315 "capability contract versions must be positive",
1316 ));
1317 }
1318 }
1319 for participant in &request.participants {
1320 validate_machine_id(&participant.capability_id)?;
1321 validate_machine_id(&participant.record_family_id)?;
1322 if participant.capability_version == 0 || participant.record_version == 0 {
1323 return Err(project_error(
1324 ProjectErrorCode::PublicationFailed,
1325 "participant contract versions must be positive",
1326 ));
1327 }
1328 }
1329 Ok(())
1330}
1331
1332fn request_metadata(
1333 request: &ProjectGenerationRequest,
1334) -> Result<(Vec<ProjectCapability>, Vec<StagedParticipant>, String), GfError> {
1335 let mut capabilities = request.capabilities.clone();
1336 capabilities.sort_by(|left, right| left.capability_id.cmp(&right.capability_id));
1337 if capabilities
1338 .windows(2)
1339 .any(|pair| pair[0].capability_id == pair[1].capability_id)
1340 {
1341 return Err(project_error(
1342 ProjectErrorCode::PublicationFailed,
1343 "duplicate capability identity",
1344 ));
1345 }
1346 if capabilities
1347 .binary_search_by(|entry| entry.capability_id.as_str().cmp("graph"))
1348 .ok()
1349 .map(|index| capabilities[index].capability_version)
1350 != Some(1)
1351 {
1352 return Err(project_error(
1353 ProjectErrorCode::PublicationFailed,
1354 "every generation must declare graph capability version 1",
1355 ));
1356 }
1357 let mut participants = Vec::with_capacity(request.participants.len());
1358 for participant in &request.participants {
1359 let content_sha256: [u8; 32] = Sha256::digest(&participant.bytes).into();
1360 participants.push(StagedParticipant {
1361 capability_id: participant.capability_id.clone(),
1362 capability_version: participant.capability_version,
1363 record_family_id: participant.record_family_id.clone(),
1364 record_version: participant.record_version,
1365 relative_path: format!(
1366 "{}/{}.{}",
1367 participant.capability_id,
1368 participant.record_family_id,
1369 participant.encoding.extension()
1370 ),
1371 encoding: participant.encoding.extension().into(),
1372 byte_length: u64::try_from(participant.bytes.len()).map_err(|_| {
1373 project_error(
1374 ProjectErrorCode::PublicationFailed,
1375 "participant byte length exceeds u64",
1376 )
1377 })?,
1378 row_count: participant.row_count,
1379 schema_fingerprint: hex_digest(participant.schema_fingerprint),
1380 content_sha256: hex_digest(content_sha256),
1381 });
1382 }
1383 participants.sort_by(|left, right| {
1384 (
1385 &left.capability_id,
1386 &left.record_family_id,
1387 &left.relative_path,
1388 )
1389 .cmp(&(
1390 &right.capability_id,
1391 &right.record_family_id,
1392 &right.relative_path,
1393 ))
1394 });
1395 if participants.windows(2).any(|pair| {
1396 pair[0].capability_id == pair[1].capability_id
1397 && pair[0].record_family_id == pair[1].record_family_id
1398 }) {
1399 return Err(project_error(
1400 ProjectErrorCode::PublicationFailed,
1401 "duplicate participant identity",
1402 ));
1403 }
1404 for participant in &participants {
1405 let capability = capabilities
1406 .binary_search_by(|entry| entry.capability_id.cmp(&participant.capability_id))
1407 .ok()
1408 .map(|index| &capabilities[index])
1409 .ok_or_else(|| {
1410 project_error(
1411 ProjectErrorCode::PublicationFailed,
1412 "participant capability is not declared",
1413 )
1414 })?;
1415 if capability.capability_version != participant.capability_version {
1416 return Err(project_error(
1417 ProjectErrorCode::PublicationFailed,
1418 "participant capability version conflicts with declaration",
1419 ));
1420 }
1421 }
1422 let fingerprint_input = RequestFingerprint {
1423 format: "graphforge-publication-request",
1424 format_version: 1,
1425 transaction_uuid: request.transaction_uuid.hyphenated().to_string(),
1426 generation_uuid: request.generation_uuid.hyphenated().to_string(),
1427 capabilities: &capabilities,
1428 participants: &participants,
1429 };
1430 let bytes = canonical_line(&fingerprint_input)?;
1431 let digest: [u8; 32] = Sha256::digest(bytes).into();
1432 Ok((capabilities, participants, hex_digest(digest)))
1433}
1434
1435fn validate_machine_id(value: &str) -> Result<(), GfError> {
1436 if value.is_empty()
1437 || value.len() > 64
1438 || !value.bytes().all(|byte| {
1439 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'-' | b'_')
1440 })
1441 {
1442 return Err(project_error(
1443 ProjectErrorCode::PublicationFailed,
1444 "machine ID must be 1-64 lowercase ASCII letters, digits, hyphens, or underscores",
1445 ));
1446 }
1447 Ok(())
1448}
1449
1450fn verify_participant_file(path: &Path, expected: &StagedParticipant) -> Result<(), GfError> {
1451 let metadata = std::fs::symlink_metadata(path).map_err(publication_io)?;
1452 if !metadata.is_file() || metadata.file_type().is_symlink() {
1453 return Err(project_error(
1454 ProjectErrorCode::PublicationFailed,
1455 "staged participant is not a regular non-link file",
1456 ));
1457 }
1458 #[cfg(unix)]
1459 {
1460 use std::os::unix::fs::MetadataExt;
1461 if metadata.nlink() != 1 {
1462 return Err(project_error(
1463 ProjectErrorCode::PublicationFailed,
1464 "staged participant is hard-linked",
1465 ));
1466 }
1467 }
1468 if metadata.len() != expected.byte_length {
1469 return Err(project_error(
1470 ProjectErrorCode::PublicationFailed,
1471 "staged participant byte length changed",
1472 ));
1473 }
1474 let mut file = File::open(path).map_err(publication_io)?;
1475 let mut hasher = Sha256::new();
1476 std::io::copy(&mut file, &mut hasher).map_err(publication_io)?;
1477 let actual: [u8; 32] = hasher.finalize().into();
1478 if hex_digest(actual) != expected.content_sha256 {
1479 return Err(project_error(
1480 ProjectErrorCode::PublicationFailed,
1481 "staged participant digest changed",
1482 ));
1483 }
1484 Ok(())
1485}
1486
1487fn sync_participant_directories(
1488 participants_root: &Path,
1489 participants: &[StagedParticipant],
1490) -> Result<(), GfError> {
1491 let mut directories: Vec<PathBuf> = participants
1492 .iter()
1493 .filter_map(|participant| {
1494 participants_root
1495 .join(&participant.relative_path)
1496 .parent()
1497 .map(Path::to_owned)
1498 })
1499 .collect();
1500 directories.sort();
1501 directories.dedup();
1502 directories.sort_by_key(|path| std::cmp::Reverse(path.components().count()));
1503 for directory in directories {
1504 sync_directory(&directory)?;
1505 }
1506 sync_directory(participants_root)
1507}
1508
1509fn write_new(path: &Path, bytes: &[u8]) -> Result<File, GfError> {
1510 let mut file = OpenOptions::new()
1511 .write(true)
1512 .create_new(true)
1513 .open(path)
1514 .map_err(publication_io)?;
1515 file.write_all(bytes).map_err(publication_io)?;
1516 Ok(file)
1517}
1518
1519fn failpoint_as_io(
1520 name: &str,
1521 transaction_uuid: Uuid,
1522 generation_uuid: Uuid,
1523 phase: &str,
1524 committed: bool,
1525) -> std::io::Result<()> {
1526 project_failpoint::hit(
1527 name,
1528 Some(transaction_uuid),
1529 Some(generation_uuid),
1530 phase,
1531 committed,
1532 )
1533 .map_err(|error| std::io::Error::other(error.to_string()))
1534}
1535
1536pub(crate) fn ensure_machine_directory(root: &Path, relative: &Path) -> Result<PathBuf, GfError> {
1537 let mut current = root.to_owned();
1538 for component in relative.components() {
1539 let std::path::Component::Normal(component) = component else {
1540 return Err(project_error(
1541 ProjectErrorCode::ProjectCorrupt,
1542 "machine directory path is not normalized",
1543 ));
1544 };
1545 current.push(component);
1546 match std::fs::symlink_metadata(¤t) {
1547 Ok(metadata) if metadata.is_dir() && !metadata.file_type().is_symlink() => {}
1548 Ok(_) => {
1549 return Err(project_error(
1550 ProjectErrorCode::ProjectCorrupt,
1551 "machine directory is linked or not a directory",
1552 ));
1553 }
1554 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1555 match std::fs::create_dir(¤t) {
1556 Ok(()) => {
1557 sync_directory(
1558 current
1559 .parent()
1560 .expect("machine directory beneath project has a parent"),
1561 )?;
1562 }
1563 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
1564 let metadata =
1565 std::fs::symlink_metadata(¤t).map_err(publication_io)?;
1566 if !metadata.is_dir() || metadata.file_type().is_symlink() {
1567 return Err(project_error(
1568 ProjectErrorCode::ProjectCorrupt,
1569 "concurrently created machine path is linked or not a directory",
1570 ));
1571 }
1572 }
1573 Err(error) => return Err(publication_io(error)),
1574 }
1575 }
1576 Err(error) => return Err(publication_io(error)),
1577 }
1578 }
1579 Ok(current)
1580}
1581
1582pub(crate) fn open_regular_lock(path: &Path) -> Result<File, GfError> {
1583 if let Ok(metadata) = std::fs::symlink_metadata(path)
1584 && (!metadata.is_file() || metadata.file_type().is_symlink())
1585 {
1586 return Err(project_error(
1587 ProjectErrorCode::ProjectCorrupt,
1588 "writer lock is linked or not a regular file",
1589 ));
1590 }
1591 let file = OpenOptions::new()
1592 .read(true)
1593 .write(true)
1594 .create(true)
1595 .truncate(false)
1596 .open(path)
1597 .map_err(publication_io)?;
1598 let metadata = file.metadata().map_err(publication_io)?;
1599 if !metadata.is_file() {
1600 return Err(project_error(
1601 ProjectErrorCode::ProjectCorrupt,
1602 "writer lock is not a regular file",
1603 ));
1604 }
1605 #[cfg(unix)]
1606 {
1607 use std::os::unix::fs::MetadataExt;
1608 if metadata.nlink() != 1 {
1609 return Err(project_error(
1610 ProjectErrorCode::ProjectCorrupt,
1611 "writer lock is hard-linked",
1612 ));
1613 }
1614 }
1615 Ok(file)
1616}
1617
1618fn verify_exact_file(path: &Path, expected: &[u8]) -> Result<(), GfError> {
1619 let mut actual = Vec::new();
1620 File::open(path)
1621 .and_then(|mut file| file.read_to_end(&mut actual))
1622 .map_err(publication_io)?;
1623 if actual != expected {
1624 return Err(project_error(
1625 ProjectErrorCode::PublicationFailed,
1626 "durable file reread did not match staged bytes",
1627 ));
1628 }
1629 Ok(())
1630}
1631
1632pub(crate) fn write_journal(path: &Path, journal: &JournalRecord) -> Result<(), GfError> {
1633 let bytes = canonical_line(journal)?;
1634 AtomicFile::new(path, AllowOverwrite)
1635 .write(|file| file.write_all(&bytes))
1636 .map_err(|error| publication_io(std::io::Error::other(error.to_string())))?;
1637 sync_directory(
1638 path.parent()
1639 .expect("transaction journal always has a parent"),
1640 )
1641}
1642
1643pub(crate) fn read_journal(path: &Path) -> Result<JournalRecord, GfError> {
1644 let metadata = std::fs::symlink_metadata(path).map_err(publication_io)?;
1645 if !metadata.is_file()
1646 || metadata.file_type().is_symlink()
1647 || metadata.len() > MAX_JOURNAL_BYTES
1648 {
1649 return Err(project_error(
1650 ProjectErrorCode::ProjectCorrupt,
1651 "transaction journal is invalid",
1652 ));
1653 }
1654 let bytes = std::fs::read(path).map_err(publication_io)?;
1655 let journal: JournalRecord = serde_json::from_slice(&bytes).map_err(|_| {
1656 project_error(
1657 ProjectErrorCode::ProjectCorrupt,
1658 "transaction journal is not canonical JSON",
1659 )
1660 })?;
1661 if canonical_line(&journal)? != bytes
1662 || journal.format != "graphforge-transaction"
1663 || journal.format_version != 1
1664 || parse_digest(&journal.request_fingerprint).is_none()
1665 || journal
1666 .operation_fingerprint
1667 .as_deref()
1668 .is_some_and(|fingerprint| parse_digest(fingerprint).is_none())
1669 {
1670 return Err(project_error(
1671 ProjectErrorCode::ProjectCorrupt,
1672 "transaction journal is not canonical",
1673 ));
1674 }
1675 Ok(journal)
1676}
1677
1678pub(crate) fn cleanup_atomicwrite_temp(path: &Path) -> Result<bool, GfError> {
1679 let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
1680 return Ok(false);
1681 };
1682 let Some(suffix) = name.strip_prefix(".atomicwrite") else {
1683 return Ok(false);
1684 };
1685 if suffix.len() != 6 || !suffix.bytes().all(|byte| byte.is_ascii_alphanumeric()) {
1686 return Ok(false);
1687 }
1688 let metadata = std::fs::symlink_metadata(path).map_err(publication_io)?;
1689 if !metadata.is_dir() || metadata.file_type().is_symlink() {
1690 return Ok(false);
1691 }
1692 let mut entries = std::fs::read_dir(path).map_err(publication_io)?;
1693 if let Some(entry) = entries.next().transpose().map_err(publication_io)? {
1694 if entries
1695 .next()
1696 .transpose()
1697 .map_err(publication_io)?
1698 .is_some()
1699 || entry.file_name() != "tmpfile.tmp"
1700 {
1701 return Ok(false);
1702 }
1703 let entry_metadata = std::fs::symlink_metadata(entry.path()).map_err(publication_io)?;
1704 if !entry_metadata.is_file() || entry_metadata.file_type().is_symlink() {
1705 return Ok(false);
1706 }
1707 #[cfg(unix)]
1708 {
1709 use std::os::unix::fs::MetadataExt;
1710 if entry_metadata.nlink() != 1 {
1711 return Ok(false);
1712 }
1713 }
1714 std::fs::remove_file(entry.path()).map_err(publication_io)?;
1715 }
1716 std::fs::remove_dir(path).map_err(publication_io)?;
1717 sync_directory(
1718 path.parent()
1719 .expect("atomic-write temporary directory always has a parent"),
1720 )?;
1721 Ok(true)
1722}
1723
1724fn canonical_line<T: Serialize>(value: &T) -> Result<Vec<u8>, GfError> {
1725 let mut bytes = serde_json::to_vec(value)
1726 .map_err(|error| GfError::Storage(format!("failed to encode project record: {error}")))?;
1727 bytes.push(b'\n');
1728 Ok(bytes)
1729}
1730
1731#[cfg(unix)]
1732pub(crate) fn sync_directory(path: &Path) -> Result<(), GfError> {
1733 File::open(path)
1734 .and_then(|directory| directory.sync_all())
1735 .map_err(publication_io)
1736}
1737
1738#[cfg(windows)]
1739pub(crate) fn sync_directory(path: &Path) -> Result<(), GfError> {
1740 use std::os::windows::fs::OpenOptionsExt;
1741
1742 const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
1745 OpenOptions::new()
1746 .write(true)
1749 .custom_flags(FILE_FLAG_BACKUP_SEMANTICS)
1750 .open(path)
1751 .and_then(|directory| directory.sync_all())
1752 .map_err(publication_io)
1753}
1754
1755#[cfg(all(not(unix), not(windows)))]
1756pub(crate) fn sync_directory(_path: &Path) -> Result<(), GfError> {
1757 Err(project_error(
1758 ProjectErrorCode::UnsupportedFilesystem,
1759 "directory durability is unsupported on this platform",
1760 ))
1761}
1762
1763fn hex_digest(bytes: [u8; 32]) -> String {
1764 let mut output = String::with_capacity(64);
1765 for byte in bytes {
1766 use std::fmt::Write as _;
1767 write!(&mut output, "{byte:02x}").expect("writing to String cannot fail");
1768 }
1769 output
1770}
1771
1772fn parse_digest(value: &str) -> Option<[u8; 32]> {
1773 if value.len() != 64
1774 || !value
1775 .bytes()
1776 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1777 {
1778 return None;
1779 }
1780 let mut digest = [0u8; 32];
1781 for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
1782 let high = hex_nibble(pair[0])?;
1783 let low = hex_nibble(pair[1])?;
1784 digest[index] = (high << 4) | low;
1785 }
1786 Some(digest)
1787}
1788
1789const fn hex_nibble(byte: u8) -> Option<u8> {
1790 match byte {
1791 b'0'..=b'9' => Some(byte - b'0'),
1792 b'a'..=b'f' => Some(byte - b'a' + 10),
1793 _ => None,
1794 }
1795}
1796
1797fn project_error(code: ProjectErrorCode, message: impl Into<String>) -> GfError {
1798 GfError::Project {
1799 code,
1800 message: message.into(),
1801 }
1802}
1803
1804fn publication_io(error: impl std::fmt::Display) -> GfError {
1805 GfError::Storage(error.to_string())
1806}
1807
1808fn transaction_conflict(request: &ProjectGenerationRequest) -> GfError {
1809 project_error(
1810 ProjectErrorCode::TransactionConflict,
1811 format!(
1812 "transaction_uuid={} generation_uuid={} phase=PREPARING committed=false cause=identity_conflict",
1813 request.transaction_uuid.hyphenated(),
1814 request.generation_uuid.hyphenated()
1815 ),
1816 )
1817}
1818
1819fn publication_error(
1820 request: &ProjectGenerationRequest,
1821 phase: &str,
1822 committed: bool,
1823 cause: &str,
1824) -> GfError {
1825 publication_error_from_parts(
1826 request.transaction_uuid,
1827 request.generation_uuid,
1828 phase,
1829 committed,
1830 cause,
1831 )
1832}
1833
1834fn publication_error_from_parts(
1835 transaction_uuid: Uuid,
1836 generation_uuid: Uuid,
1837 phase: &str,
1838 committed: bool,
1839 cause: &str,
1840) -> GfError {
1841 project_error(
1842 ProjectErrorCode::PublicationFailed,
1843 format!(
1844 "transaction_uuid={} generation_uuid={} phase={phase} committed={committed} cause={}",
1845 transaction_uuid.hyphenated(),
1846 generation_uuid.hyphenated(),
1847 safe_cause(cause)
1848 ),
1849 )
1850}
1851
1852fn safe_cause(cause: &str) -> String {
1853 cause
1854 .chars()
1855 .filter(|character| character.is_ascii_alphanumeric() || "_ -".contains(*character))
1856 .take(96)
1857 .collect()
1858}
1859
1860#[cfg(test)]
1861mod tests {
1862 use std::fs;
1863
1864 use super::*;
1865 use crate::open_or_initialize_project;
1866
1867 #[cfg(windows)]
1868 #[test]
1869 fn windows_directory_sync_uses_write_capable_handle() {
1870 let root = tempfile::tempdir().unwrap();
1871
1872 sync_directory(root.path()).unwrap();
1873 }
1874
1875 fn participant(capability: &str, family: &str, value: &[u8]) -> ProjectParticipant {
1876 ProjectParticipant {
1877 capability_id: capability.into(),
1878 capability_version: 1,
1879 record_family_id: family.into(),
1880 record_version: 1,
1881 encoding: ProjectParticipantEncoding::Parquet,
1882 schema_fingerprint: Sha256::digest(format!("{capability}/{family}")).into(),
1883 row_count: 1,
1884 bytes: value.to_vec(),
1885 }
1886 }
1887
1888 fn request(participants: Vec<ProjectParticipant>) -> ProjectGenerationRequest {
1889 let mut capabilities = vec![ProjectCapability {
1890 capability_id: "graph".into(),
1891 capability_version: 1,
1892 }];
1893 for participant in &participants {
1894 if participant.capability_id != "graph"
1895 && !capabilities
1896 .iter()
1897 .any(|entry| entry.capability_id == participant.capability_id)
1898 {
1899 capabilities.push(ProjectCapability {
1900 capability_id: participant.capability_id.clone(),
1901 capability_version: participant.capability_version,
1902 });
1903 }
1904 }
1905 ProjectGenerationRequest {
1906 transaction_uuid: Uuid::now_v7(),
1907 generation_uuid: Uuid::now_v7(),
1908 capabilities,
1909 participants,
1910 }
1911 }
1912
1913 fn project() -> tempfile::TempDir {
1914 let root = tempfile::tempdir().unwrap();
1915 open_or_initialize_project(root.path()).unwrap();
1916 root
1917 }
1918
1919 fn publish(root: &Path, request: ProjectGenerationRequest) -> ProjectPublicationReceipt {
1920 let ProjectStageOutcome::Staged(staged) = stage_project_generation(root, &request).unwrap()
1921 else {
1922 panic!("new request unexpectedly replayed");
1923 };
1924 staged
1925 .validate(|_| Ok(()), |_, _| Ok(()))
1926 .unwrap()
1927 .publish()
1928 .unwrap()
1929 }
1930
1931 fn journal_path(root: &Path, transaction_uuid: Uuid) -> PathBuf {
1932 root.join(TRANSACTIONS_DIR)
1933 .join(format!("{}.json", transaction_uuid.hyphenated()))
1934 }
1935
1936 #[test]
1937 fn publishes_graph_only_and_multi_domain_sets_atomically() {
1938 for participants in [
1939 vec![participant("graph", "nodes", b"graph")],
1940 vec![
1941 participant("graph", "nodes", b"graph"),
1942 participant("provenance", "events", b"provenance"),
1943 ],
1944 vec![
1945 participant("graph", "nodes", b"graph"),
1946 participant("provenance", "events", b"provenance"),
1947 participant("knowledge", "assertions", b"knowledge"),
1948 ],
1949 ] {
1950 let root = project();
1951 let request = request(participants);
1952 let expected = request.generation_uuid;
1953 publish(root.path(), request);
1954 let resolved = resolve_project_generation(root.path()).unwrap();
1955 assert_eq!(resolved.generation_uuid(), expected);
1956 }
1957 }
1958
1959 #[test]
1960 fn validation_failure_leaves_parent_authoritative() {
1961 let root = project();
1962 let parent = resolve_project_generation(root.path())
1963 .unwrap()
1964 .generation_uuid();
1965 let request = request(vec![participant("graph", "nodes", b"new")]);
1966 let ProjectStageOutcome::Staged(staged) =
1967 stage_project_generation(root.path(), &request).unwrap()
1968 else {
1969 panic!("new request unexpectedly replayed");
1970 };
1971
1972 let error = staged
1973 .validate(
1974 |_| Err(GfError::Validation("domain rejected".into())),
1975 |_, _| Ok(()),
1976 )
1977 .err()
1978 .expect("validation must fail");
1979
1980 assert!(matches!(error, GfError::Validation(_)));
1981 assert_eq!(
1982 resolve_project_generation(root.path())
1983 .unwrap()
1984 .generation_uuid(),
1985 parent
1986 );
1987 }
1988
1989 #[test]
1990 fn durable_install_io_failure_is_wrapped_before_current_changes() {
1991 let root = project();
1992 let parent = resolve_project_generation(root.path())
1993 .unwrap()
1994 .generation_uuid();
1995 let request = request(vec![participant("graph", "nodes", b"new")]);
1996 let ProjectStageOutcome::Staged(staged) =
1997 stage_project_generation(root.path(), &request).unwrap()
1998 else {
1999 panic!("new request unexpectedly replayed")
2000 };
2001 let validated = staged.validate(|_| Ok(()), |_, _| Ok(())).unwrap();
2002 fs::remove_file(
2003 validated
2004 .0
2005 .generation_root
2006 .join(PARTICIPANTS_DIR)
2007 .join(&validated.0.participants[0].relative_path),
2008 )
2009 .unwrap();
2010 let error = validated.publish().unwrap_err();
2011 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
2012 assert!(error.to_string().contains("phase=DURABLE committed=false"));
2013 assert_eq!(
2014 resolve_project_generation(root.path())
2015 .unwrap()
2016 .generation_uuid(),
2017 parent
2018 );
2019 }
2020
2021 #[test]
2022 fn journal_records_each_deterministic_publication_phase() {
2023 let root = project();
2024 let request = request(vec![participant("graph", "nodes", b"new")]);
2025 let journal_path = journal_path(root.path(), request.transaction_uuid);
2026 let ProjectStageOutcome::Staged(staged) =
2027 stage_project_generation(root.path(), &request).unwrap()
2028 else {
2029 panic!("new request unexpectedly replayed");
2030 };
2031 assert_eq!(
2032 read_journal(&journal_path).unwrap().phase,
2033 JournalPhase::Staged
2034 );
2035
2036 let validated = staged.validate(|_| Ok(()), |_, _| Ok(())).unwrap();
2037 assert_eq!(
2038 read_journal(&journal_path).unwrap().phase,
2039 JournalPhase::Validated
2040 );
2041
2042 validated.publish().unwrap();
2043 let published = read_journal(&journal_path).unwrap();
2044 assert_eq!(published.phase, JournalPhase::Published);
2045 assert!(published.generation_manifest_sha256.is_some());
2046 }
2047
2048 #[test]
2049 fn request_fingerprint_is_independent_of_participant_input_order() {
2050 let mut request = request(vec![
2051 participant("provenance", "events", b"provenance"),
2052 participant("graph", "nodes", b"graph"),
2053 ]);
2054 let (_, first_metadata, first_fingerprint) = request_metadata(&request).unwrap();
2055 request.participants.reverse();
2056 let (_, second_metadata, second_fingerprint) = request_metadata(&request).unwrap();
2057
2058 assert_eq!(first_metadata, second_metadata);
2059 assert_eq!(first_fingerprint, second_fingerprint);
2060 assert_eq!(first_metadata[0].capability_id, "graph");
2061 assert_eq!(first_metadata[1].capability_id, "provenance");
2062 }
2063
2064 #[test]
2065 fn machine_ids_match_the_committed_generation_reader_contract() {
2066 let root = project();
2067 let valid = request(vec![participant("graph", "node-properties", b"properties")]);
2068 let generation_uuid = valid.generation_uuid;
2069 publish(root.path(), valid);
2070 let resolved = resolve_project_generation(root.path()).unwrap();
2071 assert_eq!(resolved.generation_uuid(), generation_uuid);
2072 assert!(
2073 resolved
2074 .participant_path("graph", "node-properties")
2075 .unwrap()
2076 .is_file()
2077 );
2078
2079 let underscore = request(vec![participant(
2080 "graph_data",
2081 "node_properties",
2082 b"properties",
2083 )]);
2084 publish(root.path(), underscore);
2085 let resolved = resolve_project_generation(root.path()).unwrap();
2086 assert!(
2087 resolved
2088 .participant_path("graph_data", "node_properties")
2089 .unwrap()
2090 .is_file()
2091 );
2092
2093 let invalid = request(vec![participant("graph", "NodeProperties", b"properties")]);
2094 let error = stage_project_generation(root.path(), &invalid)
2095 .err()
2096 .expect("reader-incompatible machine ID must be rejected");
2097 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
2098 }
2099
2100 #[test]
2101 fn tampered_staged_bytes_fail_before_publication() {
2102 let root = project();
2103 let parent = resolve_project_generation(root.path())
2104 .unwrap()
2105 .generation_uuid();
2106 let initial_request = request(vec![participant("graph", "nodes", b"original")]);
2107 let ProjectStageOutcome::Staged(staged) =
2108 stage_project_generation(root.path(), &initial_request).unwrap()
2109 else {
2110 panic!("new request unexpectedly replayed");
2111 };
2112 std::fs::write(
2113 staged.generation_root.join(PARTICIPANTS_DIR).join(
2114 staged
2115 .participants
2116 .first()
2117 .expect("participant")
2118 .relative_path
2119 .as_str(),
2120 ),
2121 b"tampered",
2122 )
2123 .unwrap();
2124
2125 let error = staged
2126 .validate(|_| Ok(()), |_, _| Ok(()))
2127 .err()
2128 .expect("tampered bytes must fail validation");
2129
2130 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
2131 assert_eq!(
2132 resolve_project_generation(root.path())
2133 .unwrap()
2134 .generation_uuid(),
2135 parent
2136 );
2137
2138 let request = request(vec![participant("graph", "nodes", b"original")]);
2139 let ProjectStageOutcome::Staged(staged) =
2140 stage_project_generation(root.path(), &request).unwrap()
2141 else {
2142 panic!("new request unexpectedly replayed");
2143 };
2144 let path = staged
2145 .generation_root
2146 .join(PARTICIPANTS_DIR)
2147 .join(&staged.participants[0].relative_path);
2148 std::fs::write(path, b"short").unwrap();
2149 assert_eq!(
2150 staged
2151 .validate(|_| Ok(()), |_, _| Ok(()))
2152 .err()
2153 .expect("truncated staged bytes must fail")
2154 .code(),
2155 "GF_PUBLICATION_FAILED"
2156 );
2157 assert_eq!(
2158 resolve_project_generation(root.path())
2159 .unwrap()
2160 .generation_uuid(),
2161 parent
2162 );
2163 }
2164
2165 #[cfg(unix)]
2166 #[test]
2167 fn staged_participant_hard_link_fails_before_current_mutation() {
2168 let root = project();
2169 let parent = resolve_project_generation(root.path())
2170 .unwrap()
2171 .generation_uuid();
2172 let request = request(vec![participant("graph", "nodes", b"stable")]);
2173 let ProjectStageOutcome::Staged(staged) =
2174 stage_project_generation(root.path(), &request).unwrap()
2175 else {
2176 panic!("unexpected replay")
2177 };
2178 let path = staged
2179 .generation_root
2180 .join(PARTICIPANTS_DIR)
2181 .join(&staged.participants[0].relative_path);
2182 let external = root.path().join("external-participant");
2183 fs::rename(&path, &external).unwrap();
2184 fs::hard_link(&external, &path).unwrap();
2185
2186 assert_eq!(
2187 staged
2188 .validate(|_| Ok(()), |_, _| Ok(()))
2189 .err()
2190 .expect("hard-linked staged bytes must fail")
2191 .code(),
2192 "GF_PUBLICATION_FAILED"
2193 );
2194 assert_eq!(
2195 resolve_project_generation(root.path())
2196 .unwrap()
2197 .generation_uuid(),
2198 parent
2199 );
2200 }
2201
2202 #[test]
2203 fn identical_published_transaction_is_idempotent() {
2204 let root = project();
2205 let request = request(vec![participant("graph", "nodes", b"same")]);
2206 publish(root.path(), request.clone());
2207
2208 let ProjectStageOutcome::AlreadyPublished(receipt) =
2209 stage_project_generation(root.path(), &request).unwrap()
2210 else {
2211 panic!("identical replay was not recognized");
2212 };
2213 assert!(receipt.idempotent_replay);
2214 }
2215
2216 #[test]
2217 fn historical_published_transaction_remains_idempotent() {
2218 let root = project();
2219 let first = request(vec![participant("graph", "nodes", b"first")]);
2220 publish(root.path(), first.clone());
2221 publish(
2222 root.path(),
2223 request(vec![participant("graph", "nodes", b"second")]),
2224 );
2225
2226 let ProjectStageOutcome::AlreadyPublished(receipt) =
2227 stage_project_generation(root.path(), &first).unwrap()
2228 else {
2229 panic!("historical identical replay was not recognized");
2230 };
2231 assert!(receipt.idempotent_replay);
2232 }
2233
2234 #[test]
2235 fn changed_content_under_same_transaction_conflicts() {
2236 let root = project();
2237 let request = request(vec![participant("graph", "nodes", b"first")]);
2238 publish(root.path(), request.clone());
2239 let mut conflicting = request;
2240 conflicting.participants[0].bytes = b"different".to_vec();
2241
2242 let error = stage_project_generation(root.path(), &conflicting)
2243 .err()
2244 .expect("conflicting replay must fail");
2245
2246 assert_eq!(error.code(), "GF_IDEMPOTENCY_CONFLICT");
2247 }
2248
2249 #[test]
2250 fn interrupted_stage_and_tampered_published_replay_are_exactly_classified() {
2251 let root = project();
2252 let staged_request = request(vec![participant("graph", "nodes", b"staged")]);
2253 let ProjectStageOutcome::Staged(staged) =
2254 stage_project_generation(root.path(), &staged_request).unwrap()
2255 else {
2256 panic!("fresh transaction unexpectedly replayed")
2257 };
2258 drop(staged);
2259 let interrupted = stage_project_generation(root.path(), &staged_request)
2260 .err()
2261 .expect("interrupted stage must fail");
2262 assert_eq!(interrupted.code(), "GF_PUBLICATION_FAILED");
2263 assert!(interrupted.to_string().contains("requires recovery"));
2264
2265 let published_request = request(vec![participant("graph", "nodes", b"published")]);
2266 let receipt = publish(root.path(), published_request.clone());
2267 let manifest = root
2268 .path()
2269 .join(GENERATIONS_DIR)
2270 .join(receipt.generation_uuid.hyphenated().to_string())
2271 .join(MANIFEST_FILE);
2272 fs::write(&manifest, b"tampered\n").unwrap();
2273 assert_eq!(
2274 stage_project_generation(root.path(), &published_request)
2275 .err()
2276 .expect("tampered published replay must fail")
2277 .code(),
2278 "GF_PROJECT_CORRUPT"
2279 );
2280 }
2281
2282 #[test]
2283 fn optimistic_promotion_refuses_an_existing_generation_destination() {
2284 let root = project();
2285 let request = request(vec![participant("graph", "nodes", b"optimistic")]);
2286 let operation: [u8; 32] = Sha256::digest(b"wave7-existing-destination").into();
2287 let ProjectStageOutcome::Staged(staged) =
2288 stage_project_generation_optimistic(root.path(), &request, operation).unwrap()
2289 else {
2290 panic!("fresh optimistic transaction unexpectedly replayed")
2291 };
2292 let destination = root
2293 .path()
2294 .join(GENERATIONS_DIR)
2295 .join(request.generation_uuid.hyphenated().to_string());
2296 fs::create_dir(&destination).unwrap();
2297 let error = promote_optimistic_generation(&staged).unwrap_err();
2298 assert_eq!(error.code(), "GF_IDEMPOTENCY_CONFLICT");
2299 assert!(error.to_string().contains("generation_exists"));
2300 assert!(staged.generation_root.is_dir());
2301 }
2302
2303 #[test]
2304 fn reader_pinned_to_parent_survives_publication() {
2305 let root = project();
2306 let parent = resolve_project_generation(root.path()).unwrap();
2307 let request = request(vec![participant("graph", "nodes", b"new")]);
2308 let child = request.generation_uuid;
2309
2310 publish(root.path(), request);
2311
2312 assert_ne!(parent.generation_uuid(), child);
2313 assert_eq!(
2314 resolve_project_generation(root.path())
2315 .unwrap()
2316 .generation_uuid(),
2317 child
2318 );
2319 assert!(parent.generation_root().exists());
2320 }
2321
2322 #[test]
2323 fn optimistic_attempts_stage_concurrently_and_compare_parent_at_commit() {
2324 let root = project();
2325 let first = request(vec![participant("graph", "nodes", b"first")]);
2326 let second = request(vec![participant("graph", "nodes", b"second")]);
2327 let first_operation: [u8; 32] = Sha256::digest(b"logical-first").into();
2328 let second_operation: [u8; 32] = Sha256::digest(b"logical-second").into();
2329
2330 let ProjectStageOutcome::Staged(first_staged) =
2331 stage_project_generation_optimistic(root.path(), &first, first_operation).unwrap()
2332 else {
2333 panic!("first optimistic operation replayed unexpectedly");
2334 };
2335 let ProjectStageOutcome::Staged(second_staged) =
2336 stage_project_generation_optimistic(root.path(), &second, second_operation).unwrap()
2337 else {
2338 panic!("second optimistic operation replayed unexpectedly");
2339 };
2340
2341 let first_validated = first_staged.validate(|_| Ok(()), |_, _| Ok(())).unwrap();
2342 let second_validated = second_staged.validate(|_| Ok(()), |_, _| Ok(())).unwrap();
2343 first_validated.publish().unwrap();
2344 let error = second_validated
2345 .publish()
2346 .expect_err("stale optimistic parent must not publish");
2347 assert_eq!(error.code(), "GF_WRITE_CONFLICT");
2348 assert!(
2349 !root
2350 .path()
2351 .join(GENERATIONS_DIR)
2352 .join(second.generation_uuid.hyphenated().to_string())
2353 .exists()
2354 );
2355
2356 let mut rebased = second.clone();
2357 rebased.participants[0].bytes = b"second-rebased".to_vec();
2358 let ProjectStageOutcome::Staged(rebased) =
2359 stage_project_generation_optimistic(root.path(), &rebased, second_operation).unwrap()
2360 else {
2361 panic!("aborted optimistic operation did not permit a rebase attempt");
2362 };
2363 rebased
2364 .validate(|_| Ok(()), |_, _| Ok(()))
2365 .unwrap()
2366 .publish()
2367 .unwrap();
2368 assert_eq!(
2369 resolve_project_generation(root.path())
2370 .unwrap()
2371 .generation_uuid(),
2372 second.generation_uuid
2373 );
2374 }
2375
2376 #[test]
2377 fn optimistic_validation_conflict_aborts_only_its_own_rebase_attempt() {
2378 let root = project();
2379 let mut stale_request = request(vec![participant("graph", "nodes", b"attempt")]);
2380 let operation: [u8; 32] = Sha256::digest(b"logical-validation-rebase").into();
2381 let ProjectStageOutcome::Staged(staged) =
2382 stage_project_generation_optimistic(root.path(), &stale_request, operation).unwrap()
2383 else {
2384 panic!("optimistic operation replayed unexpectedly");
2385 };
2386
2387 publish(
2388 root.path(),
2389 request(vec![participant("graph", "nodes", b"concurrent")]),
2390 );
2391 let error = staged
2392 .validate(
2393 |_| Ok(()),
2394 |parent, _| {
2395 let current = resolve_project_generation(root.path())?;
2396 if current.generation_uuid() != parent.generation_uuid() {
2397 return Err(project_error(
2398 ProjectErrorCode::WriteConflict,
2399 "project generation changed before composite validation",
2400 ));
2401 }
2402 Ok(())
2403 },
2404 )
2405 .err()
2406 .expect("stale optimistic validation must request a rebase");
2407 assert_eq!(error.code(), "GF_WRITE_CONFLICT");
2408 assert_eq!(
2409 read_journal(&journal_path(root.path(), stale_request.transaction_uuid,))
2410 .unwrap()
2411 .phase,
2412 JournalPhase::Aborted
2413 );
2414
2415 let different_operation: [u8; 32] = Sha256::digest(b"different-operation").into();
2416 let identity_error =
2417 stage_project_generation_optimistic(root.path(), &stale_request, different_operation)
2418 .err()
2419 .expect("different logical identity must not reuse the aborted transaction");
2420 assert_eq!(identity_error.code(), "GF_IDEMPOTENCY_CONFLICT");
2421
2422 stale_request.participants[0].bytes = b"rebased-attempt".to_vec();
2423 let ProjectStageOutcome::Staged(rebased) =
2424 stage_project_generation_optimistic(root.path(), &stale_request, operation).unwrap()
2425 else {
2426 panic!("aborted optimistic operation did not permit a rebase attempt");
2427 };
2428 rebased
2429 .validate(|_| Ok(()), |_, _| Ok(()))
2430 .unwrap()
2431 .publish()
2432 .unwrap();
2433 }
2434
2435 #[test]
2436 fn aborted_optimistic_replay_fails_closed_when_generation_cleanup_is_incomplete() {
2437 let root = project();
2438 let request = request(vec![participant("graph", "nodes", b"attempt")]);
2439 let operation: [u8; 32] = Sha256::digest(b"aborted-cleanup-contract").into();
2440 let ProjectStageOutcome::Staged(staged) =
2441 stage_project_generation_optimistic(root.path(), &request, operation).unwrap()
2442 else {
2443 panic!("optimistic operation replayed unexpectedly");
2444 };
2445 abort_stale_generation(&staged).unwrap();
2446 drop(staged);
2447
2448 let leftover = root
2449 .path()
2450 .join(GENERATIONS_DIR)
2451 .join(request.generation_uuid.hyphenated().to_string());
2452 std::fs::create_dir(&leftover).unwrap();
2453 let error = stage_project_generation_optimistic(root.path(), &request, operation)
2454 .err()
2455 .expect("incomplete aborted cleanup must fail closed");
2456 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
2457 assert!(error.to_string().contains("cleanup is incomplete"));
2458 assert!(leftover.exists());
2459 assert_eq!(
2460 read_journal(&journal_path(root.path(), request.transaction_uuid))
2461 .unwrap()
2462 .phase,
2463 JournalPhase::Aborted
2464 );
2465 }
2466
2467 #[test]
2468 fn optimistic_promotion_closes_staged_handles_before_directory_rename() {
2469 let root = project();
2470 let request = request(vec![participant("graph", "nodes", b"promoted")]);
2471 let generation_uuid = request.generation_uuid;
2472 let operation: [u8; 32] = Sha256::digest(b"windows-promotion-handles").into();
2473
2474 let ProjectStageOutcome::Staged(staged) =
2475 stage_project_generation_optimistic(root.path(), &request, operation).unwrap()
2476 else {
2477 panic!("optimistic operation replayed unexpectedly");
2478 };
2479 staged
2480 .validate(|_| Ok(()), |_, _| Ok(()))
2481 .unwrap()
2482 .publish()
2483 .unwrap();
2484
2485 assert_eq!(
2486 resolve_project_generation(root.path())
2487 .unwrap()
2488 .generation_uuid(),
2489 generation_uuid
2490 );
2491 }
2492
2493 #[test]
2494 fn optimistic_transaction_identity_has_one_live_attempt() {
2495 let root = project();
2496 let request = request(vec![participant("graph", "nodes", b"attempt")]);
2497 let operation: [u8; 32] = Sha256::digest(b"logical-attempt").into();
2498 let first = stage_project_generation_optimistic(root.path(), &request, operation).unwrap();
2499
2500 let error = stage_project_generation_optimistic(root.path(), &request, operation)
2501 .err()
2502 .expect("duplicate live attempt must be rejected");
2503 assert_eq!(error.code(), "GF_WRITER_BUSY");
2504 drop(first);
2505 }
2506
2507 #[test]
2508 fn recovery_preserves_live_optimistic_attempt_then_cleans_it_after_release() {
2509 let root = project();
2510 let request = request(vec![participant("graph", "nodes", b"live")]);
2511 let operation: [u8; 32] = Sha256::digest(b"logical-live").into();
2512 let (_, _, request_fingerprint) = request_metadata(&request).unwrap();
2513 let generation_path = root
2514 .path()
2515 .join(ATTEMPTS_DIR)
2516 .join(request.transaction_uuid.hyphenated().to_string())
2517 .join(request_fingerprint);
2518 let staged = stage_project_generation_optimistic(root.path(), &request, operation).unwrap();
2519
2520 let live_report = crate::recover_project_transactions(root.path()).unwrap();
2521 assert_eq!(live_report.aborted_journals, 0);
2522 assert_eq!(live_report.removed_generations, 0);
2523 assert!(generation_path.exists());
2524
2525 drop(staged);
2526 let abandoned_report = crate::recover_project_transactions(root.path()).unwrap();
2527 assert_eq!(abandoned_report.aborted_journals, 1);
2528 assert_eq!(abandoned_report.removed_generations, 1);
2529 assert!(!generation_path.exists());
2530 }
2531
2532 #[test]
2533 fn writer_lock_is_nonblocking_and_fail_closed() {
2534 let root = project();
2535 let lock_dir = root.path().join(LOCKS_DIR);
2536 std::fs::create_dir_all(&lock_dir).unwrap();
2537 let lock = OpenOptions::new()
2538 .read(true)
2539 .write(true)
2540 .create(true)
2541 .truncate(false)
2542 .open(lock_dir.join(WRITER_LOCK_FILE))
2543 .unwrap();
2544 FileExt::lock_exclusive(&lock).unwrap();
2545 let parent = resolve_project_generation(root.path())
2546 .unwrap()
2547 .generation_uuid();
2548
2549 let error = stage_project_generation(
2550 root.path(),
2551 &request(vec![participant("graph", "nodes", b"new")]),
2552 )
2553 .err()
2554 .expect("busy writer must fail");
2555
2556 assert_eq!(error.code(), "GF_WRITER_BUSY");
2557 assert_eq!(
2558 resolve_project_generation(root.path())
2559 .unwrap()
2560 .generation_uuid(),
2561 parent
2562 );
2563 }
2564
2565 #[test]
2566 fn malformed_generation_contracts_fail_before_staging_or_current_change() {
2567 let root = project();
2568 let before = fs::read(root.path().join(CURRENT_FILE)).unwrap();
2569
2570 let mut cases = Vec::new();
2571 let mut no_capability = request(vec![]);
2572 no_capability.capabilities.clear();
2573 cases.push((no_capability, "at least one capability"));
2574
2575 let mut zero_capability = request(vec![]);
2576 zero_capability.capabilities[0].capability_version = 0;
2577 cases.push((zero_capability, "capability contract versions"));
2578
2579 let mut missing_graph = request(vec![]);
2580 missing_graph.capabilities[0].capability_id = "knowledge".into();
2581 cases.push((missing_graph, "graph capability version 1"));
2582
2583 let mut duplicate_capability = request(vec![]);
2584 duplicate_capability
2585 .capabilities
2586 .push(duplicate_capability.capabilities[0].clone());
2587 cases.push((duplicate_capability, "duplicate capability identity"));
2588
2589 let mut undeclared = request(vec![participant("knowledge", "events", b"event")]);
2590 undeclared
2591 .capabilities
2592 .retain(|entry| entry.capability_id == "graph");
2593 cases.push((undeclared, "participant capability is not declared"));
2594
2595 let mut version_mismatch = request(vec![participant("knowledge", "events", b"event")]);
2596 version_mismatch.participants[0].capability_version = 2;
2597 cases.push((version_mismatch, "version conflicts with declaration"));
2598
2599 let duplicate = participant("graph", "nodes", b"same");
2600 cases.push((
2601 request(vec![duplicate.clone(), duplicate]),
2602 "duplicate participant identity",
2603 ));
2604
2605 let mut zero_record = request(vec![participant("graph", "nodes", b"node")]);
2606 zero_record.participants[0].record_version = 0;
2607 cases.push((zero_record, "participant contract versions"));
2608
2609 let mut invalid_id = request(vec![participant("graph", "nodes", b"node")]);
2610 invalid_id.participants[0].record_family_id = "../nodes".into();
2611 cases.push((invalid_id, "machine ID"));
2612
2613 for (candidate, expected) in cases {
2614 let error = stage_project_generation(root.path(), &candidate)
2615 .err()
2616 .expect("malformed request must fail");
2617 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
2618 assert!(error.to_string().contains(expected), "{error}");
2619 assert_eq!(fs::read(root.path().join(CURRENT_FILE)).unwrap(), before);
2620 assert!(!journal_path(root.path(), candidate.transaction_uuid).exists());
2621 }
2622 }
2623
2624 #[test]
2625 fn publication_error_redacts_unsafe_cause_and_digest_parser_is_canonical() {
2626 let transaction = Uuid::now_v7();
2627 let generation = Uuid::now_v7();
2628 let error = publication_error_from_parts(
2629 transaction,
2630 generation,
2631 "STAGED",
2632 false,
2633 "bad/path:\nsecret=<value>!",
2634 );
2635 let text = error.to_string();
2636 assert!(text.contains("phase=STAGED committed=false cause=badpathsecretvalue"));
2637 assert!(!text.contains('/') && !text.contains('<') && !text.contains('!'));
2638
2639 let bytes = [0xabu8; 32];
2640 let canonical = hex_digest(bytes);
2641 assert_eq!(parse_digest(&canonical), Some(bytes));
2642 for malformed in ["", "ab", &"A".repeat(64), &"g".repeat(64)] {
2643 assert_eq!(parse_digest(malformed), None);
2644 }
2645 }
2646
2647 #[test]
2648 fn published_transaction_probe_verifies_durable_manifest_on_reopen() {
2649 let root = project();
2650 let input = request(vec![
2651 participant("graph", "nodes", b"nodes"),
2652 participant("graph", "edges", b"edges"),
2653 ]);
2654 let receipt = publish(root.path(), input.clone());
2655 let probed = published_project_transaction(root.path(), input.transaction_uuid)
2656 .unwrap()
2657 .unwrap();
2658 assert_eq!(probed.transaction_uuid, receipt.transaction_uuid);
2659 assert_eq!(probed.generation_uuid, receipt.generation_uuid);
2660 assert_eq!(
2661 probed.generation_manifest_sha256,
2662 receipt.generation_manifest_sha256
2663 );
2664 assert!(probed.idempotent_replay);
2665 assert!(
2666 published_project_transaction(root.path(), Uuid::now_v7())
2667 .unwrap()
2668 .is_none()
2669 );
2670
2671 let reopened = resolve_project_generation(root.path()).unwrap();
2672 assert_eq!(reopened.generation_uuid(), receipt.generation_uuid);
2673 reopened.validate_complete_participant_inventory().unwrap();
2674 drop(reopened);
2675
2676 let manifest = root
2677 .path()
2678 .join(GENERATIONS_DIR)
2679 .join(receipt.generation_uuid.hyphenated().to_string())
2680 .join(MANIFEST_FILE);
2681 fs::write(&manifest, b"tampered\n").unwrap();
2682 let error = published_project_transaction(root.path(), input.transaction_uuid).unwrap_err();
2683 assert_eq!(error.code(), "GF_PROJECT_CORRUPT");
2684 assert!(error.to_string().contains("does not match its journal"));
2685 }
2686
2687 #[test]
2688 fn staged_participant_file_kind_matrix_fails_before_current_mutation() {
2689 for kind in ["missing", "directory", "symlink"] {
2690 let root = project();
2691 let parent = resolve_project_generation(root.path())
2692 .unwrap()
2693 .generation_uuid();
2694 let request = request(vec![participant("graph", "nodes", b"stable")]);
2695 let ProjectStageOutcome::Staged(staged) =
2696 stage_project_generation(root.path(), &request).unwrap()
2697 else {
2698 panic!("unexpected replay")
2699 };
2700 let path = staged
2701 .generation_root
2702 .join(PARTICIPANTS_DIR)
2703 .join(&staged.participants.first().unwrap().relative_path);
2704 std::fs::remove_file(&path).unwrap();
2705 match kind {
2706 "missing" => {}
2707 "directory" => std::fs::create_dir(&path).unwrap(),
2708 "symlink" => {
2709 #[cfg(unix)]
2710 std::os::unix::fs::symlink(root.path().join(CURRENT_FILE), &path).unwrap();
2711 #[cfg(not(unix))]
2712 std::fs::create_dir(&path).unwrap();
2713 }
2714 _ => unreachable!(),
2715 }
2716 let error = match staged.validate(|_| Ok(()), |_, _| Ok(())) {
2717 Ok(_) => panic!("hostile staged participant must fail"),
2718 Err(error) => error,
2719 };
2720 let expected_code = if kind == "missing" {
2721 "GF_IO"
2722 } else {
2723 "GF_PUBLICATION_FAILED"
2724 };
2725 assert_eq!(error.code(), expected_code);
2726 assert_eq!(
2727 resolve_project_generation(root.path())
2728 .unwrap()
2729 .generation_uuid(),
2730 parent
2731 );
2732 }
2733 }
2734
2735 #[test]
2736 fn journal_decode_and_atomic_temp_cleanup_matrix_is_fail_closed() {
2737 let root = project();
2738 let journal = root.path().join(TRANSACTIONS_DIR).join("malformed.json");
2739 std::fs::create_dir_all(journal.parent().unwrap()).unwrap();
2740 for bytes in [
2741 b"not-json".as_slice(),
2742 br#"{"journal_version":999}"#,
2743 br#"{"journal_version":1,"transaction_uuid":"bad"}"#,
2744 ] {
2745 std::fs::write(&journal, bytes).unwrap();
2746 assert_eq!(
2747 read_journal(&journal).unwrap_err().code(),
2748 "GF_PROJECT_CORRUPT"
2749 );
2750 assert_eq!(std::fs::read(&journal).unwrap(), bytes);
2751 }
2752
2753 let unrelated = root.path().join("metadata.json");
2754 assert!(!cleanup_atomicwrite_temp(&unrelated).unwrap());
2755 let empty = root.path().join(".atomicwriteabc123");
2756 std::fs::create_dir(&empty).unwrap();
2757 assert!(cleanup_atomicwrite_temp(&empty).unwrap());
2758 assert!(!empty.exists());
2759
2760 let populated = root.path().join(".atomicwritedef456");
2761 std::fs::create_dir(&populated).unwrap();
2762 std::fs::write(populated.join("tmpfile.tmp"), b"abandoned").unwrap();
2763 assert!(cleanup_atomicwrite_temp(&populated).unwrap());
2764 assert!(!populated.exists());
2765
2766 let hostile = root.path().join(".atomicwriteghi789");
2767 std::fs::create_dir(&hostile).unwrap();
2768 std::fs::write(hostile.join("unexpected"), b"caller bytes").unwrap();
2769 assert!(!cleanup_atomicwrite_temp(&hostile).unwrap());
2770 assert_eq!(
2771 std::fs::read(hostile.join("unexpected")).unwrap(),
2772 b"caller bytes"
2773 );
2774 }
2775
2776 #[test]
2777 fn wave9_journal_metadata_and_lock_aliases_fail_closed() {
2778 let root = project();
2779 let journal = root.path().join(TRANSACTIONS_DIR).join("hostile.json");
2780 std::fs::create_dir_all(journal.parent().unwrap()).unwrap();
2781
2782 std::fs::create_dir(&journal).unwrap();
2783 assert_eq!(
2784 read_journal(&journal).unwrap_err().code(),
2785 "GF_PROJECT_CORRUPT"
2786 );
2787 std::fs::remove_dir(&journal).unwrap();
2788 std::fs::write(&journal, vec![b'x'; MAX_JOURNAL_BYTES as usize + 1]).unwrap();
2789 assert_eq!(
2790 read_journal(&journal).unwrap_err().code(),
2791 "GF_PROJECT_CORRUPT"
2792 );
2793
2794 #[cfg(unix)]
2795 {
2796 use std::os::unix::fs::symlink;
2797
2798 std::fs::remove_file(&journal).unwrap();
2799 let target = root.path().join("caller-journal");
2800 std::fs::write(&target, b"caller bytes").unwrap();
2801 symlink(&target, &journal).unwrap();
2802 assert_eq!(
2803 read_journal(&journal).unwrap_err().code(),
2804 "GF_PROJECT_CORRUPT"
2805 );
2806
2807 let owner = root.path().join("lock-owner");
2808 let alias = root.path().join("lock-alias");
2809 std::fs::write(&owner, b"").unwrap();
2810 std::fs::hard_link(&owner, &alias).unwrap();
2811 assert_eq!(
2812 open_regular_lock(&alias).unwrap_err().code(),
2813 "GF_PROJECT_CORRUPT"
2814 );
2815 assert!(owner.exists());
2816 }
2817 }
2818
2819 #[test]
2820 fn atomic_temp_cleanup_rejects_near_misses_links_and_multiple_entries() {
2821 let root = project();
2822 for name in [
2823 ".atomicwrite",
2824 ".atomicwrite12345",
2825 ".atomicwrite1234567",
2826 ".atomicwrite12-456",
2827 "atomicwrite123456",
2828 ] {
2829 let path = root.path().join(name);
2830 std::fs::create_dir(&path).unwrap();
2831 assert!(!cleanup_atomicwrite_temp(&path).unwrap());
2832 assert!(path.exists());
2833 }
2834
2835 let regular = root.path().join(".atomicwriteabc001");
2836 std::fs::write(®ular, b"caller").unwrap();
2837 assert!(!cleanup_atomicwrite_temp(®ular).unwrap());
2838 assert_eq!(std::fs::read(®ular).unwrap(), b"caller");
2839
2840 let multiple = root.path().join(".atomicwriteabc002");
2841 std::fs::create_dir(&multiple).unwrap();
2842 std::fs::write(multiple.join("tmpfile.tmp"), b"temporary").unwrap();
2843 std::fs::write(multiple.join("second"), b"caller").unwrap();
2844 assert!(!cleanup_atomicwrite_temp(&multiple).unwrap());
2845 assert_eq!(std::fs::read(multiple.join("second")).unwrap(), b"caller");
2846
2847 let wrong_entry = root.path().join(".atomicwriteabc003");
2848 std::fs::create_dir(&wrong_entry).unwrap();
2849 std::fs::write(wrong_entry.join("not-temp"), b"caller").unwrap();
2850 assert!(!cleanup_atomicwrite_temp(&wrong_entry).unwrap());
2851 assert_eq!(
2852 std::fs::read(wrong_entry.join("not-temp")).unwrap(),
2853 b"caller"
2854 );
2855
2856 #[cfg(unix)]
2857 {
2858 use std::os::unix::fs::symlink;
2859
2860 let linked_dir = root.path().join(".atomicwriteabc004");
2861 symlink(root.path(), &linked_dir).unwrap();
2862 assert!(!cleanup_atomicwrite_temp(&linked_dir).unwrap());
2863 assert!(
2864 linked_dir
2865 .symlink_metadata()
2866 .unwrap()
2867 .file_type()
2868 .is_symlink()
2869 );
2870
2871 let linked_entry = root.path().join(".atomicwriteabc005");
2872 std::fs::create_dir(&linked_entry).unwrap();
2873 symlink(
2874 root.path().join(CURRENT_FILE),
2875 linked_entry.join("tmpfile.tmp"),
2876 )
2877 .unwrap();
2878 assert!(!cleanup_atomicwrite_temp(&linked_entry).unwrap());
2879 assert!(
2880 linked_entry
2881 .join("tmpfile.tmp")
2882 .symlink_metadata()
2883 .unwrap()
2884 .file_type()
2885 .is_symlink()
2886 );
2887
2888 let hardlinked_entry = root.path().join(".atomicwriteabc006");
2889 std::fs::create_dir(&hardlinked_entry).unwrap();
2890 let owned = root.path().join("hardlink-owner");
2891 std::fs::write(&owned, b"caller").unwrap();
2892 std::fs::hard_link(&owned, hardlinked_entry.join("tmpfile.tmp")).unwrap();
2893 assert!(!cleanup_atomicwrite_temp(&hardlinked_entry).unwrap());
2894 assert_eq!(std::fs::read(&owned).unwrap(), b"caller");
2895 }
2896 }
2897
2898 #[test]
2899 fn machine_directory_and_lock_reject_hostile_path_components_without_replacement() {
2900 let root = tempfile::tempdir().unwrap();
2901 for relative in [
2902 Path::new("../escape"),
2903 Path::new("/absolute"),
2904 Path::new("safe/../escape"),
2905 ] {
2906 assert_eq!(
2907 ensure_machine_directory(root.path(), relative)
2908 .unwrap_err()
2909 .code(),
2910 "GF_PROJECT_CORRUPT"
2911 );
2912 }
2913 let file = root.path().join("owned");
2914 std::fs::write(&file, b"caller bytes").unwrap();
2915 assert_eq!(
2916 ensure_machine_directory(root.path(), Path::new("owned/child"))
2917 .unwrap_err()
2918 .code(),
2919 "GF_PROJECT_CORRUPT"
2920 );
2921 assert_eq!(std::fs::read(&file).unwrap(), b"caller bytes");
2922
2923 let lock = root.path().join("lock");
2924 std::fs::create_dir(&lock).unwrap();
2925 assert_eq!(
2926 open_regular_lock(&lock).unwrap_err().code(),
2927 "GF_PROJECT_CORRUPT"
2928 );
2929 assert!(lock.is_dir());
2930 }
2931
2932 #[test]
2933 fn public_transaction_probe_distinguishes_absent_staged_and_durable_publication() {
2934 let root = project();
2935 let request = request(vec![participant("graph", "nodes", b"rows")]);
2936
2937 assert!(
2938 published_project_transaction(root.path(), request.transaction_uuid)
2939 .unwrap()
2940 .is_none()
2941 );
2942 let ProjectStageOutcome::Staged(staged) =
2943 stage_project_generation(root.path(), &request).unwrap()
2944 else {
2945 panic!("fresh transaction unexpectedly replayed");
2946 };
2947 assert!(
2948 published_project_transaction(root.path(), request.transaction_uuid)
2949 .unwrap()
2950 .is_none()
2951 );
2952 let receipt = staged
2953 .validate(|_| Ok(()), |_, _| Ok(()))
2954 .unwrap()
2955 .publish()
2956 .unwrap();
2957 let reopened = published_project_transaction(root.path(), request.transaction_uuid)
2958 .unwrap()
2959 .unwrap();
2960 assert_eq!(reopened.generation_uuid, receipt.generation_uuid);
2961 assert_eq!(
2962 reopened.generation_manifest_sha256,
2963 receipt.generation_manifest_sha256
2964 );
2965 assert!(reopened.idempotent_replay);
2966 }
2967}