1use std::collections::BTreeSet;
9use std::io::{Read, Write};
10use std::path::Path;
11
12use atomicwrites::{AtomicFile, DisallowOverwrite};
13use graphforge_core::{GfError, ProjectErrorCode};
14use serde::{Deserialize, Serialize};
15use sha2::{Digest, Sha256};
16use uuid::Uuid;
17
18use crate::{
19 ProjectCapability, ProjectGenerationRequest, ProjectParticipant, ProjectParticipantEncoding,
20 ProjectPublicationReceipt, ProjectStageOutcome, ResolvedProjectGeneration,
21 open_or_initialize_project, resolve_project_generation, stage_project_generation,
22};
23
24const MAGIC: &[u8; 16] = b"graphforge-exp\0\n";
25const FORMAT: &str = "graphforge-portable-export";
26const FORMAT_VERSION: u32 = 1;
27const HEADER_LENGTH_BYTES: usize = 8;
28
29#[derive(Debug, Clone, Copy, PartialEq, Eq)]
31pub struct PortableProjectLimits {
32 pub max_envelope_bytes: u64,
34 pub max_header_bytes: u64,
36 pub max_participants: usize,
38 pub max_participant_bytes: u64,
40}
41
42impl Default for PortableProjectLimits {
43 fn default() -> Self {
44 Self {
45 max_envelope_bytes: 16 * 1024 * 1024 * 1024,
46 max_header_bytes: 4 * 1024 * 1024,
47 max_participants: 100_000,
48 max_participant_bytes: 8 * 1024 * 1024 * 1024,
49 }
50 }
51}
52
53#[derive(Debug, Clone, PartialEq, Eq)]
55pub struct PortableExportReceipt {
56 pub generation_uuid: Uuid,
58 pub envelope_sha256: [u8; 32],
60 pub byte_length: u64,
62 pub participant_count: usize,
64}
65
66#[derive(Debug, Clone, PartialEq, Eq)]
68pub struct PortableImportReceipt {
69 pub envelope_sha256: [u8; 32],
71 pub source_generation_uuid: Uuid,
73 pub publication: ProjectPublicationReceipt,
75}
76
77#[derive(Debug, Serialize, Deserialize)]
78#[serde(deny_unknown_fields)]
79struct EnvelopeHeader {
80 format: String,
81 format_version: u32,
82 source_generation_uuid: String,
83 capabilities: Vec<EnvelopeCapability>,
84 participants: Vec<EnvelopeParticipant>,
85}
86
87#[derive(Debug, Serialize, Deserialize)]
88#[serde(deny_unknown_fields)]
89struct EnvelopeCapability {
90 capability_id: String,
91 capability_version: u32,
92}
93
94#[derive(Debug, Serialize, Deserialize)]
95#[serde(deny_unknown_fields)]
96struct EnvelopeParticipant {
97 capability_id: String,
98 capability_version: u32,
99 record_family_id: String,
100 record_version: u32,
101 encoding: String,
102 schema_fingerprint: String,
103 row_count: u64,
104 byte_length: u64,
105 content_sha256: String,
106}
107
108struct ValidatedEnvelope {
109 source_generation_uuid: Uuid,
110 envelope_sha256: [u8; 32],
111 capabilities: Vec<ProjectCapability>,
112 participants: Vec<ProjectParticipant>,
113}
114
115pub fn encode_portable_project(
121 generation: &ResolvedProjectGeneration,
122 limits: PortableProjectLimits,
123) -> Result<(Vec<u8>, PortableExportReceipt), GfError> {
124 let snapshots = generation.participant_snapshots()?;
125 enforce_count(snapshots.len(), limits.max_participants)?;
126 let capabilities = generation
127 .capabilities()
128 .into_iter()
129 .map(|item| EnvelopeCapability {
130 capability_id: item.capability_id,
131 capability_version: item.capability_version,
132 })
133 .collect();
134 let mut participants = Vec::with_capacity(snapshots.len());
135 let mut body_length = 0_u64;
136 for snapshot in &snapshots {
137 let byte_length = u64::try_from(snapshot.bytes.len())
138 .map_err(|_| resource("portable participant byte length exceeds u64"))?;
139 enforce_size(
140 byte_length,
141 limits.max_participant_bytes,
142 "portable participant",
143 )?;
144 body_length = body_length
145 .checked_add(byte_length)
146 .ok_or_else(|| resource("portable envelope size overflow"))?;
147 participants.push(EnvelopeParticipant {
148 capability_id: snapshot.capability_id.clone(),
149 capability_version: snapshot.capability_version,
150 record_family_id: snapshot.record_family_id.clone(),
151 record_version: snapshot.record_version,
152 encoding: snapshot.encoding.clone(),
153 schema_fingerprint: hex(snapshot.schema_fingerprint),
154 row_count: snapshot.row_count,
155 byte_length,
156 content_sha256: hex(Sha256::digest(&snapshot.bytes).into()),
157 });
158 }
159 let header = EnvelopeHeader {
160 format: FORMAT.into(),
161 format_version: FORMAT_VERSION,
162 source_generation_uuid: generation.generation_uuid().hyphenated().to_string(),
163 capabilities,
164 participants,
165 };
166 let mut header_bytes = serde_json::to_vec(&header)
167 .map_err(|error| GfError::Storage(format!("failed to encode portable header: {error}")))?;
168 header_bytes.push(b'\n');
169 let header_length = u64::try_from(header_bytes.len())
170 .map_err(|_| resource("portable header byte length exceeds u64"))?;
171 enforce_size(header_length, limits.max_header_bytes, "portable header")?;
172 let total = u64::try_from(MAGIC.len() + HEADER_LENGTH_BYTES)
173 .expect("fixed prefix fits u64")
174 .checked_add(header_length)
175 .and_then(|value| value.checked_add(body_length))
176 .ok_or_else(|| resource("portable envelope size overflow"))?;
177 enforce_size(total, limits.max_envelope_bytes, "portable envelope")?;
178 let capacity = usize::try_from(total)
179 .map_err(|_| resource("portable envelope does not fit address space"))?;
180 let mut envelope = Vec::with_capacity(capacity);
181 envelope.extend_from_slice(MAGIC);
182 envelope.extend_from_slice(&header_length.to_be_bytes());
183 envelope.extend_from_slice(&header_bytes);
184 for snapshot in snapshots {
185 envelope.extend_from_slice(&snapshot.bytes);
186 }
187 let envelope_sha256 = Sha256::digest(&envelope).into();
188 Ok((
189 envelope,
190 PortableExportReceipt {
191 generation_uuid: generation.generation_uuid(),
192 envelope_sha256,
193 byte_length: total,
194 participant_count: header.participants.len(),
195 },
196 ))
197}
198
199pub fn export_portable_project(
201 generation: &ResolvedProjectGeneration,
202 destination: impl AsRef<Path>,
203 limits: PortableProjectLimits,
204) -> Result<PortableExportReceipt, GfError> {
205 let destination = destination.as_ref();
206 reject_export_destination(destination)?;
207 let (bytes, receipt) = encode_portable_project(generation, limits)?;
208 AtomicFile::new(destination, DisallowOverwrite)
209 .write(|file| {
210 file.write_all(&bytes)?;
211 file.sync_all()
212 })
213 .map_err(|error| GfError::Storage(format!("failed to write portable export: {error}")))?;
214 Ok(receipt)
215}
216
217pub fn import_portable_project(
224 envelope: &[u8],
225 target: impl AsRef<Path>,
226 transaction_uuid: Uuid,
227 generation_uuid: Uuid,
228 supported_capabilities: &[ProjectCapability],
229 limits: PortableProjectLimits,
230) -> Result<PortableImportReceipt, GfError> {
231 let validated = validate_envelope(envelope, supported_capabilities, limits)?;
232 let target = target.as_ref();
233 let existing_parent = prepare_import_target(target)?;
234 let initialized_parent;
235 let _parent = if let Some(parent) = existing_parent {
236 parent
237 } else {
238 initialized_parent = open_or_initialize_project(target)?;
239 initialized_parent
240 };
241 let request = ProjectGenerationRequest {
242 transaction_uuid,
243 generation_uuid,
244 capabilities: validated.capabilities,
245 participants: validated.participants,
246 };
247 let publication = match stage_project_generation(target, &request)? {
248 ProjectStageOutcome::AlreadyPublished(receipt) => receipt,
249 ProjectStageOutcome::Staged(staged) => {
250 staged.validate(|_| Ok(()), |_, _| Ok(()))?.publish()?
251 }
252 };
253 Ok(PortableImportReceipt {
254 envelope_sha256: validated.envelope_sha256,
255 source_generation_uuid: validated.source_generation_uuid,
256 publication,
257 })
258}
259
260pub fn import_portable_project_file(
266 source: impl AsRef<Path>,
267 target: impl AsRef<Path>,
268 transaction_uuid: Uuid,
269 generation_uuid: Uuid,
270 supported_capabilities: &[ProjectCapability],
271 limits: PortableProjectLimits,
272) -> Result<PortableImportReceipt, GfError> {
273 let source = source.as_ref();
274 reject_symlink_components(source, "portable import source")?;
275 let metadata = std::fs::symlink_metadata(source).map_err(|error| {
276 GfError::Storage(format!("failed to inspect portable import source: {error}"))
277 })?;
278 if metadata.file_type().is_symlink() || !metadata.is_file() {
279 return Err(project_error(
280 ProjectErrorCode::UnsupportedFilesystem,
281 "portable import source is linked or not a regular file",
282 ));
283 }
284 enforce_size(
285 metadata.len(),
286 limits.max_envelope_bytes,
287 "portable envelope",
288 )?;
289 let bounded_capacity = usize::try_from(metadata.len())
290 .map_err(|_| resource("portable envelope does not fit address space"))?;
291 let mut bytes = Vec::with_capacity(bounded_capacity);
292 let mut file = open_regular_nofollow(source).map_err(|error| {
293 GfError::Storage(format!("failed to open portable import source: {error}"))
294 })?;
295 let opened_metadata = file.metadata().map_err(|error| {
296 GfError::Storage(format!(
297 "failed to inspect opened portable import source: {error}"
298 ))
299 })?;
300 if !opened_metadata.is_file() || !same_file_identity(&metadata, &opened_metadata) {
301 return Err(project_error(
302 ProjectErrorCode::UnsupportedFilesystem,
303 "portable import source changed while it was being opened",
304 ));
305 }
306 reject_symlink_components(source, "portable import source")?;
307 Read::by_ref(&mut file)
308 .take(limits.max_envelope_bytes.saturating_add(1))
309 .read_to_end(&mut bytes)
310 .map_err(|error| {
311 GfError::Storage(format!("failed to read portable import source: {error}"))
312 })?;
313 if u64::try_from(bytes.len()).unwrap_or(u64::MAX) > limits.max_envelope_bytes {
314 return Err(resource("portable envelope exceeds limit"));
315 }
316 import_portable_project(
317 &bytes,
318 target,
319 transaction_uuid,
320 generation_uuid,
321 supported_capabilities,
322 limits,
323 )
324}
325
326fn validate_envelope(
327 envelope: &[u8],
328 supported_capabilities: &[ProjectCapability],
329 limits: PortableProjectLimits,
330) -> Result<ValidatedEnvelope, GfError> {
331 let envelope_length = u64::try_from(envelope.len())
332 .map_err(|_| resource("portable envelope byte length exceeds u64"))?;
333 enforce_size(
334 envelope_length,
335 limits.max_envelope_bytes,
336 "portable envelope",
337 )?;
338 let prefix = MAGIC.len() + HEADER_LENGTH_BYTES;
339 if envelope.len() < prefix || &envelope[..MAGIC.len()] != MAGIC {
340 return Err(corrupt("portable envelope magic is invalid"));
341 }
342 let header_length = u64::from_be_bytes(
343 envelope[MAGIC.len()..prefix]
344 .try_into()
345 .expect("fixed header length slice"),
346 );
347 enforce_size(header_length, limits.max_header_bytes, "portable header")?;
348 let header_length = usize::try_from(header_length)
349 .map_err(|_| resource("portable header does not fit address space"))?;
350 let header_end = prefix
351 .checked_add(header_length)
352 .ok_or_else(|| corrupt("portable header length overflows"))?;
353 let header_bytes = envelope
354 .get(prefix..header_end)
355 .ok_or_else(|| corrupt("portable envelope header is truncated"))?;
356 let header: EnvelopeHeader = serde_json::from_slice(header_bytes)
357 .map_err(|_| corrupt("portable envelope header is not valid canonical JSON"))?;
358 let canonical = {
359 let mut bytes = serde_json::to_vec(&header).map_err(|error| {
360 GfError::Storage(format!("failed to canonicalize portable header: {error}"))
361 })?;
362 bytes.push(b'\n');
363 bytes
364 };
365 if canonical != header_bytes {
366 return Err(corrupt("portable envelope header is not canonical"));
367 }
368 if header.format != FORMAT || header.format_version != FORMAT_VERSION {
369 return Err(project_error(
370 ProjectErrorCode::UnsupportedProjectFormat,
371 "portable envelope format or version is unsupported",
372 ));
373 }
374 let source_generation_uuid = parse_canonical_uuid(&header.source_generation_uuid)?;
375 enforce_count(header.participants.len(), limits.max_participants)?;
376 validate_capabilities(&header.capabilities, supported_capabilities)?;
377 let participants = validate_participants(&header, envelope, header_end, limits)?;
378 let capabilities = header
379 .capabilities
380 .into_iter()
381 .map(|item| ProjectCapability {
382 capability_id: item.capability_id,
383 capability_version: item.capability_version,
384 })
385 .collect();
386 Ok(ValidatedEnvelope {
387 source_generation_uuid,
388 envelope_sha256: Sha256::digest(envelope).into(),
389 capabilities,
390 participants,
391 })
392}
393
394fn validate_participants(
395 header: &EnvelopeHeader,
396 envelope: &[u8],
397 header_end: usize,
398 limits: PortableProjectLimits,
399) -> Result<Vec<ProjectParticipant>, GfError> {
400 let mut cursor = header_end;
401 let mut identities = BTreeSet::new();
402 let mut prior_identity: Option<(&str, &str)> = None;
403 let mut participants = Vec::with_capacity(header.participants.len());
404 for item in &header.participants {
405 validate_machine_id(&item.capability_id)?;
406 validate_machine_id(&item.record_family_id)?;
407 let identity = (item.capability_id.as_str(), item.record_family_id.as_str());
408 if prior_identity.is_some_and(|prior| prior >= identity) {
409 return Err(corrupt("portable participant inventory is not canonical"));
410 }
411 prior_identity = Some(identity);
412 if !identities.insert((&item.capability_id, &item.record_family_id)) {
413 return Err(corrupt(
414 "portable envelope has duplicate participant identity",
415 ));
416 }
417 enforce_size(
418 item.byte_length,
419 limits.max_participant_bytes,
420 "portable participant",
421 )?;
422 let length = usize::try_from(item.byte_length)
423 .map_err(|_| resource("portable participant does not fit address space"))?;
424 let end = cursor
425 .checked_add(length)
426 .ok_or_else(|| corrupt("portable participant length overflows"))?;
427 let bytes = envelope
428 .get(cursor..end)
429 .ok_or_else(|| corrupt("portable participant is truncated"))?;
430 if Sha256::digest(bytes).as_slice() != parse_digest(&item.content_sha256)? {
431 return Err(corrupt(
432 "portable participant content digest does not match",
433 ));
434 }
435 let encoding = match item.encoding.as_str() {
436 "parquet" => ProjectParticipantEncoding::Parquet,
437 "arrow" => ProjectParticipantEncoding::Arrow,
438 "json" => ProjectParticipantEncoding::Json,
439 _ => return Err(corrupt("portable participant encoding is unsupported")),
440 };
441 if !header.capabilities.iter().any(|capability| {
442 capability.capability_id == item.capability_id
443 && capability.capability_version == item.capability_version
444 }) {
445 return Err(corrupt("portable participant capability is not declared"));
446 }
447 participants.push(ProjectParticipant {
448 capability_id: item.capability_id.clone(),
449 capability_version: item.capability_version,
450 record_family_id: item.record_family_id.clone(),
451 record_version: item.record_version,
452 encoding,
453 schema_fingerprint: parse_digest(&item.schema_fingerprint)?,
454 row_count: item.row_count,
455 bytes: bytes.to_vec(),
456 });
457 cursor = end;
458 }
459 if cursor != envelope.len() {
460 return Err(corrupt("portable envelope has trailing bytes"));
461 }
462 Ok(participants)
463}
464
465fn validate_capabilities(
466 capabilities: &[EnvelopeCapability],
467 supported: &[ProjectCapability],
468) -> Result<(), GfError> {
469 let mut prior: Option<&str> = None;
470 for item in capabilities {
471 validate_machine_id(&item.capability_id)?;
472 if item.capability_version == 0
473 || prior.is_some_and(|value| value >= item.capability_id.as_str())
474 {
475 return Err(corrupt("portable capability inventory is not canonical"));
476 }
477 prior = Some(&item.capability_id);
478 if !supported.iter().any(|candidate| {
479 candidate.capability_id == item.capability_id
480 && candidate.capability_version == item.capability_version
481 }) {
482 return Err(project_error(
483 ProjectErrorCode::UnsupportedCapabilityVersion,
484 format!(
485 "portable capability {}@{} is unsupported",
486 item.capability_id, item.capability_version
487 ),
488 ));
489 }
490 }
491 Ok(())
492}
493
494fn prepare_import_target(target: &Path) -> Result<Option<ResolvedProjectGeneration>, GfError> {
495 reject_symlink_components(target, "portable import target")?;
496 match std::fs::symlink_metadata(target) {
497 Ok(metadata) if metadata.file_type().is_symlink() || !metadata.is_dir() => {
498 Err(project_error(
499 ProjectErrorCode::UnsupportedProjectFormat,
500 "portable import target is linked or not a directory",
501 ))
502 }
503 Ok(_) => {
504 let is_empty = std::fs::read_dir(target)
505 .map_err(|error| {
506 GfError::Storage(format!("failed to inspect portable import target: {error}"))
507 })?
508 .next()
509 .is_none();
510 if is_empty {
511 Ok(None)
512 } else if is_pristine_initialized_target(target)? {
513 resolve_project_generation(target).map(Some)
514 } else {
515 Err(project_error(
516 ProjectErrorCode::UnsupportedProjectFormat,
517 "portable import target must be empty or pristine",
518 ))
519 }
520 }
521 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
522 std::fs::create_dir(target).map_err(|error| {
523 GfError::Storage(format!("failed to create portable import target: {error}"))
524 })?;
525 Ok(None)
526 }
527 Err(error) => Err(GfError::Storage(format!(
528 "failed to inspect portable import target: {error}"
529 ))),
530 }
531}
532
533fn is_pristine_initialized_target(target: &Path) -> Result<bool, GfError> {
534 let Ok(resolved) = resolve_project_generation(target) else {
535 return Ok(false);
536 };
537 let capabilities = resolved.capabilities();
538 if capabilities.len() != 2
539 || capabilities[0].capability_id != "graph"
540 || capabilities[0].capability_version != 1
541 || capabilities[1].capability_id != "workspace"
542 || capabilities[1].capability_version != 1
543 {
544 return Ok(false);
545 }
546 let actual = resolved.participant_snapshots()?;
547 let expected = crate::workspace_participants::empty_workspace_participants()?;
548 if actual.len() != expected.len()
549 || !actual.iter().zip(&expected).all(|(actual, expected)| {
550 actual.capability_id == expected.capability_id
551 && actual.capability_version == expected.capability_version
552 && actual.record_family_id == expected.record_family_id
553 && actual.record_version == expected.record_version
554 && actual.encoding
555 == match expected.encoding {
556 ProjectParticipantEncoding::Parquet => "parquet",
557 ProjectParticipantEncoding::Arrow => "arrow",
558 ProjectParticipantEncoding::Json => "json",
559 }
560 && actual.schema_fingerprint == expected.schema_fingerprint
561 && actual.row_count == expected.row_count
562 && actual.bytes == expected.bytes
563 })
564 {
565 return Ok(false);
566 }
567 has_exact_pristine_layout(target, resolved.generation_uuid())
568}
569
570fn has_exact_pristine_layout(target: &Path, generation_uuid: Uuid) -> Result<bool, GfError> {
571 let root_names = directory_names(target)?;
572 if root_names != ["CURRENT", "FORMAT", "generations"] {
573 return Ok(false);
574 }
575 let generations = target.join("generations");
576 if directory_names(&generations)? != [generation_uuid.hyphenated().to_string()] {
577 return Ok(false);
578 }
579 let generation = generations.join(generation_uuid.hyphenated().to_string());
580 if directory_names(&generation)? != ["lease.lock", "manifest.json", "participants"] {
581 return Ok(false);
582 }
583 let participants = generation.join("participants");
584 if directory_names(&participants)? != ["workspace"] {
585 return Ok(false);
586 }
587 Ok(
588 directory_names(&participants.join("workspace"))?
589 == ["configuration.json", "ontology.json"],
590 )
591}
592
593fn directory_names(path: &Path) -> Result<Vec<String>, GfError> {
594 let mut names = std::fs::read_dir(path)
595 .map_err(|error| {
596 GfError::Storage(format!("failed to inspect portable import target: {error}"))
597 })?
598 .map(|entry| {
599 entry
600 .map_err(|error| {
601 GfError::Storage(format!("failed to inspect portable import target: {error}"))
602 })?
603 .file_name()
604 .into_string()
605 .map_err(|_| {
606 project_error(
607 ProjectErrorCode::UnsupportedProjectFormat,
608 "portable import target contains a non-UTF-8 entry",
609 )
610 })
611 })
612 .collect::<Result<Vec<_>, _>>()?;
613 names.sort();
614 Ok(names)
615}
616
617fn reject_export_destination(path: &Path) -> Result<(), GfError> {
618 reject_symlink_components(path, "portable export destination")?;
619 if std::fs::symlink_metadata(path).is_ok() {
620 return Err(project_error(
621 ProjectErrorCode::UnsupportedProjectFormat,
622 "portable export destination already exists",
623 ));
624 }
625 let parent = path.parent().ok_or_else(|| {
626 project_error(
627 ProjectErrorCode::UnsupportedFilesystem,
628 "portable export destination has no parent",
629 )
630 })?;
631 let metadata = std::fs::symlink_metadata(parent).map_err(|error| {
632 GfError::Storage(format!("failed to inspect portable export parent: {error}"))
633 })?;
634 if metadata.file_type().is_symlink() || !metadata.is_dir() {
635 return Err(project_error(
636 ProjectErrorCode::UnsupportedFilesystem,
637 "portable export parent is linked or not a directory",
638 ));
639 }
640 Ok(())
641}
642
643fn reject_symlink_components(path: &Path, name: &str) -> Result<(), GfError> {
644 let absolute = if path.is_absolute() {
645 path.to_path_buf()
646 } else {
647 std::env::current_dir()
648 .map_err(|error| {
649 GfError::Storage(format!("failed to resolve current directory: {error}"))
650 })?
651 .join(path)
652 };
653 for component in absolute.ancestors() {
654 match std::fs::symlink_metadata(component) {
655 Ok(metadata)
656 if metadata.file_type().is_symlink() && !trusted_platform_symlink(&metadata) =>
657 {
658 return Err(project_error(
659 ProjectErrorCode::UnsupportedFilesystem,
660 format!("{name} has a symbolic-link path component"),
661 ));
662 }
663 Ok(metadata)
664 if component != absolute
665 && !metadata.is_dir()
666 && !metadata.file_type().is_symlink() =>
667 {
668 return Err(project_error(
669 ProjectErrorCode::UnsupportedFilesystem,
670 format!("{name} has a non-directory ancestor"),
671 ));
672 }
673 Ok(_) => {}
674 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
675 Err(error) => {
676 return Err(GfError::Storage(format!(
677 "failed to inspect {name} path components: {error}"
678 )));
679 }
680 }
681 }
682 Ok(())
683}
684
685#[cfg(unix)]
688fn trusted_platform_symlink(metadata: &std::fs::Metadata) -> bool {
689 use std::os::unix::fs::MetadataExt;
690
691 metadata.uid() == 0
692}
693
694#[cfg(not(unix))]
695fn trusted_platform_symlink(_metadata: &std::fs::Metadata) -> bool {
696 false
697}
698
699#[cfg(unix)]
700fn open_regular_nofollow(path: &Path) -> std::io::Result<std::fs::File> {
701 use std::os::unix::fs::OpenOptionsExt;
702
703 std::fs::OpenOptions::new()
704 .read(true)
705 .custom_flags(libc::O_NOFOLLOW)
706 .open(path)
707}
708
709#[cfg(not(unix))]
710fn open_regular_nofollow(path: &Path) -> std::io::Result<std::fs::File> {
711 std::fs::File::open(path)
712}
713
714#[cfg(unix)]
715fn same_file_identity(before: &std::fs::Metadata, after: &std::fs::Metadata) -> bool {
716 use std::os::unix::fs::MetadataExt;
717
718 before.dev() == after.dev() && before.ino() == after.ino()
719}
720
721#[cfg(not(unix))]
722fn same_file_identity(before: &std::fs::Metadata, after: &std::fs::Metadata) -> bool {
723 before.len() == after.len()
724 && before.modified().ok() == after.modified().ok()
725 && before.created().ok() == after.created().ok()
726}
727
728fn validate_machine_id(value: &str) -> Result<(), GfError> {
729 if value.is_empty()
730 || value.len() > 128
731 || !value.bytes().all(|byte| {
732 byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'-' || byte == b'_'
733 })
734 {
735 return Err(corrupt(
736 "portable envelope contains an invalid machine identifier",
737 ));
738 }
739 Ok(())
740}
741
742fn parse_canonical_uuid(value: &str) -> Result<Uuid, GfError> {
743 let parsed = Uuid::parse_str(value)
744 .map_err(|_| corrupt("portable envelope generation UUID is invalid"))?;
745 if parsed.hyphenated().to_string() != value {
746 return Err(corrupt(
747 "portable envelope generation UUID is not canonical",
748 ));
749 }
750 Ok(parsed)
751}
752
753fn parse_digest(value: &str) -> Result<[u8; 32], GfError> {
754 if value.len() != 64
755 || !value
756 .bytes()
757 .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
758 {
759 return Err(corrupt("portable envelope digest is invalid"));
760 }
761 let mut digest = [0_u8; 32];
762 for (index, chunk) in value.as_bytes().chunks_exact(2).enumerate() {
763 digest[index] = (hex_nibble(chunk[0]) << 4) | hex_nibble(chunk[1]);
764 }
765 Ok(digest)
766}
767
768fn hex_nibble(byte: u8) -> u8 {
769 if byte <= b'9' {
770 byte - b'0'
771 } else {
772 byte - b'a' + 10
773 }
774}
775
776fn hex(digest: [u8; 32]) -> String {
777 const DIGITS: &[u8; 16] = b"0123456789abcdef";
778 let mut output = String::with_capacity(64);
779 for byte in digest {
780 output.push(DIGITS[usize::from(byte >> 4)] as char);
781 output.push(DIGITS[usize::from(byte & 0xf)] as char);
782 }
783 output
784}
785
786fn enforce_count(actual: usize, limit: usize) -> Result<(), GfError> {
787 if actual > limit {
788 Err(resource("portable participant count exceeds limit"))
789 } else {
790 Ok(())
791 }
792}
793
794fn enforce_size(actual: u64, limit: u64, name: &str) -> Result<(), GfError> {
795 if actual > limit {
796 Err(resource(format!("{name} exceeds limit")))
797 } else {
798 Ok(())
799 }
800}
801
802fn corrupt(message: impl Into<String>) -> GfError {
803 project_error(ProjectErrorCode::ProjectCorrupt, message)
804}
805fn resource(message: impl Into<String>) -> GfError {
806 project_error(ProjectErrorCode::ResourceLimit, message)
807}
808fn project_error(code: ProjectErrorCode, message: impl Into<String>) -> GfError {
809 GfError::Project {
810 code,
811 message: message.into(),
812 }
813}
814
815#[cfg(test)]
816mod tests {
817 use std::process::Command;
818
819 use super::*;
820
821 const ENABLE_COOKIE: &str = "graphforge-internal-subprocess-v1";
822 const IMPORT_HELPER: &str = "project_portable::tests::subprocess_portable_import_writer";
823
824 fn supported(generation: &ResolvedProjectGeneration) -> Vec<ProjectCapability> {
825 generation
826 .capabilities()
827 .into_iter()
828 .map(|item| ProjectCapability {
829 capability_id: item.capability_id,
830 capability_version: item.capability_version,
831 })
832 .collect()
833 }
834
835 #[test]
836 fn subprocess_portable_import_writer() {
837 let Ok(target) = std::env::var("GRAPHFORGE_TEST_PROJECT_ROOT") else {
838 return;
839 };
840 let envelope =
841 std::fs::read(std::env::var("GRAPHFORGE_TEST_PORTABLE_ENVELOPE").unwrap()).unwrap();
842 import_portable_project(
843 &envelope,
844 target,
845 Uuid::parse_str(&std::env::var("GRAPHFORGE_TEST_TRANSACTION_UUID").unwrap()).unwrap(),
846 Uuid::parse_str(&std::env::var("GRAPHFORGE_TEST_GENERATION_UUID").unwrap()).unwrap(),
847 &[
848 ProjectCapability {
849 capability_id: "graph".into(),
850 capability_version: 1,
851 },
852 ProjectCapability {
853 capability_id: "workspace".into(),
854 capability_version: 1,
855 },
856 ],
857 PortableProjectLimits::default(),
858 )
859 .unwrap();
860 }
861
862 #[test]
863 fn prepublication_failure_keeps_pristine_current_authoritative() {
864 let source = tempfile::tempdir().unwrap();
865 let source_generation = open_or_initialize_project(source.path()).unwrap();
866 let (envelope, _) =
867 encode_portable_project(&source_generation, PortableProjectLimits::default()).unwrap();
868 let envelope_path = source.path().join("portable.gfportable");
869 std::fs::write(&envelope_path, envelope).unwrap();
870
871 let target = tempfile::tempdir().unwrap();
872 let parent = open_or_initialize_project(target.path())
873 .unwrap()
874 .generation_uuid();
875 let status = Command::new(std::env::current_exe().unwrap())
876 .arg("--exact")
877 .arg(IMPORT_HELPER)
878 .arg("--nocapture")
879 .env("GRAPHFORGE_TEST_PROJECT_ROOT", target.path())
880 .env("GRAPHFORGE_TEST_PORTABLE_ENVELOPE", &envelope_path)
881 .env(
882 "GRAPHFORGE_TEST_TRANSACTION_UUID",
883 Uuid::new_v4().to_string(),
884 )
885 .env(
886 "GRAPHFORGE_TEST_GENERATION_UUID",
887 Uuid::new_v4().to_string(),
888 )
889 .env("GRAPHFORGE_PROJECT_FAILPOINTS", ENABLE_COOKIE)
890 .env(
891 "GRAPHFORGE_PROJECT_FAILPOINT",
892 "project.before_current_replace",
893 )
894 .status()
895 .unwrap();
896 assert_eq!(status.code(), Some(crate::project_failpoint::exit_code()));
897 assert_eq!(
898 resolve_project_generation(target.path())
899 .unwrap()
900 .generation_uuid(),
901 parent
902 );
903 }
904
905 #[test]
906 fn deterministic_round_trip_publishes_complete_new_generation() {
907 let source = tempfile::tempdir().unwrap();
908 let source_generation = open_or_initialize_project(source.path()).unwrap();
909 let expected = source_generation.participant_snapshots().unwrap();
910 let limits = PortableProjectLimits::default();
911 let (first, first_receipt) = encode_portable_project(&source_generation, limits).unwrap();
912 let (second, second_receipt) = encode_portable_project(&source_generation, limits).unwrap();
913 assert_eq!(first, second);
914 assert_eq!(first_receipt, second_receipt);
915
916 let parent = tempfile::tempdir().unwrap();
917 let target = parent.path().join("imported project");
918 let imported = import_portable_project(
919 &first,
920 &target,
921 Uuid::new_v4(),
922 Uuid::new_v4(),
923 &supported(&source_generation),
924 limits,
925 )
926 .unwrap();
927 assert_eq!(
928 imported.source_generation_uuid,
929 source_generation.generation_uuid()
930 );
931 assert_eq!(imported.envelope_sha256, first_receipt.envelope_sha256);
932 let reopened = resolve_project_generation(&target).unwrap();
933 assert_eq!(
934 reopened.generation_uuid(),
935 imported.publication.generation_uuid
936 );
937 assert_eq!(reopened.participant_snapshots().unwrap(), expected);
938 assert!(!target.join("trash").exists());
939 assert!(!target.join("cache").exists());
940 }
941
942 #[test]
943 fn pristine_initialized_target_is_importable() {
944 let source = tempfile::tempdir().unwrap();
945 let generation = open_or_initialize_project(source.path()).unwrap();
946 let limits = PortableProjectLimits::default();
947 let (envelope, _) = encode_portable_project(&generation, limits).unwrap();
948 let target = tempfile::tempdir().unwrap();
949 let prior = open_or_initialize_project(target.path()).unwrap();
950 let prior_uuid = prior.generation_uuid();
951 drop(prior);
952
953 let imported = import_portable_project(
954 &envelope,
955 target.path(),
956 Uuid::new_v4(),
957 Uuid::new_v4(),
958 &supported(&generation),
959 limits,
960 )
961 .unwrap();
962 assert_ne!(imported.publication.generation_uuid, prior_uuid);
963 assert_eq!(
964 resolve_project_generation(target.path())
965 .unwrap()
966 .generation_uuid(),
967 imported.publication.generation_uuid
968 );
969 }
970
971 #[test]
972 fn validation_failure_preserves_pristine_target_current() {
973 let source = tempfile::tempdir().unwrap();
974 let generation = open_or_initialize_project(source.path()).unwrap();
975 let limits = PortableProjectLimits::default();
976 let (envelope, _) = encode_portable_project(&generation, limits).unwrap();
977 let target = tempfile::tempdir().unwrap();
978 let prior = open_or_initialize_project(target.path()).unwrap();
979 let prior_uuid = prior.generation_uuid();
980 drop(prior);
981
982 let error = import_portable_project(
983 &envelope,
984 target.path(),
985 Uuid::new_v4(),
986 Uuid::new_v4(),
987 &[],
988 limits,
989 )
990 .unwrap_err();
991 assert_eq!(error.code(), "GF_UNSUPPORTED_CAPABILITY_VERSION");
992 assert_eq!(
993 resolve_project_generation(target.path())
994 .unwrap()
995 .generation_uuid(),
996 prior_uuid
997 );
998 }
999
1000 #[test]
1001 fn export_refuses_to_overwrite_existing_destination() {
1002 let source = tempfile::tempdir().unwrap();
1003 let generation = open_or_initialize_project(source.path()).unwrap();
1004 let destination = source.path().join("existing.gfx");
1005 std::fs::write(&destination, b"keep").unwrap();
1006
1007 let error =
1008 export_portable_project(&generation, &destination, PortableProjectLimits::default())
1009 .unwrap_err();
1010 assert_eq!(error.code(), "GF_UNSUPPORTED_PROJECT_FORMAT");
1011 assert_eq!(std::fs::read(destination).unwrap(), b"keep");
1012 }
1013
1014 #[cfg(unix)]
1015 #[test]
1016 fn linked_ancestor_is_rejected_for_source_export_and_target() {
1017 use std::os::unix::fs::symlink;
1018
1019 let source_project = tempfile::tempdir().unwrap();
1020 let generation = open_or_initialize_project(source_project.path()).unwrap();
1021 let limits = PortableProjectLimits::default();
1022 let (envelope, _) = encode_portable_project(&generation, limits).unwrap();
1023 let real = tempfile::tempdir().unwrap();
1024 let links = tempfile::tempdir().unwrap();
1025 let linked = links.path().join("linked");
1026 symlink(real.path(), &linked).unwrap();
1027
1028 let export_error =
1029 export_portable_project(&generation, linked.join("out.gfx"), limits).unwrap_err();
1030 assert_eq!(export_error.code(), "GF_UNSUPPORTED_FILESYSTEM");
1031 assert!(!real.path().join("out.gfx").exists());
1032
1033 std::fs::write(real.path().join("in.gfx"), &envelope).unwrap();
1034 let source_error = import_portable_project_file(
1035 linked.join("in.gfx"),
1036 real.path().join("unused-target"),
1037 Uuid::new_v4(),
1038 Uuid::new_v4(),
1039 &supported(&generation),
1040 limits,
1041 )
1042 .unwrap_err();
1043 assert_eq!(source_error.code(), "GF_UNSUPPORTED_FILESYSTEM");
1044 assert!(!real.path().join("unused-target").exists());
1045
1046 let input = links.path().join("input.gfx");
1047 std::fs::write(&input, envelope).unwrap();
1048 let target_error = import_portable_project_file(
1049 input,
1050 linked.join("target"),
1051 Uuid::new_v4(),
1052 Uuid::new_v4(),
1053 &supported(&generation),
1054 limits,
1055 )
1056 .unwrap_err();
1057 assert_eq!(target_error.code(), "GF_UNSUPPORTED_FILESYSTEM");
1058 assert!(!real.path().join("target").exists());
1059 }
1060
1061 #[test]
1062 fn corruption_and_trailing_bytes_fail_before_target_creation() {
1063 let source = tempfile::tempdir().unwrap();
1064 let generation = open_or_initialize_project(source.path()).unwrap();
1065 let limits = PortableProjectLimits::default();
1066 let (mut envelope, _) = encode_portable_project(&generation, limits).unwrap();
1067 *envelope.last_mut().unwrap() ^= 1;
1068 let parent = tempfile::tempdir().unwrap();
1069 let corrupt_target = parent.path().join("corrupt");
1070 let error = import_portable_project(
1071 &envelope,
1072 &corrupt_target,
1073 Uuid::new_v4(),
1074 Uuid::new_v4(),
1075 &supported(&generation),
1076 limits,
1077 )
1078 .unwrap_err();
1079 assert_eq!(error.code(), "GF_PROJECT_CORRUPT");
1080 assert!(!corrupt_target.exists());
1081
1082 let (mut envelope, _) = encode_portable_project(&generation, limits).unwrap();
1083 envelope.push(0);
1084 let trailing_target = parent.path().join("trailing");
1085 let error = import_portable_project(
1086 &envelope,
1087 &trailing_target,
1088 Uuid::new_v4(),
1089 Uuid::new_v4(),
1090 &supported(&generation),
1091 limits,
1092 )
1093 .unwrap_err();
1094 assert_eq!(error.code(), "GF_PROJECT_CORRUPT");
1095 assert!(!trailing_target.exists());
1096 }
1097
1098 #[test]
1099 fn resource_and_capability_checks_precede_mutation() {
1100 let source = tempfile::tempdir().unwrap();
1101 let generation = open_or_initialize_project(source.path()).unwrap();
1102 let limits = PortableProjectLimits::default();
1103 let (envelope, _) = encode_portable_project(&generation, limits).unwrap();
1104 let parent = tempfile::tempdir().unwrap();
1105
1106 let bounded_target = parent.path().join("bounded");
1107 let tiny = PortableProjectLimits {
1108 max_envelope_bytes: u64::try_from(envelope.len() - 1).unwrap(),
1109 ..limits
1110 };
1111 let error = import_portable_project(
1112 &envelope,
1113 &bounded_target,
1114 Uuid::new_v4(),
1115 Uuid::new_v4(),
1116 &supported(&generation),
1117 tiny,
1118 )
1119 .unwrap_err();
1120 assert_eq!(error.code(), "GF_RESOURCE_LIMIT");
1121 assert!(!bounded_target.exists());
1122
1123 let capability_target = parent.path().join("unsupported");
1124 let error = import_portable_project(
1125 &envelope,
1126 &capability_target,
1127 Uuid::new_v4(),
1128 Uuid::new_v4(),
1129 &[],
1130 limits,
1131 )
1132 .unwrap_err();
1133 assert_eq!(error.code(), "GF_UNSUPPORTED_CAPABILITY_VERSION");
1134 assert!(!capability_target.exists());
1135 }
1136
1137 #[test]
1138 fn nonempty_target_is_never_modified() {
1139 let source = tempfile::tempdir().unwrap();
1140 let generation = open_or_initialize_project(source.path()).unwrap();
1141 let limits = PortableProjectLimits::default();
1142 let (envelope, _) = encode_portable_project(&generation, limits).unwrap();
1143 let target = tempfile::tempdir().unwrap();
1144 std::fs::write(target.path().join("keep.txt"), b"keep").unwrap();
1145 let error = import_portable_project(
1146 &envelope,
1147 target.path(),
1148 Uuid::new_v4(),
1149 Uuid::new_v4(),
1150 &supported(&generation),
1151 limits,
1152 )
1153 .unwrap_err();
1154 assert_eq!(error.code(), "GF_UNSUPPORTED_PROJECT_FORMAT");
1155 assert_eq!(
1156 std::fs::read(target.path().join("keep.txt")).unwrap(),
1157 b"keep"
1158 );
1159 assert!(!target.path().join("CURRENT").exists());
1160 }
1161
1162 #[test]
1163 fn traversal_identity_is_rejected_before_target_creation() {
1164 let source = tempfile::tempdir().unwrap();
1165 let generation = open_or_initialize_project(source.path()).unwrap();
1166 let limits = PortableProjectLimits::default();
1167 let (envelope, _) = encode_portable_project(&generation, limits).unwrap();
1168 let prefix = MAGIC.len() + HEADER_LENGTH_BYTES;
1169 let header_length =
1170 u64::from_be_bytes(envelope[MAGIC.len()..prefix].try_into().unwrap()) as usize;
1171 let header_end = prefix + header_length;
1172 let mut header: EnvelopeHeader =
1173 serde_json::from_slice(&envelope[prefix..header_end]).unwrap();
1174 header.participants[0].record_family_id = "../escape".into();
1175 let mut header_bytes = serde_json::to_vec(&header).unwrap();
1176 header_bytes.push(b'\n');
1177 let mut malicious = Vec::new();
1178 malicious.extend_from_slice(MAGIC);
1179 malicious.extend_from_slice(&u64::try_from(header_bytes.len()).unwrap().to_be_bytes());
1180 malicious.extend_from_slice(&header_bytes);
1181 malicious.extend_from_slice(&envelope[header_end..]);
1182
1183 let parent = tempfile::tempdir().unwrap();
1184 let target = parent.path().join("traversal");
1185 let error = import_portable_project(
1186 &malicious,
1187 &target,
1188 Uuid::new_v4(),
1189 Uuid::new_v4(),
1190 &supported(&generation),
1191 limits,
1192 )
1193 .unwrap_err();
1194 assert_eq!(error.code(), "GF_PROJECT_CORRUPT");
1195 assert!(!target.exists());
1196 }
1197}