From ab61fe41f6b3115738bfab2c9ddd2c6c7ae6e2ce Mon Sep 17 00:00:00 2001 From: FeathBow Date: Mon, 31 Aug 2026 01:28:12 +0100 Subject: [PATCH] perf(core): bound pre-seal v3 state --- crates/degu-core/src/backend/held.rs | 591 +++++++++++++++++- crates/degu-core/src/backend/held/tests.rs | 107 +++- crates/degu-core/src/seal/sidecar.rs | 231 ++++++- crates/degu-core/src/seal/sidecar/scratch.rs | 609 ++++++++++++++++++- crates/degu-core/src/seal/wal.rs | 36 +- crates/degu-core/src/staging/rename.rs | 381 +++++++++++- crates/degu-core/src/staging/rename/tests.rs | 188 ++++++ 7 files changed, 2100 insertions(+), 43 deletions(-) diff --git a/crates/degu-core/src/backend/held.rs b/crates/degu-core/src/backend/held.rs index b55c7fc..4497bae 100644 --- a/crates/degu-core/src/backend/held.rs +++ b/crates/degu-core/src/backend/held.rs @@ -1035,6 +1035,50 @@ pub(crate) struct StreamedV3Inventory { manifest: HeldTreeFingerprint, } +/// Production schema-v3 proof collected before directory sealing. The complete +/// manifest and directory permission plan live only in authenticated private +/// storage; resident state is limited to the root context, fixed aggregate/count +/// fields, and the active manifest ancestor chain. +pub(crate) struct PreSealV3Inventory { + tree: StreamedV3Inventory, + directory_count: u64, +} + +/// Traversal result that still owns the historical BFS directory evidence only +/// until it is sealed into a private descriptor-only permission plan. +pub(crate) struct CollectedPreSealV3Inventory { + pending: PendingV3Inventory, + directories: Vec, +} + +/// Pre-seal traversal state after the directory plan has left resident memory. +pub(crate) struct PendingPreSealV3Inventory { + pending: PendingV3Inventory, + directory_count: u64, +} + +/// Bounded finalizer for a pre-seal scratch manifest. +pub(crate) struct PendingPreSealV3Finalizer { + pending: PendingV3Finalizer, + directory_count: u64, + observed_directories: u64, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub(crate) struct PreSealDirectoryPlanRecord { + ordinal: u64, + evidence: DirectoryEvidence, +} + +/// Allocation-bounded digest state for substituting the exact WAL-applied +/// directory modes into an authenticated pre-seal manifest. +pub(crate) struct PostSealExpectationBuilder { + digest: Sha256, + expected_entries: u64, + observed_entries: u64, + observed_directories: u64, +} + /// A deletion engine for an authenticated v3 manifest whose resident state is /// bounded by the root anchor, the current plan record, and its ancestor chain. pub(crate) struct StreamedV3Purger { @@ -2026,6 +2070,62 @@ impl PendingV3Inventory { }) } + /// Traverses schema v3 in the exact historical pre-seal breadth-first + /// order while emitting every complete record to private scratch. The + /// short-lived BFS evidence is immediately drained into an anonymous plan. + pub(crate) fn collect_pre_seal( + parent: HeldLocalBackendEvidence, + root_name: &OsStr, + protected_names: Vec, + limits: HeldTreeLimits, + mut emit_record: impl FnMut(&[u8]) -> Result<(), E>, + ) -> Result> { + let mut record = Vec::with_capacity(MANIFEST_V3_MAX_SEGMENT_PAYLOAD); + let walked = match traverse_v2_with_sink::>( + parent, + root_name, + protected_names, + limits, + TraversalConfiguration { + parent_admission: ParentAdmission::CurrentExclusive, + manifest_schema: CONTENT_PROOF_VERSION, + retain_entries: false, + }, + |entry| emit_forward_v3_record(entry, &mut record, &mut emit_record), + ) { + Ok(walked) => walked, + Err(TraversalSinkError::Tree(error)) => { + return Err(HeldTreeV3CollectError::Tree(error)); + } + Err(TraversalSinkError::Emit(error)) => return Err(error), + }; + debug_assert!(walked.entries.is_empty()); + // Do not overlap the historical traversal index with the short-lived + // vector that is immediately drained into the anonymous plan. + drop(walked.directory_index); + let directories = walked + .directories + .into_iter() + .map(|directory| directory.evidence) + .collect(); + Ok(CollectedPreSealV3Inventory { + pending: PendingV3Inventory { + context: ForwardV3Context { + parent: walked.parent, + root_name: walked.root_name, + root_identity: walked.root_identity, + root: walked.root, + backend: walked.backend, + mount_id: walked.mount_id, + protected_names: walked.protected_names, + limits: walked.limits, + }, + entry_count: walked.budget.entries, + }, + directories, + }) + } + pub(crate) fn entry_count(&self) -> u64 { self.entry_count } @@ -2063,6 +2163,113 @@ impl PendingV3Inventory { } } +impl CollectedPreSealV3Inventory { + pub(crate) fn directory_count(&self) -> u64 { + self.directories.len() as u64 + } + + pub(crate) fn emit_directory_plan( + self, + mut emit: impl FnMut(&[u8]) -> Result<(), E>, + ) -> Result> { + let directory_count = self.directories.len() as u64; + let mut record = Vec::new(); + for (ordinal, evidence) in self.directories.into_iter().enumerate() { + encode_pre_seal_directory_plan_record(ordinal as u64, &evidence, &mut record) + .map_err(HeldTreeV3CollectError::Codec)?; + emit(&record).map_err(HeldTreeV3CollectError::Emit)?; + } + Ok(PendingPreSealV3Inventory { + pending: self.pending, + directory_count, + }) + } + + #[cfg(test)] + pub(crate) fn resident_manifest_entries_for_test(&self) -> usize { + 0 + } + + #[cfg(test)] + pub(crate) fn resident_directory_entries_for_test(&self) -> usize { + self.directories.len() + } +} + +impl PendingPreSealV3Inventory { + pub(crate) fn entry_count(&self) -> u64 { + self.pending.entry_count() + } + + pub(crate) fn directory_count(&self) -> u64 { + self.directory_count + } + + pub(crate) fn into_finalizer( + self, + expected_manifest: DurableTreeManifest, + ) -> Result { + Ok(PendingPreSealV3Finalizer { + pending: self.pending.into_finalizer(expected_manifest)?, + directory_count: self.directory_count, + observed_directories: 0, + }) + } + + #[cfg(test)] + pub(crate) fn resident_manifest_entries_for_test(&self) -> usize { + 0 + } + + #[cfg(test)] + pub(crate) fn resident_directory_entries_for_test(&self) -> usize { + 0 + } +} + +impl PendingPreSealV3Finalizer { + pub(crate) fn observe( + &mut self, + record: ManifestV3Record<'_>, + emit_hardlink: &mut dyn FnMut(&[u8]) -> Result<(), E>, + ) -> Result<(), HeldTreeV3CollectError> { + if record.kind == ManifestV3RecordKind::Directory { + self.observed_directories = + self.observed_directories + .checked_add(1) + .ok_or(HeldTreeV3CollectError::Tree(HeldTreeError::Limit { + kind: HeldTreeLimit::Directories, + limit: self.pending.context.limits.max_directories, + }))?; + } + self.pending.observe(record, emit_hardlink) + } + + /// Finishes only after the caller's complete sorted-scratch fold returned + /// successfully. The scratch reader has already validated its count, codec, + /// aggregate fingerprint, run digests, and EOF before this can be called. + pub(crate) fn finish( + self, + regular_hard_links: RegularHardLinkTopology, + ) -> Result { + if self.observed_directories != self.directory_count { + return Err(HeldTreeError::PostChanged(PathBuf::new())); + } + let expected = self.pending.expected_manifest; + let authenticated = DurableTreeManifest { + schema_version: expected.schema_version, + entry_count: expected.entry_count, + sha256: expected.sha256, + }; + Ok(PreSealV3Inventory { + tree: self + .pending + .finish_manifest(authenticated, regular_hard_links)?, + directory_count: self.directory_count, + }) + } +} + fn path_is_beneath(directory: &Path, path: &Path) -> bool { if directory.as_os_str().is_empty() { return !path.as_os_str().is_empty(); @@ -2231,7 +2438,14 @@ impl PendingV3Finalizer { authenticated: AuthenticatedTreeManifest, regular_hard_links: RegularHardLinkTopology, ) -> Result { - let authenticated = authenticated.manifest(); + self.finish_manifest(authenticated.manifest(), regular_hard_links) + } + + fn finish_manifest( + self, + authenticated: DurableTreeManifest, + regular_hard_links: RegularHardLinkTopology, + ) -> Result { if authenticated.schema_version != self.expected_manifest.schema_version || authenticated.entry_count != self.expected_manifest.entry_count || authenticated.sha256 != self.expected_manifest.sha256 @@ -2569,6 +2783,310 @@ fn encode_v3_purge_plan_record( Ok(()) } +impl PreSealV3Inventory { + pub(crate) fn root_strong_identity(&self) -> StrongObjectIdentity { + StrongObjectIdentity::new_with_mount( + self.tree.context.root_identity.device, + self.tree.context.root_identity.inode, + crate::seal::wal::ObjectIncarnation::new(self.tree.context.root_identity.incarnation), + self.tree.context.mount_id, + ) + } + + #[cfg(test)] + pub(crate) fn resident_manifest_entries_for_test(&self) -> usize { + 0 + } + + #[cfg(test)] + pub(crate) fn resident_directory_entries_for_test(&self) -> usize { + 0 + } + + pub(crate) fn directory_count(&self) -> u64 { + self.directory_count + } + + pub(crate) fn validate_directory_plan_record( + &self, + record: &PreSealDirectoryPlanRecord, + expected_ordinal: u64, + ) -> Result<(), HeldTreeError> { + let evidence = &record.evidence; + if record.ordinal != expected_ordinal + || record.ordinal >= self.directory_count + || evidence.depth as u64 > self.tree.context.limits.max_depth as u64 + || evidence.owner_uid != rustix::process::geteuid().as_raw() + { + return Err(HeldTreeError::PostChanged(evidence.relative_path.clone())); + } + if record.ordinal == 0 { + if !evidence.relative_path.as_os_str().is_empty() + || evidence.depth != 0 + || evidence.identity != self.tree.context.root_identity + || evidence.owner_uid != self.tree.context.root.held.owner_uid() + || evidence.group_gid != self.tree.context.root.held.group_gid() + || evidence.observed_mode != self.tree.context.root.held.mode() + { + return Err(HeldTreeError::PostChanged(PathBuf::new())); + } + } else if evidence.relative_path.as_os_str().is_empty() || evidence.depth == 0 { + return Err(HeldTreeError::PostChanged(evidence.relative_path.clone())); + } + Ok(()) + } + + /// Retains only the target's root-to-directory evidence while an authenticated + /// descriptor-only plan is scanned. The vector is bounded by max_depth, never + /// by the number of directories in the tree. + pub(crate) fn consider_directory_plan_ancestor( + &self, + target: &PreSealDirectoryPlanRecord, + candidate: PreSealDirectoryPlanRecord, + chain: &mut Vec, + ) -> Result<(), HeldTreeError> { + let target_path = &target.evidence.relative_path; + let candidate_path = &candidate.evidence.relative_path; + let is_ancestor = candidate_path.as_os_str().is_empty() + || candidate_path == target_path + || path_is_beneath(candidate_path, target_path); + if !is_ancestor { + return Ok(()); + } + if candidate.evidence.depth as usize != chain.len() + || candidate.evidence.depth > target.evidence.depth + || (candidate.evidence.depth > 0 + && candidate_path.parent() + != chain + .last() + .map(|record| record.evidence.relative_path.as_path())) + { + return Err(HeldTreeError::PostChanged(target_path.clone())); + } + chain.push(candidate); + Ok(()) + } + + fn validate_directory_plan_chain( + &self, + target: &PreSealDirectoryPlanRecord, + chain: &[PreSealDirectoryPlanRecord], + ) -> Result<(), HeldTreeError> { + if chain.len() != target.evidence.depth as usize + 1 + || chain + .last() + .is_none_or(|record| record.evidence != target.evidence) + { + return Err(HeldTreeError::PostChanged( + target.evidence.relative_path.clone(), + )); + } + Ok(()) + } + + fn reopen_directory_for_transient_seal<'a>( + &'a self, + relative_path: &Path, + chain: &'a [PreSealDirectoryPlanRecord], + ) -> Result, HeldTreeError> { + reopen_directory_from_root( + &self.tree.context.root, + relative_path, + |candidate| { + chain + .iter() + .find(|record| record.evidence.relative_path == candidate) + .map(|record| &record.evidence) + }, + self.tree.context.backend, + self.tree.context.mount_id, + || { + verify_root_binding_fields( + &self.tree.context.parent, + &self.tree.context.root_name, + self.tree.context.root_identity, + self.tree.context.mount_id, + self.tree.context.backend, + ) + }, + true, + ) + } + + /// Applies one authenticated reverse-BFS directory record. The caller has + /// already authenticated the complete anonymous plan and supplies only this + /// record's bounded root-to-target chain. + #[allow(clippy::too_many_arguments)] + pub(crate) fn seal_directory_for_staging( + &mut self, + wal: &mut SealWal, + transaction: TransactionId, + source_root: &Path, + filesystem_id: &str, + mutation_id: u64, + target: &PreSealDirectoryPlanRecord, + chain: &[PreSealDirectoryPlanRecord], + ) -> Result { + self.validate_directory_plan_chain(target, chain) + .map_err(HeldTreeSealError::Tree)?; + let evidence = &target.evidence; + let relative_path = source_root.join(&evidence.relative_path); + let path = evidence.relative_path.clone(); + let result = if evidence.depth == 0 { + execute_staging_local_mode_mutation( + wal, + &mut self.tree.context.root.held, + LocalModeMutationRequest { + transaction, + mutation_id, + locator: RecoveryLocator::held_staging( + relative_path, + filesystem_id.to_owned(), + evidence.identity.incarnation, + ), + transform: LocalModeTransform::Seal { + acquire_owner_write_search: false, + }, + }, + ) + } else { + let mut reopened = self + .reopen_directory_for_transient_seal(&evidence.relative_path, chain) + .map_err(HeldTreeSealError::Tree)?; + let held = match &mut reopened.held { + ReopenedHeldDirectory::Descendant(held) => held, + ReopenedHeldDirectory::Root(_) => { + return Err(HeldTreeSealError::Mutation { + path, + source: LocalModeExecutionError::InvalidRequest( + "reopener returned retained root unexpectedly", + ), + }); + } + }; + let result = execute_staging_local_mode_mutation( + wal, + held, + LocalModeMutationRequest { + transaction, + mutation_id, + locator: RecoveryLocator::held_staging( + relative_path, + filesystem_id.to_owned(), + evidence.identity.incarnation, + ), + transform: LocalModeTransform::Seal { + acquire_owner_write_search: false, + }, + }, + ); + drop(reopened); + result + } + .map_err(|source| HeldTreeSealError::Mutation { + path: evidence.relative_path.clone(), + source, + })?; + match result { + LocalModeMutationResult::Applied { applied_mode, .. } => Ok(applied_mode), + LocalModeMutationResult::ConfirmedNotApplied => Err( + HeldTreeSealError::ConfirmedNotApplied(evidence.relative_path.clone()), + ), + } + } + + pub(crate) fn post_seal_expectation_builder(&self) -> PostSealExpectationBuilder { + let mut digest = Sha256::new(); + digest.update(MANIFEST_DOMAIN_V3); + digest.update(self.tree.manifest.entry_count.to_be_bytes()); + PostSealExpectationBuilder { + digest, + expected_entries: self.tree.manifest.entry_count, + observed_entries: 0, + observed_directories: 0, + } + } + + pub(crate) fn finish_post_seal_expectation( + self, + builder: PostSealExpectationBuilder, + ) -> Result { + if builder.observed_entries != builder.expected_entries + || builder.observed_directories != self.directory_count + { + return Err(HeldTreeError::PostChanged(PathBuf::new())); + } + Ok(PostSealManifestExpectation { + root_identity: self.tree.context.root_identity, + backend: self.tree.context.backend, + mount_id: self.tree.context.mount_id, + fingerprint: HeldTreeFingerprint { + schema_version: CONTENT_PROOF_VERSION, + entry_count: builder.expected_entries, + sha256: builder.digest.finalize().into(), + }, + }) + } +} + +impl PostSealExpectationBuilder { + pub(crate) fn observe( + &mut self, + inventory: &PreSealV3Inventory, + record: ManifestV3Record<'_>, + directory_mode: impl FnOnce(&Path, u64, u64, u64) -> Option, + ) -> Result<(), HeldTreeError> { + self.observed_entries = + self.observed_entries + .checked_add(1) + .ok_or(HeldTreeError::Limit { + kind: HeldTreeLimit::Entries, + limit: inventory.tree.context.limits.max_entries, + })?; + let mode = if record.kind == ManifestV3RecordKind::Directory { + self.observed_directories = + self.observed_directories + .checked_add(1) + .ok_or(HeldTreeError::Limit { + kind: HeldTreeLimit::Directories, + limit: inventory.tree.context.limits.max_directories, + })?; + let path = Path::new(OsStr::from_bytes(record.path)); + directory_mode(path, record.device, record.inode, record.incarnation) + .ok_or_else(|| HeldTreeError::PostChanged(path.to_path_buf()))? + } else { + record.mode + }; + update_manifest_v3_digest_with_mode(&mut self.digest, record, mode) + } +} + +fn update_manifest_v3_digest_with_mode( + digest: &mut Sha256, + record: ManifestV3Record<'_>, + mode: u32, +) -> Result<(), HeldTreeError> { + let mode_offset = 8_usize + .checked_add(record.path.len()) + .and_then(|offset| offset.checked_add(1 + 8 + 8 + 8 + 4 + 4)) + .ok_or_else(|| HeldTreeError::PostChanged(PathBuf::new()))?; + let suffix_offset = mode_offset + .checked_add(4) + .ok_or_else(|| HeldTreeError::PostChanged(PathBuf::new()))?; + let prefix = record + .encoded + .get(..mode_offset) + .ok_or_else(|| HeldTreeError::PostChanged(PathBuf::new()))?; + let suffix = record + .encoded + .get(suffix_offset..) + .ok_or_else(|| HeldTreeError::PostChanged(PathBuf::new()))?; + digest.update(prefix); + digest.update(mode.to_be_bytes()); + digest.update(suffix); + Ok(()) +} + impl StreamedV3Inventory { pub(crate) fn fingerprint(&self) -> HeldTreeFingerprint { self.manifest @@ -2863,6 +3381,77 @@ impl StreamedV3Purger { } const STRUCTURE_SCRATCH_RECORD_MAGIC: &[u8; 4] = b"DHS1"; +const PRE_SEAL_DIRECTORY_PLAN_MAGIC: &[u8; 4] = b"DHDP"; + +fn encode_pre_seal_directory_plan_record( + ordinal: u64, + evidence: &DirectoryEvidence, + encoded: &mut Vec, +) -> Result<(), ManifestV3CodecError> { + let path = evidence.relative_path.as_os_str().as_bytes(); + let record_len = 60_usize + .checked_add(path.len()) + .ok_or(ManifestV3CodecError::LengthOverflow)?; + if record_len > MANIFEST_V3_MAX_SEGMENT_PAYLOAD { + return Err(ManifestV3CodecError::RecordTooLarge); + } + encoded.clear(); + encoded.reserve(record_len); + encoded.extend_from_slice(PRE_SEAL_DIRECTORY_PLAN_MAGIC); + encoded.extend_from_slice(&ordinal.to_be_bytes()); + encoded.extend_from_slice(&(path.len() as u64).to_be_bytes()); + encoded.extend_from_slice(path); + encoded.extend_from_slice(&evidence.depth.to_be_bytes()); + encoded.extend_from_slice(&evidence.identity.device.to_be_bytes()); + encoded.extend_from_slice(&evidence.identity.inode.to_be_bytes()); + encoded.extend_from_slice(&evidence.identity.incarnation.to_be_bytes()); + encoded.extend_from_slice(&evidence.owner_uid.to_be_bytes()); + encoded.extend_from_slice(&evidence.group_gid.to_be_bytes()); + encoded.extend_from_slice(&evidence.observed_mode.to_be_bytes()); + debug_assert_eq!(encoded.len(), record_len); + Ok(()) +} + +pub(crate) fn decode_pre_seal_directory_plan_record( + mut record: &[u8], +) -> Result { + if record.len() > MANIFEST_V3_MAX_SEGMENT_PAYLOAD + || take(&mut record, PRE_SEAL_DIRECTORY_PLAN_MAGIC.len())? != PRE_SEAL_DIRECTORY_PLAN_MAGIC + { + return Err(ManifestV3CodecError::InvalidTag); + } + let ordinal = take_u64(&mut record)?; + let path_len = usize::try_from(take_u64(&mut record)?) + .map_err(|_| ManifestV3CodecError::LengthOverflow)?; + let path_bytes = take(&mut record, path_len)?; + validate_manifest_path(path_bytes, HeldTreeLimits::default().max_depth)?; + let depth = u32::from_be_bytes(take(&mut record, 4)?.try_into().unwrap()); + let path = PathBuf::from(OsStr::from_bytes(path_bytes)); + if path.components().count() != depth as usize { + return Err(ManifestV3CodecError::InvalidPath); + } + let evidence = DirectoryEvidence { + relative_path: path, + depth, + identity: NodeIdentity { + kind: NodeKind::Directory, + device: take_u64(&mut record)?, + inode: take_u64(&mut record)?, + incarnation: take_u64(&mut record)?, + }, + owner_uid: u32::from_be_bytes(take(&mut record, 4)?.try_into().unwrap()), + group_gid: u32::from_be_bytes(take(&mut record, 4)?.try_into().unwrap()), + observed_mode: u32::from_be_bytes(take(&mut record, 4)?.try_into().unwrap()), + }; + if evidence.observed_mode > 0o7777 || !record.is_empty() { + return Err(if !record.is_empty() { + ManifestV3CodecError::TrailingBytes + } else { + ManifestV3CodecError::InvalidMode + }); + } + Ok(PreSealDirectoryPlanRecord { ordinal, evidence }) +} pub(crate) fn hardlink_scratch_sentinel_record() -> &'static [u8] { &[0, 0, 0, 0, 0, 0, 0, 1, 0] diff --git a/crates/degu-core/src/backend/held/tests.rs b/crates/degu-core/src/backend/held/tests.rs index 8d35311..4192d88 100644 --- a/crates/degu-core/src/backend/held/tests.rs +++ b/crates/degu-core/src/backend/held/tests.rs @@ -1491,10 +1491,37 @@ fn v3_directory_mode_projection_hashes_in_place_without_changing_codec_bytes() { } } + let expected = fingerprint_manifest_v3(&projected); assert_eq!( inventory.fingerprint_with_directory_modes(&modes).unwrap(), - fingerprint_manifest_v3(&projected) + expected ); + + let mut streamed_digest = Sha256::new(); + streamed_digest.update(MANIFEST_DOMAIN_V3); + streamed_digest.update((inventory.manifest.len() as u64).to_be_bytes()); + let mut decoder = ManifestV3Decoder::new(inventory.manifest.len() as u64).unwrap(); + for entry in &inventory.manifest { + let mut encoded = Vec::new(); + emit_manifest_entry_v3(entry, |bytes| encoded.extend_from_slice(bytes)); + decoder + .push_segment_with(1, &encoded, |record| { + let mode = modes.get(&entry.path).copied().unwrap_or(record.mode); + update_manifest_v3_digest_with_mode(&mut streamed_digest, record, mode) + }) + .unwrap(); + } + decoder.finish().unwrap(); + assert_eq!( + HeldTreeFingerprint { + schema_version: CONTENT_PROOF_VERSION, + entry_count: inventory.manifest.len() as u64, + sha256: streamed_digest.finalize().into(), + }, + expected, + "streamed pre-seal mode substitution must remain byte-identical to the v3 codec" + ); + let mut missing = modes.clone(); missing.remove(Path::new("a")); assert!(matches!( @@ -2322,6 +2349,84 @@ fn rewalk_rejects_acl_planted_after_collect() { )); } +#[test] +fn pre_seal_v3_stream_moves_the_historical_bfs_permission_plan_out_of_resident_state() { + let (temp, root) = setup_tree(); + for path in ["z", "z/left", "a", "a/right", "z/left/deep"] { + std::fs::create_dir_all(root.join(path)).unwrap(); + } + for index in 0_u32..256 { + std::fs::write( + root.join(format!("payload-{index:04}")), + index.to_be_bytes(), + ) + .unwrap(); + } + + let baseline = collect(&temp, vec![], HeldTreeLimits::default()).unwrap(); + let expected_directories = baseline.directories.clone(); + let mut expected_records = Vec::new(); + baseline + .stream_manifest_v3_records(|record| { + expected_records.push(record.to_vec()); + Ok::<(), std::convert::Infallible>(()) + }) + .unwrap(); + + let mut streamed_records = Vec::new(); + let collected = PendingV3Inventory::collect_pre_seal( + certify_held_fd(open_directory(temp.path())).unwrap(), + OsStr::new("root"), + vec![], + HeldTreeLimits::default(), + |record| { + streamed_records.push(record.to_vec()); + Ok::<(), std::convert::Infallible>(()) + }, + ) + .unwrap(); + streamed_records.sort_by(|left, right| { + let left_len = u64::from_be_bytes(left[..8].try_into().unwrap()) as usize; + let right_len = u64::from_be_bytes(right[..8].try_into().unwrap()) as usize; + compare_manifest_paths(&left[8..8 + left_len], &right[8..8 + right_len]) + }); + + assert_eq!(collected.resident_manifest_entries_for_test(), 0); + assert_eq!( + collected.resident_directory_entries_for_test(), + expected_directories.len() + ); + let mut encoded_plan = Vec::new(); + let pending = collected + .emit_directory_plan(|record| { + encoded_plan.push(record.to_vec()); + Ok::<(), std::convert::Infallible>(()) + }) + .unwrap(); + let decoded_plan = encoded_plan + .iter() + .map(|record| decode_pre_seal_directory_plan_record(record).unwrap()) + .collect::>(); + assert_eq!(pending.resident_manifest_entries_for_test(), 0); + assert_eq!(pending.resident_directory_entries_for_test(), 0); + assert_eq!(pending.entry_count(), baseline.entry_count()); + assert_eq!(pending.directory_count(), expected_directories.len() as u64); + assert_eq!( + decoded_plan + .iter() + .map(|record| record.evidence.clone()) + .collect::>(), + expected_directories + ); + assert!( + decoded_plan + .iter() + .enumerate() + .all(|(ordinal, record)| record.ordinal == ordinal as u64) + ); + assert_eq!(streamed_records, expected_records); +} + #[test] #[allow(clippy::disallowed_methods)] fn pending_v3_streamed_finalizer_matches_resident_inventory_oracle() { diff --git a/crates/degu-core/src/seal/sidecar.rs b/crates/degu-core/src/seal/sidecar.rs index 12ef4c4..6cde074 100644 --- a/crates/degu-core/src/seal/sidecar.rs +++ b/crates/degu-core/src/seal/sidecar.rs @@ -26,6 +26,8 @@ use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicU64, Ordering}; mod scratch; +#[cfg(test)] +pub(crate) use scratch::TreeDirectoryPlan; pub(crate) use scratch::{ TreeManifestScratchBuildError, TreePurgePlan, TreeStructureScratchCursor, }; @@ -2248,7 +2250,7 @@ mod tests { expected, &mut scratch, Vec::new(), - |mut visited, record| { + |mut visited, record, _wal| { visited.push((record.path.to_vec(), record.kind)); Ok::<_, std::convert::Infallible>(visited) }, @@ -2265,6 +2267,229 @@ mod tests { assert_eq!(store.cleanup_unpublished(&mut wal).unwrap(), 1); } + #[test] + fn directory_plan_is_anonymous_authenticated_and_exactly_reverse_bfs() { + let (_temp, root, store, mut wal) = fixture(); + let transaction = tx(44); + let records = vec![ + b"root\xff".to_vec(), + b"child-a\x80".to_vec(), + b"child-b\xfe".to_vec(), + ]; + let (mut plan, output) = store + .build_directory_plan_with_output(&mut wal, transaction, records.len() as u64, |emit| { + for record in &records { + emit(record)?; + } + Ok::<_, TreeSidecarError>(91_u64) + }) + .unwrap(); + assert_eq!(output, 91); + assert_eq!(plan.link_count_for_test(), 0); + assert!( + std::fs::read_dir(&root) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .all(|name| !name.to_string_lossy().starts_with(".tree-scratch-v1-")), + "directory plan must be descriptor-only before it is returned" + ); + plan.authenticate().unwrap(); + let mut forward = Vec::new(); + plan.for_each_forward(|record| { + forward.push(record.to_vec()); + Ok::<(), std::convert::Infallible>(()) + }) + .unwrap(); + assert_eq!(forward, records); + let mut reverse = Vec::new(); + while let Some(record) = plan.next_reverse().unwrap() { + reverse.push(record); + } + assert_eq!(reverse, records.into_iter().rev().collect::>()); + plan.finish().unwrap(); + } + + #[test] + fn directory_plan_producer_failure_leaves_no_named_residue() { + #[derive(Debug, Eq, PartialEq)] + struct Stop; + + let (_temp, root, store, mut wal) = fixture(); + let transaction = tx(46); + let result = store.build_directory_plan_with_output(&mut wal, transaction, 2, |emit| { + emit(b"root").map_err(|_| Stop)?; + Err::<(), _>(Stop) + }); + assert!(matches!( + result, + Err(TreeManifestScratchBuildError::Produce(Stop)) + )); + assert!( + std::fs::read_dir(root) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .all(|name| !name.to_string_lossy().starts_with(".tree-scratch-v1-")) + ); + } + + #[test] + fn directory_plan_frame_tamper_after_preflight_fails_before_use() { + let (_temp, root, store, mut wal) = fixture(); + let transaction = tx(45); + let mut plan = store + .build_directory_plan_with_output(&mut wal, transaction, 2, |emit| { + emit(b"root")?; + emit(b"child")?; + Ok::<(), TreeSidecarError>(()) + }) + .unwrap() + .0; + plan.authenticate().unwrap(); + plan.corrupt_frame_for_test(1); + assert!(matches!( + plan.next_reverse(), + Err(TreeSidecarError::InvalidScratch( + "directory plan frame authentication failed" + )) + )); + assert!( + std::fs::read_dir(root) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .all(|name| !name.to_string_lossy().starts_with(".tree-scratch-v1-")) + ); + } + + #[test] + fn pre_seal_hardlink_fold_preserves_then_discards_only_manifest_scratch() { + let (_temp, root, store, mut wal) = fixture(); + let transaction = tx(42); + let records = vec![ + v3_regular_record(b"z", 3), + v3_directory_record(b"", 1), + v3_regular_record(b"a", 2), + ]; + let (_, expected) = sorted_v3_manifest(&records); + let mut manifest_scratch = store + .build_sorted_manifest_scratch(&mut wal, transaction, |emit| { + for record in &records { + emit(record)?; + } + Ok(()) + }) + .unwrap(); + let (hardlink_scratch, visited) = store + .build_sorted_hardlink_scratch_from_manifest( + &mut wal, + transaction, + &mut manifest_scratch, + expected, + Vec::new(), + |visited, record, _emit| { + visited.push(record.path.to_vec()); + Ok::<(), std::convert::Infallible>(()) + }, + ) + .unwrap(); + assert_eq!(visited, [b"".to_vec(), b"a".to_vec(), b"z".to_vec()]); + assert_eq!( + std::fs::read_dir(&root) + .unwrap() + .filter(|entry| { + entry + .as_ref() + .unwrap() + .file_name() + .to_string_lossy() + .starts_with(".tree-scratch-v1-") + }) + .count(), + 2 + ); + + store + .fold_sorted_hardlink_scratch_preserving_manifest( + &mut wal, + transaction, + hardlink_scratch, + 0_u64, + |_count, _record| -> Result<(), std::convert::Infallible> { + unreachable!("sentinel-only hardlink scratch has no fold record") + }, + ) + .unwrap(); + assert_eq!( + store + .fingerprint_sorted_manifest_scratch(&mut wal, transaction, &mut manifest_scratch) + .unwrap(), + expected + ); + store + .discard_sorted_manifest_scratch(&mut wal, transaction, expected, manifest_scratch) + .unwrap(); + assert_eq!(store.cleanup_unpublished(&mut wal).unwrap(), 0); + } + + #[test] + fn corrupt_pre_seal_hardlink_scratch_preserves_manifest_until_explicit_cleanup() { + let (_temp, root, store, mut wal) = fixture(); + let transaction = tx(43); + let records = vec![v3_directory_record(b"", 1), v3_regular_record(b"a", 2)]; + let (_, expected) = sorted_v3_manifest(&records); + let mut manifest_scratch = store + .build_sorted_manifest_scratch(&mut wal, transaction, |emit| { + for record in &records { + emit(record)?; + } + Ok(()) + }) + .unwrap(); + let (hardlink_scratch, ()) = store + .build_sorted_hardlink_scratch_from_manifest( + &mut wal, + transaction, + &mut manifest_scratch, + expected, + (), + |(), _record, _emit| Ok::<(), std::convert::Infallible>(()), + ) + .unwrap(); + let hardlink_run = root.join(&hardlink_scratch.run_names_for_test()[0]); + let length = std::fs::metadata(&hardlink_run).unwrap().len(); + std::fs::OpenOptions::new() + .write(true) + .open(&hardlink_run) + .unwrap() + .set_len(length - 1) + .unwrap(); + + let error = store + .fold_sorted_hardlink_scratch_preserving_manifest( + &mut wal, + transaction, + hardlink_scratch, + (), + |(), _record| Ok::<(), std::convert::Infallible>(()), + ) + .unwrap_err(); + assert!(matches!(error, TreeSidecarFoldError::Sidecar(_))); + assert_eq!( + store + .fingerprint_sorted_manifest_scratch(&mut wal, transaction, &mut manifest_scratch) + .unwrap(), + expected, + "hardlink corruption must not consume the separately authenticated manifest" + ); + assert!(store.cleanup_unpublished(&mut wal).unwrap() >= 2); + assert!( + std::fs::read_dir(root) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .all(|name| !name.to_string_lossy().starts_with(".tree-scratch-v1-")) + ); + assert_eq!(wal.tree_sidecar_commitment(transaction), None); + } + #[test] fn sorted_scratch_preserves_historical_component_path_order() { let (_temp, _root, store, mut wal) = fixture(); @@ -2414,7 +2639,7 @@ mod tests { expected, &mut scratch, (), - |(), _| { + |(), _, _wal| { visits.set(visits.get() + 1); Err(Stop) }, @@ -2459,7 +2684,7 @@ mod tests { expected, &mut scratch, Vec::>::new(), - |mut paths, record| { + |mut paths, record, _wal| { visits.set(visits.get() + 1); paths.push(record.path.to_vec()); Ok::<_, std::convert::Infallible>(paths) diff --git a/crates/degu-core/src/seal/sidecar/scratch.rs b/crates/degu-core/src/seal/sidecar/scratch.rs index f09aca1..35f77d0 100644 --- a/crates/degu-core/src/seal/sidecar/scratch.rs +++ b/crates/degu-core/src/seal/sidecar/scratch.rs @@ -7,13 +7,12 @@ //! bound digest, and private-file checks make same-UID replacement fail closed. use super::*; -#[cfg(test)] -use crate::backend::held::ManifestV3Record; use crate::backend::held::{ - ManifestV3Decoder, ManifestV3VisitError, StructureEvidence, compare_manifest_paths, - decode_structure_record, + ManifestV3Decoder, ManifestV3Record, ManifestV3VisitError, StructureEvidence, + compare_manifest_paths, decode_structure_record, }; use crate::seal::wal::DurableTreeManifest; +use std::os::unix::fs::FileExt; const SCRATCH_MAGIC: &[u8; 4] = b"DHSR"; const SCRATCH_VERSION: u16 = 1; @@ -36,6 +35,13 @@ const PURGE_PLAN_VERSION: u16 = 1; const PURGE_PLAN_HEADER_LEN: usize = 72; const PURGE_PLAN_DOMAIN: &[u8] = b"degu-held-tree-purge-plan-v1\0"; const PURGE_PLAN_FRAME_DOMAIN: &[u8] = b"degu-held-tree-purge-frame-v1\0"; +const DIRECTORY_PLAN_MAGIC: &[u8; 4] = b"DHDP"; +const DIRECTORY_PLAN_VERSION: u16 = 1; +const DIRECTORY_PLAN_HEADER_LEN: usize = 72; +const DIRECTORY_PLAN_DOMAIN: &[u8] = b"degu-held-tree-directory-plan-v1\0"; +const DIRECTORY_PLAN_FRAME_DOMAIN: &[u8] = b"degu-held-tree-directory-frame-v1\0"; +const DIRECTORY_PLAN_MAX_RECORD_BYTES: usize = MAX_SEGMENT_PAYLOAD; +const DIRECTORY_PLAN_MAX_TOTAL_BYTES: u64 = MAX_TOTAL_PAYLOAD_BYTES + MAX_RECORDS * 40; static NEXT_SCRATCH_NAME: AtomicU64 = AtomicU64::new(1); #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -130,6 +136,23 @@ pub(crate) struct TreeHardlinkScratch(TreeManifestScratch); #[derive(Debug)] pub(crate) struct TreePurgeScratch(TreeManifestScratch); +/// One-shot pre-seal directory plan. The backing file is unlinked before any +/// record is written; only this keyed descriptor can enumerate its contents. +#[derive(Debug)] +pub(crate) struct TreeDirectoryPlan { + file: File, + path: PathBuf, + transaction: TransactionId, + frame_key: [u8; 32], + expected_records: u64, + expected_payload_bytes: u64, + expected_digest: [u8; 32], + file_bytes: u64, + authenticated: bool, + reverse_remaining: u64, + reverse_offset: u64, +} + pub(crate) struct TreeStructureScratchCursor { scratch: TreeManifestScratch, readers: Vec, @@ -317,6 +340,83 @@ impl TreeSidecarStore { .map(|(scratch, output)| (TreeHardlinkScratch(scratch), output)) } + /// Folds an authenticated sorted manifest while spooling final regular-file + /// observations into a second fixed-memory identity-sorted scratch. Keeping + /// both builders inside this lease-bound method avoids any resident hardlink + /// inventory and avoids publishing a pre-seal sidecar. + pub(crate) fn build_sorted_hardlink_scratch_from_manifest( + &self, + wal: &mut SealWal, + transaction: TransactionId, + manifest_scratch: &mut TreeManifestScratch, + expected_manifest: DurableTreeManifest, + initial: A, + mut fold: F, + ) -> Result<(TreeHardlinkScratch, A), TreeSidecarFoldError> + where + F: FnMut( + &mut A, + ManifestV3Record<'_>, + &mut dyn FnMut(&[u8]) -> Result<(), TreeSidecarError>, + ) -> Result<(), E>, + { + self.require_matching_wal(wal)?; + self.revalidate_store_binding()?; + let mut builder = RunBuilder { + store: self, + transaction, + memory_bytes: SORT_MEMORY_BYTES, + order: ScratchOrder::HardlinkIdentityThenPath, + arena: Vec::with_capacity(SORT_MEMORY_BYTES), + records: Vec::new(), + runs: Vec::new(), + record_count: 0, + payload_bytes: 0, + }; + builder.push(crate::backend::held::hardlink_scratch_sentinel_record())?; + let output = self.read_sorted_manifest_scratch( + wal, + transaction, + expected_manifest, + manifest_scratch, + initial, + |mut accumulator, record, _wal| { + let mut emit = |bytes: &[u8]| builder.push(bytes); + fold(&mut accumulator, record, &mut emit)?; + Ok(accumulator) + }, + )?; + let scratch = TreeHardlinkScratch(builder.finish()?); + self.require_matching_wal(wal)?; + self.revalidate_store_binding()?; + Ok((scratch, output)) + } + + /// Seals BFS directory records into an anonymous keyed plan. The producer's + /// resident traversal vector is consumed before this method returns. + pub(crate) fn build_directory_plan_with_output( + &self, + wal: &mut SealWal, + transaction: TransactionId, + expected_records: u64, + produce: P, + ) -> Result<(TreeDirectoryPlan, T), TreeManifestScratchBuildError> + where + P: FnOnce(&mut dyn FnMut(&[u8]) -> Result<(), TreeSidecarError>) -> Result, + { + self.require_matching_wal(wal)?; + self.revalidate_store_binding()?; + let mut writer = DirectoryPlanWriter::create(self, transaction)?; + let output = { + let mut emit = |record: &[u8]| writer.push(record); + produce(&mut emit).map_err(TreeManifestScratchBuildError::Produce)? + }; + let plan = writer.finish(expected_records)?; + self.require_matching_wal(wal)?; + self.revalidate_store_binding()?; + Ok((plan, output)) + } + /// Builds private purge runs in deepest-first historical postorder. pub(crate) fn build_sorted_purge_scratch_with_output( &self, @@ -532,12 +632,60 @@ impl TreeSidecarStore { /// ordering, aggregate count, cleanup, and store/WAL binding all validate /// before a caller fold error can be returned. pub(crate) fn fold_sorted_hardlink_scratch( + &self, + wal: &mut SealWal, + transaction: TransactionId, + scratch: TreeHardlinkScratch, + accumulator: A, + fold: F, + ) -> Result> + where + F: FnMut(&mut A, &[u8]) -> Result<(), E>, + { + self.fold_sorted_hardlink_scratch_with_cleanup( + wal, + transaction, + scratch, + accumulator, + fold, + true, + ) + } + + /// Pre-seal manifest scratch must survive this fold until its exact + /// WAL-applied directory modes have been substituted. On success only the + /// hardlink runs are removed; callers still perform store-wide cleanup on + /// every error path. + pub(crate) fn fold_sorted_hardlink_scratch_preserving_manifest( + &self, + wal: &mut SealWal, + transaction: TransactionId, + scratch: TreeHardlinkScratch, + accumulator: A, + fold: F, + ) -> Result> + where + F: FnMut(&mut A, &[u8]) -> Result<(), E>, + { + self.fold_sorted_hardlink_scratch_with_cleanup( + wal, + transaction, + scratch, + accumulator, + fold, + false, + ) + } + + #[allow(clippy::too_many_arguments)] + fn fold_sorted_hardlink_scratch_with_cleanup( &self, wal: &mut SealWal, transaction: TransactionId, scratch: TreeHardlinkScratch, mut accumulator: A, mut fold: F, + cleanup_all_unpublished: bool, ) -> Result> where F: FnMut(&mut A, &[u8]) -> Result<(), E>, @@ -604,7 +752,11 @@ impl TreeSidecarStore { } Ok(accumulator) })(); - let cleanup = self.cleanup_unpublished(wal).map(|_| ()); + let cleanup = if cleanup_all_unpublished { + self.cleanup_unpublished(wal).map(|_| ()) + } else { + Ok(()) + }; match (validation, cleanup) { (Err(TreeSidecarFoldError::Sidecar(primary)), _) => { Err(TreeSidecarFoldError::Sidecar(primary)) @@ -660,13 +812,38 @@ impl TreeSidecarStore { Ok(manifest) } + /// Authenticates and removes one unpublished manifest scratch without ever + /// publishing it or making it WAL-referenceable. The containing directory is + /// synced before success so later post-seal collection starts from a clean + /// scratch namespace. + pub(crate) fn discard_sorted_manifest_scratch( + &self, + wal: &mut SealWal, + transaction: TransactionId, + expected_manifest: DurableTreeManifest, + mut scratch: TreeManifestScratch, + ) -> Result<(), TreeSidecarError> { + let actual = self.fingerprint_sorted_manifest_scratch(wal, transaction, &mut scratch)?; + if actual != expected_manifest { + return Err(TreeSidecarError::InvalidScratch( + "discarded scratch manifest fingerprint changed", + )); + } + self.require_scratch_binding(transaction, &scratch, ScratchOrder::ManifestPath)?; + self.require_matching_wal(wal)?; + self.revalidate_store_binding()?; + self.remove_runs(&scratch.runs)?; + rustix::fs::fsync(&self.directory).map_err(|error| io_error(&self.path, error.into()))?; + self.require_matching_wal(wal)?; + self.revalidate_store_binding() + } + /// Authenticates and globally decodes every sorted scratch record, folding /// borrowed typed records into owned authority-neutral data. The accumulator /// is returned only after all run identities/digests, global v3 ordering, /// aggregate fingerprint, root/parent constraints, and EOF checks pass. /// A fold error consumes partial decoder state; callers must discard any /// observations made before the error and may not treat them as evidence. - #[cfg(test)] pub(crate) fn read_sorted_manifest_scratch( &self, wal: &mut SealWal, @@ -677,7 +854,7 @@ impl TreeSidecarStore { mut fold: F, ) -> Result> where - F: FnMut(A, ManifestV3Record<'_>) -> Result, + F: FnMut(A, ManifestV3Record<'_>, &SealWal) -> Result, { self.require_scratch_binding(transaction, scratch, ScratchOrder::ManifestPath)?; if expected_manifest.entry_count != scratch.record_count { @@ -698,7 +875,7 @@ impl TreeSidecarStore { let current = accumulator .take() .expect("scratch fold accumulator is always present"); - accumulator = Some(fold(current, typed)?); + accumulator = Some(fold(current, typed, wal)?); Ok(()) }) { Ok(()) => Ok(()), @@ -1348,6 +1525,421 @@ impl TreeStructureScratchCursor { } } +fn directory_plan_digest(transaction: TransactionId) -> Sha256 { + let mut digest = Sha256::new(); + digest.update(DIRECTORY_PLAN_DOMAIN); + digest.update(transaction.0); + digest +} + +fn directory_plan_frame_tag( + key: &[u8; 32], + transaction: TransactionId, + ordinal: u64, + length: [u8; 4], + record: &[u8], +) -> [u8; 32] { + let mut inner_pad = [0x36_u8; 64]; + let mut outer_pad = [0x5c_u8; 64]; + for (pad, byte) in inner_pad.iter_mut().zip(key) { + *pad ^= byte; + } + for (pad, byte) in outer_pad.iter_mut().zip(key) { + *pad ^= byte; + } + let mut inner = Sha256::new(); + inner.update(inner_pad); + inner.update(DIRECTORY_PLAN_FRAME_DOMAIN); + inner.update(transaction.0); + inner.update(ordinal.to_be_bytes()); + inner.update(length); + inner.update(record); + let inner: [u8; 32] = inner.finalize().into(); + let mut outer = Sha256::new(); + outer.update(outer_pad); + outer.update(inner); + outer.finalize().into() +} + +fn plan_read_exact_at( + file: &File, + path: &Path, + offset: u64, + bytes: &mut [u8], +) -> Result<(), TreeSidecarError> { + file.read_exact_at(bytes, offset) + .map_err(|error| io_error(path, error)) +} + +impl TreeDirectoryPlan { + fn validate_header(&self) -> Result<(), TreeSidecarError> { + let mut header = [0_u8; DIRECTORY_PLAN_HEADER_LEN]; + plan_read_exact_at(&self.file, &self.path, 0, &mut header)?; + if &header[0..4] != DIRECTORY_PLAN_MAGIC + || u16::from_be_bytes(header[4..6].try_into().unwrap()) != DIRECTORY_PLAN_VERSION + || u16::from_be_bytes(header[6..8].try_into().unwrap()) as usize + != DIRECTORY_PLAN_HEADER_LEN + || header[8..24] != self.transaction.0 + || u64::from_be_bytes(header[24..32].try_into().unwrap()) != self.expected_records + || u64::from_be_bytes(header[32..40].try_into().unwrap()) != self.expected_payload_bytes + || header[40..72] != self.expected_digest + { + return Err(TreeSidecarError::InvalidScratch( + "directory plan header validation failed", + )); + } + Ok(()) + } + + fn read_forward_frame( + &self, + ordinal: u64, + offset: u64, + ) -> Result<(Vec, u64), TreeSidecarError> { + let mut length_bytes = [0_u8; 4]; + plan_read_exact_at(&self.file, &self.path, offset, &mut length_bytes)?; + let length = u32::from_be_bytes(length_bytes) as usize; + if length == 0 || length > DIRECTORY_PLAN_MAX_RECORD_BYTES { + return Err(TreeSidecarError::InvalidScratch( + "invalid directory plan record length", + )); + } + let mut frame_tag = [0_u8; 32]; + plan_read_exact_at(&self.file, &self.path, offset + 4, &mut frame_tag)?; + let mut record = vec![0_u8; length]; + plan_read_exact_at(&self.file, &self.path, offset + 36, &mut record)?; + let trailing_offset = + offset + .checked_add(36 + length as u64) + .ok_or(TreeSidecarError::InvalidScratch( + "directory plan frame offset overflow", + ))?; + let mut trailing = [0_u8; 4]; + plan_read_exact_at(&self.file, &self.path, trailing_offset, &mut trailing)?; + if trailing != length_bytes + || directory_plan_frame_tag( + &self.frame_key, + self.transaction, + ordinal, + length_bytes, + &record, + ) != frame_tag + { + return Err(TreeSidecarError::InvalidScratch( + "directory plan frame authentication failed", + )); + } + Ok((record, trailing_offset + 4)) + } + + /// Authenticates the entire anonymous plan before TreeSealIntent. Every frame + /// is authenticated again when scanned or consumed in reverse. + pub(crate) fn authenticate(&mut self) -> Result<(), TreeSidecarError> { + self.authenticated = false; + self.validate_header()?; + let mut offset = DIRECTORY_PLAN_HEADER_LEN as u64; + let mut payload_bytes = 0_u64; + let mut digest = directory_plan_digest(self.transaction); + for ordinal in 0..self.expected_records { + let (record, next) = self.read_forward_frame(ordinal, offset)?; + let length = (record.len() as u32).to_be_bytes(); + let tag = directory_plan_frame_tag( + &self.frame_key, + self.transaction, + ordinal, + length, + &record, + ); + digest.update(length); + digest.update(tag); + digest.update(&record); + digest.update(length); + payload_bytes = payload_bytes.checked_add(next - offset).ok_or( + TreeSidecarError::InvalidScratch("directory plan payload overflow"), + )?; + offset = next; + } + digest.update(self.expected_records.to_be_bytes()); + digest.update(self.expected_payload_bytes.to_be_bytes()); + let actual: [u8; 32] = digest.finalize().into(); + if offset != self.file_bytes + || payload_bytes != self.expected_payload_bytes + || actual != self.expected_digest + { + return Err(TreeSidecarError::InvalidScratch( + "directory plan aggregate integrity failed", + )); + } + self.authenticated = true; + self.reverse_remaining = self.expected_records; + self.reverse_offset = self.file_bytes; + Ok(()) + } + + pub(crate) fn record_count(&self) -> u64 { + self.expected_records + } + + pub(crate) fn for_each_forward( + &self, + mut visit: impl FnMut(&[u8]) -> Result<(), E>, + ) -> Result<(), TreeSidecarFoldError> { + if !self.authenticated { + return Err(TreeSidecarError::InvalidScratch( + "directory plan scan precedes authentication", + ) + .into()); + } + let mut offset = DIRECTORY_PLAN_HEADER_LEN as u64; + for ordinal in 0..self.expected_records { + let (record, next) = self.read_forward_frame(ordinal, offset)?; + visit(&record).map_err(TreeSidecarFoldError::Fold)?; + offset = next; + } + if offset != self.file_bytes { + return Err( + TreeSidecarError::InvalidScratch("directory plan scan did not reach EOF").into(), + ); + } + Ok(()) + } + + pub(crate) fn next_reverse(&mut self) -> Result>, TreeSidecarError> { + if !self.authenticated { + return Err(TreeSidecarError::InvalidScratch( + "directory plan reverse read precedes authentication", + )); + } + if self.reverse_remaining == 0 { + if self.reverse_offset != DIRECTORY_PLAN_HEADER_LEN as u64 { + return Err(TreeSidecarError::InvalidScratch( + "directory plan reverse read missed the header boundary", + )); + } + return Ok(None); + } + if self.reverse_offset < DIRECTORY_PLAN_HEADER_LEN as u64 + 40 { + return Err(TreeSidecarError::InvalidScratch( + "directory plan reverse frame is truncated", + )); + } + let mut trailing = [0_u8; 4]; + plan_read_exact_at( + &self.file, + &self.path, + self.reverse_offset - 4, + &mut trailing, + )?; + let length = u32::from_be_bytes(trailing) as u64; + if length == 0 || length as usize > DIRECTORY_PLAN_MAX_RECORD_BYTES { + return Err(TreeSidecarError::InvalidScratch( + "invalid reverse directory plan length", + )); + } + let frame_bytes = length + .checked_add(40) + .ok_or(TreeSidecarError::InvalidScratch( + "directory plan reverse frame length overflow", + ))?; + let start = self + .reverse_offset + .checked_sub(frame_bytes) + .filter(|offset| *offset >= DIRECTORY_PLAN_HEADER_LEN as u64) + .ok_or(TreeSidecarError::InvalidScratch( + "directory plan reverse frame crosses the header", + ))?; + let ordinal = self.reverse_remaining - 1; + let (record, next) = self.read_forward_frame(ordinal, start)?; + if next != self.reverse_offset { + return Err(TreeSidecarError::InvalidScratch( + "directory plan reverse frame boundary changed", + )); + } + self.reverse_offset = start; + self.reverse_remaining -= 1; + Ok(Some(record)) + } + + pub(crate) fn finish(mut self) -> Result<(), TreeSidecarError> { + while self.next_reverse()?.is_some() {} + Ok(()) + } + + #[cfg(test)] + pub(crate) fn link_count_for_test(&self) -> libc::nlink_t { + rustix::fs::fstat(&self.file).unwrap().st_nlink + } + + #[cfg(test)] + pub(crate) fn corrupt_frame_for_test(&self, target: u64) { + let mut offset = DIRECTORY_PLAN_HEADER_LEN as u64; + for ordinal in 0..self.expected_records { + let (record, next) = self.read_forward_frame(ordinal, offset).unwrap(); + if ordinal == target { + let byte_offset = offset + 36; + let byte = [record[0] ^ 0x80]; + self.file.write_all_at(&byte, byte_offset).unwrap(); + self.file.sync_data().unwrap(); + return; + } + offset = next; + } + panic!("directory plan test frame {target} is absent"); + } +} + +fn require_anonymous_plan_fd(file: &File, path: &Path) -> Result<(), TreeSidecarError> { + let stat = rustix::fs::fstat(file).map_err(|error| io_error(path, error.into()))?; + if stat.st_nlink != 0 { + return Err(TreeSidecarError::InvalidScratch( + "ephemeral plan descriptor remains named", + )); + } + Ok(()) +} + +struct DirectoryPlanWriter { + transaction: TransactionId, + path: PathBuf, + file: File, + frame_key: [u8; 32], + digest: Sha256, + payload_bytes: u64, + records: u64, +} + +impl DirectoryPlanWriter { + #[allow(clippy::disallowed_methods)] + fn create( + store: &TreeSidecarStore, + transaction: TransactionId, + ) -> Result { + let name = scratch_name(transaction); + let path = store.path.join(&name); + let fd = rustix::fs::openat(&store.directory, &name, OPEN_NEW, FILE_MODE) + .map_err(|error| io_error(&path, error.into()))?; + rustix::fs::fchmod(&fd, FILE_MODE).map_err(|error| io_error(&path, error.into()))?; + validate_file( + &store.directory, + &name, + &fd, + store.backend, + store.device, + &path, + )?; + let mut file = File::from(fd); + file.write_all(&[0_u8; DIRECTORY_PLAN_HEADER_LEN]) + .map_err(|error| io_error(&path, error))?; + rustix::fs::unlinkat(&store.directory, &name, AtFlags::empty()) + .map_err(|error| io_error(&path, error.into()))?; + require_anonymous_plan_fd(&file, &path)?; + rustix::fs::fsync(&store.directory).map_err(|error| io_error(&store.path, error.into()))?; + let mut frame_key = [0_u8; 32]; + getrandom::fill(&mut frame_key) + .map_err(|error| io_error(&path, io::Error::other(error)))?; + Ok(Self { + transaction, + path, + file, + frame_key, + digest: directory_plan_digest(transaction), + payload_bytes: 0, + records: 0, + }) + } + + fn push(&mut self, record: &[u8]) -> Result<(), TreeSidecarError> { + if record.is_empty() || record.len() > DIRECTORY_PLAN_MAX_RECORD_BYTES { + return Err(TreeSidecarError::InvalidScratch( + "invalid directory plan record", + )); + } + let length = u32::try_from(record.len()) + .map_err(|_| TreeSidecarError::InvalidScratch("directory plan length overflow"))? + .to_be_bytes(); + let tag = directory_plan_frame_tag( + &self.frame_key, + self.transaction, + self.records, + length, + record, + ); + self.file + .write_all(&length) + .and_then(|_| self.file.write_all(&tag)) + .and_then(|_| self.file.write_all(record)) + .and_then(|_| self.file.write_all(&length)) + .map_err(|error| io_error(&self.path, error))?; + self.digest.update(length); + self.digest.update(tag); + self.digest.update(record); + self.digest.update(length); + self.payload_bytes = self + .payload_bytes + .checked_add(40 + record.len() as u64) + .filter(|bytes| *bytes <= DIRECTORY_PLAN_MAX_TOTAL_BYTES) + .ok_or(TreeSidecarError::InvalidScratch( + "directory plan exceeds its payload limit", + ))?; + self.records = self + .records + .checked_add(1) + .filter(|n| *n <= MAX_RECORDS) + .ok_or(TreeSidecarError::InvalidScratch( + "directory plan record count overflow", + ))?; + Ok(()) + } + + fn finish(mut self, expected: u64) -> Result { + if expected == 0 || self.records != expected { + return Err(TreeSidecarError::InvalidScratch( + "directory plan record count changed", + )); + } + self.digest.update(self.records.to_be_bytes()); + self.digest.update(self.payload_bytes.to_be_bytes()); + let digest: [u8; 32] = self.digest.finalize().into(); + let mut header = [0_u8; DIRECTORY_PLAN_HEADER_LEN]; + header[0..4].copy_from_slice(DIRECTORY_PLAN_MAGIC); + header[4..6].copy_from_slice(&DIRECTORY_PLAN_VERSION.to_be_bytes()); + header[6..8].copy_from_slice(&(DIRECTORY_PLAN_HEADER_LEN as u16).to_be_bytes()); + header[8..24].copy_from_slice(&self.transaction.0); + header[24..32].copy_from_slice(&self.records.to_be_bytes()); + header[32..40].copy_from_slice(&self.payload_bytes.to_be_bytes()); + header[40..72].copy_from_slice(&digest); + self.file + .seek(SeekFrom::Start(0)) + .and_then(|_| self.file.write_all(&header)) + .and_then(|_| self.file.sync_all()) + .map_err(|error| io_error(&self.path, error))?; + let stat = + rustix::fs::fstat(&self.file).map_err(|error| io_error(&self.path, error.into()))?; + let file_bytes = u64::try_from(stat.st_size).map_err(|_| { + TreeSidecarError::InvalidScratch("directory plan file length is not representable") + })?; + let expected_bytes = DIRECTORY_PLAN_HEADER_LEN as u64 + self.payload_bytes; + if file_bytes != expected_bytes { + return Err(TreeSidecarError::InvalidScratch( + "directory plan file length changed", + )); + } + Ok(TreeDirectoryPlan { + file: self.file, + path: self.path, + transaction: self.transaction, + frame_key: self.frame_key, + expected_records: self.records, + expected_payload_bytes: self.payload_bytes, + expected_digest: digest, + file_bytes, + authenticated: false, + reverse_remaining: 0, + reverse_offset: 0, + }) + } +} + fn purge_plan_digest(transaction: TransactionId) -> Sha256 { let mut digest = Sha256::new(); digest.update(PURGE_PLAN_DOMAIN); @@ -1584,6 +2176,7 @@ impl PurgePlanWriter { // Unlink before any plan data is written: only this descriptor survives. rustix::fs::unlinkat(&store.directory, &name, AtFlags::empty()) .map_err(|error| io_error(&path, error.into()))?; + require_anonymous_plan_fd(&file, &path)?; let mut frame_key = [0_u8; 32]; getrandom::fill(&mut frame_key) .map_err(|error| io_error(&path, io::Error::other(error)))?; diff --git a/crates/degu-core/src/seal/wal.rs b/crates/degu-core/src/seal/wal.rs index df52584..48410ae 100644 --- a/crates/degu-core/src/seal/wal.rs +++ b/crates/degu-core/src/seal/wal.rs @@ -19,7 +19,7 @@ use std::collections::{BTreeMap, HashMap, HashSet}; use std::fs::File; use std::io::{self, Read, Seek, SeekFrom, Write}; use std::os::unix::ffi::{OsStrExt, OsStringExt}; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use std::sync::Arc; const MAGIC: &[u8; 4] = b"DSWL"; @@ -793,6 +793,40 @@ impl SealWal { }) } + /// Returns the exact executor-reported applied mode for one original tree + /// seal. This is a projection of the already-resident leased-WAL state, not + /// new authority; duplicate or mismatched evidence fails closed. + pub(crate) fn applied_tree_seal_mode( + &self, + transaction: TransactionId, + relative_path: &Path, + device: u64, + inode: u64, + incarnation: u64, + pre_mode: u32, + ) -> Option { + let mut matched = None; + for ((owner, _), permission) in &self.permissions { + if *owner != transaction + || permission.phase != TransactionState::TreeSealIntent + || permission.application != ApplicationStatus::Applied + || permission.reverses_mutation_id.is_some() + || permission.evidence.relative_path() != relative_path + || permission.evidence.device() != device + || permission.evidence.inode() != inode + || permission.evidence.generation_or_btime() != Some(incarnation) + || permission.pre_mode != pre_mode + || permission.evidence.expected_mode() != permission.expected_mode + { + continue; + } + if matched.replace(permission.expected_mode).is_some() { + return None; + } + } + matched + } + /// Allocates the next transaction-local mutation id for startup recovery. /// IDs are never reused across original seals, failed intents, or inverses. #[allow(dead_code)] // consumed by the startup recovery execution seam diff --git a/crates/degu-core/src/staging/rename.rs b/crates/degu-core/src/staging/rename.rs index 3730964..088d8a4 100644 --- a/crates/degu-core/src/staging/rename.rs +++ b/crates/degu-core/src/staging/rename.rs @@ -7,10 +7,10 @@ use crate::authority::TransactionState; use crate::backend::held::{ - HardlinkTopologyFold, HeldTreeError, HeldTreeInventory, HeldTreeLimits, HeldTreeSealError, - HeldTreeV3CollectError, ManifestV3CodecError, PendingV3Inventory, StreamedV3Inventory, - StructureEvidence, decode_hardlink_scratch_record, hardlink_scratch_sentinel_record, - structure_evidence_from_v3_record, + HardlinkTopologyFold, HeldTreeError, HeldTreeLimits, HeldTreeSealError, HeldTreeV3CollectError, + ManifestV3CodecError, PendingV3Inventory, StreamedV3Inventory, StructureEvidence, + decode_hardlink_scratch_record, decode_pre_seal_directory_plan_record, + hardlink_scratch_sentinel_record, structure_evidence_from_v3_record, }; use crate::backend::{ CertificationError, HeldLocalBackendEvidence, LocalModeRevalidationFailure, certify_held_fd, @@ -19,6 +19,10 @@ use crate::seal::executor::{ LocalModeExecutionError, LocalModeMutationRequest, LocalModeMutationResult, LocalModeTransform, RecoveryLocator, execute_staging_local_mode_mutation, }; +#[cfg(test)] +use crate::seal::sidecar::TreeDirectoryPlan; +#[cfg(test)] +type DirectoryPlanTestCallback = Box; use crate::seal::sidecar::{ TreeManifestFoldError, TreeManifestScratchBuildError, TreeSidecarCommitment, TreeSidecarError, TreeSidecarFoldError, TreeSidecarStore, TreeStructureScratchCursor, @@ -47,6 +51,12 @@ std::thread_local! { const { std::cell::RefCell::new(None) }; static AFTER_RENAME: std::cell::RefCell>> = const { std::cell::RefCell::new(None) }; + static AFTER_PRE_SEAL_DIRECTORY_PLAN_PREFLIGHT: std::cell::RefCell> = + const { std::cell::RefCell::new(None) }; + static AFTER_PRE_SEAL_SCRATCH_READY: std::cell::RefCell>> = + const { std::cell::RefCell::new(None) }; + static AFTER_PRE_SEAL_DIRECTORY_SEALS: std::cell::RefCell>> = + const { std::cell::RefCell::new(None) }; static AFTER_PRE_SEAL_INVENTORY_DROPPED: std::cell::RefCell>> = const { std::cell::RefCell::new(None) }; static AFTER_STRUCTURE_SIDECAR_PREFLIGHT: std::cell::RefCell>> = @@ -460,31 +470,356 @@ pub(crate) fn execute_prepared_rename<'a>( .map_err(classify_source_parent)?; wal.transition_staging_foundation(transaction, TransactionState::ParentSealed)?; - let mut tree = collect_source_tree(&binding)?; + let produced = + sidecars.build_sorted_manifest_scratch_with_output(wal, transaction, |emit_record| { + let parent = certify_duplicate(&binding.source_parent) + .map_err(PostSealProducerError::Binding)?; + PendingV3Inventory::collect_pre_seal( + parent, + binding.metadata.source_basename(), + crate::backend::held_tree_protected_names(), + HeldTreeLimits::default(), + emit_record, + ) + .map_err(PostSealProducerError::Collect) + }); + let (mut pre_seal_scratch, collected_pre_seal) = match produced { + Ok(produced) => produced, + Err(TreeManifestScratchBuildError::Sidecar(error)) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Sidecar(error)); + } + Err(TreeManifestScratchBuildError::Produce(PostSealProducerError::Binding(error))) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Binding(error)); + } + Err(TreeManifestScratchBuildError::Produce(PostSealProducerError::Collect(error))) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(match error { + HeldTreeV3CollectError::Tree(error) => StagingRenameError::HeldTree(error), + HeldTreeV3CollectError::Codec(error) => StagingRenameError::ManifestCodec(error), + HeldTreeV3CollectError::Emit(error) => StagingRenameError::Sidecar(error), + }); + } + }; + let directory_count = collected_pre_seal.directory_count(); + let directory_plan_build = + sidecars.build_directory_plan_with_output(wal, transaction, directory_count, |emit| { + collected_pre_seal.emit_directory_plan(emit) + }); + let (mut directory_plan, pending_pre_seal) = match directory_plan_build { + Ok(result) => result, + Err(TreeManifestScratchBuildError::Sidecar(error)) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Sidecar(error)); + } + Err(TreeManifestScratchBuildError::Produce(error)) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(match error { + HeldTreeV3CollectError::Tree(error) => StagingRenameError::HeldTree(error), + HeldTreeV3CollectError::Codec(error) => StagingRenameError::ManifestCodec(error), + HeldTreeV3CollectError::Emit(error) => StagingRenameError::Sidecar(error), + }); + } + }; + if let Err(error) = directory_plan.authenticate() { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Sidecar(error)); + } + let pre_seal_manifest = + match sidecars.fingerprint_sorted_manifest_scratch(wal, transaction, &mut pre_seal_scratch) + { + Ok(manifest) => manifest, + Err(error) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Sidecar(error)); + } + }; + let pre_seal_finalizer = match pending_pre_seal.into_finalizer(pre_seal_manifest) { + Ok(finalizer) => finalizer, + Err(error) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::HeldTree(error)); + } + }; + let hardlink_build = sidecars.build_sorted_hardlink_scratch_from_manifest( + wal, + transaction, + &mut pre_seal_scratch, + pre_seal_manifest, + pre_seal_finalizer, + |finalizer, record, emit_hardlink| finalizer.observe(record, emit_hardlink), + ); + let (pre_seal_hardlink_scratch, pre_seal_finalizer) = match hardlink_build { + Ok(result) => result, + Err(error) => { + // A complete scratch integrity pass wins over a tree/fold failure; + // no directory mutation has occurred yet. + let authenticated = sidecars.fingerprint_sorted_manifest_scratch( + wal, + transaction, + &mut pre_seal_scratch, + ); + let primary = match error { + TreeSidecarFoldError::Sidecar(error) => StagingRenameError::Sidecar(error), + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Tree(error)) => { + StagingRenameError::HeldTree(error) + } + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Codec(error)) => { + StagingRenameError::ManifestCodec(error) + } + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Emit(error)) => { + StagingRenameError::Sidecar(error) + } + }; + let _ = sidecars.cleanup_unpublished(wal); + authenticated.map_err(StagingRenameError::Sidecar)?; + return Err(primary); + } + }; + let pre_seal_hardlinks = sidecars.fold_sorted_hardlink_scratch_preserving_manifest( + wal, + transaction, + pre_seal_hardlink_scratch, + HardlinkTopologyFold::new(), + |groups, record| { + let record = decode_hardlink_scratch_record(record) + .map_err(HeldTreeV3CollectError::::Codec)? + .ok_or(HeldTreeV3CollectError::Codec( + ManifestV3CodecError::InvalidTag, + ))?; + groups.observe(record).map_err(HeldTreeV3CollectError::Tree) + }, + ); + let pre_seal_hardlinks = match pre_seal_hardlinks { + Ok(fold) => match fold.finish() { + Ok(topology) => topology, + Err(error) => { + let authenticated = sidecars.fingerprint_sorted_manifest_scratch( + wal, + transaction, + &mut pre_seal_scratch, + ); + let _ = sidecars.cleanup_unpublished(wal); + authenticated.map_err(StagingRenameError::Sidecar)?; + return Err(StagingRenameError::HeldTree(error)); + } + }, + Err(error) => { + let primary = match error { + TreeSidecarFoldError::Sidecar(error) => StagingRenameError::Sidecar(error), + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Tree(error)) => { + StagingRenameError::HeldTree(error) + } + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Codec(_)) => { + StagingRenameError::Sidecar(TreeSidecarError::InvalidScratch( + "hardlink scratch record validation failed", + )) + } + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Emit(never)) => match never {}, + }; + let authenticated = sidecars.fingerprint_sorted_manifest_scratch( + wal, + transaction, + &mut pre_seal_scratch, + ); + let _ = sidecars.cleanup_unpublished(wal); + authenticated.map_err(StagingRenameError::Sidecar)?; + return Err(primary); + } + }; + let mut tree = match pre_seal_finalizer.finish(pre_seal_hardlinks) { + Ok(tree) => tree, + Err(error) => { + let authenticated = sidecars.fingerprint_sorted_manifest_scratch( + wal, + transaction, + &mut pre_seal_scratch, + ); + let _ = sidecars.cleanup_unpublished(wal); + authenticated.map_err(StagingRenameError::Sidecar)?; + return Err(StagingRenameError::HeldTree(error)); + } + }; + #[cfg(test)] + AFTER_PRE_SEAL_SCRATCH_READY.with(|hook| { + if let Some(hook) = hook.borrow_mut().take() { + hook(); + } + }); + let mut plan_records = 0_u64; + let plan_validation = directory_plan.for_each_forward(|record| { + let record = decode_pre_seal_directory_plan_record(record) + .map_err(HeldTreeV3CollectError::::Codec)?; + tree.validate_directory_plan_record(&record, plan_records) + .map_err(HeldTreeV3CollectError::Tree)?; + plan_records = plan_records + .checked_add(1) + .ok_or(HeldTreeV3CollectError::Tree(HeldTreeError::PostChanged( + PathBuf::new(), + )))?; + Ok(()) + }); + if let Err(error) = plan_validation { + let _ = sidecars.cleanup_unpublished(wal); + return Err(match error { + TreeSidecarFoldError::Sidecar(error) => StagingRenameError::Sidecar(error), + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Tree(error)) => { + StagingRenameError::HeldTree(error) + } + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Codec(error)) => { + StagingRenameError::ManifestCodec(error) + } + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Emit(never)) => match never {}, + }); + } + if plan_records != tree.directory_count() || plan_records != directory_plan.record_count() { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::HeldTree(HeldTreeError::PostChanged( + PathBuf::new(), + ))); + } if tree.root_strong_identity() != binding.metadata.root_identity() { + let _ = sidecars.cleanup_unpublished(wal); return Err(StagingRenameError::Binding( PreparedRootError::BackendOrMountChanged, )); } - wal.transition_staging_foundation(transaction, TransactionState::TreeSealIntent)?; + if let Err(error) = + wal.transition_staging_foundation(transaction, TransactionState::TreeSealIntent) + { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Wal(error)); + } + #[cfg(test)] + AFTER_PRE_SEAL_DIRECTORY_PLAN_PREFLIGHT.with(|hook| { + if let Some(hook) = hook.borrow_mut().take() { + hook(&directory_plan); + } + }); let source_root = binding .metadata .source_parent() .relative_path() .join(binding.metadata.source_basename()); - tree.seal_directories_for_staging( + let mut mutation_id = 1_u64; + loop { + let encoded = match directory_plan.next_reverse() { + Ok(Some(record)) => record, + Ok(None) => break, + Err(error) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Sidecar(error)); + } + }; + let target = match decode_pre_seal_directory_plan_record(&encoded) { + Ok(record) => record, + Err(error) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::ManifestCodec(error)); + } + }; + let mut chain = Vec::new(); + let chain_result = directory_plan.for_each_forward(|encoded| { + let candidate = decode_pre_seal_directory_plan_record(encoded) + .map_err(HeldTreeV3CollectError::::Codec)?; + tree.consider_directory_plan_ancestor(&target, candidate, &mut chain) + .map_err(HeldTreeV3CollectError::Tree) + }); + if let Err(error) = chain_result { + let _ = sidecars.cleanup_unpublished(wal); + return Err(match error { + TreeSidecarFoldError::Sidecar(error) => StagingRenameError::Sidecar(error), + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Tree(error)) => { + StagingRenameError::HeldTree(error) + } + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Codec(error)) => { + StagingRenameError::ManifestCodec(error) + } + TreeSidecarFoldError::Fold(HeldTreeV3CollectError::Emit(never)) => match never {}, + }); + } + if let Err(error) = tree.seal_directory_for_staging( + wal, + transaction, + &source_root, + binding.metadata.filesystem_id(), + mutation_id, + &target, + &chain, + ) { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::TreeSeal(error)); + } + mutation_id = match mutation_id.checked_add(1) { + Some(next) => next, + None => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::TreeSeal( + HeldTreeSealError::MutationIdExhausted, + )); + } + }; + } + if let Err(error) = directory_plan.finish() { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Sidecar(error)); + } + #[cfg(test)] + AFTER_PRE_SEAL_DIRECTORY_SEALS.with(|hook| { + if let Some(hook) = hook.borrow_mut().take() { + hook(); + } + }); + + let mut expectation_builder = tree.post_seal_expectation_builder(); + let expectation_fold = sidecars.read_sorted_manifest_scratch( wal, transaction, - &source_root, - binding.metadata.filesystem_id(), - 1, - )?; - - // Reduce the complete pre-seal inventory to a fixed-size expected v3 - // commitment and drop it before collecting the second full content proof. - // The post-seal collection therefore cannot coexist with the pre-seal - // `Vec` in this forward path. - let post_seal_expectation = tree.into_post_seal_expectation()?; + pre_seal_manifest, + &mut pre_seal_scratch, + (), + |(), record, wal_view| { + expectation_builder.observe(&tree, record, |path, device, inode, incarnation| { + wal_view.applied_tree_seal_mode( + transaction, + &source_root.join(path), + device, + inode, + incarnation, + record.mode, + ) + })?; + Ok(()) + }, + ); + if let Err(error) = expectation_fold { + let primary = match error { + TreeSidecarFoldError::Sidecar(error) => StagingRenameError::Sidecar(error), + TreeSidecarFoldError::Fold(error) => StagingRenameError::HeldTree(error), + }; + let authenticated = + sidecars.fingerprint_sorted_manifest_scratch(wal, transaction, &mut pre_seal_scratch); + let _ = sidecars.cleanup_unpublished(wal); + authenticated.map_err(StagingRenameError::Sidecar)?; + return Err(primary); + } + let post_seal_expectation = match tree.finish_post_seal_expectation(expectation_builder) { + Ok(expectation) => expectation, + Err(error) => { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::HeldTree(error)); + } + }; + if let Err(error) = sidecars.discard_sorted_manifest_scratch( + wal, + transaction, + pre_seal_manifest, + pre_seal_scratch, + ) { + let _ = sidecars.cleanup_unpublished(wal); + return Err(StagingRenameError::Sidecar(error)); + } #[cfg(test)] AFTER_PRE_SEAL_INVENTORY_DROPPED.with(|hook| { if let Some(hook) = hook.borrow_mut().take() { @@ -939,18 +1274,6 @@ fn rewalk_streamed_v3_structure( Ok(()) } -fn collect_source_tree( - binding: &PreparedRootBinding, -) -> Result { - let parent = certify_duplicate(&binding.source_parent)?; - Ok(HeldTreeInventory::collect( - parent, - binding.metadata.source_basename(), - crate::backend::held_tree_protected_names(), - HeldTreeLimits::default(), - )?) -} - fn verify_post_rename(binding: &PreparedRootBinding) -> Result<(), PreparedRootError> { binding.source_anchor.verify_locator_binding( binding.metadata.source_parent(), diff --git a/crates/degu-core/src/staging/rename/tests.rs b/crates/degu-core/src/staging/rename/tests.rs index 6317f9d..db19040 100644 --- a/crates/degu-core/src/staging/rename/tests.rs +++ b/crates/degu-core/src/staging/rename/tests.rs @@ -199,6 +199,70 @@ fn set_group(path: &Path, gid: u32) -> std::io::Result<()> { } } +#[test] +fn anonymous_directory_plan_preserves_reverse_bfs_wal_modes_and_restart_verification() { + let Some(fixture) = Fixture::new() else { + return; + }; + let grandchild = fixture.source_root.join("child/grandchild"); + std::fs::create_dir(&grandchild).unwrap(); + set_mode(&fixture.source_root, 0o770); + set_mode(&fixture.source_root.join("child"), 0o700); + set_mode(&grandchild, 0o777); + + let transaction = TransactionId([0x72; 16]); + let (mut engine, report) = SealedStagingEngine::open(&fixture.store).unwrap(); + assert!(report.is_empty()); + drop( + engine + .stage_prepared_root(transaction, fixture.prepare()) + .unwrap(), + ); + drop(engine); + + let mut lease = fixture.store.try_lease().unwrap(); + let replay = lease.replay_and_repair().unwrap(); + let permissions = &replay.transactions[&transaction].permissions; + assert_eq!(permissions.len(), 4); + assert_eq!( + permissions + .iter() + .map(|permission| ( + permission.mutation_id, + permission.evidence.relative_path().to_path_buf(), + permission.pre_mode, + permission.expected_mode, + )) + .collect::>(), + vec![ + (0, PathBuf::from("source-parent"), 0o770, 0o750), + ( + 1, + PathBuf::from("source-parent/root/child/grandchild"), + 0o777, + 0o755, + ), + (2, PathBuf::from("source-parent/root/child"), 0o700, 0o700,), + (3, PathBuf::from("source-parent/root"), 0o770, 0o750,), + ] + ); + drop(lease); + + let (mut recovered, report) = SealedStagingEngine::open(&fixture.store).unwrap(); + let candidate = report.into_candidates().pop().unwrap(); + let capability = recovered + .prepare_startup_recovery(candidate, fixture.anchors()) + .unwrap(); + let StartupRecoveryCapability::PendingVerification(pending) = capability else { + panic!("staged plan did not resume verification") + }; + let StagedVerificationOutcome::StagedSealed(verified) = pending.verify_or_quarantine().unwrap() + else { + panic!("varied-mode tree failed exact restart verification") + }; + assert_eq!(verified.wal_state(), Some(TransactionState::StagedSealed)); +} + #[test] fn component_order_prefix_paths_stage_and_restart_verify() { let Some(fixture) = Fixture::new() else { @@ -522,6 +586,130 @@ fn consumed_pre_seal_expectation_rejects_same_size_content_drift_before_post_pro assert!(replay.transactions[&transaction].tree_sidecar.is_none()); } +#[test] +fn corrupted_anonymous_directory_plan_mutates_no_tree_directory() { + let Some(fixture) = Fixture::new() else { + return; + }; + let transaction = TransactionId([0x73; 16]); + AFTER_PRE_SEAL_DIRECTORY_PLAN_PREFLIGHT.with(|hook| { + *hook.borrow_mut() = Some(Box::new(|plan| plan.corrupt_frame_for_test(1))); + }); + let (mut engine, report) = SealedStagingEngine::open(&fixture.store).unwrap(); + assert!(report.is_empty()); + let error = match engine.stage_prepared_root(transaction, fixture.prepare()) { + Ok(_) => panic!("corrupt directory plan must fail before tree sealing"), + Err(error) => error, + }; + assert!(matches!(error, StagingRenameError::Sidecar(_))); + assert_eq!( + engine.state(transaction), + Some(TransactionState::TreeSealIntent) + ); + assert_eq!(mode(&fixture.source_parent), 0o750); + assert_eq!(mode(&fixture.source_root), 0o770); + assert_eq!(mode(&fixture.source_root.join("child")), 0o770); + assert!(fixture.source_root.join("child/data").is_file()); + assert!(!fixture.destination_root.exists()); + assert!( + std::fs::read_dir(fixture.base.join("wal-store")) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .all(|name| !name.to_string_lossy().starts_with(".tree-scratch-v1-")) + ); + drop(engine); + + let mut lease = fixture.store.try_lease().unwrap(); + let replay = lease.replay_and_repair().unwrap(); + let recovered = &replay.transactions[&transaction]; + assert_eq!(recovered.permissions.len(), 1); + assert!(recovered.tree_manifest.is_none()); + assert!(recovered.tree_sidecar.is_none()); + drop(lease); + + let (engine, report) = SealedStagingEngine::open(&fixture.store).unwrap(); + let (ready, summary) = engine + .recover_startup(report, |_, _| Ok(fixture.raw_anchors())) + .unwrap(); + assert_eq!(summary.recovered.len(), 1); + assert_eq!(ready.state(transaction), Some(TransactionState::Restored)); + assert_eq!(mode(&fixture.source_parent), 0o770); +} + +#[test] +fn pre_seal_scratch_crash_boundaries_cleanup_and_restore_exact_wal_prefix() { + for (index, (boundary, expected_state, expected_permissions)) in [ + ("scratch-ready", TransactionState::ParentSealed, 1_usize), + ( + "directories-sealed", + TransactionState::TreeSealIntent, + 3_usize, + ), + ] + .into_iter() + .enumerate() + { + let Some(fixture) = Fixture::new() else { + return; + }; + let transaction = TransactionId([0x74 + index as u8; 16]); + let (mut engine, report) = SealedStagingEngine::open(&fixture.store).unwrap(); + assert!(report.is_empty()); + match boundary { + "scratch-ready" => AFTER_PRE_SEAL_SCRATCH_READY.with(|hook| { + *hook.borrow_mut() = Some(Box::new(|| panic!("simulated pre-seal crash"))); + }), + "directories-sealed" => AFTER_PRE_SEAL_DIRECTORY_SEALS.with(|hook| { + *hook.borrow_mut() = Some(Box::new(|| panic!("simulated pre-seal crash"))); + }), + _ => unreachable!(), + } + + let crashed = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { + let _ = engine.stage_prepared_root(transaction, fixture.prepare()); + })); + assert!(crashed.is_err(), "boundary={boundary}"); + assert_eq!(engine.state(transaction), Some(expected_state)); + assert!( + std::fs::read_dir(fixture.base.join("wal-store")) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .any(|name| name.to_string_lossy().starts_with(".tree-scratch-v1-")), + "the simulated process death must strand unpublished scratch at {boundary}" + ); + drop(engine); + + let mut lease = fixture.store.try_lease().unwrap(); + let replay = lease.replay_and_repair().unwrap(); + let recovered = &replay.transactions[&transaction]; + assert_eq!(recovered.state, expected_state); + assert_eq!(recovered.permissions.len(), expected_permissions); + assert!(recovered.tree_manifest.is_none()); + assert!(recovered.tree_sidecar.is_none()); + drop(lease); + + let (engine, report) = SealedStagingEngine::open(&fixture.store).unwrap(); + assert_eq!(report.candidates().len(), 1); + assert!( + std::fs::read_dir(fixture.base.join("wal-store")) + .unwrap() + .map(|entry| entry.unwrap().file_name()) + .all(|name| !name.to_string_lossy().starts_with(".tree-scratch-v1-")), + "startup must remove unpublished pre-seal scratch at {boundary}" + ); + let (ready, summary) = engine + .recover_startup(report, |_, _| Ok(fixture.raw_anchors())) + .unwrap(); + assert_eq!(summary.recovered.len(), 1); + assert_eq!(ready.state(transaction), Some(TransactionState::Restored)); + assert_eq!(mode(&fixture.source_parent), 0o770); + assert_eq!(mode(&fixture.source_root), 0o770); + assert_eq!(mode(&fixture.source_root.join("child")), 0o770); + assert!(fixture.source_root.join("child/data").is_file()); + assert!(!fixture.destination_root.exists()); + } +} + #[test] fn streamed_structure_scratch_reports_the_first_canonical_added_path() { let Some(fixture) = Fixture::new() else {