1use std::collections::BTreeSet;
8use std::fmt::Write as _;
9use std::fs::{File, OpenOptions};
10use std::io::Read;
11use std::path::{Component, Path, PathBuf};
12use std::sync::Arc;
13
14#[cfg(windows)]
15use std::sync::{Condvar, Mutex, OnceLock};
16
17use atomicwrites::{AllowOverwrite, AtomicFile};
18use fs4::fs_std::FileExt;
19use graphforge_core::{GfError, ProjectErrorCode};
20use serde::{Deserialize, Serialize};
21use sha2::{Digest, Sha256};
22use uuid::Uuid;
23
24use crate::project_failpoint;
25
26pub const FORMAT_FILE: &str = "FORMAT";
28pub const CURRENT_FILE: &str = "CURRENT";
30pub const PROJECT_FORMAT_BYTES: &[u8] = b"graphforge-project/v1\n";
32
33const MAX_FORMAT_BYTES: u64 = 64;
34const MAX_CURRENT_BYTES: u64 = 1_024;
35const MAX_MANIFEST_BYTES: u64 = 4 * 1024 * 1024;
36const MANIFEST_FILE: &str = "manifest.json";
37const LEASE_FILE: &str = "lease.lock";
38const PARTICIPANTS_DIR: &str = "participants";
39const MAX_PARTICIPANT_INVENTORY_ENTRIES: usize = 100_000;
40
41#[derive(Debug, Clone, PartialEq, Eq)]
43pub struct ProjectCapabilityDescriptor {
44 pub capability_id: String,
46 pub capability_version: u32,
48}
49
50#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct ProjectParticipantSnapshot {
53 pub capability_id: String,
55 pub capability_version: u32,
57 pub record_family_id: String,
59 pub record_version: u32,
61 pub encoding: String,
63 pub schema_fingerprint: [u8; 32],
65 pub row_count: u64,
67 pub bytes: Vec<u8>,
69}
70
71#[derive(Debug, Clone, PartialEq, Eq)]
73pub struct ProjectParticipantDescriptor {
74 pub capability_id: String,
76 pub capability_version: u32,
78 pub record_family_id: String,
80 pub record_version: u32,
82 pub encoding: String,
84 pub schema_fingerprint: [u8; 32],
86 pub row_count: u64,
88 pub content_sha256: [u8; 32],
90}
91
92#[derive(Debug, Clone)]
98pub struct ResolvedProjectGeneration {
99 container_root: PathBuf,
100 generation_uuid: Uuid,
101 generation_root: PathBuf,
102 manifest_sha256: [u8; 32],
103 manifest: Arc<GenerationManifest>,
104 _lease_handle: Arc<GenerationLease>,
105}
106
107#[derive(Debug)]
108struct GenerationLease(File);
109
110impl Drop for GenerationLease {
111 fn drop(&mut self) {
112 let _ = FileExt::unlock(&self.0);
113 }
114}
115
116impl ResolvedProjectGeneration {
117 #[must_use]
119 pub fn container_root(&self) -> &Path {
120 &self.container_root
121 }
122
123 #[must_use]
125 pub const fn generation_uuid(&self) -> Uuid {
126 self.generation_uuid
127 }
128
129 #[must_use]
131 pub fn generation_root(&self) -> &Path {
132 &self.generation_root
133 }
134
135 #[must_use]
137 pub fn participants_root(&self) -> PathBuf {
138 self.generation_root.join(PARTICIPANTS_DIR)
139 }
140
141 #[must_use]
143 pub const fn manifest_sha256(&self) -> [u8; 32] {
144 self.manifest_sha256
145 }
146
147 #[must_use]
152 pub fn capabilities(&self) -> Vec<ProjectCapabilityDescriptor> {
153 self.manifest
154 .capabilities
155 .iter()
156 .map(|capability| ProjectCapabilityDescriptor {
157 capability_id: capability.capability_id.clone(),
158 capability_version: capability.capability_version,
159 })
160 .collect()
161 }
162
163 pub fn capability(
168 &self,
169 capability_id: &str,
170 ) -> Result<Option<ProjectCapabilityDescriptor>, GfError> {
171 validate_machine_id(capability_id)?;
172 Ok(self
173 .manifest
174 .capabilities
175 .binary_search_by(|entry| entry.capability_id.as_str().cmp(capability_id))
176 .ok()
177 .map(|index| {
178 let capability = &self.manifest.capabilities[index];
179 ProjectCapabilityDescriptor {
180 capability_id: capability.capability_id.clone(),
181 capability_version: capability.capability_version,
182 }
183 }))
184 }
185
186 pub fn require_capability(
195 &self,
196 capability_id: &str,
197 supported_version: u32,
198 ) -> Result<ProjectCapabilityDescriptor, GfError> {
199 let Some(capability) = self.capability(capability_id)? else {
200 return Err(project_error(
201 ProjectErrorCode::CapabilityDisabled,
202 format!("capability {capability_id} is not enabled"),
203 ));
204 };
205 if capability.capability_version != supported_version {
206 return Err(project_error(
207 ProjectErrorCode::UnsupportedCapabilityVersion,
208 format!(
209 "capability {capability_id}@{} is not supported",
210 capability.capability_version
211 ),
212 ));
213 }
214 Ok(capability)
215 }
216
217 #[must_use]
219 pub fn parent_generation_uuid(&self) -> Option<Uuid> {
220 self.manifest
221 .parent_generation_uuid
222 .as_deref()
223 .and_then(|value| Uuid::parse_str(value).ok())
224 }
225
226 pub fn participant_path(
233 &self,
234 capability_id: &str,
235 record_family_id: &str,
236 ) -> Result<PathBuf, GfError> {
237 validate_machine_id(capability_id)?;
238 validate_machine_id(record_family_id)?;
239 let descriptor = self
240 .manifest
241 .participants
242 .iter()
243 .find(|entry| {
244 entry.capability_id == capability_id && entry.record_family_id == record_family_id
245 })
246 .ok_or_else(|| corrupt("requested participant is absent from generation manifest"))?;
247 let relative = Path::new(&descriptor.relative_path);
248 validate_relative_path(relative)?;
249 let participants_root = self.participants_root();
250 let candidate = participants_root.join(relative);
251 reject_link_components(&participants_root, relative)?;
252 Ok(candidate)
253 }
254
255 pub fn participant_snapshots(&self) -> Result<Vec<ProjectParticipantSnapshot>, GfError> {
265 let mut snapshots = Vec::with_capacity(self.manifest.participants.len());
266 for descriptor in &self.manifest.participants {
267 snapshots.push(self.read_participant_snapshot(descriptor)?);
268 }
269 Ok(snapshots)
270 }
271
272 pub fn participant_descriptors(&self) -> Result<Vec<ProjectParticipantDescriptor>, GfError> {
277 self.manifest
278 .participants
279 .iter()
280 .map(|entry| {
281 Ok(ProjectParticipantDescriptor {
282 capability_id: entry.capability_id.clone(),
283 capability_version: entry.capability_version,
284 record_family_id: entry.record_family_id.clone(),
285 record_version: entry.record_version,
286 encoding: entry.encoding.clone(),
287 schema_fingerprint: parse_sha256(&entry.schema_fingerprint)?,
288 row_count: entry.row_count,
289 content_sha256: parse_sha256(&entry.content_sha256)?,
290 })
291 })
292 .collect()
293 }
294
295 pub fn participant_snapshot(
302 &self,
303 capability_id: &str,
304 record_family_id: &str,
305 ) -> Result<Option<ProjectParticipantSnapshot>, GfError> {
306 validate_machine_id(capability_id)?;
307 validate_machine_id(record_family_id)?;
308 self.manifest
309 .participants
310 .iter()
311 .find(|entry| {
312 entry.capability_id == capability_id && entry.record_family_id == record_family_id
313 })
314 .map(|descriptor| self.read_participant_snapshot(descriptor))
315 .transpose()
316 }
317
318 fn read_participant_snapshot(
319 &self,
320 descriptor: &ParticipantDescriptor,
321 ) -> Result<ProjectParticipantSnapshot, GfError> {
322 let path =
323 self.participant_path(&descriptor.capability_id, &descriptor.record_family_id)?;
324 let bytes = read_exact_participant(&path, descriptor.byte_length)?;
325 let digest: [u8; 32] = Sha256::digest(&bytes).into();
326 if digest != parse_sha256(&descriptor.content_sha256)? {
327 return Err(corrupt(
328 "participant content digest does not match manifest",
329 ));
330 }
331 Ok(ProjectParticipantSnapshot {
332 capability_id: descriptor.capability_id.clone(),
333 capability_version: descriptor.capability_version,
334 record_family_id: descriptor.record_family_id.clone(),
335 record_version: descriptor.record_version,
336 encoding: descriptor.encoding.clone(),
337 schema_fingerprint: parse_sha256(&descriptor.schema_fingerprint)?,
338 row_count: descriptor.row_count,
339 bytes,
340 })
341 }
342
343 pub fn validate_complete_participant_inventory(&self) -> Result<(), GfError> {
354 let expected = self
355 .manifest
356 .participants
357 .iter()
358 .map(|entry| entry.relative_path.clone())
359 .collect::<BTreeSet<_>>();
360 let root = self.participants_root();
361 let mut directories = vec![root.clone()];
362 let mut observed = BTreeSet::new();
363 let mut entry_count = 0_usize;
364 while let Some(directory) = directories.pop() {
365 for entry in std::fs::read_dir(&directory)
366 .map_err(|_| transaction_failed("participant inventory cannot be read"))?
367 {
368 entry_count = entry_count.saturating_add(1);
369 if entry_count > MAX_PARTICIPANT_INVENTORY_ENTRIES {
370 return Err(transaction_failed(
371 "participant inventory exceeds the entry limit",
372 ));
373 }
374 let entry =
375 entry.map_err(|_| transaction_failed("participant entry cannot be read"))?;
376 let file_type = entry
377 .file_type()
378 .map_err(|_| transaction_failed("participant type cannot be read"))?;
379 if file_type.is_symlink() {
380 return Err(transaction_failed("participant inventory contains a link"));
381 }
382 let path = entry.path();
383 if file_type.is_dir() {
384 directories.push(path);
385 } else if file_type.is_file() {
386 let relative = path
387 .strip_prefix(&root)
388 .map_err(|_| transaction_failed("participant path is not contained"))?;
389 validate_relative_path(relative)
390 .map_err(|_| transaction_failed("participant path is invalid"))?;
391 let relative = relative
392 .to_str()
393 .ok_or_else(|| transaction_failed("participant path is not UTF-8"))?
394 .replace(std::path::MAIN_SEPARATOR, "/");
395 observed.insert(relative);
396 } else {
397 return Err(transaction_failed(
398 "participant inventory contains a special file",
399 ));
400 }
401 }
402 }
403 if observed != expected {
404 return Err(transaction_failed(
405 "participant inventory does not match the committed manifest",
406 ));
407 }
408 Ok(())
409 }
410}
411
412#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
413#[serde(deny_unknown_fields)]
414struct CurrentRecord {
415 format: String,
416 format_version: u32,
417 generation_uuid: String,
418 generation_manifest_sha256: String,
419}
420
421#[derive(Debug, Clone, Serialize, Deserialize)]
422#[serde(deny_unknown_fields)]
423struct GenerationManifest {
424 format: String,
425 format_version: u32,
426 generation_uuid: String,
427 parent_generation_uuid: Option<String>,
428 transaction_uuid: String,
429 capabilities: Vec<CapabilityDescriptor>,
430 participants: Vec<ParticipantDescriptor>,
431}
432
433#[derive(Debug, Clone, Serialize, Deserialize)]
434#[serde(deny_unknown_fields)]
435struct CapabilityDescriptor {
436 capability_id: String,
437 capability_version: u32,
438}
439
440#[derive(Debug, Clone, Serialize, Deserialize)]
441#[serde(deny_unknown_fields)]
442struct ParticipantDescriptor {
443 capability_id: String,
444 capability_version: u32,
445 record_family_id: String,
446 record_version: u32,
447 relative_path: String,
448 encoding: String,
449 byte_length: u64,
450 row_count: u64,
451 schema_fingerprint: String,
452 content_sha256: String,
453}
454
455pub fn resolve_project_generation(
465 container_root: impl AsRef<Path>,
466) -> Result<ResolvedProjectGeneration, GfError> {
467 let supplied_root = container_root.as_ref();
468 reject_root_link(supplied_root)?;
469 let root = supplied_root
470 .canonicalize()
471 .map_err(|_| unsupported("project root does not exist or is inaccessible"))?;
472
473 let format_path = root.join(FORMAT_FILE);
474 let format_bytes = read_bounded_regular_file(&format_path, MAX_FORMAT_BYTES)
475 .map_err(|_| unsupported("project root does not contain the supported FORMAT marker"))?;
476 if format_bytes != PROJECT_FORMAT_BYTES {
477 return Err(unsupported("project FORMAT marker is not supported"));
478 }
479
480 let current_path = root.join(CURRENT_FILE);
481 loop {
482 if !current_path.exists() {
483 return Err(project_error(
484 ProjectErrorCode::ProjectUninitialized,
485 "project has no committed generation",
486 ));
487 }
488 let current_bytes = read_bounded_regular_file(¤t_path, MAX_CURRENT_BYTES)
489 .map_err(|_| corrupt("CURRENT is missing, linked, oversized, or unreadable"))?;
490 let current: CurrentRecord = parse_canonical_json_line(¤t_bytes, "CURRENT")?;
491 validate_format(¤t.format, current.format_version)?;
492 let generation_uuid = parse_canonical_uuid(¤t.generation_uuid)?;
493 let expected_manifest_digest = parse_sha256(¤t.generation_manifest_sha256)?;
494
495 let generations_dir = root.join("generations");
496 reject_exact_directory(&generations_dir)?;
497 let selected_dir = generations_dir.join(generation_uuid.hyphenated().to_string());
498 reject_exact_directory(&selected_dir)?;
499 let lease_path = selected_dir.join(LEASE_FILE);
500 let lease = acquire_generation_lease(&lease_path)?;
501
502 let confirmed_current = read_bounded_regular_file(¤t_path, MAX_CURRENT_BYTES)
505 .map_err(|_| corrupt("CURRENT changed to an invalid record during open"))?;
506 if confirmed_current != current_bytes {
507 continue;
508 }
509 reject_exact_directory(&selected_dir)?;
510
511 let manifest_path = selected_dir.join(MANIFEST_FILE);
512 let manifest_bytes = read_bounded_regular_file(&manifest_path, MAX_MANIFEST_BYTES)
513 .map_err(|_| corrupt("selected generation manifest is missing or invalid"))?;
514 let actual_digest: [u8; 32] = Sha256::digest(&manifest_bytes).into();
515 if actual_digest != expected_manifest_digest {
516 return Err(corrupt(
517 "selected generation manifest digest does not match CURRENT",
518 ));
519 }
520 let manifest: GenerationManifest =
521 parse_canonical_json_line(&manifest_bytes, "generation manifest")?;
522 validate_manifest(&manifest, generation_uuid)?;
523 reject_exact_directory(&selected_dir.join(PARTICIPANTS_DIR))?;
524
525 return Ok(ResolvedProjectGeneration {
526 container_root: root,
527 generation_uuid,
528 generation_root: selected_dir,
529 manifest_sha256: actual_digest,
530 manifest: Arc::new(manifest),
531 _lease_handle: Arc::new(lease),
532 });
533 }
534}
535
536pub(crate) fn resolve_verified_generation(
537 container_root: &Path,
538 generation_uuid: Uuid,
539 expected_manifest_digest: [u8; 32],
540) -> Result<ResolvedProjectGeneration, GfError> {
541 reject_root_link(container_root)?;
542 let root = container_root
543 .canonicalize()
544 .map_err(|_| unsupported("project root does not exist or is inaccessible"))?;
545 let format_bytes = read_bounded_regular_file(&root.join(FORMAT_FILE), MAX_FORMAT_BYTES)
546 .map_err(|_| unsupported("project root does not contain the supported FORMAT marker"))?;
547 if format_bytes != PROJECT_FORMAT_BYTES {
548 return Err(unsupported("project FORMAT marker is not supported"));
549 }
550 let generations_dir = root.join("generations");
551 reject_exact_directory(&generations_dir)?;
552 let selected_dir = generations_dir.join(generation_uuid.hyphenated().to_string());
553 reject_exact_directory(&selected_dir)?;
554 let lease = acquire_generation_lease(&selected_dir.join(LEASE_FILE))?;
555 reject_exact_directory(&selected_dir)?;
556 let manifest_bytes =
557 read_bounded_regular_file(&selected_dir.join(MANIFEST_FILE), MAX_MANIFEST_BYTES)
558 .map_err(|_| corrupt("selected generation manifest is missing or invalid"))?;
559 let actual_digest: [u8; 32] = Sha256::digest(&manifest_bytes).into();
560 if actual_digest != expected_manifest_digest {
561 return Err(corrupt(
562 "selected generation manifest digest does not match checkpoint",
563 ));
564 }
565 let manifest: GenerationManifest =
566 parse_canonical_json_line(&manifest_bytes, "generation manifest")?;
567 validate_manifest(&manifest, generation_uuid)?;
568 reject_exact_directory(&selected_dir.join(PARTICIPANTS_DIR))?;
569 Ok(ResolvedProjectGeneration {
570 container_root: root,
571 generation_uuid,
572 generation_root: selected_dir,
573 manifest_sha256: actual_digest,
574 manifest: Arc::new(manifest),
575 _lease_handle: Arc::new(lease),
576 })
577}
578
579pub fn open_or_initialize_project(
591 container_root: impl AsRef<Path>,
592) -> Result<ResolvedProjectGeneration, GfError> {
593 let root = container_root.as_ref();
594 reject_root_link(root)?;
595 let _root_lock = lock_project_root(root)?;
596 let mut entries = std::fs::read_dir(root).map_err(|error| {
597 GfError::Storage(format!("failed to inspect new project root: {error}"))
598 })?;
599 if entries
600 .next()
601 .transpose()
602 .map_err(|error| GfError::Storage(format!("failed to inspect new project root: {error}")))?
603 .is_some()
604 {
605 return match resolve_project_generation(root) {
606 Err(error) if error.code() == "GF_PROJECT_UNINITIALIZED" => {
607 let generation_uuid = validate_resumable_uninitialized_layout(root)?;
608 initialize_empty_generation(root, false, Some(generation_uuid))
609 }
610 result => result,
611 };
612 }
613 initialize_empty_generation(root, true, None)
614}
615
616fn validate_resumable_uninitialized_layout(root: &Path) -> Result<Uuid, GfError> {
617 let mut root_entries = std::fs::read_dir(root)
618 .map_err(|_| unsupported("uninitialized project layout cannot be inspected"))?;
619 let mut root_count = 0_usize;
620 for entry in &mut root_entries {
621 let entry = entry.map_err(|_| unsupported("uninitialized project layout is unreadable"))?;
622 let name = entry.file_name();
623 if crate::project_publication::cleanup_atomicwrite_temp(&entry.path())? {
624 continue;
625 }
626 root_count += 1;
627 if root_count > 2 {
628 return Err(unsupported(
629 "uninitialized project contains unknown root entries",
630 ));
631 }
632 if name != FORMAT_FILE && name != "generations" {
633 return Err(unsupported(
634 "uninitialized project contains unknown root entries",
635 ));
636 }
637 }
638 let generations = root.join("generations");
639 reject_exact_directory(&generations)
640 .map_err(|_| unsupported("uninitialized generations directory is invalid"))?;
641 let mut generation_count = 0_usize;
642 let mut generation_uuids = Vec::new();
643 for entry in std::fs::read_dir(&generations)
644 .map_err(|_| unsupported("uninitialized generations cannot be inspected"))?
645 {
646 generation_count += 1;
647 if generation_count > 16 {
648 return Err(unsupported(
649 "uninitialized project contains too many private generations",
650 ));
651 }
652 let entry =
653 entry.map_err(|_| unsupported("uninitialized generation entry is unreadable"))?;
654 let name = entry
655 .file_name()
656 .to_str()
657 .map(str::to_owned)
658 .ok_or_else(|| unsupported("uninitialized generation UUID is invalid"))?;
659 let generation_uuid = parse_canonical_uuid(&name)
660 .map_err(|_| unsupported("uninitialized generation UUID is invalid"))?;
661 validate_partial_generation(&entry.path())?;
662 generation_uuids.push(generation_uuid);
663 }
664 generation_uuids.sort_unstable();
665 let generation_uuid = generation_uuids
666 .pop()
667 .ok_or_else(|| unsupported("uninitialized project has no private generation"))?;
668 for abandoned_uuid in generation_uuids {
669 remove_partial_generation(&generations, abandoned_uuid)?;
670 }
671 reset_partial_generation(&generations, generation_uuid)?;
672 sync_directory(&generations)?;
673 Ok(generation_uuid)
674}
675
676fn validate_partial_generation(generation: &Path) -> Result<(), GfError> {
677 reject_exact_directory(generation)
678 .map_err(|_| unsupported("uninitialized generation directory is invalid"))?;
679 let mut count = 0_usize;
680 for entry in std::fs::read_dir(generation)
681 .map_err(|_| unsupported("uninitialized generation cannot be inspected"))?
682 {
683 count += 1;
684 if count > 3 {
685 return Err(unsupported(
686 "uninitialized generation contains unknown entries",
687 ));
688 }
689 let entry =
690 entry.map_err(|_| unsupported("uninitialized generation entry is unreadable"))?;
691 let name = entry.file_name();
692 if name == PARTICIPANTS_DIR {
693 reject_exact_directory(&entry.path())
694 .map_err(|_| unsupported("uninitialized participants directory is invalid"))?;
695 validate_partial_workspace_participants(&entry.path())?;
696 } else if name == LEASE_FILE || name == MANIFEST_FILE {
697 let metadata = std::fs::symlink_metadata(entry.path())
698 .map_err(|_| unsupported("uninitialized generation file is unreadable"))?;
699 if !metadata.is_file() || metadata.file_type().is_symlink() {
700 return Err(unsupported("uninitialized generation file is invalid"));
701 }
702 } else {
703 return Err(unsupported(
704 "uninitialized generation contains unknown entries",
705 ));
706 }
707 }
708 Ok(())
709}
710
711fn remove_partial_generation(generations: &Path, generation_uuid: Uuid) -> Result<(), GfError> {
712 let generation = generations.join(generation_uuid.hyphenated().to_string());
713 reset_partial_generation(generations, generation_uuid)?;
714 let participants = generation.join(PARTICIPANTS_DIR);
715 if participants.exists() {
716 std::fs::remove_dir(&participants).map_err(|error| {
717 GfError::Storage(format!(
718 "failed to remove interrupted participants directory: {error}"
719 ))
720 })?;
721 }
722 std::fs::remove_dir(&generation).map_err(|error| {
723 GfError::Storage(format!(
724 "failed to remove interrupted generation directory: {error}"
725 ))
726 })
727}
728
729fn reset_partial_generation(generations: &Path, generation_uuid: Uuid) -> Result<(), GfError> {
730 let generation = generations.join(generation_uuid.hyphenated().to_string());
731 let workspace = generation.join(PARTICIPANTS_DIR).join("workspace");
732 if workspace.exists() {
733 for family in ["configuration.json", "ontology.json"] {
734 let path = workspace.join(family);
735 if path.exists() {
736 std::fs::remove_file(&path).map_err(|error| {
737 GfError::Storage(format!("failed to reset workspace participant: {error}"))
738 })?;
739 }
740 }
741 std::fs::remove_dir(&workspace).map_err(|error| {
742 GfError::Storage(format!("failed to reset workspace directory: {error}"))
743 })?;
744 }
745 for name in [LEASE_FILE, MANIFEST_FILE] {
746 let path = generation.join(name);
747 if path.exists() {
748 std::fs::remove_file(&path).map_err(|error| {
749 GfError::Storage(format!("failed to reset interrupted generation: {error}"))
750 })?;
751 }
752 }
753 sync_directory(&generation)
754}
755
756fn validate_partial_workspace_participants(participants: &Path) -> Result<(), GfError> {
757 let entries = std::fs::read_dir(participants)
758 .map_err(|_| unsupported("uninitialized participants cannot be inspected"))?
759 .collect::<Result<Vec<_>, _>>()
760 .map_err(|_| unsupported("uninitialized participant entry is unreadable"))?;
761 if entries.is_empty() {
762 return Ok(());
763 }
764 if entries.len() != 1 || entries[0].file_name() != "workspace" {
765 return Err(unsupported(
766 "uninitialized participants contain unknown entries",
767 ));
768 }
769 reject_exact_directory(&entries[0].path())
770 .map_err(|_| unsupported("uninitialized workspace directory is invalid"))?;
771 let mut names = std::fs::read_dir(entries[0].path())
772 .map_err(|_| unsupported("uninitialized workspace cannot be inspected"))?
773 .map(|entry| {
774 let entry =
775 entry.map_err(|_| unsupported("uninitialized workspace entry is unreadable"))?;
776 let metadata = std::fs::symlink_metadata(entry.path())
777 .map_err(|_| unsupported("uninitialized workspace entry is unreadable"))?;
778 if !metadata.is_file() || metadata.file_type().is_symlink() {
779 return Err(unsupported("uninitialized workspace entry is invalid"));
780 }
781 entry
782 .file_name()
783 .into_string()
784 .map_err(|_| unsupported("uninitialized workspace entry name is invalid"))
785 })
786 .collect::<Result<Vec<_>, _>>()?;
787 names.sort();
788 if names != ["configuration.json", "ontology.json"] && names != ["configuration.json"] {
789 return Err(unsupported(
790 "uninitialized workspace contains unknown entries",
791 ));
792 }
793 Ok(())
794}
795
796fn initialize_empty_generation(
797 root: &Path,
798 write_format: bool,
799 generation_uuid: Option<Uuid>,
800) -> Result<ResolvedProjectGeneration, GfError> {
801 let generation_uuid = generation_uuid.unwrap_or_else(Uuid::now_v7);
802 let transaction_uuid = Uuid::now_v7();
803 let generation_root = root
804 .join("generations")
805 .join(generation_uuid.hyphenated().to_string());
806 let participants_root = generation_root.join(PARTICIPANTS_DIR);
807 std::fs::create_dir_all(&participants_root)
808 .map_err(|error| GfError::Storage(format!("failed to create project layout: {error}")))?;
809 if write_format {
810 let format_path = root.join(FORMAT_FILE);
811 write_new_synced(&format_path, PROJECT_FORMAT_BYTES, "project format")?;
812 project_failpoint::hit(
813 "project.after_format_fsync",
814 Some(transaction_uuid),
815 Some(generation_uuid),
816 "FORMAT",
817 false,
818 )?;
819 }
820 let participant_descriptors = install_empty_workspace_participants(&participants_root)?;
821 sync_directory(&participants_root)?;
822 sync_directory(&generation_root)?;
823 sync_directory(&root.join("generations"))?;
824 sync_directory(root)?;
825 if write_format {
826 project_failpoint::hit(
827 "project.after_container_dir_fsync",
828 Some(transaction_uuid),
829 Some(generation_uuid),
830 "CONTAINER",
831 false,
832 )?;
833 }
834 write_new_synced(&generation_root.join(LEASE_FILE), &[], "generation lease")?;
835 let manifest = GenerationManifest {
836 format: "graphforge-generation".into(),
837 format_version: 1,
838 generation_uuid: generation_uuid.hyphenated().to_string(),
839 parent_generation_uuid: None,
840 transaction_uuid: transaction_uuid.hyphenated().to_string(),
841 capabilities: vec![
842 CapabilityDescriptor {
843 capability_id: "graph".into(),
844 capability_version: 1,
845 },
846 CapabilityDescriptor {
847 capability_id: "workspace".into(),
848 capability_version: 1,
849 },
850 ],
851 participants: participant_descriptors,
852 };
853 let mut manifest_bytes = serde_json::to_vec(&manifest)
854 .map_err(|error| GfError::Storage(format!("failed to encode generation: {error}")))?;
855 manifest_bytes.push(b'\n');
856 write_new_synced(
857 &generation_root.join(MANIFEST_FILE),
858 &manifest_bytes,
859 "generation",
860 )?;
861 sync_directory(&generation_root)?;
862 sync_directory(&root.join("generations"))?;
863 let digest: [u8; 32] = Sha256::digest(&manifest_bytes).into();
864 let current = CurrentRecord {
865 format: "graphforge-project".into(),
866 format_version: 1,
867 generation_uuid: generation_uuid.hyphenated().to_string(),
868 generation_manifest_sha256: sha256_hex(digest),
869 };
870 let mut current_bytes = serde_json::to_vec(¤t)
871 .map_err(|error| GfError::Storage(format!("failed to encode CURRENT: {error}")))?;
872 current_bytes.push(b'\n');
873 AtomicFile::new(root.join(CURRENT_FILE), AllowOverwrite)
874 .write(|file| {
875 use std::io::Write as _;
876
877 file.write_all(¤t_bytes)?;
878 file.sync_all()
879 })
880 .map_err(|error| GfError::Storage(format!("failed to write CURRENT: {error}")))?;
881 sync_directory(root)?;
882 resolve_project_generation(root)
883}
884
885fn install_empty_workspace_participants(
886 participants_root: &Path,
887) -> Result<Vec<ParticipantDescriptor>, GfError> {
888 let workspace_participants = crate::workspace_participants::empty_workspace_participants()?;
889 let workspace_root = participants_root.join("workspace");
890 std::fs::create_dir_all(&workspace_root)
891 .map_err(|error| GfError::Storage(format!("failed to create project layout: {error}")))?;
892 let mut descriptors = Vec::with_capacity(workspace_participants.len());
893 for participant in workspace_participants {
894 let relative_path = format!(
895 "{}/{}.json",
896 participant.capability_id, participant.record_family_id
897 );
898 write_new_synced(
899 &participants_root.join(&relative_path),
900 &participant.bytes,
901 "workspace participant",
902 )?;
903 descriptors.push(ParticipantDescriptor {
904 capability_id: participant.capability_id,
905 capability_version: participant.capability_version,
906 record_family_id: participant.record_family_id,
907 record_version: participant.record_version,
908 relative_path,
909 encoding: "json".into(),
910 byte_length: participant.bytes.len() as u64,
911 row_count: participant.row_count,
912 schema_fingerprint: sha256_hex(participant.schema_fingerprint),
913 content_sha256: sha256_hex(Sha256::digest(&participant.bytes).into()),
914 });
915 }
916 descriptors.sort_by(|left, right| {
917 (
918 &left.capability_id,
919 &left.record_family_id,
920 &left.relative_path,
921 )
922 .cmp(&(
923 &right.capability_id,
924 &right.record_family_id,
925 &right.relative_path,
926 ))
927 });
928 sync_directory(&workspace_root)?;
929 Ok(descriptors)
930}
931
932fn write_new_synced(path: &Path, bytes: &[u8], name: &str) -> Result<(), GfError> {
933 use std::io::Write as _;
934
935 let mut file = OpenOptions::new()
936 .write(true)
937 .create_new(true)
938 .open(path)
939 .map_err(|error| GfError::Storage(format!("failed to create {name}: {error}")))?;
940 file.write_all(bytes)
941 .and_then(|()| file.sync_all())
942 .map_err(|error| GfError::Storage(format!("failed to write {name}: {error}")))
943}
944
945#[cfg(unix)]
946fn lock_project_root(path: &Path) -> Result<File, GfError> {
947 let root = File::open(path)
948 .map_err(|error| GfError::Storage(format!("failed to open project root: {error}")))?;
949 FileExt::lock_exclusive(&root)
950 .map_err(|error| GfError::Storage(format!("failed to lock project root: {error}")))?;
951 Ok(root)
952}
953
954#[cfg(windows)]
955static LOCAL_PROJECT_ROOT_LOCKS: OnceLock<(Mutex<BTreeSet<String>>, Condvar)> = OnceLock::new();
956
957#[cfg(windows)]
958struct LocalProjectRootLock {
959 name: String,
960}
961
962#[cfg(windows)]
963impl Drop for LocalProjectRootLock {
964 fn drop(&mut self) {
965 let (active, available) =
966 LOCAL_PROJECT_ROOT_LOCKS.get_or_init(|| (Mutex::new(BTreeSet::new()), Condvar::new()));
967 let mut active = active
968 .lock()
969 .unwrap_or_else(|poisoned| poisoned.into_inner());
970 active.remove(&self.name);
971 available.notify_all();
972 }
973}
974
975#[cfg(windows)]
976struct WindowsProjectRootLock {
977 kernel: Option<named_lock::NamedLockGuard>,
978 _local: LocalProjectRootLock,
979}
980
981#[cfg(windows)]
982impl Drop for WindowsProjectRootLock {
983 fn drop(&mut self) {
984 drop(self.kernel.take());
987 }
988}
989
990#[cfg(windows)]
991fn acquire_local_project_root_lock(name: &str) -> LocalProjectRootLock {
992 let (active, available) =
993 LOCAL_PROJECT_ROOT_LOCKS.get_or_init(|| (Mutex::new(BTreeSet::new()), Condvar::new()));
994 let mut active = active
995 .lock()
996 .unwrap_or_else(|poisoned| poisoned.into_inner());
997 while active.contains(name) {
998 active = available
999 .wait(active)
1000 .unwrap_or_else(|poisoned| poisoned.into_inner());
1001 }
1002 active.insert(name.to_owned());
1003 LocalProjectRootLock {
1004 name: name.to_owned(),
1005 }
1006}
1007
1008#[cfg(windows)]
1009fn lock_project_root(path: &Path) -> Result<WindowsProjectRootLock, GfError> {
1010 let canonical = std::fs::canonicalize(path)
1011 .map_err(|error| GfError::Storage(format!("failed to lock project root: {error}")))?;
1012 let identity = canonical.as_os_str().to_string_lossy().to_lowercase();
1013 let digest: [u8; 32] = Sha256::digest(identity.as_bytes()).into();
1014 let name = format!("GraphForge.ProjectRoot.{}", sha256_hex(digest));
1015 let local = acquire_local_project_root_lock(&name);
1016 let lock = named_lock::NamedLock::create(&name)
1017 .map_err(|error| GfError::Storage(format!("failed to lock project root: {error}")))?;
1018
1019 match lock.try_lock() {
1020 Ok(kernel) => Ok(WindowsProjectRootLock {
1021 kernel: Some(kernel),
1022 _local: local,
1023 }),
1024 Err(named_lock::Error::WouldBlock) => Err(project_error(
1025 ProjectErrorCode::WriterBusy,
1026 "phase=PROJECT_ROOT_LOCK committed=false cause=busy",
1027 )),
1028 Err(error) => Err(GfError::Storage(format!(
1029 "failed to lock project root: {error}"
1030 ))),
1031 }
1032}
1033
1034#[cfg(all(not(unix), not(windows)))]
1035fn lock_project_root(_path: &Path) -> Result<File, GfError> {
1036 Err(GfError::Storage(
1037 "project-root locks are unsupported on this platform".into(),
1038 ))
1039}
1040
1041fn sync_directory(path: &Path) -> Result<(), GfError> {
1042 crate::project_publication::sync_directory(path).map_err(|error| match error {
1043 GfError::Storage(message) => {
1044 GfError::Storage(format!("failed to sync project directory: {message}"))
1045 }
1046 other => other,
1047 })
1048}
1049
1050fn validate_format(format: &str, version: u32) -> Result<(), GfError> {
1051 if format != "graphforge-project" || version != 1 {
1052 return Err(unsupported("project format version is not supported"));
1053 }
1054 Ok(())
1055}
1056
1057fn validate_manifest(manifest: &GenerationManifest, expected: Uuid) -> Result<(), GfError> {
1058 if manifest.format != "graphforge-generation" || manifest.format_version != 1 {
1059 return Err(corrupt("generation manifest format is invalid"));
1060 }
1061 if parse_canonical_uuid(&manifest.generation_uuid)? != expected {
1062 return Err(corrupt("generation manifest UUID does not match CURRENT"));
1063 }
1064 parse_canonical_uuid(&manifest.transaction_uuid)?;
1065 if let Some(parent) = &manifest.parent_generation_uuid {
1066 parse_canonical_uuid(parent)?;
1067 }
1068 if !manifest
1069 .capabilities
1070 .windows(2)
1071 .all(|pair| pair[0].capability_id < pair[1].capability_id)
1072 {
1073 return Err(corrupt(
1074 "generation capabilities are not in canonical order",
1075 ));
1076 }
1077 for capability in &manifest.capabilities {
1078 validate_machine_id(&capability.capability_id)?;
1079 if capability.capability_version == 0 {
1080 return Err(corrupt("capability version must be positive"));
1081 }
1082 }
1083 if manifest
1084 .capabilities
1085 .binary_search_by(|entry| entry.capability_id.as_str().cmp("graph"))
1086 .ok()
1087 .map(|index| manifest.capabilities[index].capability_version)
1088 != Some(1)
1089 {
1090 return Err(corrupt(
1091 "generation manifest must declare graph capability version 1",
1092 ));
1093 }
1094 if !manifest.participants.windows(2).all(|pair| {
1095 (
1096 &pair[0].capability_id,
1097 &pair[0].record_family_id,
1098 &pair[0].relative_path,
1099 ) < (
1100 &pair[1].capability_id,
1101 &pair[1].record_family_id,
1102 &pair[1].relative_path,
1103 )
1104 }) {
1105 return Err(corrupt(
1106 "generation participants are not in canonical order",
1107 ));
1108 }
1109 for participant in &manifest.participants {
1110 validate_machine_id(&participant.capability_id)?;
1111 validate_machine_id(&participant.record_family_id)?;
1112 if participant.capability_version == 0 || participant.record_version == 0 {
1113 return Err(corrupt("participant contract versions must be positive"));
1114 }
1115 let capability = manifest
1116 .capabilities
1117 .binary_search_by(|entry| entry.capability_id.cmp(&participant.capability_id))
1118 .ok()
1119 .map(|index| &manifest.capabilities[index])
1120 .ok_or_else(|| corrupt("participant capability is not declared"))?;
1121 if capability.capability_version != participant.capability_version {
1122 return Err(corrupt(
1123 "participant capability version does not match its declaration",
1124 ));
1125 }
1126 validate_relative_path(Path::new(&participant.relative_path))?;
1127 parse_sha256(&participant.content_sha256)?;
1128 parse_sha256(&participant.schema_fingerprint)?;
1129 }
1130 Ok(())
1131}
1132
1133pub(crate) fn validated_generation_parent(
1134 root: &Path,
1135 generation_uuid: Uuid,
1136) -> Result<Option<Uuid>, GfError> {
1137 validated_generation_metadata(root, generation_uuid).map(|(parent, _)| parent)
1138}
1139
1140pub(crate) fn validated_generation_manifest_sha256(
1141 root: &Path,
1142 generation_uuid: Uuid,
1143) -> Result<[u8; 32], GfError> {
1144 validated_generation_metadata(root, generation_uuid).map(|(_, digest)| digest)
1145}
1146
1147fn validated_generation_metadata(
1148 root: &Path,
1149 generation_uuid: Uuid,
1150) -> Result<(Option<Uuid>, [u8; 32]), GfError> {
1151 let generation_root = root
1152 .join("generations")
1153 .join(generation_uuid.hyphenated().to_string());
1154 reject_exact_directory(&generation_root)?;
1155 let manifest_bytes =
1156 read_bounded_regular_file(&generation_root.join(MANIFEST_FILE), MAX_MANIFEST_BYTES)
1157 .map_err(|_| corrupt("retained ancestor manifest is missing or invalid"))?;
1158 let manifest: GenerationManifest =
1159 parse_canonical_json_line(&manifest_bytes, "retained ancestor manifest")?;
1160 validate_manifest(&manifest, generation_uuid)?;
1161 let parent = manifest
1162 .parent_generation_uuid
1163 .as_deref()
1164 .map(parse_canonical_uuid)
1165 .transpose()?;
1166 Ok((parent, Sha256::digest(&manifest_bytes).into()))
1167}
1168
1169fn parse_canonical_json_line<T>(bytes: &[u8], name: &str) -> Result<T, GfError>
1170where
1171 T: serde::de::DeserializeOwned + Serialize,
1172{
1173 if !bytes.ends_with(b"\n") || bytes[..bytes.len().saturating_sub(1)].contains(&b'\n') {
1174 return Err(corrupt(format!("{name} is not one canonical JSON line")));
1175 }
1176 let parsed: T =
1177 serde_json::from_slice(bytes).map_err(|_| corrupt(format!("{name} is invalid JSON")))?;
1178 let mut canonical =
1179 serde_json::to_vec(&parsed).map_err(|_| corrupt(format!("{name} cannot be encoded")))?;
1180 canonical.push(b'\n');
1181 if canonical != bytes {
1182 return Err(corrupt(format!("{name} is not canonical JSON")));
1183 }
1184 Ok(parsed)
1185}
1186
1187fn parse_canonical_uuid(value: &str) -> Result<Uuid, GfError> {
1188 let parsed = Uuid::parse_str(value).map_err(|_| corrupt("UUID is invalid"))?;
1189 if parsed.hyphenated().to_string() != value {
1190 return Err(corrupt("UUID is not lowercase canonical hyphenated form"));
1191 }
1192 Ok(parsed)
1193}
1194
1195fn parse_sha256(value: &str) -> Result<[u8; 32], GfError> {
1196 if value.len() != 64
1197 || !value
1198 .bytes()
1199 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1200 {
1201 return Err(corrupt("SHA-256 is not 64 lowercase hexadecimal bytes"));
1202 }
1203 let mut output = [0_u8; 32];
1204 for (index, pair) in value.as_bytes().chunks_exact(2).enumerate() {
1205 output[index] = (hex_nibble(pair[0]) << 4) | hex_nibble(pair[1]);
1206 }
1207 Ok(output)
1208}
1209
1210const fn hex_nibble(byte: u8) -> u8 {
1211 match byte {
1212 b'0'..=b'9' => byte - b'0',
1213 b'a'..=b'f' => byte - b'a' + 10,
1214 _ => 0,
1215 }
1216}
1217
1218fn sha256_hex(bytes: [u8; 32]) -> String {
1219 bytes
1220 .iter()
1221 .fold(String::with_capacity(64), |mut output, byte| {
1222 write!(output, "{byte:02x}").expect("writing to String is infallible");
1223 output
1224 })
1225}
1226
1227fn validate_machine_id(value: &str) -> Result<(), GfError> {
1228 if value.is_empty()
1229 || value.len() > 64
1230 || !value.bytes().all(|byte| {
1231 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'-' | b'_')
1232 })
1233 {
1234 return Err(corrupt("machine identifier is not lowercase ASCII"));
1235 }
1236 Ok(())
1237}
1238
1239fn validate_relative_path(path: &Path) -> Result<(), GfError> {
1240 if path.as_os_str().is_empty()
1241 || path.is_absolute()
1242 || path.components().any(|component| {
1243 !matches!(component, Component::Normal(_)) || component.as_os_str().to_str().is_none()
1244 })
1245 {
1246 return Err(corrupt(
1247 "participant path is not a normalized relative path",
1248 ));
1249 }
1250 Ok(())
1251}
1252
1253fn reject_root_link(path: &Path) -> Result<(), GfError> {
1254 let metadata = std::fs::symlink_metadata(path)
1255 .map_err(|_| unsupported("project root does not exist or is inaccessible"))?;
1256 if metadata.file_type().is_symlink() || !metadata.is_dir() {
1257 return Err(unsupported(
1258 "project root must be a real local directory, not a link",
1259 ));
1260 }
1261 Ok(())
1262}
1263
1264fn reject_exact_directory(path: &Path) -> Result<(), GfError> {
1265 let metadata =
1266 std::fs::symlink_metadata(path).map_err(|_| corrupt("required directory is missing"))?;
1267 if metadata.file_type().is_symlink() || !metadata.is_dir() {
1268 return Err(corrupt("required directory is linked or not a directory"));
1269 }
1270 Ok(())
1271}
1272
1273fn reject_link_components(root: &Path, relative: &Path) -> Result<(), GfError> {
1274 let mut candidate = root.to_path_buf();
1275 for component in relative.components() {
1276 let Component::Normal(part) = component else {
1277 return Err(corrupt("participant path escapes its generation"));
1278 };
1279 candidate.push(part);
1280 let metadata = std::fs::symlink_metadata(&candidate)
1281 .map_err(|_| corrupt("participant path is missing"))?;
1282 if metadata.file_type().is_symlink() {
1283 return Err(corrupt("participant path contains a link"));
1284 }
1285 }
1286 Ok(())
1287}
1288
1289fn open_regular_file(path: &Path) -> Result<File, std::io::Error> {
1290 let metadata = std::fs::symlink_metadata(path)?;
1291 if metadata.file_type().is_symlink() || !metadata.is_file() {
1292 return Err(std::io::Error::other("not a regular non-link file"));
1293 }
1294 #[cfg(unix)]
1295 {
1296 use std::os::unix::fs::MetadataExt;
1297 if metadata.nlink() != 1 {
1298 return Err(std::io::Error::other("hard-linked project file"));
1299 }
1300 }
1301 OpenOptions::new().read(true).open(path)
1302}
1303
1304fn acquire_generation_lease(path: &Path) -> Result<GenerationLease, GfError> {
1305 let handle = open_regular_file(path)
1306 .map_err(|_| corrupt("selected generation lease is missing or invalid"))?;
1307 FileExt::lock_shared(&handle)
1308 .map_err(|_| corrupt("selected generation lease cannot be acquired"))?;
1309 Ok(GenerationLease(handle))
1310}
1311
1312fn read_bounded_regular_file(path: &Path, maximum: u64) -> Result<Vec<u8>, std::io::Error> {
1313 let file = open_regular_file(path)?;
1314 let metadata = file.metadata()?;
1315 if metadata.len() > maximum {
1316 return Err(std::io::Error::other("file exceeds bounded read limit"));
1317 }
1318 let capacity = usize::try_from(metadata.len())
1319 .map_err(|_| std::io::Error::other("file length does not fit address space"))?;
1320 let mut bytes = Vec::with_capacity(capacity);
1321 file.take(maximum + 1).read_to_end(&mut bytes)?;
1322 if bytes.len() as u64 > maximum {
1323 return Err(std::io::Error::other("file exceeds bounded read limit"));
1324 }
1325 Ok(bytes)
1326}
1327
1328fn read_exact_participant(path: &Path, expected_length: u64) -> Result<Vec<u8>, GfError> {
1329 let file = open_regular_file(path)
1330 .map_err(|_| corrupt("participant is missing, linked, or unreadable"))?;
1331 let metadata = file
1332 .metadata()
1333 .map_err(|_| corrupt("participant metadata is unreadable"))?;
1334 if metadata.len() != expected_length {
1335 return Err(corrupt("participant byte length does not match manifest"));
1336 }
1337 let capacity = usize::try_from(expected_length)
1338 .map_err(|_| corrupt("participant byte length exceeds address space"))?;
1339 let mut bytes = Vec::with_capacity(capacity);
1340 file.take(expected_length.saturating_add(1))
1341 .read_to_end(&mut bytes)
1342 .map_err(|_| corrupt("participant cannot be read"))?;
1343 if u64::try_from(bytes.len()).ok() != Some(expected_length) {
1344 return Err(corrupt("participant byte length changed while reading"));
1345 }
1346 Ok(bytes)
1347}
1348
1349fn unsupported(message: impl Into<String>) -> GfError {
1350 project_error(ProjectErrorCode::UnsupportedProjectFormat, message)
1351}
1352
1353fn corrupt(message: impl Into<String>) -> GfError {
1354 project_error(ProjectErrorCode::ProjectCorrupt, message)
1355}
1356
1357fn transaction_failed(message: impl Into<String>) -> GfError {
1358 project_error(ProjectErrorCode::TransactionFailed, message)
1359}
1360
1361fn project_error(code: ProjectErrorCode, message: impl Into<String>) -> GfError {
1362 GfError::Project {
1363 code,
1364 message: message.into(),
1365 }
1366}
1367
1368#[cfg(test)]
1369mod tests {
1370 use super::*;
1371 use std::fs;
1372 use std::sync::{Arc, Barrier};
1373
1374 #[cfg(windows)]
1375 const ROOT_LOCK_CONTENDER: &str =
1376 "project_generation::tests::windows_project_root_lock_contender";
1377 #[cfg(windows)]
1378 const ROOT_LOCK_ABANDONER: &str =
1379 "project_generation::tests::windows_project_root_lock_abandoner";
1380
1381 #[cfg(windows)]
1382 fn wait_for_lock_subprocess(mut child: std::process::Child, proof: &str) {
1383 use std::time::Duration;
1384 use wait_timeout::ChildExt as _;
1385
1386 let status = child
1387 .wait_timeout(Duration::from_secs(30))
1388 .unwrap_or_else(|error| panic!("failed to wait for {proof}: {error}"));
1389 let Some(status) = status else {
1390 child.kill().expect("failed to kill timed-out subprocess");
1391 child.wait().expect("failed to reap timed-out subprocess");
1392 panic!("{proof} timed out after 30 seconds");
1393 };
1394 assert!(status.success(), "{proof} failed with {status}");
1395 }
1396
1397 fn canonical_line<T: Serialize>(value: &T) -> Vec<u8> {
1398 let mut bytes = serde_json::to_vec(value).unwrap();
1399 bytes.push(b'\n');
1400 bytes
1401 }
1402
1403 fn install_generation(root: &Path, generation_uuid: Uuid) -> [u8; 32] {
1404 let generation_root = root
1405 .join("generations")
1406 .join(generation_uuid.hyphenated().to_string());
1407 fs::create_dir_all(generation_root.join(PARTICIPANTS_DIR)).unwrap();
1408 fs::write(generation_root.join(LEASE_FILE), []).unwrap();
1409 let manifest = GenerationManifest {
1410 format: "graphforge-generation".into(),
1411 format_version: 1,
1412 generation_uuid: generation_uuid.hyphenated().to_string(),
1413 parent_generation_uuid: None,
1414 transaction_uuid: Uuid::now_v7().hyphenated().to_string(),
1415 capabilities: vec![CapabilityDescriptor {
1416 capability_id: "graph".into(),
1417 capability_version: 1,
1418 }],
1419 participants: vec![],
1420 };
1421 let bytes = canonical_line(&manifest);
1422 fs::write(generation_root.join(MANIFEST_FILE), &bytes).unwrap();
1423 Sha256::digest(&bytes).into()
1424 }
1425
1426 fn write_current(root: &Path, generation_uuid: Uuid, digest: [u8; 32]) {
1427 let current = CurrentRecord {
1428 format: "graphforge-project".into(),
1429 format_version: 1,
1430 generation_uuid: generation_uuid.hyphenated().to_string(),
1431 generation_manifest_sha256: sha256_hex(digest),
1432 };
1433 fs::write(root.join(CURRENT_FILE), canonical_line(¤t)).unwrap();
1434 }
1435
1436 fn project() -> (tempfile::TempDir, Uuid) {
1437 let root = tempfile::tempdir().unwrap();
1438 fs::write(root.path().join(FORMAT_FILE), PROJECT_FORMAT_BYTES).unwrap();
1439 fs::create_dir(root.path().join("generations")).unwrap();
1440 let generation_uuid = Uuid::now_v7();
1441 let digest = install_generation(root.path(), generation_uuid);
1442 write_current(root.path(), generation_uuid, digest);
1443 (root, generation_uuid)
1444 }
1445
1446 fn open_generation_lease(root: &Path, generation: Uuid) -> File {
1447 OpenOptions::new()
1448 .read(true)
1449 .write(true)
1450 .open(
1451 root.join("generations")
1452 .join(generation.hyphenated().to_string())
1453 .join(LEASE_FILE),
1454 )
1455 .unwrap()
1456 }
1457
1458 fn assert_code(error: GfError, expected: &str) {
1459 assert_eq!(error.code(), expected, "{error}");
1460 }
1461
1462 #[test]
1463 fn initial_generation_declares_graph_and_workspace_capabilities() {
1464 let root = tempfile::tempdir().unwrap();
1465
1466 let resolved = open_or_initialize_project(root.path()).unwrap();
1467
1468 assert_eq!(
1469 resolved.capabilities(),
1470 vec![
1471 ProjectCapabilityDescriptor {
1472 capability_id: "graph".into(),
1473 capability_version: 1,
1474 },
1475 ProjectCapabilityDescriptor {
1476 capability_id: "workspace".into(),
1477 capability_version: 1,
1478 },
1479 ]
1480 );
1481 let ontology = resolved
1482 .participant_snapshot("workspace", "ontology")
1483 .unwrap()
1484 .unwrap();
1485 assert_eq!(
1486 crate::WorkspaceOntology::from_canonical_json(&ontology.bytes)
1487 .unwrap()
1488 .mode,
1489 crate::WorkspaceOntologyMode::None
1490 );
1491 let configuration = resolved
1492 .participant_snapshot("workspace", "configuration")
1493 .unwrap()
1494 .unwrap();
1495 assert_eq!(
1496 crate::WorkspaceConfiguration::from_canonical_json(&configuration.bytes)
1497 .unwrap()
1498 .ontology_mode,
1499 crate::WorkspaceOntologyMode::None
1500 );
1501 assert_eq!(
1502 resolved.capability("graph").unwrap(),
1503 Some(ProjectCapabilityDescriptor {
1504 capability_id: "graph".into(),
1505 capability_version: 1,
1506 })
1507 );
1508 assert_eq!(resolved.capability("knowledge").unwrap(), None);
1509 assert_eq!(
1510 resolved
1511 .require_capability("knowledge", 1)
1512 .unwrap_err()
1513 .code(),
1514 "GF_CAPABILITY_DISABLED"
1515 );
1516 assert_eq!(
1517 resolved.require_capability("graph", 2).unwrap_err().code(),
1518 "GF_UNSUPPORTED_CAPABILITY_VERSION"
1519 );
1520 }
1521
1522 #[test]
1523 fn resolves_only_the_generation_named_by_current() {
1524 let (root, expected) = project();
1525 let abandoned = Uuid::now_v7();
1526 install_generation(root.path(), abandoned);
1527
1528 let resolved = resolve_project_generation(root.path()).unwrap();
1529
1530 assert_eq!(resolved.generation_uuid(), expected);
1531 assert_ne!(resolved.generation_uuid(), abandoned);
1532 }
1533
1534 #[test]
1535 fn concurrent_first_openers_converge_on_one_generation() {
1536 let root = tempfile::tempdir().unwrap();
1537 let barrier = Arc::new(Barrier::new(3));
1538 let mut workers = Vec::new();
1539 for _ in 0..2 {
1540 let root = root.path().to_owned();
1541 let barrier = Arc::clone(&barrier);
1542 workers.push(std::thread::spawn(move || {
1543 barrier.wait();
1544 open_or_initialize_project(root).unwrap().generation_uuid()
1545 }));
1546 }
1547
1548 barrier.wait();
1549 let first = workers.remove(0).join().unwrap();
1550 let second = workers.remove(0).join().unwrap();
1551
1552 assert_eq!(first, second);
1553 assert_eq!(
1554 fs::read_dir(root.path().join("generations"))
1555 .unwrap()
1556 .count(),
1557 1
1558 );
1559 }
1560
1561 #[test]
1562 fn project_root_lock_is_released_for_reopen() {
1563 let root = tempfile::tempdir().unwrap();
1564
1565 let first = open_or_initialize_project(root.path()).unwrap();
1566 let generation = first.generation_uuid();
1567 let reopened = open_or_initialize_project(root.path()).unwrap();
1568
1569 assert_eq!(reopened.generation_uuid(), generation);
1570 assert_eq!(first.generation_uuid(), generation);
1571 }
1572
1573 #[cfg(windows)]
1574 #[test]
1575 fn windows_project_root_lock_allows_owner_to_inspect_directory() {
1576 let root = tempfile::tempdir().unwrap();
1577 let _owner = lock_project_root(root.path()).unwrap();
1578
1579 assert_eq!(fs::read_dir(root.path()).unwrap().count(), 0);
1580 }
1581
1582 #[cfg(windows)]
1583 #[test]
1584 fn windows_project_directory_sync_uses_write_capable_handle() {
1585 let root = tempfile::tempdir().unwrap();
1586
1587 sync_directory(root.path()).unwrap();
1588 }
1589
1590 #[cfg(windows)]
1591 #[test]
1592 fn windows_project_root_lock_fails_closed_and_releases() {
1593 use std::process::Command;
1594
1595 let root = tempfile::tempdir().unwrap();
1596 let owner = lock_project_root(root.path()).unwrap();
1597
1598 let contender = Command::new(std::env::current_exe().unwrap())
1599 .arg("--exact")
1600 .arg(ROOT_LOCK_CONTENDER)
1601 .arg("--nocapture")
1602 .env("GRAPHFORGE_TEST_PROJECT_ROOT", root.path())
1603 .spawn()
1604 .unwrap();
1605 wait_for_lock_subprocess(contender, "subprocess contention proof");
1606
1607 drop(owner);
1608 open_or_initialize_project(root.path()).unwrap();
1609 }
1610
1611 #[cfg(windows)]
1612 #[test]
1613 fn windows_project_root_lock_contender() {
1614 let Ok(root) = std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT") else {
1615 return;
1616 };
1617
1618 let contention = open_or_initialize_project(PathBuf::from(root)).unwrap_err();
1619 assert_code(contention, "GF_WRITER_BUSY");
1620 }
1621
1622 #[cfg(windows)]
1623 #[test]
1624 fn windows_project_root_lock_recovers_abandoned_owner() {
1625 use std::process::Command;
1626
1627 let root = tempfile::tempdir().unwrap();
1628 let abandoner = Command::new(std::env::current_exe().unwrap())
1629 .arg("--exact")
1630 .arg(ROOT_LOCK_ABANDONER)
1631 .arg("--nocapture")
1632 .env("GRAPHFORGE_TEST_PROJECT_ROOT", root.path())
1633 .spawn()
1634 .unwrap();
1635 wait_for_lock_subprocess(abandoner, "subprocess abandonment proof");
1636
1637 let _recovered = lock_project_root(root.path()).unwrap();
1638 }
1639
1640 #[cfg(windows)]
1641 #[test]
1642 fn windows_project_root_lock_abandoner() {
1643 let Ok(root) = std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT") else {
1644 return;
1645 };
1646
1647 let root = PathBuf::from(root);
1648 let _owner = lock_project_root(&root).unwrap();
1649 std::process::exit(0);
1650 }
1651
1652 #[test]
1653 fn repeated_interrupted_initializations_reuse_one_validated_generation() {
1654 let root = tempfile::tempdir().unwrap();
1655 fs::write(root.path().join(FORMAT_FILE), PROJECT_FORMAT_BYTES).unwrap();
1656 fs::create_dir(root.path().join("generations")).unwrap();
1657 let mut generations = Vec::new();
1658 for _ in 0..16 {
1659 let generation = Uuid::now_v7();
1660 install_generation(root.path(), generation);
1661 generations.push(generation);
1662 }
1663 generations.sort_unstable();
1664
1665 let resolved = open_or_initialize_project(root.path()).unwrap();
1666
1667 assert_eq!(resolved.generation_uuid(), *generations.last().unwrap());
1668 assert_eq!(
1669 fs::read_dir(root.path().join("generations"))
1670 .unwrap()
1671 .count(),
1672 1
1673 );
1674 }
1675
1676 #[test]
1677 fn interrupted_current_temp_is_removed_before_initialization_resumes() {
1678 let root = tempfile::tempdir().unwrap();
1679 fs::write(root.path().join(FORMAT_FILE), PROJECT_FORMAT_BYTES).unwrap();
1680 fs::create_dir(root.path().join("generations")).unwrap();
1681 install_generation(root.path(), Uuid::now_v7());
1682 let writer_temp = root.path().join(".atomicwriteAb12Cd");
1683 fs::create_dir(&writer_temp).unwrap();
1684 fs::write(writer_temp.join("tmpfile.tmp"), b"partial CURRENT").unwrap();
1685
1686 open_or_initialize_project(root.path()).unwrap();
1687
1688 assert!(!writer_temp.exists());
1689 }
1690
1691 #[test]
1692 fn interrupted_initialization_rejects_unknown_layout_entries_without_mutation() {
1693 fn partial_project() -> (tempfile::TempDir, PathBuf) {
1694 let root = tempfile::tempdir().unwrap();
1695 fs::write(root.path().join(FORMAT_FILE), PROJECT_FORMAT_BYTES).unwrap();
1696 let generation = root
1697 .path()
1698 .join("generations")
1699 .join(Uuid::now_v7().hyphenated().to_string());
1700 fs::create_dir_all(&generation).unwrap();
1701 (root, generation)
1702 }
1703
1704 let (root, _) = partial_project();
1705 fs::write(root.path().join("caller-data"), b"preserve").unwrap();
1706 assert_code(
1707 open_or_initialize_project(root.path()).unwrap_err(),
1708 "GF_UNSUPPORTED_PROJECT_FORMAT",
1709 );
1710 assert_eq!(
1711 fs::read(root.path().join("caller-data")).unwrap(),
1712 b"preserve"
1713 );
1714
1715 let (root, generation) = partial_project();
1716 fs::write(generation.join("caller-data"), b"preserve").unwrap();
1717 assert_code(
1718 open_or_initialize_project(root.path()).unwrap_err(),
1719 "GF_UNSUPPORTED_PROJECT_FORMAT",
1720 );
1721 assert_eq!(
1722 fs::read(generation.join("caller-data")).unwrap(),
1723 b"preserve"
1724 );
1725
1726 let (root, generation) = partial_project();
1727 let participants = generation.join(PARTICIPANTS_DIR);
1728 fs::create_dir(&participants).unwrap();
1729 fs::write(participants.join("unknown"), b"preserve").unwrap();
1730 assert_code(
1731 open_or_initialize_project(root.path()).unwrap_err(),
1732 "GF_UNSUPPORTED_PROJECT_FORMAT",
1733 );
1734 assert_eq!(fs::read(participants.join("unknown")).unwrap(), b"preserve");
1735
1736 let (root, generation) = partial_project();
1737 let workspace = generation.join(PARTICIPANTS_DIR).join("workspace");
1738 fs::create_dir_all(&workspace).unwrap();
1739 fs::write(workspace.join("configuration.json"), b"{}").unwrap();
1740 fs::write(workspace.join("unknown.json"), b"preserve").unwrap();
1741 assert_code(
1742 open_or_initialize_project(root.path()).unwrap_err(),
1743 "GF_UNSUPPORTED_PROJECT_FORMAT",
1744 );
1745 assert_eq!(
1746 fs::read(workspace.join("unknown.json")).unwrap(),
1747 b"preserve"
1748 );
1749 }
1750
1751 #[test]
1752 fn wave9_interrupted_initialization_enforces_bounded_private_generation_layout() {
1753 let root = tempfile::tempdir().unwrap();
1754 fs::write(root.path().join(FORMAT_FILE), PROJECT_FORMAT_BYTES).unwrap();
1755 let generations = root.path().join("generations");
1756 fs::create_dir(&generations).unwrap();
1757 for _ in 0..17 {
1758 fs::create_dir(generations.join(Uuid::now_v7().hyphenated().to_string())).unwrap();
1759 }
1760 assert_code(
1761 open_or_initialize_project(root.path()).unwrap_err(),
1762 "GF_UNSUPPORTED_PROJECT_FORMAT",
1763 );
1764
1765 let root = tempfile::tempdir().unwrap();
1766 fs::write(root.path().join(FORMAT_FILE), PROJECT_FORMAT_BYTES).unwrap();
1767 let generation = root
1768 .path()
1769 .join("generations")
1770 .join(Uuid::now_v7().hyphenated().to_string());
1771 fs::create_dir_all(&generation).unwrap();
1772 for name in [PARTICIPANTS_DIR, LEASE_FILE, MANIFEST_FILE] {
1773 let path = generation.join(name);
1774 if name == PARTICIPANTS_DIR {
1775 fs::create_dir(path).unwrap();
1776 } else {
1777 fs::write(path, b"partial").unwrap();
1778 }
1779 }
1780 fs::write(generation.join("fourth-entry"), b"caller").unwrap();
1781 assert_code(
1782 open_or_initialize_project(root.path()).unwrap_err(),
1783 "GF_UNSUPPORTED_PROJECT_FORMAT",
1784 );
1785 assert_eq!(
1786 fs::read(generation.join("fourth-entry")).unwrap(),
1787 b"caller"
1788 );
1789 }
1790
1791 #[cfg(unix)]
1792 #[test]
1793 fn wave9_participant_inventory_rejects_links_and_special_files() {
1794 use std::os::unix::fs::symlink;
1795 use std::os::unix::net::UnixListener;
1796
1797 for kind in ["link", "socket"] {
1798 let root = tempfile::Builder::new()
1799 .prefix("gf")
1800 .tempdir_in("/tmp")
1801 .unwrap();
1802 let resolved = open_or_initialize_project(root.path()).unwrap();
1803 let participants = resolved.participants_root();
1804 let hostile = participants.join(format!("hostile-{kind}"));
1805 if kind == "link" {
1806 symlink(root.path().join(CURRENT_FILE), &hostile).unwrap();
1807 } else {
1808 let _listener = UnixListener::bind(&hostile).unwrap();
1809 }
1810 assert_code(
1811 resolved
1812 .validate_complete_participant_inventory()
1813 .unwrap_err(),
1814 "GF_TRANSACTION_FAILED",
1815 );
1816 }
1817 }
1818
1819 #[test]
1820 fn existing_resolution_remains_pinned_after_current_changes() {
1821 let (root, first) = project();
1822 let old_reader = resolve_project_generation(root.path()).unwrap();
1823 let second = Uuid::now_v7();
1824 let digest = install_generation(root.path(), second);
1825 write_current(root.path(), second, digest);
1826
1827 let new_reader = resolve_project_generation(root.path()).unwrap();
1828
1829 assert_eq!(old_reader.generation_uuid(), first);
1830 assert_eq!(new_reader.generation_uuid(), second);
1831 assert_ne!(old_reader.generation_root(), new_reader.generation_root());
1832 }
1833
1834 #[test]
1835 fn resolved_reader_holds_the_generation_lease() {
1836 let (root, generation) = project();
1837 let resolved = resolve_project_generation(root.path()).unwrap();
1838 let lease = open_generation_lease(root.path(), generation);
1839
1840 assert!(!FileExt::try_lock_exclusive(&lease).unwrap());
1841 drop(resolved);
1842 assert!(FileExt::try_lock_exclusive(&lease).unwrap());
1843 }
1844
1845 #[test]
1846 fn cloned_reader_holds_the_generation_lease_until_the_last_drop() {
1847 let (root, generation) = project();
1848 let resolved = resolve_project_generation(root.path()).unwrap();
1849 let cloned = resolved.clone();
1850 let lease = open_generation_lease(root.path(), generation);
1851
1852 drop(resolved);
1853 assert!(!FileExt::try_lock_exclusive(&lease).unwrap());
1854 drop(cloned);
1855 assert!(FileExt::try_lock_exclusive(&lease).unwrap());
1856 }
1857
1858 #[cfg(unix)]
1859 #[test]
1860 fn final_reader_drop_unlocks_a_duplicated_lease_description() {
1861 let (root, generation) = project();
1862 let resolved = resolve_project_generation(root.path()).unwrap();
1863 let duplicated_lease = resolved._lease_handle.0.try_clone().unwrap();
1864 let lease = open_generation_lease(root.path(), generation);
1865
1866 assert!(!FileExt::try_lock_exclusive(&lease).unwrap());
1867 drop(resolved);
1868 assert!(duplicated_lease.metadata().is_ok());
1869 assert!(FileExt::try_lock_exclusive(&lease).unwrap());
1870 }
1871
1872 #[test]
1873 fn rejects_pre_v1_root_without_mutation() {
1874 let root = tempfile::tempdir().unwrap();
1875 fs::create_dir(root.path().join("topology")).unwrap();
1876 fs::write(root.path().join("topology/nodes.parquet"), b"old").unwrap();
1877 let before = fs::read(root.path().join("topology/nodes.parquet")).unwrap();
1878
1879 let error = resolve_project_generation(root.path()).unwrap_err();
1880
1881 assert_code(error, "GF_UNSUPPORTED_PROJECT_FORMAT");
1882 assert_eq!(
1883 fs::read(root.path().join("topology/nodes.parquet")).unwrap(),
1884 before
1885 );
1886 assert!(!root.path().join(FORMAT_FILE).exists());
1887 }
1888
1889 #[test]
1890 fn exact_format_without_current_is_uninitialized() {
1891 let root = tempfile::tempdir().unwrap();
1892 fs::write(root.path().join(FORMAT_FILE), PROJECT_FORMAT_BYTES).unwrap();
1893
1894 let error = resolve_project_generation(root.path()).unwrap_err();
1895
1896 assert_code(error, "GF_PROJECT_UNINITIALIZED");
1897 }
1898
1899 #[test]
1900 fn future_format_is_unsupported() {
1901 let root = tempfile::tempdir().unwrap();
1902 fs::write(root.path().join(FORMAT_FILE), b"graphforge-project/v2\n").unwrap();
1903
1904 let error = resolve_project_generation(root.path()).unwrap_err();
1905
1906 assert_code(error, "GF_UNSUPPORTED_PROJECT_FORMAT");
1907 }
1908
1909 #[test]
1910 fn noncanonical_current_is_corrupt() {
1911 let (root, _) = project();
1912 let bytes = fs::read(root.path().join(CURRENT_FILE)).unwrap();
1913 let spaced = String::from_utf8(bytes).unwrap().replace("\":\"", "\": \"");
1914 fs::write(root.path().join(CURRENT_FILE), spaced).unwrap();
1915
1916 let error = resolve_project_generation(root.path()).unwrap_err();
1917
1918 assert_code(error, "GF_PROJECT_CORRUPT");
1919 }
1920
1921 #[test]
1922 fn manifest_digest_mismatch_is_corrupt() {
1923 let (root, generation) = project();
1924 write_current(root.path(), generation, [0_u8; 32]);
1925
1926 let error = resolve_project_generation(root.path()).unwrap_err();
1927
1928 assert_code(error, "GF_PROJECT_CORRUPT");
1929 let lease = open_generation_lease(root.path(), generation);
1930 assert!(FileExt::try_lock_exclusive(&lease).unwrap());
1931 }
1932
1933 #[test]
1934 fn participant_path_rejects_traversal_from_manifest() {
1935 let (root, generation) = project();
1936 let generation_root = root
1937 .path()
1938 .join("generations")
1939 .join(generation.hyphenated().to_string());
1940 let manifest = GenerationManifest {
1941 format: "graphforge-generation".into(),
1942 format_version: 1,
1943 generation_uuid: generation.hyphenated().to_string(),
1944 parent_generation_uuid: None,
1945 transaction_uuid: Uuid::now_v7().hyphenated().to_string(),
1946 capabilities: vec![CapabilityDescriptor {
1947 capability_id: "graph".into(),
1948 capability_version: 1,
1949 }],
1950 participants: vec![ParticipantDescriptor {
1951 capability_id: "graph".into(),
1952 capability_version: 1,
1953 record_family_id: "topology".into(),
1954 record_version: 1,
1955 relative_path: "../outside.parquet".into(),
1956 encoding: "parquet".into(),
1957 byte_length: 0,
1958 row_count: 0,
1959 schema_fingerprint: "0".repeat(64),
1960 content_sha256: "0".repeat(64),
1961 }],
1962 };
1963 let bytes = canonical_line(&manifest);
1964 fs::write(generation_root.join(MANIFEST_FILE), &bytes).unwrap();
1965 write_current(root.path(), generation, Sha256::digest(&bytes).into());
1966
1967 let error = resolve_project_generation(root.path()).unwrap_err();
1968
1969 assert_code(error, "GF_PROJECT_CORRUPT");
1970 }
1971
1972 #[test]
1973 fn resolution_does_not_open_unrequested_participant_tables() {
1974 let (root, generation) = project();
1975 let generation_root = root
1976 .path()
1977 .join("generations")
1978 .join(generation.hyphenated().to_string());
1979 let manifest = GenerationManifest {
1980 format: "graphforge-generation".into(),
1981 format_version: 1,
1982 generation_uuid: generation.hyphenated().to_string(),
1983 parent_generation_uuid: None,
1984 transaction_uuid: Uuid::now_v7().hyphenated().to_string(),
1985 capabilities: vec![
1986 CapabilityDescriptor {
1987 capability_id: "graph".into(),
1988 capability_version: 1,
1989 },
1990 CapabilityDescriptor {
1991 capability_id: "knowledge".into(),
1992 capability_version: 1,
1993 },
1994 ],
1995 participants: vec![ParticipantDescriptor {
1996 capability_id: "knowledge".into(),
1997 capability_version: 1,
1998 record_family_id: "assertions".into(),
1999 record_version: 1,
2000 relative_path: "knowledge/assertions.parquet".into(),
2001 encoding: "parquet".into(),
2002 byte_length: 99,
2003 row_count: 1,
2004 schema_fingerprint: "0".repeat(64),
2005 content_sha256: "0".repeat(64),
2006 }],
2007 };
2008 let bytes = canonical_line(&manifest);
2009 fs::write(generation_root.join(MANIFEST_FILE), &bytes).unwrap();
2010 write_current(root.path(), generation, Sha256::digest(&bytes).into());
2011 assert!(
2014 !generation_root
2015 .join("participants/knowledge/assertions.parquet")
2016 .exists()
2017 );
2018
2019 let resolved = resolve_project_generation(root.path()).unwrap();
2020
2021 assert_eq!(resolved.generation_uuid(), generation);
2022 }
2023
2024 #[cfg(unix)]
2025 #[test]
2026 fn selected_generation_symlink_is_corrupt() {
2027 use std::os::unix::fs::symlink;
2028
2029 let (root, generation) = project();
2030 let generation_root = root
2031 .path()
2032 .join("generations")
2033 .join(generation.hyphenated().to_string());
2034 fs::remove_file(generation_root.join(LEASE_FILE)).unwrap();
2035 symlink("/dev/null", generation_root.join(LEASE_FILE)).unwrap();
2036
2037 let error = resolve_project_generation(root.path()).unwrap_err();
2038
2039 assert_code(error, "GF_PROJECT_CORRUPT");
2040 }
2041
2042 #[test]
2043 fn canonical_identity_and_manifest_validation_matrix_is_total() {
2044 let expected = Uuid::now_v7();
2045 assert_eq!(
2046 parse_canonical_uuid(&expected.hyphenated().to_string()).unwrap(),
2047 expected
2048 );
2049 for value in [
2050 "not-a-uuid".to_owned(),
2051 expected.simple().to_string(),
2052 expected.hyphenated().to_string().to_uppercase(),
2053 ] {
2054 assert_code(
2055 parse_canonical_uuid(&value).unwrap_err(),
2056 "GF_PROJECT_CORRUPT",
2057 );
2058 }
2059 assert_eq!(parse_sha256(&"00".repeat(32)).unwrap(), [0; 32]);
2060 for value in ["00".to_owned(), "AA".repeat(32), "gg".repeat(32)] {
2061 assert_code(parse_sha256(&value).unwrap_err(), "GF_PROJECT_CORRUPT");
2062 }
2063 for value in ["", "Upper", "has/slash", "has space", ".", ".."] {
2064 assert_code(
2065 validate_machine_id(value).unwrap_err(),
2066 "GF_PROJECT_CORRUPT",
2067 );
2068 }
2069 assert!(validate_machine_id("graph_data-1").is_ok());
2070
2071 let base = GenerationManifest {
2072 format: "graphforge-generation".into(),
2073 format_version: 1,
2074 generation_uuid: expected.hyphenated().to_string(),
2075 parent_generation_uuid: None,
2076 transaction_uuid: Uuid::now_v7().hyphenated().to_string(),
2077 capabilities: vec![CapabilityDescriptor {
2078 capability_id: "graph".into(),
2079 capability_version: 1,
2080 }],
2081 participants: vec![],
2082 };
2083 assert!(validate_manifest(&base, expected).is_ok());
2084 let mutations: Vec<Box<dyn Fn(&mut GenerationManifest)>> = vec![
2085 Box::new(|manifest| manifest.format = "future".into()),
2086 Box::new(|manifest| manifest.format_version = 2),
2087 Box::new(|manifest| manifest.generation_uuid = Uuid::now_v7().to_string()),
2088 Box::new(|manifest| manifest.transaction_uuid = "bad".into()),
2089 Box::new(|manifest| manifest.capabilities[0].capability_id = "Upper".into()),
2090 Box::new(|manifest| manifest.capabilities[0].capability_version = 0),
2091 Box::new(|manifest| manifest.capabilities.push(manifest.capabilities[0].clone())),
2092 ];
2093 for mutate in mutations {
2094 let mut manifest = base.clone();
2095 mutate(&mut manifest);
2096 assert!(validate_manifest(&manifest, expected).is_err());
2097 }
2098
2099 let participant = ParticipantDescriptor {
2100 capability_id: "graph".into(),
2101 capability_version: 1,
2102 record_family_id: "topology".into(),
2103 record_version: 1,
2104 relative_path: "topology/nodes.parquet".into(),
2105 encoding: "parquet".into(),
2106 byte_length: 0,
2107 row_count: 0,
2108 schema_fingerprint: "0".repeat(64),
2109 content_sha256: "0".repeat(64),
2110 };
2111 let participant_mutations: Vec<Box<dyn Fn(&mut ParticipantDescriptor)>> = vec![
2112 Box::new(|entry| entry.capability_id = "missing".into()),
2113 Box::new(|entry| entry.capability_version = 2),
2114 Box::new(|entry| entry.record_family_id = "Upper".into()),
2115 Box::new(|entry| entry.record_version = 0),
2116 Box::new(|entry| entry.relative_path = "/absolute".into()),
2117 Box::new(|entry| entry.relative_path = "../escape".into()),
2118 Box::new(|entry| entry.schema_fingerprint = "short".into()),
2119 Box::new(|entry| entry.content_sha256 = "GG".repeat(32)),
2120 ];
2121 for mutate in participant_mutations {
2122 let mut manifest = base.clone();
2123 let mut entry = participant.clone();
2124 mutate(&mut entry);
2125 manifest.participants.push(entry);
2126 assert!(validate_manifest(&manifest, expected).is_err());
2127 }
2128
2129 let mut valid = base.clone();
2130 valid.participants.push(participant.clone());
2131 assert!(validate_manifest(&valid, expected).is_ok());
2132 valid.participants.push(participant);
2133 assert!(validate_manifest(&valid, expected).is_err());
2134 }
2135
2136 #[test]
2137 fn bounded_regular_file_rejects_missing_directory_and_oversize_without_mutation() {
2138 let root = tempfile::tempdir().unwrap();
2139 let missing = root.path().join("missing");
2140 assert!(read_bounded_regular_file(&missing, 4).is_err());
2141 let directory = root.path().join("directory");
2142 std::fs::create_dir(&directory).unwrap();
2143 assert!(read_bounded_regular_file(&directory, 4).is_err());
2144 let oversized = root.path().join("oversized");
2145 std::fs::write(&oversized, b"12345").unwrap();
2146 assert!(read_bounded_regular_file(&oversized, 4).is_err());
2147 assert_eq!(std::fs::read(&oversized).unwrap(), b"12345");
2148 assert_eq!(read_bounded_regular_file(&oversized, 5).unwrap(), b"12345");
2149 }
2150
2151 #[test]
2152 fn verified_reopen_rejects_wrong_digest_and_missing_generation_without_current_mutation() {
2153 let (root, generation_uuid) = project();
2154 let current = resolve_project_generation(root.path()).unwrap();
2155 let current_bytes = std::fs::read(root.path().join(CURRENT_FILE)).unwrap();
2156
2157 assert_code(
2158 resolve_verified_generation(root.path(), generation_uuid, [0x55; 32]).unwrap_err(),
2159 "GF_PROJECT_CORRUPT",
2160 );
2161 assert_code(
2162 resolve_verified_generation(root.path(), Uuid::now_v7(), current.manifest_sha256())
2163 .unwrap_err(),
2164 "GF_PROJECT_CORRUPT",
2165 );
2166 assert_eq!(
2167 std::fs::read(root.path().join(CURRENT_FILE)).unwrap(),
2168 current_bytes
2169 );
2170 assert_eq!(
2171 resolve_project_generation(root.path())
2172 .unwrap()
2173 .generation_uuid(),
2174 generation_uuid
2175 );
2176 }
2177
2178 #[test]
2179 fn participant_snapshot_rejects_same_length_content_tampering_after_reopen() {
2180 let root = tempfile::tempdir().unwrap();
2181 let resolved = open_or_initialize_project(root.path()).unwrap();
2182 let generation_uuid = resolved.generation_uuid();
2183 let snapshot = resolved
2184 .participant_snapshot("workspace", "configuration")
2185 .unwrap()
2186 .unwrap();
2187 let path = resolved
2188 .participant_path("workspace", "configuration")
2189 .unwrap();
2190 let mut tampered = snapshot.bytes.clone();
2191 tampered[0] ^= 1;
2192 fs::write(&path, tampered).unwrap();
2193 drop(resolved);
2194
2195 let reopened = resolve_project_generation(root.path()).unwrap();
2196 assert_eq!(reopened.generation_uuid(), generation_uuid);
2197 assert_code(
2198 reopened
2199 .participant_snapshot("workspace", "configuration")
2200 .unwrap_err(),
2201 "GF_PROJECT_CORRUPT",
2202 );
2203 }
2204}