This is an automated email from the ASF dual-hosted git repository.

spetz pushed a commit to branch persisted_mode_perf
in repository https://gitbox.apache.org/repos/asf/iggy.git

commit c5cb3ec68a2af85bb57a4a7a173c73b60a0d5cd5
Author: spetz <[email protected]>
AuthorDate: Fri Sep 11 21:15:36 2026 +0200

    perf(partitions): cut persisted-mode disk traffic per acknowledgment
---
 core/journal/src/partition_journal.rs          | 461 ++++++++++++++++++++++---
 core/journal/src/partition_journal/segments.rs |  18 +-
 core/partitions/src/iggy_index_writer.rs       |  23 +-
 core/partitions/src/iggy_partition.rs          |  46 ++-
 core/partitions/src/install_backup.rs          |  21 ++
 core/partitions/src/persistence.rs             |  11 +-
 core/server/src/partition_helpers.rs           |   7 +-
 core/server/src/segment_recovery.rs            |  13 +-
 core/server/src/server_error.rs                |   7 +-
 core/simulator/src/storage/tests.rs            |  93 ++++-
 10 files changed, 602 insertions(+), 98 deletions(-)

diff --git a/core/journal/src/partition_journal.rs 
b/core/journal/src/partition_journal.rs
index 76d5b8223..ae11f3b1b 100644
--- a/core/journal/src/partition_journal.rs
+++ b/core/journal/src/partition_journal.rs
@@ -42,6 +42,24 @@ pub const PARTITION_WAL_BYTES_MAX: u64 = 256 * 1024 * 1024;
 pub const PARTITION_WAL_CAPACITY_MIN: u64 = 2 * (64 * 1024 * 1024 + 4096);
 pub const PARTITION_WAL_CAPACITY_MAX: u64 = 4 * 1024 * 1024 * 1024;
 const RECORD_PREFIX: usize = 32;
+/// The published frontier. Rewritten IN PLACE, so anything that freezes a
+/// partition's files must copy this one rather than retain it by hard link.
+pub const FRONTIER_FILE_NAME: &str = "frontier";
+/// Scratch name the complete slot file is built under before it is renamed 
over
+/// the frontier. Only [`PartitionPrepareJournal::install_frontier`] uses it, 
so
+/// it appears once per open and never on an acknowledgment path.
+const FRONTIER_TEMPORARY_NAME: &str = "frontier.tmp";
+/// Fixed slots the frontier alternates between, so a publication overwrites 
the
+/// older copy in place instead of creating and renaming a temporary file.
+const FRONTIER_SLOTS: usize = 2;
+const FRONTIER_BYTES: usize = FRONTIER_SLOTS * PARTITION_WAL_BLOCK_SIZE;
+/// Publication counter, placed directly after the last field the frontier 
block
+/// encodes so that adding a field cannot silently move it onto this one. 
Inside
+/// the range the block's own checksum covers, and reserved in every frontier
+/// this build and its predecessor wrote, so a block from before the two-slot
+/// layout reads back as sequence zero.
+const FRONTIER_SEQUENCE_OFFSET: usize = SEALED_STATE_MAGIC_OFFSET + 
size_of::<u64>();
+const _: () = assert!(FRONTIER_SEQUENCE_OFFSET + size_of::<u64>() <= 
PARTITION_WAL_BLOCK_SIZE);
 pub const PREPARE_BYTES_MAX: usize = 64 * 1024 * 1024;
 const STATE_MAGIC: &[u8; 8] = b"IGGYWAL1";
 const REFERENCE_STATE_MAGIC: &[u8; 8] = b"IGGYWAL2";
