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).map_err(|error| match error {
205 GfError::Storage(message) => publication_error(request, "STAGE", false, &message),
206 other => other,
207 })
208}
209
210pub fn stage_project_generation_optimistic(
222 container_root: impl AsRef<Path>,
223 request: &ProjectGenerationRequest,
224 operation_fingerprint: [u8; 32],
225) -> Result<ProjectStageOutcome, GfError> {
226 stage_project_generation_optimistic_inner(
227 container_root.as_ref(),
228 request,
229 operation_fingerprint,
230 )
231 .map_err(|error| match error {
232 GfError::Storage(message) => publication_error(request, "STAGE", false, &message),
233 other => other,
234 })
235}
236
237fn stage_project_generation_inner(
238 container_root: &Path,
239 request: &ProjectGenerationRequest,
240) -> Result<ProjectStageOutcome, GfError> {
241 let root = canonical_supported_root(container_root)?;
242 let writer_lock = acquire_writer_lock(&root, request)?;
243 project_failpoint::hit(
244 "project.after_writer_lock",
245 Some(request.transaction_uuid),
246 Some(request.generation_uuid),
247 "WRITER_LOCK",
248 false,
249 )?;
250 let parent = resolve_project_generation(&root)?;
251 stage_project_generation_with_lock(root, writer_lock, parent, request, None)
252}
253
254fn stage_project_generation_optimistic_inner(
255 container_root: &Path,
256 request: &ProjectGenerationRequest,
257 operation_fingerprint: [u8; 32],
258) -> Result<ProjectStageOutcome, GfError> {
259 let root = canonical_supported_root(container_root)?;
260 let transaction_lock = acquire_transaction_lock(&root, request)?;
261 let parent = resolve_project_generation(&root)?;
262 stage_project_generation_inner_with_locks(
263 root,
264 PublicationLock::Optimistic(transaction_lock),
265 parent,
266 request,
267 None,
268 Some(operation_fingerprint),
269 )
270}
271
272pub(crate) fn stage_project_generation_with_lock(
275 root: PathBuf,
276 writer_lock: File,
277 parent: ResolvedProjectGeneration,
278 request: &ProjectGenerationRequest,
279 revert: Option<RevertJournalExtension>,
280) -> Result<ProjectStageOutcome, GfError> {
281 stage_project_generation_inner_with_locks(
282 root,
283 PublicationLock::Exclusive(writer_lock),
284 parent,
285 request,
286 revert,
287 None,
288 )
289}
290
291fn stage_project_generation_inner_with_locks(
292 root: PathBuf,
293 publication_lock: PublicationLock,
294 parent: ResolvedProjectGeneration,
295 request: &ProjectGenerationRequest,
296 revert: Option<RevertJournalExtension>,
297 operation_fingerprint: Option<[u8; 32]>,
298) -> Result<ProjectStageOutcome, GfError> {
299 validate_request(request)?;
300 let (capabilities, participants, request_fingerprint) = request_metadata(request)?;
301 let operation_fingerprint =
302 operation_fingerprint.map_or_else(|| request_fingerprint.clone(), hex_digest);
303 let transactions_dir = ensure_machine_directory(&root, Path::new(TRANSACTIONS_DIR))?;
304 sync_directory(&root)?;
305 let journal_path =
306 transactions_dir.join(format!("{}.json", request.transaction_uuid.hyphenated()));
307 if journal_path.exists()
308 && let Some(outcome) = handle_existing_journal(
309 &root,
310 request,
311 &request_fingerprint,
312 &operation_fingerprint,
313 revert.as_ref(),
314 &journal_path,
315 )?
316 {
317 return Ok(outcome);
318 }
319
320 let requires_promotion = matches!(publication_lock, PublicationLock::Optimistic(_));
321 let generation_root =
322 prepare_generation_directory(&root, request, &request_fingerprint, requires_promotion)?;
323 write_journal(
324 &journal_path,
325 &JournalRecord::new(
326 request,
327 Some(parent.generation_uuid()),
328 JournalPhase::Preparing,
329 (request_fingerprint.clone(), operation_fingerprint.clone()),
330 &participants,
331 None,
332 revert.clone(),
333 ),
334 )?;
335 project_failpoint::hit(
336 "project.after_journal_preparing",
337 Some(request.transaction_uuid),
338 Some(request.generation_uuid),
339 "PREPARING",
340 false,
341 )?;
342
343 stage_participant_files(request, &generation_root, &participants)?;
344 sync_participant_directories(&generation_root.join(PARTICIPANTS_DIR), &participants)?;
345 project_failpoint::hit(
346 "project.after_participant_dir_fsync",
347 Some(request.transaction_uuid),
348 Some(request.generation_uuid),
349 "STAGED",
350 false,
351 )?;
352 write_journal(
353 &journal_path,
354 &JournalRecord::new(
355 request,
356 Some(parent.generation_uuid()),
357 JournalPhase::Staged,
358 (request_fingerprint.clone(), operation_fingerprint.clone()),
359 &participants,
360 None,
361 revert.clone(),
362 ),
363 )?;
364 project_failpoint::hit(
365 "project.after_journal_staged",
366 Some(request.transaction_uuid),
367 Some(request.generation_uuid),
368 "STAGED",
369 false,
370 )?;
371
372 Ok(ProjectStageOutcome::Staged(Box::new(
373 StagedProjectGeneration {
374 root,
375 publication_lock,
376 parent,
377 transaction_uuid: request.transaction_uuid,
378 generation_uuid: request.generation_uuid,
379 generation_root,
380 requires_promotion,
381 request_fingerprint,
382 operation_fingerprint,
383 capabilities,
384 participants,
385 revert,
386 },
387 )))
388}
389
390fn acquire_writer_lock(root: &Path, request: &ProjectGenerationRequest) -> Result<File, GfError> {
391 acquire_writer_lock_for_parts(root, request.transaction_uuid, request.generation_uuid)
392}
393
394fn acquire_writer_lock_for_parts(
395 root: &Path,
396 transaction_uuid: Uuid,
397 generation_uuid: Uuid,
398) -> Result<File, GfError> {
399 let lock_dir = ensure_machine_directory(root, Path::new(LOCKS_DIR))?;
400 sync_directory(root)?;
401 let writer_lock = open_regular_lock(&lock_dir.join(WRITER_LOCK_FILE))?;
402 if !FileExt::try_lock_exclusive(&writer_lock).map_err(publication_io)? {
403 return Err(project_error(
404 ProjectErrorCode::WriterBusy,
405 format!(
406 "transaction_uuid={} generation_uuid={} phase=WRITER_LOCK committed=false cause=busy",
407 transaction_uuid.hyphenated(),
408 generation_uuid.hyphenated()
409 ),
410 ));
411 }
412 Ok(writer_lock)
413}
414
415fn wait_for_writer_lock(root: &Path) -> Result<File, GfError> {
416 let lock_dir = ensure_machine_directory(root, Path::new(LOCKS_DIR))?;
417 sync_directory(root)?;
418 let writer_lock = open_regular_lock(&lock_dir.join(WRITER_LOCK_FILE))?;
419 FileExt::lock_exclusive(&writer_lock).map_err(publication_io)?;
420 Ok(writer_lock)
421}
422
423fn acquire_transaction_lock(
424 root: &Path,
425 request: &ProjectGenerationRequest,
426) -> Result<File, GfError> {
427 let lock = open_transaction_lock(root, request.transaction_uuid)?;
428 if !FileExt::try_lock_exclusive(&lock).map_err(publication_io)? {
429 return Err(project_error(
430 ProjectErrorCode::WriterBusy,
431 format!(
432 "transaction_uuid={} generation_uuid={} phase=TRANSACTION_LOCK committed=false cause=busy",
433 request.transaction_uuid.hyphenated(),
434 request.generation_uuid.hyphenated()
435 ),
436 ));
437 }
438 Ok(lock)
439}
440
441pub(crate) fn open_transaction_lock(root: &Path, transaction_uuid: Uuid) -> Result<File, GfError> {
442 let lock_dir =
443 ensure_machine_directory(root, &Path::new(LOCKS_DIR).join(TRANSACTION_LOCKS_DIR))?;
444 open_regular_lock(&lock_dir.join(format!("{}.lock", transaction_uuid.hyphenated())))
445}
446
447fn handle_existing_journal(
448 root: &Path,
449 request: &ProjectGenerationRequest,
450 request_fingerprint: &str,
451 operation_fingerprint: &str,
452 expected_revert: Option<&RevertJournalExtension>,
453 journal_path: &Path,
454) -> Result<Option<ProjectStageOutcome>, GfError> {
455 let journal = read_journal(journal_path)?;
456 if journal.operation_fingerprint() != operation_fingerprint
457 || journal.generation_uuid != request.generation_uuid.hyphenated().to_string()
458 || journal.revert.as_ref() != expected_revert
459 {
460 return Err(transaction_conflict(request));
461 }
462 if journal.phase == JournalPhase::Aborted {
463 let generation_name = request.generation_uuid.hyphenated().to_string();
464 if root.join(GENERATIONS_DIR).join(&generation_name).exists()
465 || root.join("trash").join(&generation_name).exists()
466 {
467 return Err(publication_error(
468 request,
469 "ABORTED",
470 false,
471 "aborted transaction cleanup is incomplete; run recovery again",
472 ));
473 }
474 cleanup_aborted_attempts(root, request.transaction_uuid)?;
475 return Ok(None);
476 }
477 if journal.request_fingerprint != request_fingerprint
478 && journal.phase != JournalPhase::Published
479 {
480 return Err(transaction_conflict(request));
481 }
482 if journal.phase != JournalPhase::Published {
483 return Err(publication_error(
484 request,
485 "PREPARING",
486 false,
487 "an interrupted transaction requires recovery",
488 ));
489 }
490 let digest = journal
491 .generation_manifest_sha256
492 .as_deref()
493 .and_then(parse_digest)
494 .ok_or_else(|| {
495 project_error(
496 ProjectErrorCode::ProjectCorrupt,
497 "published transaction journal has no valid manifest digest",
498 )
499 })?;
500 let manifest_path = root
501 .join(GENERATIONS_DIR)
502 .join(request.generation_uuid.hyphenated().to_string())
503 .join(MANIFEST_FILE);
504 let manifest_bytes = std::fs::read(&manifest_path).map_err(publication_io)?;
505 let actual: [u8; 32] = Sha256::digest(&manifest_bytes).into();
506 if actual != digest {
507 return Err(project_error(
508 ProjectErrorCode::ProjectCorrupt,
509 "published transaction manifest does not match its journal",
510 ));
511 }
512 Ok(Some(ProjectStageOutcome::AlreadyPublished(
513 ProjectPublicationReceipt {
514 transaction_uuid: request.transaction_uuid,
515 generation_uuid: request.generation_uuid,
516 generation_manifest_sha256: digest,
517 idempotent_replay: true,
518 },
519 )))
520}
521
522fn cleanup_aborted_attempts(root: &Path, transaction_uuid: Uuid) -> Result<(), GfError> {
523 let attempts_root = root.join(ATTEMPTS_DIR);
524 let transaction_root = attempts_root.join(transaction_uuid.hyphenated().to_string());
525 let metadata = match std::fs::symlink_metadata(&transaction_root) {
526 Ok(metadata) => metadata,
527 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(()),
528 Err(error) => return Err(publication_io(error)),
529 };
530 if !metadata.is_dir() || metadata.file_type().is_symlink() {
531 return Err(project_error(
532 ProjectErrorCode::ProjectCorrupt,
533 "aborted transaction attempt path is linked or not a directory",
534 ));
535 }
536 std::fs::remove_dir_all(&transaction_root).map_err(publication_io)?;
537 sync_directory(&attempts_root)
538}
539
540pub(crate) fn load_revert_journal_extension(
542 root: &Path,
543 transaction_uuid: Uuid,
544) -> Result<Option<RevertJournalExtension>, GfError> {
545 let path = root
546 .join(TRANSACTIONS_DIR)
547 .join(format!("{}.json", transaction_uuid.hyphenated()));
548 if !path.exists() {
549 return Ok(None);
550 }
551 Ok(read_journal(&path)?.revert)
552}
553
554pub(crate) fn load_published_revert(
556 root: &Path,
557 transaction_uuid: Uuid,
558) -> Result<Option<(RevertJournalExtension, ProjectPublicationReceipt)>, GfError> {
559 let path = root
560 .join(TRANSACTIONS_DIR)
561 .join(format!("{}.json", transaction_uuid.hyphenated()));
562 if !path.exists() {
563 return Ok(None);
564 }
565 let journal = read_journal(&path)?;
566 let Some(revert) = journal.revert else {
567 return Ok(None);
568 };
569 if journal.phase != JournalPhase::Published {
570 return Ok(None);
571 }
572 let generation_uuid = Uuid::parse_str(&journal.generation_uuid).map_err(|_| {
573 project_error(
574 ProjectErrorCode::ProjectCorrupt,
575 "published journal has an invalid generation UUID",
576 )
577 })?;
578 let digest = journal
579 .generation_manifest_sha256
580 .as_deref()
581 .and_then(parse_digest)
582 .ok_or_else(|| {
583 project_error(
584 ProjectErrorCode::ProjectCorrupt,
585 "published journal has an invalid manifest digest",
586 )
587 })?;
588 Ok(Some((
589 revert,
590 ProjectPublicationReceipt {
591 transaction_uuid,
592 generation_uuid,
593 generation_manifest_sha256: digest,
594 idempotent_replay: true,
595 },
596 )))
597}
598
599pub fn published_project_transaction(
604 root: &Path,
605 transaction_uuid: Uuid,
606) -> Result<Option<ProjectPublicationReceipt>, GfError> {
607 let path = root
608 .join(TRANSACTIONS_DIR)
609 .join(format!("{}.json", transaction_uuid.hyphenated()));
610 if !path.exists() {
611 return Ok(None);
612 }
613 let journal = read_journal(&path)?;
614 if journal.phase != JournalPhase::Published {
615 return Ok(None);
616 }
617 let generation_uuid = Uuid::parse_str(&journal.generation_uuid).map_err(|_| {
618 project_error(
619 ProjectErrorCode::ProjectCorrupt,
620 "published journal has an invalid generation UUID",
621 )
622 })?;
623 let digest = journal
624 .generation_manifest_sha256
625 .as_deref()
626 .and_then(parse_digest)
627 .ok_or_else(|| {
628 project_error(
629 ProjectErrorCode::ProjectCorrupt,
630 "published journal has an invalid manifest digest",
631 )
632 })?;
633 let manifest_path = root
634 .join(GENERATIONS_DIR)
635 .join(generation_uuid.hyphenated().to_string())
636 .join(MANIFEST_FILE);
637 let actual: [u8; 32] =
638 Sha256::digest(std::fs::read(manifest_path).map_err(publication_io)?).into();
639 if actual != digest {
640 return Err(project_error(
641 ProjectErrorCode::ProjectCorrupt,
642 "published transaction manifest does not match its journal",
643 ));
644 }
645 Ok(Some(ProjectPublicationReceipt {
646 transaction_uuid,
647 generation_uuid,
648 generation_manifest_sha256: digest,
649 idempotent_replay: true,
650 }))
651}
652
653fn prepare_generation_directory(
654 root: &Path,
655 request: &ProjectGenerationRequest,
656 request_fingerprint: &str,
657 requires_promotion: bool,
658) -> Result<PathBuf, GfError> {
659 let relative_root = if requires_promotion {
660 Path::new(ATTEMPTS_DIR)
661 .join(request.transaction_uuid.hyphenated().to_string())
662 .join(request_fingerprint)
663 } else {
664 Path::new(GENERATIONS_DIR).join(request.generation_uuid.hyphenated().to_string())
665 };
666 let generation_root = root.join(&relative_root);
667 if generation_root.exists() {
668 return Err(transaction_conflict(request));
669 }
670 ensure_machine_directory(root, &relative_root.join(PARTICIPANTS_DIR))?;
671 Ok(generation_root)
672}
673
674fn stage_participant_files(
675 request: &ProjectGenerationRequest,
676 generation_root: &Path,
677 participants: &[StagedParticipant],
678) -> Result<(), GfError> {
679 for metadata in participants {
680 let input = request
681 .participants
682 .iter()
683 .find(|candidate| {
684 candidate.capability_id == metadata.capability_id
685 && candidate.record_family_id == metadata.record_family_id
686 })
687 .expect("validated canonical metadata has one source participant");
688 let destination = generation_root
689 .join(PARTICIPANTS_DIR)
690 .join(&metadata.relative_path);
691 let parent_dir = destination
692 .parent()
693 .expect("machine-derived participant path has a parent");
694 let relative_parent = parent_dir
695 .strip_prefix(generation_root)
696 .expect("machine-derived participant parent is contained");
697 ensure_machine_directory(generation_root, relative_parent)?;
698 let mut file = OpenOptions::new()
699 .write(true)
700 .create_new(true)
701 .open(&destination)
702 .map_err(publication_io)?;
703 file.write_all(&input.bytes).map_err(publication_io)?;
704 project_failpoint::hit(
705 "project.after_participant_write",
706 Some(request.transaction_uuid),
707 Some(request.generation_uuid),
708 "STAGED",
709 false,
710 )?;
711 file.sync_all().map_err(publication_io)?;
712 project_failpoint::hit(
713 "project.after_participant_fsync",
714 Some(request.transaction_uuid),
715 Some(request.generation_uuid),
716 "STAGED",
717 false,
718 )?;
719 verify_participant_file(&destination, metadata)?;
720 }
721 Ok(())
722}
723
724impl StagedProjectGeneration {
725 #[must_use]
727 pub fn participants(&self) -> &[StagedParticipant] {
728 &self.participants
729 }
730
731 pub fn validate<D, C>(
736 self,
737 domain_validation: D,
738 composite_validation: C,
739 ) -> Result<ValidatedProjectGeneration, GfError>
740 where
741 D: FnOnce(&[StagedParticipant]) -> Result<(), GfError>,
742 C: FnOnce(&ResolvedProjectGeneration, &[StagedParticipant]) -> Result<(), GfError>,
743 {
744 for participant in &self.participants {
745 verify_participant_file(
746 &self
747 .generation_root
748 .join(PARTICIPANTS_DIR)
749 .join(&participant.relative_path),
750 participant,
751 )?;
752 }
753 domain_validation(&self.participants)?;
754 project_failpoint::hit(
755 "project.after_domain_validation",
756 Some(self.transaction_uuid),
757 Some(self.generation_uuid),
758 "VALIDATED",
759 false,
760 )?;
761 composite_validation(&self.parent, &self.participants)?;
762 project_failpoint::hit(
763 "project.after_composite_validation",
764 Some(self.transaction_uuid),
765 Some(self.generation_uuid),
766 "VALIDATED",
767 false,
768 )?;
769 write_journal(
770 &self.journal_path(),
771 &self.journal(JournalPhase::Validated, None),
772 )?;
773 project_failpoint::hit(
774 "project.after_journal_validated",
775 Some(self.transaction_uuid),
776 Some(self.generation_uuid),
777 "VALIDATED",
778 false,
779 )?;
780 Ok(ValidatedProjectGeneration(self))
781 }
782
783 fn journal_path(&self) -> PathBuf {
784 self.root
785 .join(TRANSACTIONS_DIR)
786 .join(format!("{}.json", self.transaction_uuid.hyphenated()))
787 }
788
789 fn journal(&self, phase: JournalPhase, manifest_sha256: Option<String>) -> JournalRecord {
790 JournalRecord {
791 format: "graphforge-transaction".into(),
792 format_version: 1,
793 transaction_uuid: self.transaction_uuid.hyphenated().to_string(),
794 generation_uuid: self.generation_uuid.hyphenated().to_string(),
795 parent_generation_uuid: Some(self.parent.generation_uuid().hyphenated().to_string()),
796 phase,
797 request_fingerprint: self.request_fingerprint.clone(),
798 operation_fingerprint: Some(self.operation_fingerprint.clone()),
799 participant_paths: self
800 .participants
801 .iter()
802 .map(|participant| participant.relative_path.clone())
803 .collect(),
804 generation_manifest_sha256: manifest_sha256,
805 revert: self.revert.clone(),
806 }
807 }
808}
809
810impl ValidatedProjectGeneration {
811 pub fn publish(self) -> Result<ProjectPublicationReceipt, GfError> {
817 let commit_lock = self.prepare_commit_lock()?;
818 let result = self.publish_inner().map_err(|error| {
819 if matches!(error, GfError::Project { .. }) {
820 error
821 } else {
822 publication_error_from_parts(
823 self.0.transaction_uuid,
824 self.0.generation_uuid,
825 "DURABLE",
826 false,
827 &error.to_string(),
828 )
829 }
830 });
831 drop(commit_lock);
832 result
833 }
834
835 fn prepare_commit_lock(&self) -> Result<Option<File>, GfError> {
836 if matches!(self.0.publication_lock, PublicationLock::Exclusive(_)) {
837 return Ok(None);
838 }
839 let staged = &self.0;
840 let writer_lock = wait_for_writer_lock(&staged.root)?;
841 project_failpoint::hit(
842 "project.after_optimistic_commit_lock",
843 Some(staged.transaction_uuid),
844 Some(staged.generation_uuid),
845 "COMMIT_LOCK",
846 false,
847 )?;
848 let current = resolve_project_generation(&staged.root)?;
849 if current.generation_uuid() != staged.parent.generation_uuid() {
850 abort_stale_generation(staged)?;
851 return Err(project_error(
852 ProjectErrorCode::WriteConflict,
853 format!(
854 "transaction_uuid={} generation_uuid={} phase=COMMIT_LOCK committed=false cause=stale_parent expected_parent={} actual_parent={}",
855 staged.transaction_uuid.hyphenated(),
856 staged.generation_uuid.hyphenated(),
857 staged.parent.generation_uuid().hyphenated(),
858 current.generation_uuid().hyphenated()
859 ),
860 ));
861 }
862 Ok(Some(writer_lock))
863 }
864
865 fn publish_inner(&self) -> Result<ProjectPublicationReceipt, GfError> {
866 let staged = &self.0;
867 let manifest_sha256 = make_generation_durable(staged)?;
868 replace_current(staged, manifest_sha256)?;
869 finish_published_generation(staged, manifest_sha256)?;
870 Ok(ProjectPublicationReceipt {
871 transaction_uuid: staged.transaction_uuid,
872 generation_uuid: staged.generation_uuid,
873 generation_manifest_sha256: manifest_sha256,
874 idempotent_replay: false,
875 })
876 }
877}
878
879fn abort_stale_generation(staged: &StagedProjectGeneration) -> Result<(), GfError> {
880 write_journal(
881 &staged.journal_path(),
882 &staged.journal(JournalPhase::Aborted, None),
883 )?;
884 if staged.generation_root.exists() {
885 std::fs::remove_dir_all(&staged.generation_root).map_err(publication_io)?;
886 sync_directory(
887 staged
888 .generation_root
889 .parent()
890 .expect("machine attempt path has a parent"),
891 )?;
892 }
893 Ok(())
894}
895
896fn make_generation_durable(staged: &StagedProjectGeneration) -> Result<[u8; 32], GfError> {
897 let lease_path = staged.generation_root.join(LEASE_FILE);
898 let lease = OpenOptions::new()
899 .write(true)
900 .create_new(true)
901 .open(&lease_path)
902 .map_err(publication_io)?;
903 lease.sync_all().map_err(publication_io)?;
904 drop(lease);
908
909 let manifest = GenerationManifestRecord {
910 format: "graphforge-generation".into(),
911 format_version: 1,
912 generation_uuid: staged.generation_uuid.hyphenated().to_string(),
913 parent_generation_uuid: Some(staged.parent.generation_uuid().hyphenated().to_string()),
914 transaction_uuid: staged.transaction_uuid.hyphenated().to_string(),
915 capabilities: staged
916 .capabilities
917 .iter()
918 .map(|capability| CapabilityRecord {
919 capability_id: capability.capability_id.clone(),
920 capability_version: capability.capability_version,
921 })
922 .collect(),
923 participants: staged.participants.clone(),
924 };
925 let manifest_bytes = canonical_line(&manifest)?;
926 let manifest_path = staged.generation_root.join(MANIFEST_FILE);
927 let manifest_file = write_new(&manifest_path, &manifest_bytes)?;
928 project_failpoint::hit(
929 "project.after_manifest_write",
930 Some(staged.transaction_uuid),
931 Some(staged.generation_uuid),
932 "DURABLE",
933 false,
934 )?;
935 manifest_file.sync_all().map_err(publication_io)?;
936 drop(manifest_file);
939 project_failpoint::hit(
940 "project.after_manifest_fsync",
941 Some(staged.transaction_uuid),
942 Some(staged.generation_uuid),
943 "DURABLE",
944 false,
945 )?;
946 let manifest_sha256: [u8; 32] = Sha256::digest(&manifest_bytes).into();
947 for participant in &staged.participants {
948 verify_participant_file(
949 &staged
950 .generation_root
951 .join(PARTICIPANTS_DIR)
952 .join(&participant.relative_path),
953 participant,
954 )?;
955 }
956 verify_exact_file(&manifest_path, &manifest_bytes)?;
957 sync_participant_directories(
958 &staged.generation_root.join(PARTICIPANTS_DIR),
959 &staged.participants,
960 )?;
961 sync_directory(&staged.generation_root)?;
962 project_failpoint::hit(
963 "project.after_generation_dir_fsync",
964 Some(staged.transaction_uuid),
965 Some(staged.generation_uuid),
966 "DURABLE",
967 false,
968 )?;
969 if staged.requires_promotion {
970 promote_optimistic_generation(staged)?;
971 } else {
972 sync_directory(&staged.root.join(GENERATIONS_DIR))?;
973 }
974 write_journal(
975 &staged.journal_path(),
976 &staged.journal(JournalPhase::Durable, Some(hex_digest(manifest_sha256))),
977 )?;
978 project_failpoint::hit(
979 "project.after_journal_durable",
980 Some(staged.transaction_uuid),
981 Some(staged.generation_uuid),
982 "DURABLE",
983 false,
984 )?;
985 Ok(manifest_sha256)
986}
987
988fn promote_optimistic_generation(staged: &StagedProjectGeneration) -> Result<(), GfError> {
989 let generations_root = staged.root.join(GENERATIONS_DIR);
990 let destination = generations_root.join(staged.generation_uuid.hyphenated().to_string());
991 if destination.exists() {
992 return Err(project_error(
993 ProjectErrorCode::TransactionConflict,
994 format!(
995 "transaction_uuid={} generation_uuid={} phase=PROMOTE committed=false cause=generation_exists",
996 staged.transaction_uuid.hyphenated(),
997 staged.generation_uuid.hyphenated()
998 ),
999 ));
1000 }
1001 std::fs::rename(&staged.generation_root, &destination).map_err(publication_io)?;
1002 let transaction_attempt_root = staged
1003 .generation_root
1004 .parent()
1005 .expect("machine attempt path has a parent");
1006 sync_directory(transaction_attempt_root)?;
1007 std::fs::remove_dir(transaction_attempt_root).map_err(publication_io)?;
1008 sync_directory(&staged.root.join(ATTEMPTS_DIR))?;
1009 sync_directory(&generations_root)?;
1010 project_failpoint::hit(
1011 "project.after_optimistic_promotion",
1012 Some(staged.transaction_uuid),
1013 Some(staged.generation_uuid),
1014 "DURABLE",
1015 false,
1016 )
1017}
1018
1019fn replace_current(
1020 staged: &StagedProjectGeneration,
1021 manifest_sha256: [u8; 32],
1022) -> Result<(), GfError> {
1023 let current = CurrentRecord {
1024 format: "graphforge-project".into(),
1025 format_version: 1,
1026 generation_uuid: staged.generation_uuid.hyphenated().to_string(),
1027 generation_manifest_sha256: hex_digest(manifest_sha256),
1028 };
1029 let current_bytes = canonical_line(¤t)?;
1030 let current_path = staged.root.join(CURRENT_FILE);
1031 AtomicFile::new(¤t_path, AllowOverwrite)
1032 .write(|file| {
1033 file.write_all(¤t_bytes)?;
1034 failpoint_as_io(
1035 "project.after_current_temp_write",
1036 staged.transaction_uuid,
1037 staged.generation_uuid,
1038 "CURRENT",
1039 false,
1040 )?;
1041 file.sync_all()?;
1042 failpoint_as_io(
1043 "project.after_current_temp_fsync",
1044 staged.transaction_uuid,
1045 staged.generation_uuid,
1046 "CURRENT",
1047 false,
1048 )?;
1049 failpoint_as_io(
1050 "project.before_current_replace",
1051 staged.transaction_uuid,
1052 staged.generation_uuid,
1053 "CURRENT",
1054 false,
1055 )
1056 })
1057 .map_err(|error| {
1058 publication_error_from_parts(
1059 staged.transaction_uuid,
1060 staged.generation_uuid,
1061 "CURRENT",
1062 false,
1063 &error.to_string(),
1064 )
1065 })?;
1066 project_failpoint::hit(
1067 "project.after_current_replace",
1068 Some(staged.transaction_uuid),
1069 Some(staged.generation_uuid),
1070 "CURRENT",
1071 true,
1072 )
1073}
1074
1075fn finish_published_generation(
1076 staged: &StagedProjectGeneration,
1077 manifest_sha256: [u8; 32],
1078) -> Result<(), GfError> {
1079 sync_directory(&staged.root).map_err(|error| {
1082 publication_error_from_parts(
1083 staged.transaction_uuid,
1084 staged.generation_uuid,
1085 "CURRENT",
1086 true,
1087 &error.to_string(),
1088 )
1089 })?;
1090 project_failpoint::hit(
1091 "project.after_root_fsync",
1092 Some(staged.transaction_uuid),
1093 Some(staged.generation_uuid),
1094 "PUBLISHED",
1095 true,
1096 )?;
1097 write_journal(
1098 &staged.journal_path(),
1099 &staged.journal(JournalPhase::Published, Some(hex_digest(manifest_sha256))),
1100 )
1101 .map_err(|error| {
1102 publication_error_from_parts(
1103 staged.transaction_uuid,
1104 staged.generation_uuid,
1105 "PUBLISHED",
1106 true,
1107 &error.to_string(),
1108 )
1109 })?;
1110 project_failpoint::hit(
1111 "project.after_journal_published",
1112 Some(staged.transaction_uuid),
1113 Some(staged.generation_uuid),
1114 "PUBLISHED",
1115 true,
1116 )?;
1117
1118 let resolved = resolve_project_generation(&staged.root).map_err(|error| {
1119 publication_error_from_parts(
1120 staged.transaction_uuid,
1121 staged.generation_uuid,
1122 "PUBLISHED",
1123 true,
1124 &error.to_string(),
1125 )
1126 })?;
1127 if resolved.generation_uuid() != staged.generation_uuid
1128 || resolved.manifest_sha256() != manifest_sha256
1129 {
1130 return Err(publication_error_from_parts(
1131 staged.transaction_uuid,
1132 staged.generation_uuid,
1133 "PUBLISHED",
1134 true,
1135 "published CURRENT did not resolve to exact generation bytes",
1136 ));
1137 }
1138 Ok(())
1139}
1140
1141impl Drop for StagedProjectGeneration {
1142 fn drop(&mut self) {
1143 match &self.publication_lock {
1144 PublicationLock::Exclusive(lock) | PublicationLock::Optimistic(lock) => {
1145 let _ = FileExt::unlock(lock);
1146 }
1147 }
1148 }
1149}
1150
1151#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1152#[serde(rename_all = "SCREAMING_SNAKE_CASE")]
1153pub(crate) enum JournalPhase {
1154 Preparing,
1155 Staged,
1156 Validated,
1157 Durable,
1158 Published,
1159 Aborted,
1160}
1161
1162#[derive(Debug, Serialize, Deserialize)]
1163#[serde(deny_unknown_fields)]
1164pub(crate) struct JournalRecord {
1165 pub(crate) format: String,
1166 pub(crate) format_version: u32,
1167 pub(crate) transaction_uuid: String,
1168 pub(crate) generation_uuid: String,
1169 pub(crate) parent_generation_uuid: Option<String>,
1170 pub(crate) phase: JournalPhase,
1171 pub(crate) request_fingerprint: String,
1172 #[serde(default, skip_serializing_if = "Option::is_none")]
1173 pub(crate) operation_fingerprint: Option<String>,
1174 pub(crate) participant_paths: Vec<String>,
1175 pub(crate) generation_manifest_sha256: Option<String>,
1176 #[serde(default, skip_serializing_if = "Option::is_none")]
1177 pub(crate) revert: Option<RevertJournalExtension>,
1178}
1179
1180impl JournalRecord {
1181 fn new(
1182 request: &ProjectGenerationRequest,
1183 parent: Option<Uuid>,
1184 phase: JournalPhase,
1185 fingerprints: (String, String),
1186 participants: &[StagedParticipant],
1187 generation_manifest_sha256: Option<String>,
1188 revert: Option<RevertJournalExtension>,
1189 ) -> Self {
1190 let (request_fingerprint, operation_fingerprint) = fingerprints;
1191 Self {
1192 format: "graphforge-transaction".into(),
1193 format_version: 1,
1194 transaction_uuid: request.transaction_uuid.hyphenated().to_string(),
1195 generation_uuid: request.generation_uuid.hyphenated().to_string(),
1196 parent_generation_uuid: parent.map(|uuid| uuid.hyphenated().to_string()),
1197 phase,
1198 request_fingerprint,
1199 operation_fingerprint: Some(operation_fingerprint),
1200 participant_paths: participants
1201 .iter()
1202 .map(|participant| participant.relative_path.clone())
1203 .collect(),
1204 generation_manifest_sha256,
1205 revert,
1206 }
1207 }
1208
1209 pub(crate) fn operation_fingerprint(&self) -> &str {
1210 self.operation_fingerprint
1211 .as_deref()
1212 .unwrap_or(&self.request_fingerprint)
1213 }
1214}
1215
1216#[derive(Debug, Serialize)]
1217struct RequestFingerprint<'a> {
1218 format: &'static str,
1219 format_version: u32,
1220 transaction_uuid: String,
1221 generation_uuid: String,
1222 capabilities: &'a [ProjectCapability],
1223 participants: &'a [StagedParticipant],
1224}
1225
1226#[derive(Debug, Serialize)]
1227struct CurrentRecord {
1228 format: String,
1229 format_version: u32,
1230 generation_uuid: String,
1231 generation_manifest_sha256: String,
1232}
1233
1234#[derive(Debug, Serialize)]
1235struct GenerationManifestRecord {
1236 format: String,
1237 format_version: u32,
1238 generation_uuid: String,
1239 parent_generation_uuid: Option<String>,
1240 transaction_uuid: String,
1241 capabilities: Vec<CapabilityRecord>,
1242 participants: Vec<StagedParticipant>,
1243}
1244
1245#[derive(Debug, Serialize)]
1246struct CapabilityRecord {
1247 capability_id: String,
1248 capability_version: u32,
1249}
1250
1251impl Serialize for StagedParticipant {
1252 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
1253 where
1254 S: serde::Serializer,
1255 {
1256 #[derive(Serialize)]
1257 struct Ordered<'a> {
1258 capability_id: &'a str,
1259 capability_version: u32,
1260 record_family_id: &'a str,
1261 record_version: u32,
1262 relative_path: &'a str,
1263 encoding: &'a str,
1264 byte_length: u64,
1265 row_count: u64,
1266 schema_fingerprint: &'a str,
1267 content_sha256: &'a str,
1268 }
1269 Ordered {
1270 capability_id: &self.capability_id,
1271 capability_version: self.capability_version,
1272 record_family_id: &self.record_family_id,
1273 record_version: self.record_version,
1274 relative_path: &self.relative_path,
1275 encoding: &self.encoding,
1276 byte_length: self.byte_length,
1277 row_count: self.row_count,
1278 schema_fingerprint: &self.schema_fingerprint,
1279 content_sha256: &self.content_sha256,
1280 }
1281 .serialize(serializer)
1282 }
1283}
1284
1285fn canonical_supported_root(root: &Path) -> Result<PathBuf, GfError> {
1286 resolve_project_generation(root).map(|resolved| resolved.container_root().to_owned())
1288}
1289
1290fn validate_request(request: &ProjectGenerationRequest) -> Result<(), GfError> {
1291 if request.capabilities.is_empty() {
1292 return Err(project_error(
1293 ProjectErrorCode::PublicationFailed,
1294 "a generation must declare at least one capability",
1295 ));
1296 }
1297 for capability in &request.capabilities {
1298 validate_machine_id(&capability.capability_id)?;
1299 if capability.capability_version == 0 {
1300 return Err(project_error(
1301 ProjectErrorCode::PublicationFailed,
1302 "capability contract versions must be positive",
1303 ));
1304 }
1305 }
1306 for participant in &request.participants {
1307 validate_machine_id(&participant.capability_id)?;
1308 validate_machine_id(&participant.record_family_id)?;
1309 if participant.capability_version == 0 || participant.record_version == 0 {
1310 return Err(project_error(
1311 ProjectErrorCode::PublicationFailed,
1312 "participant contract versions must be positive",
1313 ));
1314 }
1315 }
1316 Ok(())
1317}
1318
1319fn request_metadata(
1320 request: &ProjectGenerationRequest,
1321) -> Result<(Vec<ProjectCapability>, Vec<StagedParticipant>, String), GfError> {
1322 let mut capabilities = request.capabilities.clone();
1323 capabilities.sort_by(|left, right| left.capability_id.cmp(&right.capability_id));
1324 if capabilities
1325 .windows(2)
1326 .any(|pair| pair[0].capability_id == pair[1].capability_id)
1327 {
1328 return Err(project_error(
1329 ProjectErrorCode::PublicationFailed,
1330 "duplicate capability identity",
1331 ));
1332 }
1333 if capabilities
1334 .binary_search_by(|entry| entry.capability_id.as_str().cmp("graph"))
1335 .ok()
1336 .map(|index| capabilities[index].capability_version)
1337 != Some(1)
1338 {
1339 return Err(project_error(
1340 ProjectErrorCode::PublicationFailed,
1341 "every generation must declare graph capability version 1",
1342 ));
1343 }
1344 let mut participants = Vec::with_capacity(request.participants.len());
1345 for participant in &request.participants {
1346 let content_sha256: [u8; 32] = Sha256::digest(&participant.bytes).into();
1347 participants.push(StagedParticipant {
1348 capability_id: participant.capability_id.clone(),
1349 capability_version: participant.capability_version,
1350 record_family_id: participant.record_family_id.clone(),
1351 record_version: participant.record_version,
1352 relative_path: format!(
1353 "{}/{}.{}",
1354 participant.capability_id,
1355 participant.record_family_id,
1356 participant.encoding.extension()
1357 ),
1358 encoding: participant.encoding.extension().into(),
1359 byte_length: u64::try_from(participant.bytes.len()).map_err(|_| {
1360 project_error(
1361 ProjectErrorCode::PublicationFailed,
1362 "participant byte length exceeds u64",
1363 )
1364 })?,
1365 row_count: participant.row_count,
1366 schema_fingerprint: hex_digest(participant.schema_fingerprint),
1367 content_sha256: hex_digest(content_sha256),
1368 });
1369 }
1370 participants.sort_by(|left, right| {
1371 (
1372 &left.capability_id,
1373 &left.record_family_id,
1374 &left.relative_path,
1375 )
1376 .cmp(&(
1377 &right.capability_id,
1378 &right.record_family_id,
1379 &right.relative_path,
1380 ))
1381 });
1382 if participants.windows(2).any(|pair| {
1383 pair[0].capability_id == pair[1].capability_id
1384 && pair[0].record_family_id == pair[1].record_family_id
1385 }) {
1386 return Err(project_error(
1387 ProjectErrorCode::PublicationFailed,
1388 "duplicate participant identity",
1389 ));
1390 }
1391 for participant in &participants {
1392 let capability = capabilities
1393 .binary_search_by(|entry| entry.capability_id.cmp(&participant.capability_id))
1394 .ok()
1395 .map(|index| &capabilities[index])
1396 .ok_or_else(|| {
1397 project_error(
1398 ProjectErrorCode::PublicationFailed,
1399 "participant capability is not declared",
1400 )
1401 })?;
1402 if capability.capability_version != participant.capability_version {
1403 return Err(project_error(
1404 ProjectErrorCode::PublicationFailed,
1405 "participant capability version conflicts with declaration",
1406 ));
1407 }
1408 }
1409 let fingerprint_input = RequestFingerprint {
1410 format: "graphforge-publication-request",
1411 format_version: 1,
1412 transaction_uuid: request.transaction_uuid.hyphenated().to_string(),
1413 generation_uuid: request.generation_uuid.hyphenated().to_string(),
1414 capabilities: &capabilities,
1415 participants: &participants,
1416 };
1417 let bytes = canonical_line(&fingerprint_input)?;
1418 let digest: [u8; 32] = Sha256::digest(bytes).into();
1419 Ok((capabilities, participants, hex_digest(digest)))
1420}
1421
1422fn validate_machine_id(value: &str) -> Result<(), GfError> {
1423 if value.is_empty()
1424 || value.len() > 64
1425 || !value.bytes().all(|byte| {
1426 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'-' | b'_')
1427 })
1428 {
1429 return Err(project_error(
1430 ProjectErrorCode::PublicationFailed,
1431 "machine ID must be 1-64 lowercase ASCII letters, digits, hyphens, or underscores",
1432 ));
1433 }
1434 Ok(())
1435}
1436
1437fn verify_participant_file(path: &Path, expected: &StagedParticipant) -> Result<(), GfError> {
1438 let metadata = std::fs::symlink_metadata(path).map_err(publication_io)?;
1439 if !metadata.is_file() || metadata.file_type().is_symlink() {
1440 return Err(project_error(
1441 ProjectErrorCode::PublicationFailed,
1442 "staged participant is not a regular non-link file",
1443 ));
1444 }
1445 #[cfg(unix)]
1446 {
1447 use std::os::unix::fs::MetadataExt;
1448 if metadata.nlink() != 1 {
1449 return Err(project_error(
1450 ProjectErrorCode::PublicationFailed,
1451 "staged participant is hard-linked",
1452 ));
1453 }
1454 }
1455 if metadata.len() != expected.byte_length {
1456 return Err(project_error(
1457 ProjectErrorCode::PublicationFailed,
1458 "staged participant byte length changed",
1459 ));
1460 }
1461 let mut file = File::open(path).map_err(publication_io)?;
1462 let mut hasher = Sha256::new();
1463 std::io::copy(&mut file, &mut hasher).map_err(publication_io)?;
1464 let actual: [u8; 32] = hasher.finalize().into();
1465 if hex_digest(actual) != expected.content_sha256 {
1466 return Err(project_error(
1467 ProjectErrorCode::PublicationFailed,
1468 "staged participant digest changed",
1469 ));
1470 }
1471 Ok(())
1472}
1473
1474fn sync_participant_directories(
1475 participants_root: &Path,
1476 participants: &[StagedParticipant],
1477) -> Result<(), GfError> {
1478 let mut directories: Vec<PathBuf> = participants
1479 .iter()
1480 .filter_map(|participant| {
1481 participants_root
1482 .join(&participant.relative_path)
1483 .parent()
1484 .map(Path::to_owned)
1485 })
1486 .collect();
1487 directories.sort();
1488 directories.dedup();
1489 directories.sort_by_key(|path| std::cmp::Reverse(path.components().count()));
1490 for directory in directories {
1491 sync_directory(&directory)?;
1492 }
1493 sync_directory(participants_root)
1494}
1495
1496fn write_new(path: &Path, bytes: &[u8]) -> Result<File, GfError> {
1497 let mut file = OpenOptions::new()
1498 .write(true)
1499 .create_new(true)
1500 .open(path)
1501 .map_err(publication_io)?;
1502 file.write_all(bytes).map_err(publication_io)?;
1503 Ok(file)
1504}
1505
1506fn failpoint_as_io(
1507 name: &str,
1508 transaction_uuid: Uuid,
1509 generation_uuid: Uuid,
1510 phase: &str,
1511 committed: bool,
1512) -> std::io::Result<()> {
1513 project_failpoint::hit(
1514 name,
1515 Some(transaction_uuid),
1516 Some(generation_uuid),
1517 phase,
1518 committed,
1519 )
1520 .map_err(|error| std::io::Error::other(error.to_string()))
1521}
1522
1523pub(crate) fn ensure_machine_directory(root: &Path, relative: &Path) -> Result<PathBuf, GfError> {
1524 let mut current = root.to_owned();
1525 for component in relative.components() {
1526 let std::path::Component::Normal(component) = component else {
1527 return Err(project_error(
1528 ProjectErrorCode::ProjectCorrupt,
1529 "machine directory path is not normalized",
1530 ));
1531 };
1532 current.push(component);
1533 match std::fs::symlink_metadata(¤t) {
1534 Ok(metadata) if metadata.is_dir() && !metadata.file_type().is_symlink() => {}
1535 Ok(_) => {
1536 return Err(project_error(
1537 ProjectErrorCode::ProjectCorrupt,
1538 "machine directory is linked or not a directory",
1539 ));
1540 }
1541 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1542 match std::fs::create_dir(¤t) {
1543 Ok(()) => {
1544 sync_directory(
1545 current
1546 .parent()
1547 .expect("machine directory beneath project has a parent"),
1548 )?;
1549 }
1550 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
1551 let metadata =
1552 std::fs::symlink_metadata(¤t).map_err(publication_io)?;
1553 if !metadata.is_dir() || metadata.file_type().is_symlink() {
1554 return Err(project_error(
1555 ProjectErrorCode::ProjectCorrupt,
1556 "concurrently created machine path is linked or not a directory",
1557 ));
1558 }
1559 }
1560 Err(error) => return Err(publication_io(error)),
1561 }
1562 }
1563 Err(error) => return Err(publication_io(error)),
1564 }
1565 }
1566 Ok(current)
1567}
1568
1569pub(crate) fn open_regular_lock(path: &Path) -> Result<File, GfError> {
1570 if let Ok(metadata) = std::fs::symlink_metadata(path)
1571 && (!metadata.is_file() || metadata.file_type().is_symlink())
1572 {
1573 return Err(project_error(
1574 ProjectErrorCode::ProjectCorrupt,
1575 "writer lock is linked or not a regular file",
1576 ));
1577 }
1578 let file = OpenOptions::new()
1579 .read(true)
1580 .write(true)
1581 .create(true)
1582 .truncate(false)
1583 .open(path)
1584 .map_err(publication_io)?;
1585 let metadata = file.metadata().map_err(publication_io)?;
1586 if !metadata.is_file() {
1587 return Err(project_error(
1588 ProjectErrorCode::ProjectCorrupt,
1589 "writer lock is not a regular file",
1590 ));
1591 }
1592 #[cfg(unix)]
1593 {
1594 use std::os::unix::fs::MetadataExt;
1595 if metadata.nlink() != 1 {
1596 return Err(project_error(
1597 ProjectErrorCode::ProjectCorrupt,
1598 "writer lock is hard-linked",
1599 ));
1600 }
1601 }
1602 Ok(file)
1603}
1604
1605fn verify_exact_file(path: &Path, expected: &[u8]) -> Result<(), GfError> {
1606 let mut actual = Vec::new();
1607 File::open(path)
1608 .and_then(|mut file| file.read_to_end(&mut actual))
1609 .map_err(publication_io)?;
1610 if actual != expected {
1611 return Err(project_error(
1612 ProjectErrorCode::PublicationFailed,
1613 "durable file reread did not match staged bytes",
1614 ));
1615 }
1616 Ok(())
1617}
1618
1619pub(crate) fn write_journal(path: &Path, journal: &JournalRecord) -> Result<(), GfError> {
1620 let bytes = canonical_line(journal)?;
1621 AtomicFile::new(path, AllowOverwrite)
1622 .write(|file| file.write_all(&bytes))
1623 .map_err(|error| publication_io(std::io::Error::other(error.to_string())))?;
1624 sync_directory(
1625 path.parent()
1626 .expect("transaction journal always has a parent"),
1627 )
1628}
1629
1630pub(crate) fn read_journal(path: &Path) -> Result<JournalRecord, GfError> {
1631 let metadata = std::fs::symlink_metadata(path).map_err(publication_io)?;
1632 if !metadata.is_file()
1633 || metadata.file_type().is_symlink()
1634 || metadata.len() > MAX_JOURNAL_BYTES
1635 {
1636 return Err(project_error(
1637 ProjectErrorCode::ProjectCorrupt,
1638 "transaction journal is invalid",
1639 ));
1640 }
1641 let bytes = std::fs::read(path).map_err(publication_io)?;
1642 let journal: JournalRecord = serde_json::from_slice(&bytes).map_err(|_| {
1643 project_error(
1644 ProjectErrorCode::ProjectCorrupt,
1645 "transaction journal is not canonical JSON",
1646 )
1647 })?;
1648 if canonical_line(&journal)? != bytes
1649 || journal.format != "graphforge-transaction"
1650 || journal.format_version != 1
1651 || parse_digest(&journal.request_fingerprint).is_none()
1652 || journal
1653 .operation_fingerprint
1654 .as_deref()
1655 .is_some_and(|fingerprint| parse_digest(fingerprint).is_none())
1656 {
1657 return Err(project_error(
1658 ProjectErrorCode::ProjectCorrupt,
1659 "transaction journal is not canonical",
1660 ));
1661 }
1662 Ok(journal)
1663}
1664
1665pub(crate) fn cleanup_atomicwrite_temp(path: &Path) -> Result<bool, GfError> {
1666 let Some(name) = path.file_name().and_then(|name| name.to_str()) else {
1667 return Ok(false);
1668 };
1669 let Some(suffix) = name.strip_prefix(".atomicwrite") else {
1670 return Ok(false);
1671 };
1672 if suffix.len() != 6 || !suffix.bytes().all(|byte| byte.is_ascii_alphanumeric()) {
1673 return Ok(false);
1674 }
1675 let metadata = std::fs::symlink_metadata(path).map_err(publication_io)?;
1676 if !metadata.is_dir() || metadata.file_type().is_symlink() {
1677 return Ok(false);
1678 }
1679 let mut entries = std::fs::read_dir(path).map_err(publication_io)?;
1680 if let Some(entry) = entries.next().transpose().map_err(publication_io)? {
1681 if entries
1682 .next()
1683 .transpose()
1684 .map_err(publication_io)?
1685 .is_some()
1686 || entry.file_name() != "tmpfile.tmp"
1687 {
1688 return Ok(false);
1689 }
1690 let entry_metadata = std::fs::symlink_metadata(entry.path()).map_err(publication_io)?;
1691 if !entry_metadata.is_file() || entry_metadata.file_type().is_symlink() {
1692 return Ok(false);
1693 }
1694 #[cfg(unix)]
1695 {
1696 use std::os::unix::fs::MetadataExt;
1697 if entry_metadata.nlink() != 1 {
1698 return Ok(false);
1699 }
1700 }
1701 std::fs::remove_file(entry.path()).map_err(publication_io)?;
1702 }
1703 std::fs::remove_dir(path).map_err(publication_io)?;
1704 sync_directory(
1705 path.parent()
1706 .expect("atomic-write temporary directory always has a parent"),
1707 )?;
1708 Ok(true)
1709}
1710
1711fn canonical_line<T: Serialize>(value: &T) -> Result<Vec<u8>, GfError> {
1712 let mut bytes = serde_json::to_vec(value)
1713 .map_err(|error| GfError::Storage(format!("failed to encode project record: {error}")))?;
1714 bytes.push(b'\n');
1715 Ok(bytes)
1716}
1717
1718#[cfg(unix)]
1719pub(crate) fn sync_directory(path: &Path) -> Result<(), GfError> {
1720 File::open(path)
1721 .and_then(|directory| directory.sync_all())
1722 .map_err(publication_io)
1723}
1724
1725#[cfg(windows)]
1726pub(crate) fn sync_directory(path: &Path) -> Result<(), GfError> {
1727 use std::os::windows::fs::OpenOptionsExt;
1728
1729 const FILE_FLAG_BACKUP_SEMANTICS: u32 = 0x0200_0000;
1732 OpenOptions::new()
1733 .write(true)
1736 .custom_flags(FILE_FLAG_BACKUP_SEMANTICS)
1737 .open(path)
1738 .and_then(|directory| directory.sync_all())
1739 .map_err(publication_io)
1740}
1741
1742#[cfg(all(not(unix), not(windows)))]
1743pub(crate) fn sync_directory(_path: &Path) -> Result<(), GfError> {
1744 Err(project_error(
1745 ProjectErrorCode::UnsupportedFilesystem,
1746 "directory durability is unsupported on this platform",
1747 ))
1748}
1749
1750fn hex_digest(bytes: [u8; 32]) -> String {
1751 let mut output = String::with_capacity(64);
1752 for byte in bytes {
1753 use std::fmt::Write as _;
1754 write!(&mut output, "{byte:02x}").expect("writing to String cannot fail");
1755 }
1756 output
1757}
1758
1759fn parse_digest(value: &str) -> Option<[u8; 32]> {
1760 if value.len() != 64
1761 || !value
1762 .bytes()
1763 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1764 {
1765 return None;
1766 }
1767 let mut digest = [0u8; 32];
1768 for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
1769 let high = hex_nibble(pair[0])?;
1770 let low = hex_nibble(pair[1])?;
1771 digest[index] = (high << 4) | low;
1772 }
1773 Some(digest)
1774}
1775
1776const fn hex_nibble(byte: u8) -> Option<u8> {
1777 match byte {
1778 b'0'..=b'9' => Some(byte - b'0'),
1779 b'a'..=b'f' => Some(byte - b'a' + 10),
1780 _ => None,
1781 }
1782}
1783
1784fn project_error(code: ProjectErrorCode, message: impl Into<String>) -> GfError {
1785 GfError::Project {
1786 code,
1787 message: message.into(),
1788 }
1789}
1790
1791fn publication_io(error: impl std::fmt::Display) -> GfError {
1792 GfError::Storage(error.to_string())
1793}
1794
1795fn transaction_conflict(request: &ProjectGenerationRequest) -> GfError {
1796 project_error(
1797 ProjectErrorCode::TransactionConflict,
1798 format!(
1799 "transaction_uuid={} generation_uuid={} phase=PREPARING committed=false cause=identity_conflict",
1800 request.transaction_uuid.hyphenated(),
1801 request.generation_uuid.hyphenated()
1802 ),
1803 )
1804}
1805
1806fn publication_error(
1807 request: &ProjectGenerationRequest,
1808 phase: &str,
1809 committed: bool,
1810 cause: &str,
1811) -> GfError {
1812 publication_error_from_parts(
1813 request.transaction_uuid,
1814 request.generation_uuid,
1815 phase,
1816 committed,
1817 cause,
1818 )
1819}
1820
1821fn publication_error_from_parts(
1822 transaction_uuid: Uuid,
1823 generation_uuid: Uuid,
1824 phase: &str,
1825 committed: bool,
1826 cause: &str,
1827) -> GfError {
1828 project_error(
1829 ProjectErrorCode::PublicationFailed,
1830 format!(
1831 "transaction_uuid={} generation_uuid={} phase={phase} committed={committed} cause={}",
1832 transaction_uuid.hyphenated(),
1833 generation_uuid.hyphenated(),
1834 safe_cause(cause)
1835 ),
1836 )
1837}
1838
1839fn safe_cause(cause: &str) -> String {
1840 cause
1841 .chars()
1842 .filter(|character| character.is_ascii_alphanumeric() || "_ -".contains(*character))
1843 .take(96)
1844 .collect()
1845}
1846
1847#[cfg(test)]
1848mod tests {
1849 use super::*;
1850 use crate::open_or_initialize_project;
1851
1852 #[cfg(windows)]
1853 #[test]
1854 fn windows_directory_sync_uses_write_capable_handle() {
1855 let root = tempfile::tempdir().unwrap();
1856
1857 sync_directory(root.path()).unwrap();
1858 }
1859
1860 fn participant(capability: &str, family: &str, value: &[u8]) -> ProjectParticipant {
1861 ProjectParticipant {
1862 capability_id: capability.into(),
1863 capability_version: 1,
1864 record_family_id: family.into(),
1865 record_version: 1,
1866 encoding: ProjectParticipantEncoding::Parquet,
1867 schema_fingerprint: Sha256::digest(format!("{capability}/{family}")).into(),
1868 row_count: 1,
1869 bytes: value.to_vec(),
1870 }
1871 }
1872
1873 fn request(participants: Vec<ProjectParticipant>) -> ProjectGenerationRequest {
1874 let mut capabilities = vec![ProjectCapability {
1875 capability_id: "graph".into(),
1876 capability_version: 1,
1877 }];
1878 for participant in &participants {
1879 if participant.capability_id != "graph"
1880 && !capabilities
1881 .iter()
1882 .any(|entry| entry.capability_id == participant.capability_id)
1883 {
1884 capabilities.push(ProjectCapability {
1885 capability_id: participant.capability_id.clone(),
1886 capability_version: participant.capability_version,
1887 });
1888 }
1889 }
1890 ProjectGenerationRequest {
1891 transaction_uuid: Uuid::now_v7(),
1892 generation_uuid: Uuid::now_v7(),
1893 capabilities,
1894 participants,
1895 }
1896 }
1897
1898 fn project() -> tempfile::TempDir {
1899 let root = tempfile::tempdir().unwrap();
1900 open_or_initialize_project(root.path()).unwrap();
1901 root
1902 }
1903
1904 fn publish(root: &Path, request: ProjectGenerationRequest) -> ProjectPublicationReceipt {
1905 let ProjectStageOutcome::Staged(staged) = stage_project_generation(root, &request).unwrap()
1906 else {
1907 panic!("new request unexpectedly replayed");
1908 };
1909 staged
1910 .validate(|_| Ok(()), |_, _| Ok(()))
1911 .unwrap()
1912 .publish()
1913 .unwrap()
1914 }
1915
1916 fn journal_path(root: &Path, transaction_uuid: Uuid) -> PathBuf {
1917 root.join(TRANSACTIONS_DIR)
1918 .join(format!("{}.json", transaction_uuid.hyphenated()))
1919 }
1920
1921 #[test]
1922 fn publishes_graph_only_and_multi_domain_sets_atomically() {
1923 for participants in [
1924 vec![participant("graph", "nodes", b"graph")],
1925 vec![
1926 participant("graph", "nodes", b"graph"),
1927 participant("provenance", "events", b"provenance"),
1928 ],
1929 vec![
1930 participant("graph", "nodes", b"graph"),
1931 participant("provenance", "events", b"provenance"),
1932 participant("knowledge", "assertions", b"knowledge"),
1933 ],
1934 ] {
1935 let root = project();
1936 let request = request(participants);
1937 let expected = request.generation_uuid;
1938 publish(root.path(), request);
1939 let resolved = resolve_project_generation(root.path()).unwrap();
1940 assert_eq!(resolved.generation_uuid(), expected);
1941 }
1942 }
1943
1944 #[test]
1945 fn validation_failure_leaves_parent_authoritative() {
1946 let root = project();
1947 let parent = resolve_project_generation(root.path())
1948 .unwrap()
1949 .generation_uuid();
1950 let request = request(vec![participant("graph", "nodes", b"new")]);
1951 let ProjectStageOutcome::Staged(staged) =
1952 stage_project_generation(root.path(), &request).unwrap()
1953 else {
1954 panic!("new request unexpectedly replayed");
1955 };
1956
1957 let error = staged
1958 .validate(
1959 |_| Err(GfError::Validation("domain rejected".into())),
1960 |_, _| Ok(()),
1961 )
1962 .err()
1963 .expect("validation must fail");
1964
1965 assert!(matches!(error, GfError::Validation(_)));
1966 assert_eq!(
1967 resolve_project_generation(root.path())
1968 .unwrap()
1969 .generation_uuid(),
1970 parent
1971 );
1972 }
1973
1974 #[test]
1975 fn journal_records_each_deterministic_publication_phase() {
1976 let root = project();
1977 let request = request(vec![participant("graph", "nodes", b"new")]);
1978 let journal_path = journal_path(root.path(), request.transaction_uuid);
1979 let ProjectStageOutcome::Staged(staged) =
1980 stage_project_generation(root.path(), &request).unwrap()
1981 else {
1982 panic!("new request unexpectedly replayed");
1983 };
1984 assert_eq!(
1985 read_journal(&journal_path).unwrap().phase,
1986 JournalPhase::Staged
1987 );
1988
1989 let validated = staged.validate(|_| Ok(()), |_, _| Ok(())).unwrap();
1990 assert_eq!(
1991 read_journal(&journal_path).unwrap().phase,
1992 JournalPhase::Validated
1993 );
1994
1995 validated.publish().unwrap();
1996 let published = read_journal(&journal_path).unwrap();
1997 assert_eq!(published.phase, JournalPhase::Published);
1998 assert!(published.generation_manifest_sha256.is_some());
1999 }
2000
2001 #[test]
2002 fn request_fingerprint_is_independent_of_participant_input_order() {
2003 let mut request = request(vec![
2004 participant("provenance", "events", b"provenance"),
2005 participant("graph", "nodes", b"graph"),
2006 ]);
2007 let (_, first_metadata, first_fingerprint) = request_metadata(&request).unwrap();
2008 request.participants.reverse();
2009 let (_, second_metadata, second_fingerprint) = request_metadata(&request).unwrap();
2010
2011 assert_eq!(first_metadata, second_metadata);
2012 assert_eq!(first_fingerprint, second_fingerprint);
2013 assert_eq!(first_metadata[0].capability_id, "graph");
2014 assert_eq!(first_metadata[1].capability_id, "provenance");
2015 }
2016
2017 #[test]
2018 fn machine_ids_match_the_committed_generation_reader_contract() {
2019 let root = project();
2020 let valid = request(vec![participant("graph", "node-properties", b"properties")]);
2021 let generation_uuid = valid.generation_uuid;
2022 publish(root.path(), valid);
2023 let resolved = resolve_project_generation(root.path()).unwrap();
2024 assert_eq!(resolved.generation_uuid(), generation_uuid);
2025 assert!(
2026 resolved
2027 .participant_path("graph", "node-properties")
2028 .unwrap()
2029 .is_file()
2030 );
2031
2032 let underscore = request(vec![participant(
2033 "graph_data",
2034 "node_properties",
2035 b"properties",
2036 )]);
2037 publish(root.path(), underscore);
2038 let resolved = resolve_project_generation(root.path()).unwrap();
2039 assert!(
2040 resolved
2041 .participant_path("graph_data", "node_properties")
2042 .unwrap()
2043 .is_file()
2044 );
2045
2046 let invalid = request(vec![participant("graph", "NodeProperties", b"properties")]);
2047 let error = stage_project_generation(root.path(), &invalid)
2048 .err()
2049 .expect("reader-incompatible machine ID must be rejected");
2050 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
2051 }
2052
2053 #[test]
2054 fn tampered_staged_bytes_fail_before_publication() {
2055 let root = project();
2056 let parent = resolve_project_generation(root.path())
2057 .unwrap()
2058 .generation_uuid();
2059 let request = request(vec![participant("graph", "nodes", b"original")]);
2060 let ProjectStageOutcome::Staged(staged) =
2061 stage_project_generation(root.path(), &request).unwrap()
2062 else {
2063 panic!("new request unexpectedly replayed");
2064 };
2065 std::fs::write(
2066 staged.generation_root.join(PARTICIPANTS_DIR).join(
2067 staged
2068 .participants
2069 .first()
2070 .expect("participant")
2071 .relative_path
2072 .as_str(),
2073 ),
2074 b"tampered",
2075 )
2076 .unwrap();
2077
2078 let error = staged
2079 .validate(|_| Ok(()), |_, _| Ok(()))
2080 .err()
2081 .expect("tampered bytes must fail validation");
2082
2083 assert_eq!(error.code(), "GF_PUBLICATION_FAILED");
2084 assert_eq!(
2085 resolve_project_generation(root.path())
2086 .unwrap()
2087 .generation_uuid(),
2088 parent
2089 );
2090 }
2091
2092 #[test]
2093 fn identical_published_transaction_is_idempotent() {
2094 let root = project();
2095 let request = request(vec![participant("graph", "nodes", b"same")]);
2096 publish(root.path(), request.clone());
2097
2098 let ProjectStageOutcome::AlreadyPublished(receipt) =
2099 stage_project_generation(root.path(), &request).unwrap()
2100 else {
2101 panic!("identical replay was not recognized");
2102 };
2103 assert!(receipt.idempotent_replay);
2104 }
2105
2106 #[test]
2107 fn historical_published_transaction_remains_idempotent() {
2108 let root = project();
2109 let first = request(vec![participant("graph", "nodes", b"first")]);
2110 publish(root.path(), first.clone());
2111 publish(
2112 root.path(),
2113 request(vec![participant("graph", "nodes", b"second")]),
2114 );
2115
2116 let ProjectStageOutcome::AlreadyPublished(receipt) =
2117 stage_project_generation(root.path(), &first).unwrap()
2118 else {
2119 panic!("historical identical replay was not recognized");
2120 };
2121 assert!(receipt.idempotent_replay);
2122 }
2123
2124 #[test]
2125 fn changed_content_under_same_transaction_conflicts() {
2126 let root = project();
2127 let request = request(vec![participant("graph", "nodes", b"first")]);
2128 publish(root.path(), request.clone());
2129 let mut conflicting = request;
2130 conflicting.participants[0].bytes = b"different".to_vec();
2131
2132 let error = stage_project_generation(root.path(), &conflicting)
2133 .err()
2134 .expect("conflicting replay must fail");
2135
2136 assert_eq!(error.code(), "GF_IDEMPOTENCY_CONFLICT");
2137 }
2138
2139 #[test]
2140 fn reader_pinned_to_parent_survives_publication() {
2141 let root = project();
2142 let parent = resolve_project_generation(root.path()).unwrap();
2143 let request = request(vec![participant("graph", "nodes", b"new")]);
2144 let child = request.generation_uuid;
2145
2146 publish(root.path(), request);
2147
2148 assert_ne!(parent.generation_uuid(), child);
2149 assert_eq!(
2150 resolve_project_generation(root.path())
2151 .unwrap()
2152 .generation_uuid(),
2153 child
2154 );
2155 assert!(parent.generation_root().exists());
2156 }
2157
2158 #[test]
2159 fn optimistic_attempts_stage_concurrently_and_compare_parent_at_commit() {
2160 let root = project();
2161 let first = request(vec![participant("graph", "nodes", b"first")]);
2162 let second = request(vec![participant("graph", "nodes", b"second")]);
2163 let first_operation: [u8; 32] = Sha256::digest(b"logical-first").into();
2164 let second_operation: [u8; 32] = Sha256::digest(b"logical-second").into();
2165
2166 let ProjectStageOutcome::Staged(first_staged) =
2167 stage_project_generation_optimistic(root.path(), &first, first_operation).unwrap()
2168 else {
2169 panic!("first optimistic operation replayed unexpectedly");
2170 };
2171 let ProjectStageOutcome::Staged(second_staged) =
2172 stage_project_generation_optimistic(root.path(), &second, second_operation).unwrap()
2173 else {
2174 panic!("second optimistic operation replayed unexpectedly");
2175 };
2176
2177 let first_validated = first_staged.validate(|_| Ok(()), |_, _| Ok(())).unwrap();
2178 let second_validated = second_staged.validate(|_| Ok(()), |_, _| Ok(())).unwrap();
2179 first_validated.publish().unwrap();
2180 let error = second_validated
2181 .publish()
2182 .expect_err("stale optimistic parent must not publish");
2183 assert_eq!(error.code(), "GF_WRITE_CONFLICT");
2184 assert!(
2185 !root
2186 .path()
2187 .join(GENERATIONS_DIR)
2188 .join(second.generation_uuid.hyphenated().to_string())
2189 .exists()
2190 );
2191
2192 let mut rebased = second.clone();
2193 rebased.participants[0].bytes = b"second-rebased".to_vec();
2194 let ProjectStageOutcome::Staged(rebased) =
2195 stage_project_generation_optimistic(root.path(), &rebased, second_operation).unwrap()
2196 else {
2197 panic!("aborted optimistic operation did not permit a rebase attempt");
2198 };
2199 rebased
2200 .validate(|_| Ok(()), |_, _| Ok(()))
2201 .unwrap()
2202 .publish()
2203 .unwrap();
2204 assert_eq!(
2205 resolve_project_generation(root.path())
2206 .unwrap()
2207 .generation_uuid(),
2208 second.generation_uuid
2209 );
2210 }
2211
2212 #[test]
2213 fn optimistic_promotion_closes_staged_handles_before_directory_rename() {
2214 let root = project();
2215 let request = request(vec![participant("graph", "nodes", b"promoted")]);
2216 let generation_uuid = request.generation_uuid;
2217 let operation: [u8; 32] = Sha256::digest(b"windows-promotion-handles").into();
2218
2219 let ProjectStageOutcome::Staged(staged) =
2220 stage_project_generation_optimistic(root.path(), &request, operation).unwrap()
2221 else {
2222 panic!("optimistic operation replayed unexpectedly");
2223 };
2224 staged
2225 .validate(|_| Ok(()), |_, _| Ok(()))
2226 .unwrap()
2227 .publish()
2228 .unwrap();
2229
2230 assert_eq!(
2231 resolve_project_generation(root.path())
2232 .unwrap()
2233 .generation_uuid(),
2234 generation_uuid
2235 );
2236 }
2237
2238 #[test]
2239 fn optimistic_transaction_identity_has_one_live_attempt() {
2240 let root = project();
2241 let request = request(vec![participant("graph", "nodes", b"attempt")]);
2242 let operation: [u8; 32] = Sha256::digest(b"logical-attempt").into();
2243 let first = stage_project_generation_optimistic(root.path(), &request, operation).unwrap();
2244
2245 let error = stage_project_generation_optimistic(root.path(), &request, operation)
2246 .err()
2247 .expect("duplicate live attempt must be rejected");
2248 assert_eq!(error.code(), "GF_WRITER_BUSY");
2249 drop(first);
2250 }
2251
2252 #[test]
2253 fn recovery_preserves_live_optimistic_attempt_then_cleans_it_after_release() {
2254 let root = project();
2255 let request = request(vec![participant("graph", "nodes", b"live")]);
2256 let operation: [u8; 32] = Sha256::digest(b"logical-live").into();
2257 let (_, _, request_fingerprint) = request_metadata(&request).unwrap();
2258 let generation_path = root
2259 .path()
2260 .join(ATTEMPTS_DIR)
2261 .join(request.transaction_uuid.hyphenated().to_string())
2262 .join(request_fingerprint);
2263 let staged = stage_project_generation_optimistic(root.path(), &request, operation).unwrap();
2264
2265 let live_report = crate::recover_project_transactions(root.path()).unwrap();
2266 assert_eq!(live_report.aborted_journals, 0);
2267 assert_eq!(live_report.removed_generations, 0);
2268 assert!(generation_path.exists());
2269
2270 drop(staged);
2271 let abandoned_report = crate::recover_project_transactions(root.path()).unwrap();
2272 assert_eq!(abandoned_report.aborted_journals, 1);
2273 assert_eq!(abandoned_report.removed_generations, 1);
2274 assert!(!generation_path.exists());
2275 }
2276
2277 #[test]
2278 fn writer_lock_is_nonblocking_and_fail_closed() {
2279 let root = project();
2280 let lock_dir = root.path().join(LOCKS_DIR);
2281 std::fs::create_dir_all(&lock_dir).unwrap();
2282 let lock = OpenOptions::new()
2283 .read(true)
2284 .write(true)
2285 .create(true)
2286 .truncate(false)
2287 .open(lock_dir.join(WRITER_LOCK_FILE))
2288 .unwrap();
2289 FileExt::lock_exclusive(&lock).unwrap();
2290 let parent = resolve_project_generation(root.path())
2291 .unwrap()
2292 .generation_uuid();
2293
2294 let error = stage_project_generation(
2295 root.path(),
2296 &request(vec![participant("graph", "nodes", b"new")]),
2297 )
2298 .err()
2299 .expect("busy writer must fail");
2300
2301 assert_eq!(error.code(), "GF_WRITER_BUSY");
2302 assert_eq!(
2303 resolve_project_generation(root.path())
2304 .unwrap()
2305 .generation_uuid(),
2306 parent
2307 );
2308 }
2309}