From c3cb0617e4d139e3930f3cfb7b2729f6280c1f2f Mon Sep 17 00:00:00 2001 From: longjin Date: Mon, 21 Sep 2026 13:57:10 +0000 Subject: [PATCH 1/2] fix(vfs): batch filesystem sync across page-cache mappings Global filesystem synchronization previously waited for each mapping's durable data completion before starting the next mapping. With ext4's synchronous journal requests, this repeatedly sealed tiny transactions and paid the complete journal/checkpoint barrier cost for small files. Start bounded windows using the existing asynchronous writeback budget, then request synchronous progress and drain every mapping in the window. Submit the first claimed Token batch on the sync caller's stack so its acceptance is not deferred behind the window's forced commit. Reuse the existing continuation for successors; keep Legacy and ordinary range WRITE dispatch asynchronous and preserve journal durability barriers. Retain the canonical inode and page cache across data completion and metadata writeback. Preserve the frozen range after a partial dispatch failure, sample mapping errseq before starting, and report errors only after draining the accepted work. Replace the completion scan's temporary allocation with a bounded array so memory pressure cannot interrupt drain. Add PageCache success/submission-error/deferred-error selftests and ext4 multi-window, mixed-size, concurrent-redirty, and remount coverage. Validation: - make kernel and focused rustfmt/diff checks passed. - KVM guests with both one and two vCPUs passed 46 ext4, one PageCache accounting (including kernel selftests), and four sync_file_range tests. - Repeated DragonOS repository clone/checkout/sync/fsck validation passed. - The measured checkout sync barrier count fell from about 20,500 to about 2,700 without changing the journal protocol or single-file fsync. - Three-role adversarial review found no unresolved blocking findings. Signed-off-by: longjin --- kernel/src/filesystem/page_cache.rs | 10 +- kernel/src/filesystem/page_cache/selftest.rs | 75 ++- kernel/src/filesystem/page_cache/writeback.rs | 626 +++++++++++------- kernel/src/filesystem/vfs/mount/mod.rs | 46 +- .../suites/normal/ext4_inode_identity.cc | 136 ++++ 5 files changed, 630 insertions(+), 263 deletions(-) diff --git a/kernel/src/filesystem/page_cache.rs b/kernel/src/filesystem/page_cache.rs index e1a1bf070..779017e9a 100644 --- a/kernel/src/filesystem/page_cache.rs +++ b/kernel/src/filesystem/page_cache.rs @@ -58,8 +58,8 @@ pub use read_batch::{PageCacheReadBatchCompletion, PageCacheReadBatchRequest}; pub use read_dma::PageCacheReadDmaReservation; pub(crate) use selftest::{run_accounting_debug_selftest, run_completion_domain_debug_selftest}; pub(crate) use writeback::{ - async_writeback_progress_snapshot, wait_for_async_writeback_progress, - PageCacheWritebackDispatchOutcome, + async_writeback_progress_snapshot, wait_for_async_writeback_progress, PageCacheDomainWriteback, + PageCacheWritebackDispatchOutcome, DOMAIN_WRITEBACK_WINDOW, }; use writeback::{ run_async_writeback_budget_retry_selftest, run_submitted_writeback_selftest, @@ -386,7 +386,7 @@ pub struct PageCacheWritebackDomain { pub(crate) trait PageCacheDomainMember: Send + Sync { fn has_dirty_or_writeback_work(&self) -> bool; - fn sync_for_domain(&self) -> Result<(), SystemError>; + fn start_sync_for_domain(&self) -> Result; fn retire_for_domain(&self) -> Result<(), SystemError>; } @@ -3948,8 +3948,8 @@ impl PageCacheDomainMember for PageCache { self.has_dirty_or_writeback_pages() } - fn sync_for_domain(&self) -> Result<(), SystemError> { - self.manager.sync() + fn start_sync_for_domain(&self) -> Result { + self.manager.start_sync_for_domain() } fn retire_for_domain(&self) -> Result<(), SystemError> { diff --git a/kernel/src/filesystem/page_cache/selftest.rs b/kernel/src/filesystem/page_cache/selftest.rs index 041b15719..c836e7f2d 100644 --- a/kernel/src/filesystem/page_cache/selftest.rs +++ b/kernel/src/filesystem/page_cache/selftest.rs @@ -565,13 +565,14 @@ impl PageCacheWritebackSubmission for PageCacheSubmissionSelftestToken { return Err(SystemError::EIO); } if self.state.defer_next_submit.swap(false, Ordering::AcqRel) { - self.state - .deferred_after_submit - .fetch_add(1, Ordering::Relaxed); let progress = self.state.deferred_progress(); let previous = self.state.private_claims.fetch_sub(1, Ordering::AcqRel); assert_eq!(previous, 1, "submission token private claim underflow"); self.resolved = true; + self.state + .deferred_after_submit + .fetch_add(1, Ordering::Release); + self.state.progress_wait.wake_all(); return Ok(PageCacheWritebackSubmitResult::Deferred(progress)); } if self.state.fail_next_submit.swap(false, Ordering::AcqRel) { @@ -1563,6 +1564,70 @@ fn run_remote_dirty_publish_selftest() -> Result { && removal_after_token_release) } +/// Exercise domain start separately from its completion/error boundary. +fn run_domain_writeback_start_selftest() -> Result { + use crate::filesystem::vfs::{FileType, InodeMode}; + + // Reuse the token fixture, but supply a real inode for domain retention + // and stable size. The separate shmem cache needs no mounted filesystem. + for (fail_inline, fail_after_defer) in [(false, false), (true, false), (false, true)] { + let fs = crate::filesystem::ramfs::RamFS::new(); + let inode = fs.root_inode().create( + "domain-writeback", + FileType::File, + InodeMode::from_bits_truncate(0o600), + )?; + inode.resize(MMArch::PAGE_SIZE)?; + let state = Arc::new(PageCacheSubmissionSelftestState::default()); + let backend: Arc = Arc::new(PageCacheSubmissionSelftestBackend { + state: state.clone(), + admission_order: PageCacheWritebackAdmissionOrder::AdmissionBeforeInvalidate, + snapshot_phase: PageCacheWritebackSnapshotPhase::WithinAdmission, + }); + let cache = PageCache::new_shmem(Some(Arc::downgrade(&inode)), Some(backend)); + let page = cache.get_or_create_page_zero(0)?; + { + let mut locked = page.write(); + locked.add_flags(PageFlags::PG_DIRTY); + cache.mark_page_dirty_page_locked(0, &locked)?; + } + state.fail_next_submit.store(fail_inline, Ordering::Release); + state + .defer_next_submit + .store(fail_after_defer, Ordering::Release); + let pending = cache.manager().start_sync_for_domain()?; + if fail_after_defer { + // Global writeback-budget contention may legitimately postpone + // the first attempt. Wait for the actual producer event, not an + // assumption that this machine was idle when start was called. + state.progress_wait.wait_until(|| { + (state.deferred_after_submit.load(Ordering::Acquire) != 0).then_some(()) + }); + state.fail_next_submit.store(true, Ordering::Release); + state.release_deferred_progress(); + } + let result = pending.finish(); + let expected_result = if fail_inline || fail_after_defer { + result == Err(SystemError::EIO) && state.failed_submissions.load(Ordering::Acquire) == 1 + } else { + result.is_ok() && state.submitted.load(Ordering::Acquire) == 1 + }; + let resolved = state.private_claims.load(Ordering::Acquire) == 0 + && state.submitted_while_admitted.load(Ordering::Acquire) == 0 + && state.admission_depth.load(Ordering::Acquire) == 0 + && state.fallback_writes.load(Ordering::Acquire) == 0 + && cache.inner.lock().writeback_pages.is_empty(); + let removed = cache.manager.remove_page(0)?.is_some(); + let paddr = page.phys_address(); + page_manager_lock().remove_page(&paddr); + let _ = page_reclaimer_lock().remove_page(&paddr); + if !expected_result || !resolved || !removed { + return Ok(false); + } + } + Ok(true) +} + /// Exercise the descriptor half of the front-dirty certificate independently /// from the future ext4 consumer. In particular, g1 must remain frozen in /// its already-bound descriptor while a writer creates g2 during g1's @@ -2014,6 +2079,10 @@ pub(crate) fn run_accounting_debug_selftest() -> Result { + frozen: T, + range: PageCacheWritebackRange, + error: Option, +} + +impl StartedWriteback { + fn into_result(self) -> Result<(T, PageCacheWritebackRange), SystemError> { + match self.error { + Some(error) => Err(error), + None => Ok((self.frozen, self.range)), + } + } +} + +/// Pins the canonical inode through data completion and metadata writeback. +/// A range alone only owns a Weak; the last dirty-page completion +/// may otherwise release the inode's last retention owner before write_inode. +pub(crate) struct PageCacheDomainWriteback { + cache: Arc, + inode: Arc, + _retention: InodeRetentionGuard, + started: StartedWriteback<()>, + error_since: crate::libs::errseq::ErrSeqValue, +} + +impl PageCacheDomainWriteback { + pub(crate) fn finish(self) -> Result<(), SystemError> { + let waited = self.started.range.wait_for_completion(); + let error = self + .started + .error + .or_else(|| waited.err()) + .or_else(|| self.cache.check_writeback_error_since(self.error_since)); + if let Some(error) = error { + return Err(error); + } + self.inode + .write_inode(&WritebackControl::sync_all_for_sync()) + .map_err(|error| { + self.cache + .record_writeback_error_with_superblock(error.clone()); + error + }) + } +} static PAGECACHE_WRITEBACK_RR: AtomicUsize = AtomicUsize::new(0); static ASYNC_WRITEBACK_BATCHES: AtomicUsize = AtomicUsize::new(0); static ASYNC_WRITEBACK_COMPLETIONS: AtomicU64 = AtomicU64::new(0); @@ -1329,6 +1379,30 @@ impl PageCacheWritebackRange { } impl PageCacheManager { + pub(crate) fn start_sync_for_domain(&self) -> Result { + let cache = self.upgrade()?; + let inode = cache + .inode() + .and_then(|inode| inode.upgrade()) + .ok_or(SystemError::EIO)?; + let retention = InodeRetentionGuard::new(inode.clone(), InodeRetentionKind::AsyncWork)?; + let error_since = cache.sample_writeback_error(); + let started = self.start_writeback_range_impl( + 0, + usize::MAX, + || Ok(()), + crate::sched::sched_yield, + true, + )?; + Ok(PageCacheDomainWriteback { + cache, + inode, + _retention: retention, + started, + error_since, + }) + } + fn lock_writeback_protocol( cache: &PageCache, protocol: PageCacheWritebackProtocol, @@ -2640,10 +2714,12 @@ impl PageCacheManager { } let (entries, last_scanned) = { let inner = cache.inner.lock(); - let mut entries = Vec::new(); - entries - .try_reserve_exact(WAIT_BATCH_ENTRIES) - .map_err(|_| SystemError::ENOMEM)?; + // Completion must remain drainable under memory pressure. + // In particular a domain-sync handle must not lose earlier + // untagged I/O merely because its wait cannot allocate. + let mut entries: [Option<(Arc, u64)>; WAIT_BATCH_ENTRIES] = + [const { None }; WAIT_BATCH_ENTRIES]; + let mut entry_count = 0; let mut last_scanned = None; let mut scanned = 0usize; for index in inner.page_indices.range(cursor..=end_index) { @@ -2659,9 +2735,10 @@ impl PageCacheManager { if entry.state() == PageState::Writeback && frontier.is_none_or(|frontier| incarnation <= frontier) { - entries.push((entry.clone(), incarnation)); + entries[entry_count] = Some((entry.clone(), incarnation)); + entry_count += 1; } - if entries.len() == WAIT_BATCH_ENTRIES || scanned == WAIT_SCAN_INDICES { + if entry_count == WAIT_BATCH_ENTRIES || scanned == WAIT_SCAN_INDICES { break; } } @@ -2670,7 +2747,7 @@ impl PageCacheManager { let Some(last_scanned) = last_scanned else { break; }; - for (entry, incarnation) in entries { + for (entry, incarnation) in entries.into_iter().flatten() { if let Err(error) = Self::wait_writeback_entry_incarnation(entry, incarnation) { first_error.get_or_insert(error); } @@ -3612,6 +3689,30 @@ impl PageCacheManager { last_index: usize, permit: AsyncWritebackPermit, batch: ClaimedWritebackBatch, + ) { + let work_state = Mutex::new(Some((permit, batch))); + schedule_pagecache_writeback(Work::new(move || { + let Some((permit, batch)) = work_state.lock().take() else { + return; + }; + Self::run_tagged_writeback_submission( + cache.clone(), + inode.clone(), + continuation, + last_index, + permit, + batch, + ); + })); + } + + fn run_tagged_writeback_submission( + cache: Weak, + inode: Weak, + continuation: TaggedWritebackCursor, + last_index: usize, + permit: AsyncWritebackPermit, + batch: ClaimedWritebackBatch, ) { let TaggedWritebackCursor { start_index, @@ -3619,68 +3720,62 @@ impl PageCacheManager { epoch, cursor: retry_cursor, } = continuation; - let work_state = Mutex::new(Some((permit, batch))); - schedule_pagecache_writeback(Work::new(move || { - let Some((permit, batch)) = work_state.lock().take() else { - return; - }; - let outcome = Self::submit_writeback_batch(batch); - drop(permit); - match outcome { - Ok(WritebackSubmitOutcome::Completed | WritebackSubmitOutcome::Submitted) => { - if last_index == usize::MAX { - if let Some(cache) = cache.upgrade() { - Self::notify_tagged_writeback_progress(&cache); - } - return; + let outcome = Self::submit_writeback_batch(batch); + drop(permit); + match outcome { + Ok(WritebackSubmitOutcome::Completed | WritebackSubmitOutcome::Submitted) => { + if last_index == usize::MAX { + if let Some(cache) = cache.upgrade() { + Self::notify_tagged_writeback_progress(&cache); } - Self::schedule_tagged_writeback_drain_with_permit( - cache.clone(), - inode.clone(), + return; + } + Self::schedule_tagged_writeback_drain_with_permit( + cache.clone(), + inode.clone(), + start_index, + frozen_end, + epoch, + last_index + 1, + None, + ); + } + Ok(WritebackSubmitOutcome::Deferred(progress)) => { + Self::schedule_tagged_writeback_retry( + progress, + cache.clone(), + inode.clone(), + start_index, + frozen_end, + epoch, + retry_cursor, + ); + } + Ok(WritebackSubmitOutcome::Failed(_)) => { + if let Some(cache) = cache.upgrade() { + // Batch completion already recorded the error. Do not + // leave later tags waiting for a continuation that + // cannot be scheduled after a terminal failure. + Self::retire_tagged_writeback_generation( + &cache, start_index, frozen_end, epoch, - last_index + 1, - None, ); } - Ok(WritebackSubmitOutcome::Deferred(progress)) => { - Self::schedule_tagged_writeback_retry( - progress, - cache.clone(), - inode.clone(), + } + Err(error) => { + if let Some(cache) = cache.upgrade() { + Self::abandon_tagged_writeback_generation( + &cache, start_index, frozen_end, epoch, - retry_cursor, + error, ); } - Ok(WritebackSubmitOutcome::Failed(_)) => { - if let Some(cache) = cache.upgrade() { - // Batch completion already recorded the error. Do not - // leave later tags waiting for a continuation that - // cannot be scheduled after a terminal failure. - Self::retire_tagged_writeback_generation( - &cache, - start_index, - frozen_end, - epoch, - ); - } - } - Err(error) => { - if let Some(cache) = cache.upgrade() { - Self::abandon_tagged_writeback_generation( - &cache, - start_index, - frozen_end, - epoch, - error, - ); - } - } } - })); + } } fn schedule_tagged_writeback_retry( @@ -4076,8 +4171,24 @@ impl PageCacheManager { start_index: usize, end_index: usize, freeze: F, - mut chunk_released: C, + chunk_released: C, ) -> Result<(T, PageCacheWritebackRange), SystemError> + where + F: FnOnce() -> Result, + C: FnMut(), + { + self.start_writeback_range_impl(start_index, end_index, freeze, chunk_released, false)? + .into_result() + } + + fn start_writeback_range_impl( + &self, + start_index: usize, + end_index: usize, + freeze: F, + mut chunk_released: C, + submit_first_token_inline: bool, + ) -> Result, SystemError> where F: FnOnce() -> Result, C: FnMut(), @@ -4091,16 +4202,17 @@ impl PageCacheManager { let mut invalidate = cache.invalidate_write(); let frozen_filesystem_state = freeze()?; if start_index > end_index { - return Ok(( - frozen_filesystem_state, - PageCacheWritebackRange { + return Ok(StartedWriteback { + frozen: frozen_filesystem_state, + range: PageCacheWritebackRange { cache: Arc::downgrade(&cache), start_index, frozen_end: None, writeback_frontier: 0, epoch: 0, }, - )); + error: None, + }); } // Freeze the caller-visible dirty set with an epoch tag, equivalent to @@ -4137,16 +4249,17 @@ impl PageCacheManager { epoch }; let Some(frozen_end) = dirty_end.into_iter().chain(preexisting_writeback_end).max() else { - return Ok(( - frozen_filesystem_state, - PageCacheWritebackRange { + return Ok(StartedWriteback { + frozen: frozen_filesystem_state, + range: PageCacheWritebackRange { cache: Arc::downgrade(&cache), start_index, frozen_end: None, writeback_frontier: 0, epoch, }, - )); + error: None, + }); }; let mut tagged_new = false; let mut tag_cursor = start_index; @@ -4209,220 +4322,247 @@ impl PageCacheManager { // Backend claim/admission takes invalidate-read. Release the writer // only after both metadata and page tags have been frozen. drop(invalidate); + drop(_tag_scan); if !tagged_new { - return Ok((frozen_filesystem_state, operation)); + return Ok(StartedWriteback { + frozen: frozen_filesystem_state, + range: operation, + error: None, + }); } - let inode = match cache.inode().and_then(|inode| inode.upgrade()) { - Some(inode) => inode, - None => { - Self::abandon_tagged_writeback_generation( - &cache, - start_index, - frozen_end, - epoch, - SystemError::EIO, - ); - return Err(SystemError::EIO); - } - }; - let _domain_io = match cache.try_acquire_domain_io_classified() { - Ok(permit) => permit, - Err(super::PageCacheDomainIoAdmissionError::Closed) => { - Self::retire_tagged_writeback_generation(&cache, start_index, frozen_end, epoch); - return Err(SystemError::ESTALE); - } - Err(super::PageCacheDomainIoAdmissionError::Unavailable(error)) => { - Self::abandon_tagged_writeback_generation( - &cache, - start_index, - frozen_end, - epoch, - error.clone(), - ); - return Err(error); - } - }; - - // `SYNC_FILE_RANGE_WRITE` is an asynchronous writeout starter: it - // publishes Dirty -> Writeback and queues legacy I/O, but must not - // run a synchronous backend write_pages() to completion on the - // syscall stack. Tagged batches additionally register a precise - // submission-boundary record. WAIT_BEFORE|WRITE waits for workers to - // cross that record, while ordinary WRITE returns after dispatch. - // - // A token queue stops after its first claimed head. Its submit result - // or defer continuation is the only path allowed to inspect a - // successor, preserving delayed-allocation head-first order. - let mut cursor = start_index; - loop { - let (target_index, target_entry, tagged_end) = - match Self::find_tagged_writeback_target(&cache, cursor, frozen_end, epoch) { - TaggedWritebackSearch::Done => { - Self::notify_tagged_writeback_progress(&cache); - break; - } - TaggedWritebackSearch::Advance(next_cursor) => { - cursor = next_cursor; - crate::sched::sched_yield(); - continue; - } - TaggedWritebackSearch::Target { index, entry, end } => (index, entry, end), - }; - - let Some(permit) = AsyncWritebackPermit::try_acquire() else { - // WRITE is an asynchronous starter. A saturated batch budget - // leaves the next frozen tag Dirty and registers a one-shot - // drain retry; it must not wait for an older write_pages() - // call to finish on the syscall stack. - Self::schedule_tagged_writeback_budget_retry( - &cache, - &inode, - start_index, - frozen_end, - epoch, - cursor, - ); - break; - }; - let claim = match self.claim_tagged_batch_with_admission( - &cache, - &inode, - WritebackBatchRange::new(target_index, tagged_end), - &target_entry, - epoch, - _domain_io.as_ref(), - ) { - Ok(claim) => claim, - Err(error) => { - drop(permit); + // Keep the frozen range even on a partially failed dispatch. Domain + // sync must wait earlier accepted batches under its completion guard. + let dispatch = (|| -> Result<(), SystemError> { + let inode = match cache.inode().and_then(|inode| inode.upgrade()) { + Some(inode) => inode, + None => { Self::abandon_tagged_writeback_generation( &cache, start_index, frozen_end, epoch, - error.clone(), + SystemError::EIO, ); - return Err(error); + return Err(SystemError::EIO); } }; - match claim { - WritebackClaimOutcome::Deferred(progress) => { - drop(permit); - // The callback is registered before WRITE returns, so a - // claim-time defer has a producer-owned progress edge - // rather than relying on an unrelated later reclaim. - Self::schedule_tagged_writeback_retry( - progress, - Arc::downgrade(&cache), - Arc::downgrade(&inode), + let _domain_io = match cache.try_acquire_domain_io_classified() { + Ok(permit) => permit, + Err(super::PageCacheDomainIoAdmissionError::Closed) => { + Self::retire_tagged_writeback_generation( + &cache, start_index, frozen_end, epoch, - target_index, ); - break; + return Err(SystemError::ESTALE); } - WritebackClaimOutcome::FailedRecorded(error) => { - drop(permit); - // The claimed batch has already been completed back to - // Dirty and its error recorded by PageCache. Only retire - // the remaining frozen tags here. - Self::retire_tagged_writeback_generation( + Err(super::PageCacheDomainIoAdmissionError::Unavailable(error)) => { + Self::abandon_tagged_writeback_generation( &cache, start_index, frozen_end, epoch, + error.clone(), ); return Err(error); } - WritebackClaimOutcome::NoBatch => { - drop(permit); - let retry_incarnation = { - let inner = cache.inner.lock(); - inner.pages.get(&target_index).and_then(|current| { - (Arc::ptr_eq(current, &target_entry) - && inner.dirty_pages.contains(&target_index) - && current.writeback_tag() == epoch) - .then(|| current.writeback_incarnation.load(Ordering::Acquire)) - }) + }; + + // `SYNC_FILE_RANGE_WRITE` is an asynchronous writeout starter: it + // publishes Dirty -> Writeback and queues legacy I/O, but must not + // run a synchronous backend write_pages() to completion on the + // syscall stack. Tagged batches additionally register a precise + // submission-boundary record. WAIT_BEFORE|WRITE waits for workers to + // cross that record, while ordinary WRITE returns after dispatch. + // + // A token queue stops after its first claimed head. Its submit result + // or defer continuation is the only path allowed to inspect a + // successor, preserving delayed-allocation head-first order. + let mut cursor = start_index; + loop { + let (target_index, target_entry, tagged_end) = + match Self::find_tagged_writeback_target(&cache, cursor, frozen_end, epoch) { + TaggedWritebackSearch::Done => { + Self::notify_tagged_writeback_progress(&cache); + break; + } + TaggedWritebackSearch::Advance(next_cursor) => { + cursor = next_cursor; + crate::sched::sched_yield(); + continue; + } + TaggedWritebackSearch::Target { index, entry, end } => (index, entry, end), }; - if let Some(observed_incarnation) = retry_incarnation { - if let Err(error) = Self::register_tagged_writeback_incarnation_retry( + + let Some(permit) = AsyncWritebackPermit::try_acquire() else { + // WRITE is an asynchronous starter. A saturated batch budget + // leaves the next frozen tag Dirty and registers a one-shot + // drain retry; it must not wait for an older write_pages() + // call to finish on the syscall stack. + Self::schedule_tagged_writeback_budget_retry( + &cache, + &inode, + start_index, + frozen_end, + epoch, + cursor, + ); + break; + }; + let claim = match self.claim_tagged_batch_with_admission( + &cache, + &inode, + WritebackBatchRange::new(target_index, tagged_end), + &target_entry, + epoch, + _domain_io.as_ref(), + ) { + Ok(claim) => claim, + Err(error) => { + drop(permit); + Self::abandon_tagged_writeback_generation( &cache, + start_index, + frozen_end, + epoch, + error.clone(), + ); + return Err(error); + } + }; + match claim { + WritebackClaimOutcome::Deferred(progress) => { + drop(permit); + // The callback is registered before WRITE returns, so a + // claim-time defer has a producer-owned progress edge + // rather than relying on an unrelated later reclaim. + Self::schedule_tagged_writeback_retry( + progress, + Arc::downgrade(&cache), + Arc::downgrade(&inode), + start_index, + frozen_end, + epoch, + target_index, + ); + break; + } + WritebackClaimOutcome::FailedRecorded(error) => { + drop(permit); + // The claimed batch has already been completed back to + // Dirty and its error recorded by PageCache. Only retire + // the remaining frozen tags here. + Self::retire_tagged_writeback_generation( + &cache, + start_index, + frozen_end, + epoch, + ); + return Err(error); + } + WritebackClaimOutcome::NoBatch => { + drop(permit); + let retry_incarnation = { + let inner = cache.inner.lock(); + inner.pages.get(&target_index).and_then(|current| { + (Arc::ptr_eq(current, &target_entry) + && inner.dirty_pages.contains(&target_index) + && current.writeback_tag() == epoch) + .then(|| current.writeback_incarnation.load(Ordering::Acquire)) + }) + }; + if let Some(observed_incarnation) = retry_incarnation { + if let Err(error) = Self::register_tagged_writeback_incarnation_retry( + &cache, + Arc::downgrade(&inode), + target_entry, + observed_incarnation, + TaggedWritebackCursor { + start_index, + frozen_end, + epoch, + cursor: target_index, + }, + ) { + Self::abandon_tagged_writeback_generation( + &cache, + start_index, + frozen_end, + epoch, + error.clone(), + ); + return Err(error); + } + break; + } + if target_index == usize::MAX { + Self::notify_tagged_writeback_progress(&cache); + break; + } + cursor = target_index + 1; + } + WritebackClaimOutcome::Claimed(batch) => { + let last_index = batch + .entries + .last() + .map(|(index, _, _)| *index) + .unwrap_or(target_index); + if batch.submission.is_none() { + // Legacy has no Deferred outcome, so it may retain + // established parallel background writeback. The + // worker clears the per-generation submission record + // only after PageCache completion and errseq + // publication make the batch result observable. + let work_state = Mutex::new(Some((permit, batch))); + schedule_pagecache_writeback(Work::new(move || { + let Some((permit, batch)) = work_state.lock().take() else { + return; + }; + let _permit = permit; + let _ = Self::submit_writeback_batch(batch); + })); + if last_index == usize::MAX { + break; + } + cursor = last_index + 1; + continue; + } + + let submit = if submit_first_token_inline { + Self::run_tagged_writeback_submission + } else { + Self::schedule_tagged_writeback_submission + }; + // Only the first Token head is run inline, after all + // scan/invalidate/admission locks have been released. + // Its successor uses the same asynchronous continuation. + submit( + Arc::downgrade(&cache), Arc::downgrade(&inode), - target_entry, - observed_incarnation, TaggedWritebackCursor { start_index, frozen_end, epoch, cursor: target_index, }, - ) { - Self::abandon_tagged_writeback_generation( - &cache, - start_index, - frozen_end, - epoch, - error.clone(), - ); - return Err(error); - } - break; - } - if target_index == usize::MAX { - Self::notify_tagged_writeback_progress(&cache); + last_index, + permit, + batch, + ); break; } - cursor = target_index + 1; - } - WritebackClaimOutcome::Claimed(batch) => { - let last_index = batch - .entries - .last() - .map(|(index, _, _)| *index) - .unwrap_or(target_index); - if batch.submission.is_none() { - // Legacy has no Deferred outcome, so it may retain - // established parallel background writeback. The - // worker clears the per-generation submission record - // only after PageCache completion and errseq - // publication make the batch result observable. - let work_state = Mutex::new(Some((permit, batch))); - schedule_pagecache_writeback(Work::new(move || { - let Some((permit, batch)) = work_state.lock().take() else { - return; - }; - let _permit = permit; - let _ = Self::submit_writeback_batch(batch); - })); - if last_index == usize::MAX { - break; - } - cursor = last_index + 1; - continue; - } - - Self::schedule_tagged_writeback_submission( - Arc::downgrade(&cache), - Arc::downgrade(&inode), - TaggedWritebackCursor { - start_index, - frozen_end, - epoch, - cursor: target_index, - }, - last_index, - permit, - batch, - ); - break; } } - } - Ok((frozen_filesystem_state, operation)) + Ok(()) + })(); + Ok(StartedWriteback { + frozen: frozen_filesystem_state, + range: operation, + error: dispatch.err(), + }) } /// Schedule one bounded reclaimer batch without doing page/MM work on the diff --git a/kernel/src/filesystem/vfs/mount/mod.rs b/kernel/src/filesystem/vfs/mount/mod.rs index 8b8ca2af7..e0cc51dc8 100644 --- a/kernel/src/filesystem/vfs/mount/mod.rs +++ b/kernel/src/filesystem/vfs/mount/mod.rs @@ -3192,13 +3192,41 @@ impl MountFS { return Ok(()); }; let mut last_err = Ok(()); - for page_cache in domain.snapshot() { - if !page_cache.has_dirty_or_writeback_work() { - continue; + let mut members = domain + .snapshot() + .into_iter() + .filter(|cache| cache.has_dirty_or_writeback_work()); + let mut pending = Vec::new(); + pending + .try_reserve_exact(crate::filesystem::page_cache::DOMAIN_WRITEBACK_WINDOW) + .map_err(|_| SystemError::ENOMEM)?; + loop { + let mut count = 0; + for page_cache in members + .by_ref() + .take(crate::filesystem::page_cache::DOMAIN_WRITEBACK_WINDOW) + { + count += 1; + match page_cache.start_sync_for_domain() { + Ok(writeback) => pending.push(writeback), + Err(error) => { + self.record_wb_error(error.clone()); + last_err = Err(error); + } + } } - if let Err(e) = page_cache.sync_for_domain() { - log::warn!("sync_inodes_of_mount: page cache sync failed: {:?}", e); - last_err = Err(e); + if count == 0 { + break; + } + // Submit across mappings before forcing durable completion. An + // outer guard would seal each small file before the next starts. + let _sync_request = self.inner_filesystem.begin_sync_writeback(); + for writeback in pending.drain(..) { + if let Err(error) = writeback.finish() { + log::warn!("sync_inodes_of_mount: page cache sync failed: {:?}", error); + self.record_wb_error(error.clone()); + last_err = Err(error); + } } } last_err @@ -3212,7 +3240,6 @@ impl MountFS { return Ok(()); } - let _sync_request = self.inner_filesystem.begin_sync_writeback(); self.sync_inodes_of_mount() } @@ -3276,11 +3303,6 @@ impl MountFS { return Ok(()); } - // Accepted filesystem metadata may complete the PageCache I/O below. - // Keep one request across this whole bounded sync, including later - // inode submissions, rather than allocating one guard per mapping. - let _sync_request = self.inner_filesystem.begin_sync_writeback(); - // writeback_inodes_sb(sb) — void let mut last_err = self.sync_inodes_of_mount(); // sync_fs(sb, 0) diff --git a/user/apps/tests/dunitest/suites/normal/ext4_inode_identity.cc b/user/apps/tests/dunitest/suites/normal/ext4_inode_identity.cc index 720624d0c..37a297ad7 100644 --- a/user/apps/tests/dunitest/suites/normal/ext4_inode_identity.cc +++ b/user/apps/tests/dunitest/suites/normal/ext4_inode_identity.cc @@ -728,6 +728,142 @@ TEST(Ext4InodeIdentity, LargeAppendBatchesAndPartialTailSurviveRemount) { ASSERT_NO_FATAL_FAILURE(fs.Unmount()); } +TEST(Ext4InodeIdentity, DomainSyncMixedFilesSurviveRemount) { + LoopExt4 fs; + ASSERT_NO_FATAL_FAILURE(fs.SetUp()); + ASSERT_NO_FATAL_FAILURE(fs.Mount()); + constexpr int kFiles = 18; // More than two eight-inode submission windows. + for (int file = 0; file < kFiles; ++file) { + const size_t size = (file == kFiles - 1 ? 129 : 3) * 4096 + 37; + const std::string data(size, static_cast('A' + file)); + const std::string path = fs.mount_point() + "/domain_" + std::to_string(file); + int fd = open(path.c_str(), O_CREAT | O_EXCL | O_WRONLY, 0600); + ASSERT_GE(fd, 0) << strerror(errno); + ASSERT_NO_FATAL_FAILURE(WriteAll(fd, data.data(), data.size())); + ASSERT_EQ(0, close(fd)); // No per-file fsync: exercise domain synchronization. + } + int directory = open(fs.mount_point().c_str(), O_RDONLY | O_DIRECTORY); + ASSERT_GE(directory, 0) << strerror(errno); + ASSERT_EQ(0, syscall(__NR_syncfs, directory)) << strerror(errno); + ASSERT_EQ(0, close(directory)); + ASSERT_NO_FATAL_FAILURE(fs.Unmount()); + ASSERT_NO_FATAL_FAILURE(fs.Mount()); + // Remount validates recovery and ownership, not a crash-at-syncfs boundary: + // unmount itself is allowed to perform additional synchronization. + for (int file = 0; file < kFiles; ++file) { + const size_t size = (file == kFiles - 1 ? 129 : 3) * 4096 + 37; + const std::string path = fs.mount_point() + "/domain_" + std::to_string(file); + int fd = open(path.c_str(), O_RDONLY); + ASSERT_GE(fd, 0) << strerror(errno); + struct stat st = {}; + ASSERT_EQ(0, fstat(fd, &st)); + ASSERT_EQ(static_cast(size), st.st_size); + std::string data(size, '\0'); + size_t done = 0; + while (done < size) { + ssize_t count = read(fd, data.data() + done, size - done); + if (count < 0 && errno == EINTR) { + continue; + } + ASSERT_GT(count, 0) << strerror(errno); + done += static_cast(count); + } + EXPECT_EQ(std::string(size, static_cast('A' + file)), data); + char tail; + EXPECT_EQ(0, read(fd, &tail, 1)); + ASSERT_EQ(0, close(fd)); + } + ASSERT_NO_FATAL_FAILURE(fs.Unmount()); +} + +TEST(Ext4InodeIdentity, ConcurrentDomainSyncAndRedirtyComplete) { + LoopExt4 fs; + ASSERT_NO_FATAL_FAILURE(fs.SetUp()); + ASSERT_NO_FATAL_FAILURE(fs.Mount()); + constexpr int kFiles = 17; + constexpr int kRounds = 16; + for (int file = 0; file < kFiles; ++file) { + const std::string path = fs.mount_point() + "/redirty_" + std::to_string(file); + int fd = open(path.c_str(), O_CREAT | O_EXCL | O_WRONLY, 0600); + ASSERT_GE(fd, 0) << strerror(errno); + const std::string data(4096, 'I'); + ASSERT_NO_FATAL_FAILURE(WriteAll(fd, data.data(), data.size())); + ASSERT_EQ(0, close(fd)); + } + int directory = open(fs.mount_point().c_str(), O_RDONLY | O_DIRECTORY); + ASSERT_GE(directory, 0) << strerror(errno); + std::atomic start{false}; + std::atomic first_error{0}; + auto record_error = [&](int error) { + int expected = 0; + first_error.compare_exchange_strong(expected, error); + }; + std::thread writer([&] { + while (!start.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + for (int round = 0; round < kRounds; ++round) { + for (int file = 0; file < kFiles; ++file) { + const std::string path = fs.mount_point() + "/redirty_" + std::to_string(file); + int fd = open(path.c_str(), O_WRONLY); + if (fd < 0) { + record_error(errno); + return; + } + const std::string data(4096, static_cast(1 + round + file)); + size_t done = 0; + while (done < data.size()) { + ssize_t count = pwrite(fd, data.data() + done, data.size() - done, done); + if (count < 0 && errno == EINTR) { + continue; + } + if (count <= 0) { + record_error(count < 0 ? errno : EIO); + close(fd); + return; + } + done += static_cast(count); + } + if (close(fd) != 0) { + record_error(errno); + return; + } + } + } + }); + std::thread syncer([&] { + while (!start.load(std::memory_order_acquire)) { + std::this_thread::yield(); + } + for (int round = 0; round < kRounds; ++round) { + if (syscall(__NR_syncfs, directory) != 0) { + record_error(errno); + return; + } + } + }); + start.store(true, std::memory_order_release); + writer.join(); + syncer.join(); + const int final_sync = syscall(__NR_syncfs, directory); + const int final_error = errno; + ASSERT_EQ(0, close(directory)); + ASSERT_EQ(0, first_error.load()) << strerror(first_error.load()); + ASSERT_EQ(0, final_sync) << strerror(final_error); + ASSERT_NO_FATAL_FAILURE(fs.Unmount()); + ASSERT_NO_FATAL_FAILURE(fs.Mount()); + for (int file = 0; file < kFiles; ++file) { + const std::string path = fs.mount_point() + "/redirty_" + std::to_string(file); + int fd = open(path.c_str(), O_RDONLY); + ASSERT_GE(fd, 0) << strerror(errno); + std::string data(4096, '\0'); + ASSERT_EQ(static_cast(data.size()), pread(fd, data.data(), data.size(), 0)); + EXPECT_EQ(std::string(4096, static_cast(kRounds + file)), data); + ASSERT_EQ(0, close(fd)); + } + ASSERT_NO_FATAL_FAILURE(fs.Unmount()); +} + TEST(Ext4InodeIdentity, ConcurrentDelallocInodesCompleteAndRecover) { constexpr int kWriters = 4; constexpr size_t kPageSize = 4096; From 0b3e3b9fe7eda44a8593f7052d659aa708a7e1bc Mon Sep 17 00:00:00 2001 From: longjin Date: Mon, 21 Sep 2026 14:03:29 +0000 Subject: [PATCH 2/2] fix(vfs): use inspect_err for writeback error reporting PageCache domain completion only records the write_inode error and passes it through unchanged. Use inspect_err rather than map_err to express that side effect and satisfy the denied clippy::manual_inspect lint invoked by make fmt. Validated with make fmt, FMT_CHECK=1 make fmt (the CI command), and make kernel. The error value, superblock reporting, and writeback completion ordering are unchanged. Signed-off-by: longjin --- kernel/src/filesystem/page_cache/writeback.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/kernel/src/filesystem/page_cache/writeback.rs b/kernel/src/filesystem/page_cache/writeback.rs index 31b637efa..2161ca949 100644 --- a/kernel/src/filesystem/page_cache/writeback.rs +++ b/kernel/src/filesystem/page_cache/writeback.rs @@ -52,10 +52,9 @@ impl PageCacheDomainWriteback { } self.inode .write_inode(&WritebackControl::sync_all_for_sync()) - .map_err(|error| { + .inspect_err(|error| { self.cache .record_writeback_error_with_superblock(error.clone()); - error }) } }