@@ -77,6 +95,8 @@ pub trait DurableAppend {
 pub struct PartitionPrepareJournal<S: DurableStorage = DiskStorage> {
     directory: PathBuf,
     file: S::File,
+    frontier: S::File,
+    frontier_sequence: u64,
     storage: S,
     capacity: u64,
     state: JournalState,
@@ -183,19 +203,17 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
         }
         storage.create_directories(directory).await?;
         storage.sync_directory(parent).await?;
-        let state_path = directory.join("frontier");
-        let existing = match storage.open(&state_path, OpenMode::Read).await {
-            Ok(file) => {
-                let bytes = file.read(0, PARTITION_WAL_BLOCK_SIZE).await?;
-                let state = JournalState::decode(&bytes)?;
-                if state.group != group || state.incarnation != incarnation {
-                    return Err(invalid("partition WAL identity mismatch"));
-                }
-                Some(state)
-            }
-            Err(error) if error.kind() == io::ErrorKind::NotFound => None,
-            Err(error) => return Err(error),
-        };
+        let state_path = directory.join(FRONTIER_FILE_NAME);
+        let (published, slots) = Self::open_frontier(&storage, 
&state_path).await?;
+        let (existing, sequence, verified) = 
Self::read_frontier(published.as_ref(), slots).await?;
+        if slots > 0 && verified == 0 {
+            return Err(invalid("unreadable partition WAL frontier"));
+        }
+        if let Some(state) = existing
+            && (state.group != group || state.incarnation != incarnation)
+        {
+            return Err(invalid("partition WAL identity mismatch"));
+        }
         if existing.is_none() {
             Self::validate_unpublished_history(&storage, directory).await?;
         }
@@ -216,9 +234,26 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
         if file.length().await? < state.length {
             return Err(invalid("partition WAL lost acknowledged bytes"));
         }
+        // After the data file it names, never before: a frontier is visible 
the
+        // moment its rename lands, and one naming a generation that does not
+        // exist is unrecoverable history rather than a fresh journal.
+        let (frontier, frontier_sequence) = match published {
+            Some(file) if slots == FRONTIER_SLOTS => (file, sequence),
+            _ => {
+                let sequence =
+                    Self::install_frontier(&storage, directory, &state_path, 
state, sequence)
+                        .await?;
+                (
+                    storage.open(&state_path, OpenMode::ReadWrite).await?,
+                    sequence,
+                )
+            }
+        };
         let mut journal = Self {
             directory: directory.to_path_buf(),
             file,
+            frontier,
+            frontier_sequence,
             storage,
             capacity,
             state,
@@ -234,27 +269,121 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
             preallocate_segments,
             retained_bytes: 0,
         };
-        journal.recover_entries().await?;
-        // Only bytes covered by the durable frontier could have released an 
ack.
-        journal.file.truncate(state.length).await?;
+        // A slot that did not verify may have been the newest, so the frontier
+        // this open read can name less than an acknowledgment already covered.
+        // Only then does recovery walk past the frontier.
+        journal.recover_entries(verified < slots).await?;
+        // Only verified bytes could have released an ack. Recovery adopted 
every
+        // record past the published frontier whose envelope, chain and body 
all
+        // verify; what follows them is a tail no barrier ever covered.
+        let recovered = journal.state;
+        journal.durable_head = recovered.head;
+        journal.file.truncate(recovered.length).await?;
         journal.file.sync().await?;
         journal.storage.sync_directory(directory).await?;
-        if existing.is_none() {
-            journal.publish(state).await?;
-        }
-        journal.discover_obsolete().await?;
-        loop {
-            let remaining = journal.obsolete.len();
-            journal.cleanup_obsolete().await;
-            if journal.obsolete.is_empty() || journal.obsolete.len() == 
remaining {
-                break;
-            }
+        if recovered != state {
+            journal.publish(recovered).await?;
         }
+        journal.remove_obsolete_history().await?;
         journal.recover_segment_files().await?;
         journal.migrate_segment_prepares().await?;
         Ok(journal)
     }
 
+    /// Open the frontier slot file without creating it, with the number of 
whole
+    /// slots it holds.
+    ///
+    /// A visible frontier always holds whole slots: it is only ever created, 
or
+    /// grown to its full slot count, by [`Self::install_frontier`], through a
+    /// rename. Publication then overwrites one slot of a file it never 
resizes.
+    /// A length in between belongs to no protocol this build can read.
+    async fn open_frontier(storage: &S, path: &Path) -> 
io::Result<(Option<S::File>, usize)> {
+        let published = match storage.open(path, OpenMode::ReadWrite).await {
+            Ok(file) => Some(file),
+            Err(error) if error.kind() == io::ErrorKind::NotFound => None,
+            Err(error) => return Err(error),
+        };
+        let length = match &published {
+            Some(file) => file.length().await?,
+            None => 0,
+        };
+        if length > FRONTIER_BYTES as u64 || 
!length.is_multiple_of(PARTITION_WAL_BLOCK_SIZE as u64)
+        {
+            return Err(invalid("unknown partition WAL frontier size"));
+        }
+        let slots = usize::try_from(length)
+            .map_err(|_| invalid("unknown partition WAL frontier size"))?
+            / PARTITION_WAL_BLOCK_SIZE;
+        Ok((published, slots))
+    }
+
+    /// Install a complete slot file, atomically, and return its newest 
sequence.
+    ///
+    /// Runs once per open: for a journal that has no frontier yet, and for one
+    /// whose frontier predates the second slot. Both slots carry `state`, so 
the
+    /// first in-place publication always has an intact partner to fall back 
on,
+    /// and a visible frontier never holds a slot no publication completed. 
That
+    /// is what lets [`Self::open_frontier`] trust the file it finds: a torn
+    /// write inside this temporary name never becomes the frontier, and a torn
+    /// write afterwards can only damage the slot being published.
+    async fn install_frontier(
+        storage: &S,
+        directory: &Path,
+        path: &Path,
+        state: JournalState,
+        sequence: u64,
+    ) -> io::Result<u64> {
+        let temporary = directory.join(FRONTIER_TEMPORARY_NAME);
+        let mut file = storage.open(&temporary, OpenMode::Create).await?;
+        let mut newest = sequence;
+        for _ in 0..FRONTIER_SLOTS {
+            newest = newest
+                .checked_add(1)
+                .ok_or_else(|| invalid("WAL frontier sequence exhausted"))?;
+            file.write(frontier_offset(newest), state.encode(newest))
+                .await?;
+        }
+        file.sync().await?;
+        storage.rename(&temporary, path).await?;
+        storage.sync_directory(directory).await?;
+        Ok(newest)
+    }
+
+    /// The newest intact slot wins. A slot that fails verification is either 
the
+    /// interrupted half of the publication in flight, which never released an
+    /// acknowledgment, or a copy the newer one superseded; either way its
+    /// partner is the frontier. Losing the NEWEST slot rolls the published
+    /// prefix back one publication, which [`Self::recover_unpublished_tail`]
+    /// then re-adopts from the records themselves.
+    ///
+    /// Returns the newest intact slot, its publication sequence, and how many
+    /// of the `slots` whole blocks verified. A slot that does not verify is
+    /// either the interrupted half of a publication in flight or a copy the
+    /// newer one superseded, so the count is what tells the caller whether the
+    /// frontier it got could be older than an acknowledgment already covered.
+    async fn read_frontier(
+        file: Option<&S::File>,
+        slots: usize,
+    ) -> io::Result<(Option<JournalState>, u64, usize)> {
+        let mut newest = None;
+        let mut sequence = 0;
+        let mut verified = 0;
+        let Some(file) = file.filter(|_| slots > 0) else {
+            return Ok((newest, sequence, verified));
+        };
+        let bytes = file.read(0, slots * PARTITION_WAL_BLOCK_SIZE).await?;
+        for slot in bytes.as_chunks::<PARTITION_WAL_BLOCK_SIZE>().0 {
+            if let Ok((state, slot_sequence)) = JournalState::decode(slot) {
+                verified += 1;
+                if newest.is_none() || slot_sequence > sequence {
+                    newest = Some(state);
+                    sequence = slot_sequence;
+                }
+            }
+        }
+        Ok((newest, sequence, verified))
+    }
+
     pub const fn certified_log_view(&self) -> Option<u32> {
         self.state.certified_log_view
     }
@@ -350,7 +479,8 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
     pub async fn prepares(&self) -> io::Result<Vec<Message<PrepareHeader>>> {
         let mut prepares = Vec::with_capacity(self.entries.len());
         for entry in self.entries.values() {
-            let (_, length, prepare, _) = 
self.read_record(entry.position).await?;
+            let (_, length, prepare, _) =
+                self.read_record(entry.position, self.state.length).await?;
             if length != entry.length {
                 return Err(invalid("partition WAL index length mismatch"));
             }
@@ -534,6 +664,11 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
 
     /// Make every buffered predecessor recoverable with one frontier 
publication.
     ///
+    /// The body barrier and the WAL barrier cover different inodes, so they 
run
+    /// under one `join` and a batch pays one barrier latency rather than two.
+    /// Their completion order does not matter, because neither is the
+    /// acknowledgment point: [`Self::publish`] is, and it follows both.
+    ///
     /// # Errors
     /// Returns an error unless the buffered prefix and its frontier are 
durable.
     pub async fn sync(&mut self) -> io::Result<()> {
@@ -542,8 +677,12 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
             return Ok(());
         }
         self.poisoned = true;
-        self.sync_segment_files().await?;
-        self.file.sync().await?;
+        let (bodies, records) =
+            futures::future::join(self.segment_barrier(), 
self.file.sync()).await;
+        bodies?;
+        records?;
+        self.segment_files_dirty = false;
+        self.segment_links_dirty = false;
         self.publish(self.state).await?;
         self.durable_head = self.state.head;
         self.retain_active_segment_file();
@@ -741,13 +880,17 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
         Ok(())
     }
 
-    async fn recover_entries(&mut self) -> io::Result<()> {
+    /// Rebuild the entry index from the published prefix, which must verify in
+    /// full. `lost_publication` extends the walk past that prefix through
+    /// [`Self::recover_unpublished_tail`].
+    async fn recover_entries(&mut self, lost_publication: bool) -> 
io::Result<()> {
         let state = self.state;
         let mut position = 0;
         let mut previous = state.checkpoint;
         let mut checksum = state.checkpoint_checksum;
         while position < state.length {
-            let (header, length, prepare, reference) = 
self.read_record(position).await?;
+            let (header, length, prepare, reference) =
+                self.read_record(position, state.length).await?;
             let next_offset = if state.segment_storage.is_some() && 
reference.is_some() {
                 Some(segments::batch_next_offset(prepare.as_slice())?)
             } else {
@@ -790,9 +933,110 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
             .values()
             .map(|entry| entry.retained_bytes)
             .sum();
+        if lost_publication {
+            self.recover_unpublished_tail(previous, checksum).await
+        } else {
+            Ok(())
+        }
+    }
+
+    /// Recover the records a frontier publication this open could not read had
+    /// already covered.
+    ///
+    /// Runs only when a slot failed to verify, which can happen only while a
+    /// publication was in flight. The writer is serial and publishes after 
both
+    /// data barriers, so a publication in flight proves every record before it
+    /// was already durable: between the surviving frontier and the end of the
+    /// data file the records are complete, and every one of them is adopted.
+    ///
+    /// That is why this walk REFUSES instead of stopping. Once a slot is lost,
+    /// nothing left on disk says how far acknowledgment had reached, so a 
record
+    /// that does not verify there cannot be dismissed as an unwritten tail. It
+    /// is damage to history that may have been acknowledged, and the partition
+    /// has to fence and rebuild from its peers rather than open a truncated 
log.
+    /// When both slots verify the frontier is exact and the tail is discarded 
as
+    /// it always was.
+    async fn recover_unpublished_tail(
+        &mut self,
+        mut previous: u64,
+        mut checksum: u128,
+    ) -> io::Result<()> {
+        let limit = self.file.length().await?;
+        if limit > self.capacity {
+            return Err(invalid("partition WAL recovered tail exceeds 
capacity"));
+        }
+        let mut position = self.state.length;
+        while position < limit {
+            let (header, length, prepare, reference) = 
self.read_record(position, limit).await?;
+            let next_op = previous
+                .checked_add(1)
+                .ok_or_else(|| invalid("WAL op overflow"))?;
+            if header.op != next_op || (self.state.anchor_known && 
header.parent != checksum) {
+                return Err(invalid("partition WAL prepare chain is broken"));
+            }
+            let body_bytes = record_length(header.size as usize)? as u64;
+            if self.retained_bytes.saturating_add(body_bytes) > self.capacity {
+                return Err(invalid("partition WAL recovered tail exceeds 
capacity"));
+            }
+            let mut state = self.state;
+            let next_offset = if let Some(segments) = &mut 
state.segment_storage {
+                let (reserved, next_offset) = segments.reserve(&header, 
prepare.as_slice())?;
+                if reserved != reference || !segments.valid() {
+                    return Err(invalid("partition WAL recovered segment 
boundary mismatch"));
+                }
+                next_offset
+            } else {
+                None
+            };
+            state.length = position + length as u64;
+            state.head = header.op;
+            state.head_checksum = header.checksum;
+            state.segment_references |= reference.is_some();
+            if state
+                .certified_log_view
+                .is_some_and(|view| header.view > view)
+            {
+                state.certified_log_view = None;
+            }
+            if !state.anchor_known {
+                state.checkpoint_checksum = header.parent;
+                state.anchor_known = true;
+            }
+            self.recovered_prepares.push(prepare);
+            self.entries.insert(
+                header.op,
+                StoredPrepare {
+                    position,
+                    length,
+                    checksum: header.checksum,
+                    reference,
+                    next_offset,
+                    retained_bytes: body_bytes,
+                },
+            );
+            self.retained_bytes += body_bytes;
+            self.state = state;
+            previous = header.op;
+            checksum = header.checksum;
+            position += length as u64;
+        }
         Ok(())
     }
 
+    /// Drop every file the recovered history does not retain, repeating while
+    /// the queue shrinks: `cleanup_obsolete` removes a bounded batch per call
+    /// and re-queues what it could not remove.
+    async fn remove_obsolete_history(&mut self) -> io::Result<()> {
+        self.discover_obsolete().await?;
+        loop {
+            let remaining = self.obsolete.len();
+            self.cleanup_obsolete().await;
+            if self.obsolete.is_empty() || self.obsolete.len() == remaining {
+                return Ok(());
+            }
+        }
+    }
+
     async fn discover_obsolete(&mut self) -> io::Result<()> {
         let retained = self.retained_segment_paths(&self.entries, 
self.state.segment_storage);
         for entry in self.storage.entries(&self.directory).await? {
@@ -805,7 +1049,7 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
                 .and_then(|value| value.parse::<u64>().ok());
             if !entry.directory
                 && (generation.is_some_and(|generation| generation != 
self.state.generation)
-                    || name == "frontier.tmp"
+                    || name == FRONTIER_TEMPORARY_NAME
                     || (is_retained_segment_name(name)
                         && 
!retained.contains(&self.directory.join(&entry.name))))
             {
@@ -842,7 +1086,7 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
 
     async fn validate_unpublished_history(storage: &S, directory: &Path) -> 
io::Result<()> {
         for entry in storage.entries(directory).await? {
-            if entry.directory || entry.name == "frontier.tmp" {
+            if entry.directory || entry.name == FRONTIER_TEMPORARY_NAME {
                 continue;
             }
             // An interrupted first open can leave only its empty 
generation-zero file.
@@ -876,13 +1120,15 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
     async fn read_record(
         &self,
         position: u64,
+        limit: u64,
     ) -> io::Result<(
         PrepareHeader,
         usize,
         Message<PrepareHeader>,
         Option<SegmentReference>,
     )> {
-        let (mut buffer, frame_length, reference) = 
self.read_encoded_record(position).await?;
+        let (mut buffer, frame_length, reference) =
+            self.read_encoded_record(position, limit).await?;
         let length = buffer.as_slice().len();
         if let Some(reference) = reference {
             let payload = &buffer.as_slice()[RECORD_PREFIX..RECORD_PREFIX + 
frame_length];
@@ -935,6 +1181,7 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
     async fn read_encoded_record(
         &self,
         position: u64,
+        limit: u64,
     ) -> io::Result<(Owned<4096>, usize, Option<SegmentReference>)> {
         let prefix = self
             .file
@@ -948,7 +1195,7 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
         let length = record_length(frame_length)?;
         if position
             .checked_add(length as u64)
-            .is_none_or(|end| end > self.state.length)
+            .is_none_or(|end| end > limit)
         {
             return Err(invalid("partition WAL record crosses durable 
frontier"));
         }
@@ -1018,21 +1265,41 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
         Ok((buffer, frame_length, reference))
     }
 
-    async fn publish(&self, state: JournalState) -> io::Result<()> {
+    /// Publish the frontier by overwriting the older of two fixed slots.
+    ///
+    /// Both slots exist and hold a complete record from the moment the journal
+    /// opens ([`Self::install_frontier`]), so a publication is one 4096-byte
+    /// overwrite plus `fdatasync`: no create, no truncate, no size change, no
+    /// rename and no directory barrier. On a journaling filesystem that is the
+    /// difference between zero metadata transactions per acknowledgment and
+    /// roughly two, which at thousands of acknowledgments per second per node
+    /// costs more than the prepare bytes the batch carries.
+    ///
+    /// A torn slot fails its checksum and its partner still holds the previous
+    /// publication, the guarantee the temporary-file rename used to provide.
+    /// Unlike the rename, an unreadable NEWEST slot leaves the published 
prefix
+    /// one publication behind; [`Self::recover_unpublished_tail`] recovers it
+    /// from the records themselves.
+    ///
+    /// The slot follows the sequence's parity, so a publication never 
overwrites
+    /// the copy it would have to fall back on.
+    async fn publish(&mut self, state: JournalState) -> io::Result<()> {
         if state
             .segment_storage
             .is_some_and(|segments| !segments.valid())
         {
             return Err(invalid("invalid durable segment boundaries"));
         }
-        let temporary = self.directory.join("frontier.tmp");
-        let mut file = self.storage.open(&temporary, OpenMode::Create).await?;
-        file.write(0, state.encode()).await?;
-        file.sync().await?;
-        self.storage
-            .rename(&temporary, &self.directory.join("frontier"))
+        let sequence = self
+            .frontier_sequence
+            .checked_add(1)
+            .ok_or_else(|| invalid("WAL frontier sequence exhausted"))?;
+        self.frontier
+            .write(frontier_offset(sequence), state.encode(sequence))
             .await?;
-        self.storage.sync_directory(&self.directory).await
+        self.frontier.sync().await?;
+        self.frontier_sequence = sequence;
+        Ok(())
     }
 
     #[allow(clippy::too_many_lines)]
@@ -1113,8 +1380,9 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
             if op < checkpoint || truncate.is_some_and(|from| op >= from) {
                 continue;
             }
-            let (record, payload_length, mut reference) =
-                self.read_encoded_record(entry.position).await?;
+            let (record, payload_length, mut reference) = self
+                .read_encoded_record(entry.position, self.state.length)
+                .await?;
             let mut next_offset = entry.next_offset;
             // Purge removes polled data, but these operations still 
participate
             // in repair. Inline their bodies before releasing whole segment 
inodes.
@@ -1122,7 +1390,8 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
                 && reference.is_some()
                 && op <= state.purge_floor
             {
-                let (_, _, prepare, _) = 
self.read_record(entry.position).await?;
+                let (_, _, prepare, _) =
+                    self.read_record(entry.position, self.state.length).await?;
                 reference = None;
                 next_offset = None;
                 Some(prepare)
@@ -1231,8 +1500,10 @@ impl<S: DurableStorage> DurableAppend for 
PartitionPrepareJournal<S> {
 }
 
 impl JournalState {
-    fn encode(self) -> Vec<u8> {
+    fn encode(self, sequence: u64) -> Vec<u8> {
         let mut bytes = vec![0; PARTITION_WAL_BLOCK_SIZE];
+        bytes[FRONTIER_SEQUENCE_OFFSET..FRONTIER_SEQUENCE_OFFSET + 
size_of::<u64>()]
+            .copy_from_slice(&sequence.to_le_bytes());
         bytes[..8].copy_from_slice(if self.segment_references {
             REFERENCE_STATE_MAGIC
         } else {
@@ -1271,7 +1542,7 @@ impl JournalState {
         bytes
     }
 
-    fn decode(bytes: &[u8]) -> io::Result<Self> {
+    fn decode(bytes: &[u8]) -> io::Result<(Self, u64)> {
         if bytes.len() != PARTITION_WAL_BLOCK_SIZE
             || (&bytes[..8] != STATE_MAGIC && &bytes[..8] != 
REFERENCE_STATE_MAGIC)
         {
@@ -1348,10 +1619,16 @@ impl JournalState {
         {
             return Err(invalid("invalid partition WAL frontier bounds"));
         }
-        Ok(state)
+        Ok((state, read_u64(FRONTIER_SEQUENCE_OFFSET)?))
     }
 }
 
+/// Byte offset of the slot a publication sequence owns. Alternating by parity 
is
+/// what keeps a publication off the copy it would have to fall back on.
+const fn frontier_offset(sequence: u64) -> u64 {
+    (sequence % FRONTIER_SLOTS as u64) * PARTITION_WAL_BLOCK_SIZE as u64
+}
+
 /// Padded size of a prepare record, including its envelope.
 ///
 /// # Errors
@@ -1916,6 +2193,76 @@ mod tests {
         );
     }
 
+    #[compio::test]
+    async fn a_lost_frontier_slot_recovers_the_acknowledged_tail() {
+        let directory = tempdir().unwrap();
+        let mut journal = PartitionPrepareJournal::open(directory.path(), 42, 
7)
+            .await
+            .unwrap();
+        let first = prepare(1, 0);
+        let second = prepare(2, first.header().checksum);
+        journal.append(first.into_frozen()).await.unwrap();
+        journal.append(second.into_frozen()).await.unwrap();
+        // Publication alternates two slots of one file in place: no temporary
+        // name is created and renamed per acknowledgment.
+        let path = directory.path().join("frontier");
+        assert_eq!(
+            std::fs::metadata(&path).unwrap().len(),
+            FRONTIER_BYTES as u64
+        );
+        assert!(!directory.path().join("frontier.tmp").exists());
+        let newest = usize::try_from(journal.frontier_sequence % 
FRONTIER_SLOTS as u64).unwrap()
+            * PARTITION_WAL_BLOCK_SIZE;
+        drop(journal);
+        let mut bytes = std::fs::read(&path).unwrap();
+        bytes[newest..newest + PARTITION_WAL_BLOCK_SIZE].fill(0);
+        std::fs::write(&path, bytes).unwrap();
+        // The surviving slot names op 1, but op 2 was acknowledged: recovery
+        // walks past the frontier it could read and adopts the verified 
record.
+        let journal = PartitionPrepareJournal::open(directory.path(), 42, 7)
+            .await
+            .unwrap();
+        assert_eq!(journal.head(), 2);
+        assert_eq!(journal.prepares().await.unwrap().len(), 2);
+    }
+
+    #[compio::test]
+    async fn a_lost_frontier_slot_does_not_hide_damaged_acknowledged_records() 
{
+        for damage in ["zeroed", "shortened"] {
+            let directory = tempdir().unwrap();
+            let mut journal = PartitionPrepareJournal::open(directory.path(), 
42, 7)
+                .await
+                .unwrap();
+            let first = prepare(1, 0);
+            let second = prepare(2, first.header().checksum);
+            journal.append(first.into_frozen()).await.unwrap();
+            journal.append(second.into_frozen()).await.unwrap();
+            let newest = 
usize::try_from(frontier_offset(journal.frontier_sequence)).unwrap();
+            let position = 
usize::try_from(journal.entries[&2].position).unwrap();
+            let data = data_path(directory.path(), journal.state.generation);
+            drop(journal);
+            let frontier = directory.path().join(FRONTIER_FILE_NAME);
+            let mut slots = std::fs::read(&frontier).unwrap();
+            slots[newest..newest + PARTITION_WAL_BLOCK_SIZE].fill(0);
+            std::fs::write(&frontier, &slots).unwrap();
+            let mut records = std::fs::read(&data).unwrap();
+            match damage {
+                "zeroed" => records[position..].fill(0),
+                "shortened" => records.truncate(position + 
PARTITION_WAL_BLOCK_SIZE / 2),
+                _ => unreachable!(),
+            }
+            std::fs::write(&data, &records).unwrap();
+            assert!(
+                PartitionPrepareJournal::open(directory.path(), 42, 7)
+                    .await
+                    .is_err(),
+                "{damage}: recovery must refuse uncertain acknowledged history"
+            );
+            assert_eq!(std::fs::read(&data).unwrap(), records);
+            assert_eq!(std::fs::read(&frontier).unwrap(), slots);
+        }
+    }
+
     #[compio::test]
     async fn durable_append_covers_buffered_predecessors() {
         let directory = tempdir().unwrap();
@@ -2184,8 +2531,8 @@ mod tests {
             incarnation: 7,
             ..JournalState::default()
         };
-        let mut bytes = state.encode();
-        assert_eq!(JournalState::decode(&bytes).unwrap(), state);
+        let mut bytes = state.encode(1);
+        assert_eq!(JournalState::decode(&bytes).unwrap().0, state);
         bytes[24] ^= 1;
         assert!(JournalState::decode(&bytes).is_err());
     }
@@ -2457,7 +2804,7 @@ mod tests {
                 }
                 _ => unreachable!(),
             }
-            let encoded = state.encode();
+            let encoded = state.encode(1);
             assert!(JournalState::decode(&encoded).is_err(), "{invalid}");
             std::fs::write(directory.join("frontier"), encoded).unwrap();
             drop(journal);
@@ -2509,7 +2856,7 @@ mod tests {
             .unwrap();
         let mut state = journal.state;
         state.segment_storage.as_mut().unwrap().max_size = 
iggy_common::MAX_TOPIC_SEGMENT_SIZE + 1;
-        let encoded = state.encode();
+        let encoded = state.encode(1);
         assert!(JournalState::decode(&encoded).is_err());
         std::fs::write(directory.join("frontier"), encoded).unwrap();
         drop(journal);
@@ -2537,8 +2884,8 @@ mod tests {
                 segment_references,
                 ..JournalState::default()
             };
-            let mut encoded = state.encode();
-            assert_eq!(JournalState::decode(&encoded).unwrap(), state);
+            let mut encoded = state.encode(1);
+            assert_eq!(JournalState::decode(&encoded).unwrap().0, state);
             // The old reader hashes the same payload and ignores reserved 
bytes.
             assert_eq!(
                 u64::from_le_bytes(encoded[8..16].try_into().unwrap()),
@@ -2550,11 +2897,11 @@ mod tests {
                 REFERENCE_STATE_MAGIC
             });
             assert!(JournalState::decode(&encoded).is_err());
-            let mut legacy = state.encode();
+            let mut legacy = state.encode(1);
             legacy[SEALED_STATE_MAGIC_OFFSET..SEALED_STATE_MAGIC_OFFSET + 
8].fill(0);
             let checksum = XxHash3_64::oneshot(&legacy[16..]);
             legacy[8..16].copy_from_slice(&checksum.to_le_bytes());
-            assert_eq!(JournalState::decode(&legacy).unwrap(), state);
+            assert_eq!(JournalState::decode(&legacy).unwrap().0, state);
         }
     }
 
diff --git a/core/journal/src/partition_journal/segments.rs 
b/core/journal/src/partition_journal/segments.rs
index 6ff31cbb8..84ca7e57b 100644
--- a/core/journal/src/partition_journal/segments.rs
+++ b/core/journal/src/partition_journal/segments.rs
@@ -335,7 +335,8 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
             if entry.reference.is_some() {
                 continue;
             }
-            let (header, _, prepare, _) = 
self.read_record(entry.position).await?;
+            let (header, _, prepare, _) =
+                self.read_record(entry.position, self.state.length).await?;
             if header.operation == Operation::SendMessages
                 && decode_batch(prepare.as_slice())?.base_offset
                     >= segments.tail.position.next_offset
@@ -425,25 +426,34 @@ impl<S: DurableStorage> PartitionPrepareJournal<S> {
             .ok_or_else(|| invalid("segment writing handle is absent"))
     }
 
-    pub(super) async fn sync_segment_files(&mut self) -> io::Result<()> {
+    /// Barrier over every body and link this journal wrote, leaving the dirty
+    /// flags alone. Split out of [`Self::sync_segment_files`] so that
+    /// `PartitionPrepareJournal::sync` can overlap it with the WAL file's own
+    /// barrier and clear both flags once the pair has completed.
+    pub(super) async fn segment_barrier(&self) -> io::Result<()> {
         if self.segment_files_dirty {
             for file in self.segment_files.values() {
                 // Keep the writing handle: reopening after an errseq 
writeback error
                 // could turn a failed body barrier into a successful 
acknowledgment.
                 file.sync().await?;
             }
-            self.segment_files_dirty = false;
         }
         if self.segment_links_dirty {
             self.storage.sync_directory(&self.directory).await?;
             self.storage
                 .sync_directory(self.segment_directory()?)
                 .await?;
-            self.segment_links_dirty = false;
         }
         Ok(())
     }
 
+    pub(super) async fn sync_segment_files(&mut self) -> io::Result<()> {
+        self.segment_barrier().await?;
+        self.segment_files_dirty = false;
+        self.segment_links_dirty = false;
+        Ok(())
+    }
+
     pub(super) fn retain_active_segment_file(&mut self) {
         // Buffered rotations must retain their original writers until 
publication.
         if let Some(segments) = self.state.segment_storage {
diff --git a/core/partitions/src/iggy_index_writer.rs 
b/core/partitions/src/iggy_index_writer.rs
index 78ab2c0eb..ddd6f294d 100644
--- a/core/partitions/src/iggy_index_writer.rs
+++ b/core/partitions/src/iggy_index_writer.rs
@@ -106,6 +106,14 @@ impl IggyIndexWriter {
     ///
     /// Returns an error if the index bytes cannot be written or synced to 
disk.
     pub(crate) async fn save_indexes(&self, indexes: Vec<u8>) -> Result<u64, 
IggyError> {
+        let saved = self.save_indexes_buffered(indexes).await?;
+        if saved > 0 && self.fsync {
+            self.fsync().await?;
+        }
+        Ok(saved)
+    }
+
+    pub(crate) async fn save_indexes_buffered(&self, indexes: Vec<u8>) -> 
Result<u64, IggyError> {
         if indexes.is_empty() {
             return Ok(0);
         }
@@ -119,10 +127,6 @@ impl IggyIndexWriter {
             .0
             .map_err(|_| IggyError::CannotSaveIndexToSegment)?;
 
-        if self.fsync {
-            self.fsync().await?;
-        }
-
         trace!(
             target: "iggy.partitions.storage",
             file = self.file_path.as_str(),
@@ -162,6 +166,17 @@ impl IggyIndexWriter {
 mod tests {
     use super::*;
 
+    #[cfg(target_os = "linux")]
+    #[compio::test]
+    async fn 
buffered_indexes_defer_sync_errors_to_the_original_writer_barrier() {
+        let writer = IggyIndexWriter::new("/dev/null", 
Rc::new(AtomicU64::new(0)), true, false)
+            .await
+            .unwrap();
+        assert_eq!(writer.save_indexes_buffered(vec![1; 32]).await.unwrap(), 
32);
+        assert!(writer.fsync().await.is_err());
+        assert!(writer.save_indexes(vec![1; 32]).await.is_err());
+    }
+
     #[compio::test]
     async fn 
given_seeded_size_diverging_from_disk_when_opening_existing_file_should_return_size_mismatch_error()
      {
diff --git a/core/partitions/src/iggy_partition.rs 
b/core/partitions/src/iggy_partition.rs
index b4849701b..f70e18f2c 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -975,6 +975,17 @@ where
                 return;
             }
         }
+        if let Some(writer) = 
self.log.index_writers().last().and_then(Option::as_ref)
+            && let Err(error) = writer.fsync().await
+        {
+            error!(%error, namespace_raw = self.namespace().inner(), 
"partition checkpoint index sync failed");
+            self.fatal = Some(FatalCommit {
+                namespace_raw: self.namespace().inner(),
+                op: through_op,
+                operation: Operation::SendMessages,
+            });
+            return;
+        }
         let (files, directories) = self.persistence_checkpoint_files(config);
         persistence.checkpoint_files(through_op, files, directories);
         self.start_persistence();
@@ -6583,7 +6594,7 @@ where
                 .last()
                 .and_then(|writer| writer.as_ref())
                 .ok_or(IggyError::CannotWriteToFile)?;
-            let saved_indexes = index_writer.save_indexes(index_bytes).await?;
+            let saved_indexes = 
index_writer.save_indexes_buffered(index_bytes).await?;
             index_writer.advance(saved_indexes);
             if let Some(writer) = self
                 .log
@@ -6747,6 +6758,11 @@ where
                 "a plant at {start_offset} leaves a gap past {sealed_end} with 
no anchor"
             );
         }
+        if self.persistence.is_some()
+            && let Some(writer) = &self.log.index_writers()[sealed_index]
+        {
+            writer.fsync().await?;
+        }
         self.log.active_segment_mut().sealed = true;
         self.install_empty_segment(config, start_offset).await?;
         self.stats.increment_segments_count(1);
@@ -8694,6 +8710,34 @@ mod tests {
         assert!(!directory.path().join("prepares-0").exists());
     }
 
+    #[cfg(target_os = "linux")]
+    #[compio::test]
+    async fn 
checkpoint_index_sync_failure_fences_before_reclaiming_wal_history() {
+        let directory = tempfile::tempdir().unwrap();
+        let (mut partition, _) = recording_partition_at(0, 3);
+        
partition.set_partition_dir(directory.path().to_string_lossy().into_owned());
+        partition.runtime_options.durability = 
iggy_common::Durability::Persisted;
+        partition.open_persistence().await.unwrap();
+        let persistence = Rc::clone(partition.persistence.as_ref().unwrap());
+        let prepare = checksummed_segment_prepare(1, 0, 0, b"durable");
+        persistence.append(prepare.into_frozen(), true).unwrap();
+        assert!(persistence.start());
+        Rc::clone(&persistence).run().await;
+        assert!(persistence.is_durable_through(1));
+        partition.consensus.restore_commit_state(1, 1);
+        let writer = IggyIndexWriter::new("/dev/null", 
Rc::new(AtomicU64::new(0)), true, false)
+            .await
+            .unwrap();
+        assert_eq!(writer.save_indexes_buffered(vec![1; 32]).await.unwrap(), 
32);
+        let active = partition.log.index_writers().len() - 1;
+        partition.log.index_writers_mut()[active] = Some(Rc::new(writer));
+        persistence.request_checkpoint();
+        partition.checkpoint_persistence(&repair_config()).await;
+        assert!(partition.fatal().is_some());
+        assert!(!persistence.checkpoint_pending());
+        assert!(persistence.is_durable_through(1));
+    }
+
     #[compio::test]
     async fn 
pending_wal_prefix_keeps_pipeline_replies_and_does_not_partially_flush() {
         for durability in [
diff --git a/core/partitions/src/install_backup.rs 
b/core/partitions/src/install_backup.rs
index 37a9e2bb0..c07c3c771 100644
--- a/core/partitions/src/install_backup.rs
+++ b/core/partitions/src/install_backup.rs
@@ -16,6 +16,7 @@
 // under the License.
 
 use journal::durable_storage::{DiskStorage, DurableFile, DurableStorage, 
OpenMode};
+use journal::partition_journal::FRONTIER_FILE_NAME;
 use std::io;
 use std::path::Path;
 
@@ -124,6 +125,12 @@ async fn link_tree<S: DurableStorage>(
             if entry.directory {
                 storage.create_directories(&destination).await?;
                 pending.push((source.join(&name), destination));
+            } else if name == FRONTIER_FILE_NAME {
+                // The partition WAL publishes its frontier by overwriting one 
of
+                // two slots in place, so a hard link would not freeze it: the
+                // WAL reset this snapshot exists to roll back would rewrite 
the
+                // snapshot's own bytes. Two blocks, copied once per install.
+                copy_file(&source.join(&name), &destination, storage).await?;
             } else {
                 // Transfer unlinks or atomically replaces these frozen files.
                 // Hard links retain the old bytes without copying segment 
data.
@@ -142,6 +149,20 @@ async fn link_tree<S: DurableStorage>(
     Ok(())
 }
 
+async fn copy_file<S: DurableStorage>(
+    source: &Path,
+    destination: &Path,
+    storage: &S,
+) -> io::Result<()> {
+    let original = storage.open(source, OpenMode::Read).await?;
+    let length = usize::try_from(original.length().await?)
+        .map_err(|_| io::Error::other("partition WAL frontier is too large to 
copy"))?;
+    let bytes = original.read(0, length).await?;
+    let mut copy = storage.open(destination, OpenMode::Create).await?;
+    copy.write(0, bytes).await?;
+    copy.sync().await
+}
+
 fn is_scratch(name: &str) -> bool {
     matches!(name, BACKUP | BUILDING | RETIRED)
         || Path::new(name)
diff --git a/core/partitions/src/persistence.rs 
b/core/partitions/src/persistence.rs
index 6a1e4b3dc..31ccf9ae6 100644
--- a/core/partitions/src/persistence.rs
+++ b/core/partitions/src/persistence.rs
@@ -35,8 +35,15 @@ use std::time::Duration;
 #[cfg(unix)]
 use nix::sys::resource::{Resource, getrlimit};
 
-const APPEND_BATCH_BYTES_MAX: u64 = 1024 * 1024;
-const APPEND_BATCH_OPS_MAX: usize = 64;
+// Group commit bounds, not throughput bounds. Every prepare in a group is
+// already queued and waiting, so widening the group moves work off the barrier
+// and onto a buffered memcpy: one body write and one durability barrier serve
+// the whole group instead of each prepare paying its own. The byte budget is
+// charged against the padded BODY size even when the WAL stores a segment
+// reference and writes 4096 bytes per record, so a tight budget caps grouping
+// far below what the write itself costs.
+const APPEND_BATCH_BYTES_MAX: u64 = 8 * 1024 * 1024;
+const APPEND_BATCH_OPS_MAX: usize = 256;
 const CHECKPOINT_DIRTY_FILES_MAX: usize = 1024;
 #[cfg(unix)]
 const OFFSET_FILES_TOTAL_MAX: usize = 1024;
diff --git a/core/server/src/partition_helpers.rs 
b/core/server/src/partition_helpers.rs
index cebe04091..136ff1c92 100644
--- a/core/server/src/partition_helpers.rs
+++ b/core/server/src/partition_helpers.rs
@@ -1749,7 +1749,12 @@ mod tests {
         drop(store);
         let frontier = Path::new(&directory).join("prepares-0/frontier");
         let mut corrupt = std::fs::read(&frontier).unwrap();
-        corrupt[0] ^= u8::MAX;
+        // Every slot: the frontier alternates between two of them, and one
+        // damaged copy is recoverable by design, so damaging a single slot
+        // would open the partition instead of fencing it.
+        for slot in 
corrupt.chunks_mut(journal::partition_journal::PARTITION_WAL_BLOCK_SIZE) {
+            slot[0] ^= u8::MAX;
+        }
         std::fs::write(&frontier, &corrupt).unwrap();
         let partitions = solo_partitions();
         let metadata = Partition::new(0, namespace.inner(), 
IggyTimestamp::now(), 0, 0);
diff --git a/core/server/src/segment_recovery.rs 
b/core/server/src/segment_recovery.rs
index ca82b42bf..d637dfa8c 100644
--- a/core/server/src/segment_recovery.rs
+++ b/core/server/src/segment_recovery.rs
@@ -78,13 +78,12 @@ const INDEX_SCAN_YIELD_STRIDE: u64 = 1024;
 
 /// Index entries the log may legitimately fail to back under 
`durable_segments`.
 /// Persistence writes exactly one entry per flush chunk and chunks never
-/// overlap. The two halves fdatasync concurrently WITHIN one flush, but
-/// flushes are serialized, and the log's fdatasync covers the whole file: an
-/// entry existing above entry N therefore proves the log was synced through
-/// chunk N. Only the chunk in flight when the process died can leave an entry
-/// the log never backed. See [`PartitionRecoveryRefusal::FsyncedLogLoss`] for
-/// why a deeper step-back is evidence about the log rather than about the
-/// index.
+/// overlap. The WAL makes a body durable before it acknowledges it, and the
+/// flush that indexes that body runs later still, so an entry existing on disk
+/// proves the log bytes it names were already fdatasynced. Only the chunk in
+/// flight when the process died can leave an entry the log never backed. See
+/// [`PartitionRecoveryRefusal::FsyncedLogLoss`] for why a deeper step-back is
+/// evidence about the log rather than about the index.
 const MAX_FSYNCED_INDEX_STEP_BACK_ENTRIES: u64 = 1;
 
 /// Index entries the backward anchor search probes before giving up, at one
diff --git a/core/server/src/server_error.rs b/core/server/src/server_error.rs
index b32a8ab7f..5e9bc6188 100644
--- a/core/server/src/server_error.rs
+++ b/core/server/src/server_error.rs
@@ -436,9 +436,10 @@ pub enum PartitionRecoveryRefusal {
     },
     /// The sparse index of a topic running under `persisted durability` 
outruns its
     /// log by more than the one entry a crash can legitimately strand there.
-    /// Persistence writes exactly one entry per flush chunk, chunks never
-    /// overlap, and flushes are serialized, so every entry below the last one
-    /// names a chunk whose log bytes completed their fdatasync. A completed
+    /// Persistence writes exactly one entry per flush chunk and chunks never
+    /// overlap. The WAL makes a body durable before acknowledging it and the
+    /// flush indexes it later, so every entry on disk names a chunk whose log
+    /// bytes completed their fdatasync. A completed
     /// chunk can contain batches acknowledged before the flush threshold was
     /// reached, while the in-flight chunk can do so too. Reply timing is not
     /// the proof. Only the chunk in flight when the process died can have an
diff --git a/core/simulator/src/storage/tests.rs 
b/core/simulator/src/storage/tests.rs
index c64ecb5dd..454598486 100644
--- a/core/simulator/src/storage/tests.rs
+++ b/core/simulator/src/storage/tests.rs
@@ -974,27 +974,22 @@ fn 
queued_prepares_share_a_barrier_and_survive_power_loss_together() {
         Rc::clone(&persistence).run().await;
         assert!(persistence.failure().is_none());
         let trace = storage.trace();
-        assert_eq!(
-            trace
-                .iter()
-                .filter(|operation| **operation == StorageOperation::FileSync)
-                .count(),
-            4
-        );
-        assert_eq!(
-            trace
-                .iter()
-                .filter(|operation| **operation == 
StorageOperation::DirectorySync)
-                .count(),
-            2
-        );
-        assert_eq!(
+        let count = |wanted: StorageOperation| {
             trace
                 .iter()
-                .filter(|operation| **operation == StorageOperation::Write)
-                .count(),
-            4
-        );
+                .filter(|operation| **operation == wanted)
+                .count()
+        };
+        // One group: the WAL extent and the frontier slot, one barrier each.
+        assert_eq!(count(StorageOperation::Write), 2);
+        assert_eq!(count(StorageOperation::FileSync), 2);
+        // Publication overwrites a pre-existing slot in place, so an
+        // acknowledgment creates no file, renames nothing and leaves no
+        // directory to make durable. Those are the filesystem metadata
+        // transactions this path must never pay per batch.
+        assert_eq!(count(StorageOperation::Create), 0);
+        assert_eq!(count(StorageOperation::Rename), 0);
+        assert_eq!(count(StorageOperation::DirectorySync), 0);
         assert!(persistence.is_durable_through(65));
         storage.crash(Crash::PowerLoss);
         let recovered = 
PartitionPrepareJournal::open_with_storage(Path::new(WAL), 42, 7, storage)
@@ -1004,6 +999,66 @@ fn 
queued_prepares_share_a_barrier_and_survive_power_loss_together() {
     });
 }
 
+#[test]
+fn 
queued_owned_prepares_share_three_file_barriers_without_directory_mutations() {
+    block_on(async {
+        for count in [65, 256, 257] {
+            let storage = storage_for_partition().await;
+            let (persistence, _) =
+                PartitionPersistence::open_with_storage(Path::new(WAL), 42, 7, 
storage.clone())
+                    .await
+                    .unwrap();
+            persistence.enable_segment_storage(SegmentPosition::default(), 64 
* 1024 * 1024);
+            let first = owned_prepare(1, 0, 0);
+            let mut parent = first.header().checksum;
+            persistence.append(first.into_frozen(), true).unwrap();
+            assert!(persistence.start());
+            Rc::clone(&persistence).run().await;
+            assert!(persistence.failure().is_none());
+            persistence.take_metrics();
+            for index in 1..=count {
+                let prepare = owned_prepare(1, parent, index).transmute_header(
+                    |original, header: &mut PrepareHeader| {
+                        *header = original;
+                        header.op = index + 1;
+                        header.checksum = header.identity_checksum();
+                    },
+                );
+                parent = prepare.header().checksum;
+                persistence.append(prepare.into_frozen(), true).unwrap();
+            }
+            storage.clear_trace();
+            assert!(persistence.start());
+            Rc::clone(&persistence).run().await;
+            assert!(persistence.failure().is_none());
+            let trace = storage.trace();
+            let operations = |wanted: StorageOperation| {
+                trace
+                    .iter()
+                    .filter(|operation| **operation == wanted)
+                    .count() as u64
+            };
+            let groups = count.div_ceil(256);
+            assert_eq!(operations(StorageOperation::Write), 3 * groups);
+            assert_eq!(operations(StorageOperation::FileSync), 3 * groups);
+            assert_eq!(operations(StorageOperation::Create), 0);
+            assert_eq!(operations(StorageOperation::Rename), 0);
+            assert_eq!(operations(StorageOperation::DirectorySync), 0);
+            let metrics = persistence.take_metrics();
+            assert_eq!(metrics.completed_batches, groups);
+            assert_eq!(metrics.batched_prepares, count);
+            assert!(persistence.is_durable_through(count + 1));
+            storage.crash(Crash::PowerLoss);
+            let recovered =
+                PartitionPrepareJournal::open_with_storage(Path::new(WAL), 42, 
7, storage)
+                    .await
+                    .unwrap();
+            assert_eq!(recovered.head(), count + 1);
+            assert_eq!(recovered.prepares().await.unwrap().len() as u64, count 
+ 1);
+        }
+    });
+}
+
 #[test]
 fn failed_group_barrier_never_acknowledges_a_partial_batch() {
     block_on(async {

Reply via email to