hubcio commented on code in PR #4092:
URL: https://github.com/apache/iggy/pull/4092#discussion_r3959455117


##########
core/partitions/src/iggy_partition.rs:
##########
@@ -711,6 +723,368 @@ where
         partition
     }
 
+    pub fn set_persistence_notifier(&self, notifier: PersistenceNotifier) {
+        if let Some(persistence) = &self.persistence {
+            persistence.set_notifier(notifier);
+        }
+    }
+
+    /// # Errors
+    /// Returns an error if durable prepare history cannot be opened or 
replayed.
+    pub async fn open_persistence(&mut self) -> Result<(), IggyError> {
+        
self.open_persistence_with_capacity(journal::partition_journal::PARTITION_WAL_BYTES_MAX)
+            .await
+    }
+
+    /// # Errors
+    /// Returns an error if durable prepare history cannot be opened or 
replayed.
+    pub async fn open_persistence_with_capacity(&mut self, capacity: u64) -> 
Result<(), IggyError> {
+        if self.consensus.replica_count() > 1
+            && let Some(directory) = &self.partition_dir
+        {
+            self.materialization_missing =
+                crate::state_transfer::materialization_is_missing(directory, 
self.created_revision)
+                    .await
+                    .map_err(|_| IggyError::CannotReadFile)?;
+            self.ensure_materialization_recovery();
+        }
+        if self.consensus.replica_count() == 1
+            || !(self.durability().is_persisted()
+                || self.consumer_offset_durability().is_persisted())
+        {
+            return Ok(());
+        }
+        let directory = self
+            .partition_dir
+            .as_ref()
+            .ok_or(IggyError::CannotReadFile)?;
+        let directory =
+            std::path::Path::new(directory).join(format!("prepares-{}", 
self.created_revision));
+        let (persistence, prepares) = PartitionPersistence::open_with_capacity(
+            &directory,
+            self.namespace().inner(),
+            self.created_revision,
+            journal::durable_storage::DiskStorage,
+            capacity,
+        )
+        .await
+        .map_err(|error| {
+            warn!(%error, "cannot open partition prepare WAL");
+            IggyError::CannotReadFile
+        })?;
+        if self.materialization_missing {
+            self.persistence = Some(persistence);
+            return Ok(());
+        }
+        let (purge_generation, purge_floor) = persistence.purge_marker();
+        if purge_generation <= self.applied_purge_generation {
+            self.purge_floor_op = self.purge_floor_op.max(purge_floor);
+        }
+        let checkpoint = persistence.checkpoint_op();
+        let head = persistence.head();
+        let mut commit = checkpoint;
+        for message in prepares {
+            let header = *message.header();
+            if header.op == checkpoint {
+                self.log
+                    .journal()
+                    .inner
+                    .restore_checkpoint_prepare(checkpoint, 
message.into_frozen());
+                continue;
+            }
+            commit = commit.max(header.commit.min(head));
+            if header.operation == Operation::SendMessages {
+                self.append_repaired_send_messages(message).await?;
+            } else {
+                self.apply_replicated_operation(message).await?;
+            }
+        }
+        if head > 0 {

Review Comment:
   critical: a crash after saving a newer `log_view` but before its WAL history 
arrives can discard acknowledged operations during election. persist the log 
and view together, or prevent incomplete replicas from voting.



##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4301,6 +4768,27 @@ where
         }
     }
 
+    fn drain_persistable_commits(&self) -> Vec<PipelineEntry> {
+        let Some(persistence) = &self.persistence else {
+            return drain_committable_prefix(self.consensus());
+        };
+        self.persist_repaired_prefix();
+        let through = self.consensus.commit_max().min(persistence.head());

Review Comment:
   critical: promotion counts the local ACK before its prepare is durable, so a 
persisted quorum can contain only one stable copy. count that ACK only after 
the matching prepare is durable in the WAL.



##########
helm/charts/iggy/templates/_helpers.tpl:
##########
@@ -290,9 +290,9 @@ Secret owns it.
   {{- $server := .Values.server }}
   {{- $generated := include "iggy.secretName" . }}
   {{- if $server.encryption.enabled }}
-- name: IGGY_SYSTEM_ENCRYPTION_ENABLED
+- name: IGGY_ENCRYPTION_ENABLED

Review Comment:
   critical: enabling encryption emits variables the default `0.9.0-edge.6` 
image ignores, leaving payload encryption disabled. align the image and 
variable names, then verify encryption with that image.



##########
core/journal/src/partition_journal.rs:
##########
@@ -0,0 +1,1133 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#![allow(clippy::future_not_send)]
+
+use crate::durable_storage::{DiskStorage, DurableFile, DurableStorage, 
OpenMode};
+use futures::TryStreamExt;
+use iggy_binary_protocol::{Command, PrepareHeader};
+use server_common::{
+    Message,
+    iobuf::{Frozen, Owned},
+};
+use std::collections::{BTreeMap, VecDeque};
+use std::io;
+use std::path::{Path, PathBuf};
+use twox_hash::XxHash3_64;
+
+pub const PARTITION_WAL_BLOCK_SIZE: usize = 4096;
+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;
+pub const PREPARE_BYTES_MAX: usize = 64 * 1024 * 1024;
+const STATE_MAGIC: &[u8; 8] = b"IGGYWAL2";
+
+pub trait DurableAppend {
+    /// # Errors
+    /// Returns an error if persistence fails or the prepare does not extend 
the journal.
+    fn append(&mut self, prepare: Frozen<4096>) -> impl Future<Output = 
io::Result<()>>;
+}
+
+pub struct PartitionPrepareJournal<S: DurableStorage = DiskStorage> {
+    directory: PathBuf,
+    file: S::File,
+    storage: S,
+    capacity: u64,
+    state: JournalState,
+    entries: BTreeMap<u64, StoredPrepare>,
+    poisoned: bool,
+    durable_head: u64,
+    obsolete: VecDeque<PathBuf>,
+    cleanup_directory_dirty: bool,
+    recovered_prepares: Vec<Message<PrepareHeader>>,
+}
+
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+struct JournalState {
+    group: u64,
+    incarnation: u64,
+    generation: u64,
+    length: u64,
+    checkpoint: u64,
+    checkpoint_checksum: u128,
+    head: u64,
+    head_checksum: u128,
+    anchor_known: bool,
+    checkpoint_prepare: bool,
+    purge_generation: u64,
+    purge_floor: u64,
+}
+
+#[derive(Clone, Copy)]
+struct StoredPrepare {
+    position: u64,
+    length: usize,
+    checksum: u128,
+}
+
+impl PartitionPrepareJournal {
+    /// # Errors
+    /// Returns an error on I/O failure or invalid durable history.
+    pub async fn open(directory: &Path, group: u64, incarnation: u64) -> 
io::Result<Self> {
+        Self::open_with_storage(directory, group, incarnation, 
DiskStorage).await
+    }
+}
+
+impl<S: DurableStorage> PartitionPrepareJournal<S> {
+    /// Open and verify the durably published partition history.
+    /// The caller must first durably materialize the parent directory.
+    ///
+    /// # Errors
+    /// Returns an error on I/O failure, invalid history, or a poisoned 
journal.
+    pub async fn open_with_storage(
+        directory: &Path,
+        group: u64,
+        incarnation: u64,
+        storage: S,
+    ) -> io::Result<Self> {
+        Self::open_with_storage_and_capacity(
+            directory,
+            group,
+            incarnation,
+            storage,
+            PARTITION_WAL_BYTES_MAX,
+        )
+        .await
+    }
+
+    /// Open history independently of the current admission capacity.
+    ///
+    /// # Errors
+    /// Returns an error for invalid capacity or unverifiable durable history.
+    pub async fn open_with_storage_and_capacity(
+        directory: &Path,
+        group: u64,
+        incarnation: u64,
+        storage: S,
+        capacity: u64,
+    ) -> io::Result<Self> {
+        if 
!(PARTITION_WAL_CAPACITY_MIN..=PARTITION_WAL_CAPACITY_MAX).contains(&capacity)
+            || !capacity.is_multiple_of(PARTITION_WAL_BLOCK_SIZE as u64)
+        {
+            return Err(invalid(
+                "partition WAL capacity is out of bounds or unaligned",
+            ));
+        }
+        let parent = directory
+            .parent()
+            .filter(|parent| !parent.as_os_str().is_empty())
+            .unwrap_or_else(|| Path::new("."));
+        if !storage.exists(parent).await? {
+            return Err(invalid(
+                "partition WAL parent must already be durably materialized",
+            ));
+        }
+        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),
+        };
+        if existing.is_none() {
+            Self::validate_unpublished_history(&storage, directory).await?;
+        }
+        let state = existing.unwrap_or_else(|| JournalState {
+            group,
+            incarnation,
+            ..JournalState::default()
+        });
+        let mode = if existing.is_none() {
+            OpenMode::Create
+        } else {
+            OpenMode::ReadWrite
+        };
+        let file = storage
+            .open(&data_path(directory, state.generation), mode)
+            .await?;
+        if file.length().await? < state.length {
+            return Err(invalid("partition WAL lost acknowledged bytes"));
+        }
+        let mut journal = Self {
+            directory: directory.to_path_buf(),
+            file,
+            storage,
+            capacity,
+            state,
+            entries: BTreeMap::new(),
+            poisoned: false,
+            durable_head: state.head,
+            obsolete: VecDeque::new(),
+            cleanup_directory_dirty: false,
+            recovered_prepares: Vec::new(),
+        };
+        journal.recover_entries().await?;
+        // Only bytes covered by the durable frontier could have released an 
ack.
+        journal.file.truncate(state.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;
+            }
+        }
+        Ok(journal)
+    }
+
+    #[must_use]
+    pub const fn durable_op(&self) -> u64 {
+        self.durable_head
+    }
+
+    #[must_use]
+    pub const fn head(&self) -> u64 {
+        self.state.head
+    }
+
+    #[must_use]
+    pub const fn checkpoint_op(&self) -> u64 {
+        self.state.checkpoint
+    }
+
+    #[must_use]
+    pub const fn checkpoint_checksum(&self) -> Option<u128> {
+        if self.state.anchor_known {
+            Some(self.state.checkpoint_checksum)
+        } else {
+            None
+        }
+    }
+
+    #[must_use]
+    pub const fn generation(&self) -> u64 {
+        self.state.generation
+    }
+
+    #[must_use]
+    pub const fn size_bytes(&self) -> u64 {
+        self.state.length
+    }
+
+    #[must_use]
+    pub fn contains(&self, header: &PrepareHeader) -> bool {
+        if header.op > self.durable_head {
+            return false;
+        }
+        self.entries
+            .get(&header.op)
+            .is_some_and(|entry| entry.checksum == header.checksum)
+    }
+
+    pub fn take_recovered_prepares(&mut self) -> Vec<Message<PrepareHeader>> {
+        std::mem::take(&mut self.recovered_prepares)
+    }
+
+    /// Read the retained prepares in operation order.
+    ///
+    /// # Errors
+    /// Returns an error on I/O failure, invalid history, or a poisoned 
journal.
+    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?;
+            if length != entry.length {
+                return Err(invalid("partition WAL index length mismatch"));
+            }
+            prepares.push(prepare);
+        }
+        Ok(prepares)
+    }
+
+    /// Durably replace the uncommitted suffix.
+    ///
+    /// # Errors
+    /// Returns an error on I/O failure, invalid history, or a poisoned 
journal.
+    pub async fn truncate_from(&mut self, from_op: u64) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if from_op <= self.state.checkpoint {
+            return Err(invalid("cannot truncate checkpointed partition 
operations"));
+        }
+        if from_op > self.state.head {
+            return Ok(());
+        }
+        self.rewrite(
+            self.state.checkpoint,
+            self.state.checkpoint_checksum,
+            Some(from_op),
+        )
+        .await
+    }
+
+    /// Synchronize required materialized files before removing their WAL 
coverage.
+    /// Authorized deletions are excluded by the caller. A missing listed path
+    /// does not prove deletion was authorized and cannot permit WAL 
reclamation.
+    ///
+    /// # Errors
+    /// Returns an error if any file or directory barrier fails.
+    pub async fn checkpoint_files(
+        &mut self,
+        through_op: u64,
+        files: &[std::path::PathBuf],
+        directories: &[std::path::PathBuf],
+    ) -> io::Result<()> {
+        futures::stream::iter(files.iter().map(Ok::<_, io::Error>))
+            .try_for_each_concurrent(16, |path| async {
+                self.storage.open(path, OpenMode::Read).await?.sync().await

Review Comment:
   critical: closing offset writers lets writeback errors disappear after inode 
eviction, so checkpointing can discard WAL before offsets are durable. sync the 
original descriptors before closing them, or retain them through checkpoint.



##########
core/partitions/src/iggy_partition.rs:
##########
@@ -5059,6 +5565,23 @@ where
             }
         }
 
+        if messages_committed
+            && self.consensus.replica_count() == 1
+            && self.durability().is_persisted()
+            && let Some(directory) = &self.partition_dir
+            && let Err(error) = 
crate::state_transfer::fsync_dir(directory).await

Review Comment:
   warning: each persisted singleton message batch reopens and syncs the 
directory even when its file names are already durable. track pending name 
publication and sync only while needed, including creation, rotation, and 
recovery.



##########
core/server/src/dispatch/partition.rs:
##########
@@ -104,6 +104,16 @@ where
             return;
         };
         let partitions = shard.plane.partitions();
+        if partitions.with_partition(
+            &namespace,
+            partitions::IggyPartition::requires_state_transfer,
+        ) == Some(true)
+        {
+            let _ = reply.try_send(PartitionReadReply::Rejected(

Review Comment:
   warning: existing fallbacks turn recovery rejection into empty success, a 
no-op deletion, or the wrong HTTP error. propagate `Rejected(error)` through 
the offset and deletion callers.



##########
core/partitions/src/state_transfer.rs:
##########
@@ -3143,6 +3266,16 @@ where
         // before the swap.
         self.purge_deferred = false;
 
+        if let Some(persistence) = &self.persistence {
+            persistence.reset(commit_op, offsets_wire.prepare_checksum);

Review Comment:
   warning: this existing checkpoint gap leaves transferred replicas unable to 
elect when the source disappears before another commit. retain the checkpoint 
prepare during transfer, or let elections use verified materialized checkpoints.



##########
scripts/ci/test-helm.sh:
##########
@@ -455,7 +455,7 @@ server:
       value: "0.0.0.0:8080"
     - name: IGGY_WEBSOCKET_ADDRESS
       value: "0.0.0.0:8092"
-    - name: IGGY_SYSTEM_SHARDING_CPU_ALLOCATION
+    - name: IGGY_SHARDING_CPU_ALLOCATION

Review Comment:
   warning: the default smoke image ignores `IGGY_SHARDING_CPU_ALLOCATION`, so 
the requested shard limit is lost. use the variable supported by the selected 
image.



##########
core/partitions/src/iggy_partition.rs:
##########
@@ -711,6 +723,368 @@ where
         partition
     }
 
+    pub fn set_persistence_notifier(&self, notifier: PersistenceNotifier) {
+        if let Some(persistence) = &self.persistence {
+            persistence.set_notifier(notifier);
+        }
+    }
+
+    /// # Errors
+    /// Returns an error if durable prepare history cannot be opened or 
replayed.
+    pub async fn open_persistence(&mut self) -> Result<(), IggyError> {
+        
self.open_persistence_with_capacity(journal::partition_journal::PARTITION_WAL_BYTES_MAX)
+            .await
+    }
+
+    /// # Errors
+    /// Returns an error if durable prepare history cannot be opened or 
replayed.
+    pub async fn open_persistence_with_capacity(&mut self, capacity: u64) -> 
Result<(), IggyError> {
+        if self.consensus.replica_count() > 1
+            && let Some(directory) = &self.partition_dir
+        {
+            self.materialization_missing =
+                crate::state_transfer::materialization_is_missing(directory, 
self.created_revision)
+                    .await
+                    .map_err(|_| IggyError::CannotReadFile)?;
+            self.ensure_materialization_recovery();
+        }
+        if self.consensus.replica_count() == 1
+            || !(self.durability().is_persisted()
+                || self.consumer_offset_durability().is_persisted())
+        {
+            return Ok(());
+        }
+        let directory = self
+            .partition_dir
+            .as_ref()
+            .ok_or(IggyError::CannotReadFile)?;
+        let directory =
+            std::path::Path::new(directory).join(format!("prepares-{}", 
self.created_revision));
+        let (persistence, prepares) = PartitionPersistence::open_with_capacity(
+            &directory,
+            self.namespace().inner(),
+            self.created_revision,
+            journal::durable_storage::DiskStorage,
+            capacity,
+        )
+        .await
+        .map_err(|error| {
+            warn!(%error, "cannot open partition prepare WAL");
+            IggyError::CannotReadFile
+        })?;
+        if self.materialization_missing {
+            self.persistence = Some(persistence);
+            return Ok(());
+        }
+        let (purge_generation, purge_floor) = persistence.purge_marker();
+        if purge_generation <= self.applied_purge_generation {
+            self.purge_floor_op = self.purge_floor_op.max(purge_floor);
+        }
+        let checkpoint = persistence.checkpoint_op();
+        let head = persistence.head();
+        let mut commit = checkpoint;
+        for message in prepares {
+            let header = *message.header();
+            if header.op == checkpoint {
+                self.log
+                    .journal()
+                    .inner
+                    .restore_checkpoint_prepare(checkpoint, 
message.into_frozen());
+                continue;
+            }
+            commit = commit.max(header.commit.min(head));
+            if header.operation == Operation::SendMessages {
+                self.append_repaired_send_messages(message).await?;
+            } else {
+                self.apply_replicated_operation(message).await?;
+            }
+        }
+        if head > 0 {
+            self.consensus.sequencer().set_sequence(head);
+            if let Some(checksum) = persistence.checksum(head) {
+                self.consensus.set_last_prepare_checksum(checksum);
+            }
+            self.consensus.restore_commit_state(checkpoint, commit);
+        }
+        // Recovery can reuse completed writes whose last barrier was 
interrupted.
+        // Include them once before reclaiming any recovered WAL history.
+        for segment in self.log.segments() {
+            persistence.mark_segment_dirty(segment.start_offset);
+        }
+        for kind in [ConsumerKind::Consumer, ConsumerKind::ConsumerGroup] {
+            self.durable_consumer_offsets.with_entries(kind, |entries| {
+                for consumer_id in entries.keys() {
+                    persistence.mark_offset_dirty(
+                        crate::state_transfer::consumer_kind_index(kind),
+                        *consumer_id,
+                        true,
+                    );
+                }
+            });
+        }
+        self.persistence = Some(persistence);
+        Ok(())
+    }
+
+    pub const fn requires_state_transfer(&self) -> bool {
+        self.materialization_missing
+    }
+
+    pub fn ensure_materialization_recovery(&self) {
+        if self.materialization_missing
+            && (self.consensus.state_transfer_stage() == 
consensus::StateTransferStage::Idle
+                || self.consensus.status() == consensus::Status::ViewChange)
+        {
+            self.consensus.begin_view_probe();
+            if self.consensus.state_transfer_stage() == 
consensus::StateTransferStage::Idle {
+                self.consensus.begin_state_transfer_await();
+            }
+        }
+    }
+
+    pub async fn on_persistence_completed(&mut self, completion: 
PersistenceCompletion) {
+        if !self
+            .persistence
+            .as_ref()
+            .is_some_and(|persistence| 
persistence.accepts_completion(completion))
+        {
+            return;
+        }
+        self.drive_persistence().await;
+    }
+
+    pub fn needs_persistence_checkpoint(&self) -> bool {
+        self.persistence.as_ref().is_some_and(|persistence| {
+            persistence.needs_checkpoint()
+                && self.consensus.commit_min().min(persistence.head()) > 
persistence.checkpoint_op()
+        })
+    }
+
+    pub async fn checkpoint_persistence(&mut self, config: &PartitionsConfig) {
+        let Some(persistence) = self
+            .persistence
+            .as_ref()
+            .filter(|persistence| persistence.needs_checkpoint())
+            .cloned()
+        else {
+            return;
+        };
+        let through_op = self.consensus.commit_min().min(persistence.head());
+        if through_op <= persistence.checkpoint_op() {
+            return;
+        }
+        if let Err(error) = self.flush_committed_messages(config).await {
+            error!(%error, namespace_raw = self.namespace().inner(), 
"partition checkpoint 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();
+    }
+
+    fn persistence_checkpoint_files(
+        &self,
+        config: &PartitionsConfig,
+    ) -> (Vec<std::path::PathBuf>, Vec<std::path::PathBuf>) {
+        let namespace = self.namespace();
+        let Some(persistence) = &self.persistence else {
+            return (Vec::new(), Vec::new());
+        };
+        let (segments, offsets) = persistence.take_dirty_files();
+        let mut paths = Vec::with_capacity(
+            segments.len() * 2
+                + offsets
+                    .iter()
+                    .map(std::collections::BTreeSet::len)
+                    .sum::<usize>(),
+        );
+        for start_offset in segments {
+            // Retention may already have removed a dirty sealed segment.
+            if self
+                .log
+                .segments()
+                .binary_search_by_key(&start_offset, |segment| 
segment.start_offset)
+                .is_err()
+            {
+                continue;
+            }
+            paths.push(config.get_messages_path(
+                namespace.stream_id(),
+                namespace.topic_id(),
+                namespace.partition_id(),
+                start_offset,
+            ));
+            paths.push(config.get_index_path(
+                namespace.stream_id(),
+                namespace.topic_id(),
+                namespace.partition_id(),
+                start_offset,
+            ));
+        }
+        for (kind, consumers) in [ConsumerKind::Consumer, 
ConsumerKind::ConsumerGroup]
+            .into_iter()
+            .zip(offsets)
+        {
+            for consumer_id in consumers {
+                if self.durable_consumer_offsets.contains(kind, consumer_id)
+                    && let Some(path) = self.persisted_offset_path(kind, 
consumer_id)
+                {
+                    paths.push(path);
+                }
+            }
+        }
+        let mut directories = Vec::with_capacity(4);
+        for directory in [
+            &self.consumer_offsets_path,
+            &self.consumer_group_offsets_path,
+        ]
+        .into_iter()
+        .flatten()
+        {
+            let path = std::path::PathBuf::from(directory);
+            directories.push(path.clone());
+            if let Some(parent) = path.parent()
+                && !directories.iter().any(|existing| existing == parent)
+            {
+                directories.push(parent.to_path_buf());
+            }
+        }
+        if let Some(directory) = &self.partition_dir {
+            directories.push(std::path::PathBuf::from(directory));
+        }
+        (
+            paths.into_iter().map(std::path::PathBuf::from).collect(),
+            directories,
+        )
+    }
+
+    pub fn take_persistence_metrics(&self) -> 
Option<crate::persistence::PersistenceMetrics> {
+        self.persistence
+            .as_ref()
+            .map(|persistence| persistence.take_metrics())
+    }
+
+    fn persistence_checkpoint_pending(&self) -> bool {
+        self.persistence
+            .as_ref()
+            .is_some_and(|persistence| persistence.checkpoint_pending())
+    }
+
+    pub async fn drive_persistence(&mut self) {
+        let Some(persistence) = self.persistence.as_ref() else {
+            return;
+        };
+        if let Some(error) = persistence.failure() {
+            error!(%error, namespace_raw = self.namespace().inner(), 
"partition prepare persistence failed");
+            if self.fatal.is_none() {
+                self.fatal = Some(FatalCommit {
+                    namespace_raw: self.namespace().inner(),
+                    op: persistence.head(),
+                    operation: Operation::SendMessages,
+                });
+            }
+            return;
+        }
+        let durable_op = persistence.durable_op();
+        loop {
+            let ready = self
+                .pending_persisted_acks
+                .borrow()
+                .first_key_value()
+                .is_some_and(|(op, _)| *op <= durable_op);
+            if !ready {
+                break;
+            }
+            let Some((_, header)) = 
self.pending_persisted_acks.borrow_mut().pop_first() else {
+                break;
+            };
+            if persistence.checksum(header.op) == Some(header.checksum) {
+                self.send_prepare_ok(&header).await;
+            }
+        }
+    }
+
+    pub async fn acknowledge_prepare(&self, op: u64) {
+        let Some(prepare) = self.log.journal().inner.repair_entry(op) else {
+            return;
+        };
+        let Ok(header) = bytemuck::checked::try_from_bytes::<PrepareHeader>(
+            &prepare.as_slice()[..size_of::<PrepareHeader>()],
+        ) else {
+            return;
+        };
+        let header = *header;
+        self.persist_repaired_prefix();
+        self.send_prepare_ok(&header).await;
+    }
+
+    const fn requires_persistence(&self, operation: Operation) -> bool {
+        match operation {
+            Operation::SendMessages => self.durability().is_persisted(),
+            Operation::StoreConsumerOffset | Operation::DeleteConsumerOffset 
=> {
+                self.consumer_offset_durability().is_persisted()
+            }
+            _ => false,
+        }
+    }
+
+    fn submit_prepare_persistence(&self, prepare: Frozen<4096>, operation: 
Operation) -> bool {
+        let Some(persistence) = &self.persistence else {
+            return self.consensus.replica_count() == 1 || 
!self.requires_persistence(operation);
+        };
+        if let Err(error) = persistence.append(prepare, 
self.requires_persistence(operation)) {
+            warn!(%error, namespace_raw = self.namespace().inner(), "partition 
WAL refused prepare");
+            if error.kind() != std::io::ErrorKind::WouldBlock {
+                persistence.fail(error);
+            }
+            return false;
+        }
+        self.start_persistence();
+        true
+    }
+
+    fn persist_repaired_prefix(&self) {
+        let Some(persistence) = &self.persistence else {
+            return;
+        };
+        while let Some(op) = persistence.head().checked_add(1) {
+            let Some(prepare) = self.log.journal().inner.repair_entry(op) else 
{

Review Comment:
   warning: a caught-up WAL commit drain still scans the repair ring for a 
missing next operation. skip repair when the WAL covers the journal, or reject 
the lookup before scanning.



##########
core/journal/src/partition_journal.rs:
##########
@@ -0,0 +1,1133 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#![allow(clippy::future_not_send)]
+
+use crate::durable_storage::{DiskStorage, DurableFile, DurableStorage, 
OpenMode};
+use futures::TryStreamExt;
+use iggy_binary_protocol::{Command, PrepareHeader};
+use server_common::{
+    Message,
+    iobuf::{Frozen, Owned},
+};
+use std::collections::{BTreeMap, VecDeque};
+use std::io;
+use std::path::{Path, PathBuf};
+use twox_hash::XxHash3_64;
+
+pub const PARTITION_WAL_BLOCK_SIZE: usize = 4096;
+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;
+pub const PREPARE_BYTES_MAX: usize = 64 * 1024 * 1024;
+const STATE_MAGIC: &[u8; 8] = b"IGGYWAL2";
+
+pub trait DurableAppend {
+    /// # Errors
+    /// Returns an error if persistence fails or the prepare does not extend 
the journal.
+    fn append(&mut self, prepare: Frozen<4096>) -> impl Future<Output = 
io::Result<()>>;
+}
+
+pub struct PartitionPrepareJournal<S: DurableStorage = DiskStorage> {
+    directory: PathBuf,
+    file: S::File,
+    storage: S,
+    capacity: u64,
+    state: JournalState,
+    entries: BTreeMap<u64, StoredPrepare>,
+    poisoned: bool,
+    durable_head: u64,
+    obsolete: VecDeque<PathBuf>,
+    cleanup_directory_dirty: bool,
+    recovered_prepares: Vec<Message<PrepareHeader>>,
+}
+
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+struct JournalState {
+    group: u64,
+    incarnation: u64,
+    generation: u64,
+    length: u64,
+    checkpoint: u64,
+    checkpoint_checksum: u128,
+    head: u64,
+    head_checksum: u128,
+    anchor_known: bool,
+    checkpoint_prepare: bool,
+    purge_generation: u64,
+    purge_floor: u64,
+}
+
+#[derive(Clone, Copy)]
+struct StoredPrepare {
+    position: u64,
+    length: usize,
+    checksum: u128,
+}
+
+impl PartitionPrepareJournal {
+    /// # Errors
+    /// Returns an error on I/O failure or invalid durable history.
+    pub async fn open(directory: &Path, group: u64, incarnation: u64) -> 
io::Result<Self> {
+        Self::open_with_storage(directory, group, incarnation, 
DiskStorage).await
+    }
+}
+
+impl<S: DurableStorage> PartitionPrepareJournal<S> {
+    /// Open and verify the durably published partition history.
+    /// The caller must first durably materialize the parent directory.
+    ///
+    /// # Errors
+    /// Returns an error on I/O failure, invalid history, or a poisoned 
journal.
+    pub async fn open_with_storage(
+        directory: &Path,
+        group: u64,
+        incarnation: u64,
+        storage: S,
+    ) -> io::Result<Self> {
+        Self::open_with_storage_and_capacity(
+            directory,
+            group,
+            incarnation,
+            storage,
+            PARTITION_WAL_BYTES_MAX,
+        )
+        .await
+    }
+
+    /// Open history independently of the current admission capacity.
+    ///
+    /// # Errors
+    /// Returns an error for invalid capacity or unverifiable durable history.
+    pub async fn open_with_storage_and_capacity(
+        directory: &Path,
+        group: u64,
+        incarnation: u64,
+        storage: S,
+        capacity: u64,
+    ) -> io::Result<Self> {
+        if 
!(PARTITION_WAL_CAPACITY_MIN..=PARTITION_WAL_CAPACITY_MAX).contains(&capacity)
+            || !capacity.is_multiple_of(PARTITION_WAL_BLOCK_SIZE as u64)
+        {
+            return Err(invalid(
+                "partition WAL capacity is out of bounds or unaligned",
+            ));
+        }
+        let parent = directory
+            .parent()
+            .filter(|parent| !parent.as_os_str().is_empty())
+            .unwrap_or_else(|| Path::new("."));
+        if !storage.exists(parent).await? {
+            return Err(invalid(
+                "partition WAL parent must already be durably materialized",
+            ));
+        }
+        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),
+        };
+        if existing.is_none() {
+            Self::validate_unpublished_history(&storage, directory).await?;
+        }
+        let state = existing.unwrap_or_else(|| JournalState {
+            group,
+            incarnation,
+            ..JournalState::default()
+        });
+        let mode = if existing.is_none() {
+            OpenMode::Create
+        } else {
+            OpenMode::ReadWrite
+        };
+        let file = storage
+            .open(&data_path(directory, state.generation), mode)
+            .await?;
+        if file.length().await? < state.length {
+            return Err(invalid("partition WAL lost acknowledged bytes"));
+        }
+        let mut journal = Self {
+            directory: directory.to_path_buf(),
+            file,
+            storage,
+            capacity,
+            state,
+            entries: BTreeMap::new(),
+            poisoned: false,
+            durable_head: state.head,
+            obsolete: VecDeque::new(),
+            cleanup_directory_dirty: false,
+            recovered_prepares: Vec::new(),
+        };
+        journal.recover_entries().await?;
+        // Only bytes covered by the durable frontier could have released an 
ack.
+        journal.file.truncate(state.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;
+            }
+        }
+        Ok(journal)
+    }
+
+    #[must_use]
+    pub const fn durable_op(&self) -> u64 {
+        self.durable_head
+    }
+
+    #[must_use]
+    pub const fn head(&self) -> u64 {
+        self.state.head
+    }
+
+    #[must_use]
+    pub const fn checkpoint_op(&self) -> u64 {
+        self.state.checkpoint
+    }
+
+    #[must_use]
+    pub const fn checkpoint_checksum(&self) -> Option<u128> {
+        if self.state.anchor_known {
+            Some(self.state.checkpoint_checksum)
+        } else {
+            None
+        }
+    }
+
+    #[must_use]
+    pub const fn generation(&self) -> u64 {
+        self.state.generation
+    }
+
+    #[must_use]
+    pub const fn size_bytes(&self) -> u64 {
+        self.state.length
+    }
+
+    #[must_use]
+    pub fn contains(&self, header: &PrepareHeader) -> bool {
+        if header.op > self.durable_head {
+            return false;
+        }
+        self.entries
+            .get(&header.op)
+            .is_some_and(|entry| entry.checksum == header.checksum)
+    }
+
+    pub fn take_recovered_prepares(&mut self) -> Vec<Message<PrepareHeader>> {
+        std::mem::take(&mut self.recovered_prepares)
+    }
+
+    /// Read the retained prepares in operation order.
+    ///
+    /// # Errors
+    /// Returns an error on I/O failure, invalid history, or a poisoned 
journal.
+    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?;
+            if length != entry.length {
+                return Err(invalid("partition WAL index length mismatch"));
+            }
+            prepares.push(prepare);
+        }
+        Ok(prepares)
+    }
+
+    /// Durably replace the uncommitted suffix.
+    ///
+    /// # Errors
+    /// Returns an error on I/O failure, invalid history, or a poisoned 
journal.
+    pub async fn truncate_from(&mut self, from_op: u64) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if from_op <= self.state.checkpoint {
+            return Err(invalid("cannot truncate checkpointed partition 
operations"));
+        }
+        if from_op > self.state.head {
+            return Ok(());
+        }
+        self.rewrite(
+            self.state.checkpoint,
+            self.state.checkpoint_checksum,
+            Some(from_op),
+        )
+        .await
+    }
+
+    /// Synchronize required materialized files before removing their WAL 
coverage.
+    /// Authorized deletions are excluded by the caller. A missing listed path
+    /// does not prove deletion was authorized and cannot permit WAL 
reclamation.
+    ///
+    /// # Errors
+    /// Returns an error if any file or directory barrier fails.
+    pub async fn checkpoint_files(
+        &mut self,
+        through_op: u64,
+        files: &[std::path::PathBuf],
+        directories: &[std::path::PathBuf],
+    ) -> io::Result<()> {
+        futures::stream::iter(files.iter().map(Ok::<_, io::Error>))
+            .try_for_each_concurrent(16, |path| async {
+                self.storage.open(path, OpenMode::Read).await?.sync().await
+            })
+            .await?;
+        for path in directories {
+            self.storage.sync_directory(path).await?;
+        }
+        self.checkpoint(through_op).await
+    }
+
+    /// The caller has durably materialized every operation through this point.
+    /// Reclaim a prefix already materialized durably by the caller.
+    ///
+    /// # Errors
+    /// Returns an error on I/O failure, invalid history, or a poisoned 
journal.
+    pub async fn checkpoint(&mut self, through_op: u64) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if through_op <= self.state.checkpoint {
+            return Ok(());
+        }
+        let checksum = self
+            .entries
+            .get(&through_op)
+            .ok_or_else(|| invalid("unknown WAL checkpoint"))?
+            .checksum;
+        self.state.anchor_known = true;
+        self.rewrite(through_op, checksum, None).await
+    }
+
+    /// Install an already durable replacement state, such as a completed 
transfer.
+    /// Replace the journal with an already durable state-transfer checkpoint.
+    ///
+    /// # Errors
+    /// Returns an error on I/O failure, invalid history, or a poisoned 
journal.
+    pub async fn reset(&mut self, op: u64, checksum: Option<u128>) -> 
io::Result<()> {
+        self.ensure_healthy()?;
+        self.state.anchor_known = checksum.is_some();
+        self.rewrite(op, checksum.unwrap_or(0), Some(0)).await
+    }
+
+    #[must_use]
+    pub const fn purge_marker(&self) -> (u64, u64) {
+        (self.state.purge_generation, self.state.purge_floor)
+    }
+
+    /// # Errors
+    /// Returns an error if the purge marker cannot be published durably.
+    pub async fn mark_purge(&mut self, generation: u64, floor: u64) -> 
io::Result<()> {
+        self.ensure_healthy()?;
+        if generation <= self.state.purge_generation {

Review Comment:
   critical: retrying a purge in the same generation after the sequencer 
advances keeps the old durable floor, letting deleted messages return after 
restart. persist the updated floor before deleting on retry.



##########
core/partitions/src/iggy_partition.rs:
##########
@@ -711,6 +723,368 @@ where
         partition
     }
 
+    pub fn set_persistence_notifier(&self, notifier: PersistenceNotifier) {
+        if let Some(persistence) = &self.persistence {
+            persistence.set_notifier(notifier);
+        }
+    }
+
+    /// # Errors
+    /// Returns an error if durable prepare history cannot be opened or 
replayed.
+    pub async fn open_persistence(&mut self) -> Result<(), IggyError> {
+        
self.open_persistence_with_capacity(journal::partition_journal::PARTITION_WAL_BYTES_MAX)
+            .await
+    }
+
+    /// # Errors
+    /// Returns an error if durable prepare history cannot be opened or 
replayed.
+    pub async fn open_persistence_with_capacity(&mut self, capacity: u64) -> 
Result<(), IggyError> {
+        if self.consensus.replica_count() > 1
+            && let Some(directory) = &self.partition_dir
+        {
+            self.materialization_missing =
+                crate::state_transfer::materialization_is_missing(directory, 
self.created_revision)
+                    .await
+                    .map_err(|_| IggyError::CannotReadFile)?;
+            self.ensure_materialization_recovery();
+        }
+        if self.consensus.replica_count() == 1
+            || !(self.durability().is_persisted()
+                || self.consumer_offset_durability().is_persisted())
+        {
+            return Ok(());
+        }
+        let directory = self
+            .partition_dir
+            .as_ref()
+            .ok_or(IggyError::CannotReadFile)?;
+        let directory =
+            std::path::Path::new(directory).join(format!("prepares-{}", 
self.created_revision));
+        let (persistence, prepares) = PartitionPersistence::open_with_capacity(
+            &directory,
+            self.namespace().inner(),
+            self.created_revision,
+            journal::durable_storage::DiskStorage,
+            capacity,
+        )
+        .await
+        .map_err(|error| {
+            warn!(%error, "cannot open partition prepare WAL");
+            IggyError::CannotReadFile
+        })?;
+        if self.materialization_missing {
+            self.persistence = Some(persistence);
+            return Ok(());
+        }
+        let (purge_generation, purge_floor) = persistence.purge_marker();
+        if purge_generation <= self.applied_purge_generation {
+            self.purge_floor_op = self.purge_floor_op.max(purge_floor);
+        }
+        let checkpoint = persistence.checkpoint_op();
+        let head = persistence.head();
+        let mut commit = checkpoint;
+        for message in prepares {
+            let header = *message.header();
+            if header.op == checkpoint {
+                self.log
+                    .journal()
+                    .inner
+                    .restore_checkpoint_prepare(checkpoint, 
message.into_frozen());
+                continue;
+            }
+            commit = commit.max(header.commit.min(head));
+            if header.operation == Operation::SendMessages {
+                self.append_repaired_send_messages(message).await?;
+            } else {
+                self.apply_replicated_operation(message).await?;
+            }
+        }
+        if head > 0 {
+            self.consensus.sequencer().set_sequence(head);
+            if let Some(checksum) = persistence.checksum(head) {
+                self.consensus.set_last_prepare_checksum(checksum);
+            }
+            self.consensus.restore_commit_state(checkpoint, commit);
+        }
+        // Recovery can reuse completed writes whose last barrier was 
interrupted.
+        // Include them once before reclaiming any recovered WAL history.
+        for segment in self.log.segments() {
+            persistence.mark_segment_dirty(segment.start_offset);
+        }
+        for kind in [ConsumerKind::Consumer, ConsumerKind::ConsumerGroup] {
+            self.durable_consumer_offsets.with_entries(kind, |entries| {
+                for consumer_id in entries.keys() {
+                    persistence.mark_offset_dirty(
+                        crate::state_transfer::consumer_kind_index(kind),
+                        *consumer_id,
+                        true,
+                    );
+                }
+            });
+        }
+        self.persistence = Some(persistence);
+        Ok(())
+    }
+
+    pub const fn requires_state_transfer(&self) -> bool {
+        self.materialization_missing
+    }
+
+    pub fn ensure_materialization_recovery(&self) {
+        if self.materialization_missing
+            && (self.consensus.state_transfer_stage() == 
consensus::StateTransferStage::Idle
+                || self.consensus.status() == consensus::Status::ViewChange)
+        {
+            self.consensus.begin_view_probe();
+            if self.consensus.state_transfer_stage() == 
consensus::StateTransferStage::Idle {
+                self.consensus.begin_state_transfer_await();
+            }
+        }
+    }
+
+    pub async fn on_persistence_completed(&mut self, completion: 
PersistenceCompletion) {
+        if !self
+            .persistence
+            .as_ref()
+            .is_some_and(|persistence| 
persistence.accepts_completion(completion))
+        {
+            return;
+        }
+        self.drive_persistence().await;
+    }
+
+    pub fn needs_persistence_checkpoint(&self) -> bool {
+        self.persistence.as_ref().is_some_and(|persistence| {
+            persistence.needs_checkpoint()
+                && self.consensus.commit_min().min(persistence.head()) > 
persistence.checkpoint_op()
+        })
+    }
+
+    pub async fn checkpoint_persistence(&mut self, config: &PartitionsConfig) {
+        let Some(persistence) = self
+            .persistence
+            .as_ref()
+            .filter(|persistence| persistence.needs_checkpoint())
+            .cloned()
+        else {
+            return;
+        };
+        let through_op = self.consensus.commit_min().min(persistence.head());
+        if through_op <= persistence.checkpoint_op() {
+            return;
+        }
+        if let Err(error) = self.flush_committed_messages(config).await {
+            error!(%error, namespace_raw = self.namespace().inner(), 
"partition checkpoint 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();
+    }
+
+    fn persistence_checkpoint_files(
+        &self,
+        config: &PartitionsConfig,
+    ) -> (Vec<std::path::PathBuf>, Vec<std::path::PathBuf>) {
+        let namespace = self.namespace();
+        let Some(persistence) = &self.persistence else {
+            return (Vec::new(), Vec::new());
+        };
+        let (segments, offsets) = persistence.take_dirty_files();
+        let mut paths = Vec::with_capacity(
+            segments.len() * 2
+                + offsets
+                    .iter()
+                    .map(std::collections::BTreeSet::len)
+                    .sum::<usize>(),
+        );
+        for start_offset in segments {
+            // Retention may already have removed a dirty sealed segment.
+            if self
+                .log
+                .segments()
+                .binary_search_by_key(&start_offset, |segment| 
segment.start_offset)
+                .is_err()
+            {
+                continue;
+            }
+            paths.push(config.get_messages_path(
+                namespace.stream_id(),
+                namespace.topic_id(),
+                namespace.partition_id(),
+                start_offset,
+            ));
+            paths.push(config.get_index_path(
+                namespace.stream_id(),
+                namespace.topic_id(),
+                namespace.partition_id(),
+                start_offset,
+            ));
+        }
+        for (kind, consumers) in [ConsumerKind::Consumer, 
ConsumerKind::ConsumerGroup]
+            .into_iter()
+            .zip(offsets)
+        {
+            for consumer_id in consumers {
+                if self.durable_consumer_offsets.contains(kind, consumer_id)
+                    && let Some(path) = self.persisted_offset_path(kind, 
consumer_id)
+                {
+                    paths.push(path);
+                }
+            }
+        }
+        let mut directories = Vec::with_capacity(4);
+        for directory in [
+            &self.consumer_offsets_path,
+            &self.consumer_group_offsets_path,
+        ]
+        .into_iter()
+        .flatten()
+        {
+            let path = std::path::PathBuf::from(directory);
+            directories.push(path.clone());
+            if let Some(parent) = path.parent()
+                && !directories.iter().any(|existing| existing == parent)
+            {
+                directories.push(parent.to_path_buf());
+            }
+        }
+        if let Some(directory) = &self.partition_dir {
+            directories.push(std::path::PathBuf::from(directory));
+        }
+        (
+            paths.into_iter().map(std::path::PathBuf::from).collect(),
+            directories,
+        )
+    }
+
+    pub fn take_persistence_metrics(&self) -> 
Option<crate::persistence::PersistenceMetrics> {
+        self.persistence
+            .as_ref()
+            .map(|persistence| persistence.take_metrics())
+    }
+
+    fn persistence_checkpoint_pending(&self) -> bool {
+        self.persistence
+            .as_ref()
+            .is_some_and(|persistence| persistence.checkpoint_pending())
+    }
+
+    pub async fn drive_persistence(&mut self) {
+        let Some(persistence) = self.persistence.as_ref() else {
+            return;
+        };
+        if let Some(error) = persistence.failure() {
+            error!(%error, namespace_raw = self.namespace().inner(), 
"partition prepare persistence failed");
+            if self.fatal.is_none() {
+                self.fatal = Some(FatalCommit {
+                    namespace_raw: self.namespace().inner(),
+                    op: persistence.head(),
+                    operation: Operation::SendMessages,
+                });
+            }
+            return;
+        }
+        let durable_op = persistence.durable_op();
+        loop {
+            let ready = self
+                .pending_persisted_acks
+                .borrow()
+                .first_key_value()
+                .is_some_and(|(op, _)| *op <= durable_op);
+            if !ready {
+                break;
+            }
+            let Some((_, header)) = 
self.pending_persisted_acks.borrow_mut().pop_first() else {
+                break;
+            };
+            if persistence.checksum(header.op) == Some(header.checksum) {
+                self.send_prepare_ok(&header).await;
+            }
+        }
+    }
+
+    pub async fn acknowledge_prepare(&self, op: u64) {
+        let Some(prepare) = self.log.journal().inner.repair_entry(op) else {
+            return;
+        };
+        let Ok(header) = bytemuck::checked::try_from_bytes::<PrepareHeader>(
+            &prepare.as_slice()[..size_of::<PrepareHeader>()],
+        ) else {
+            return;
+        };
+        let header = *header;
+        self.persist_repaired_prefix();
+        self.send_prepare_ok(&header).await;
+    }
+
+    const fn requires_persistence(&self, operation: Operation) -> bool {
+        match operation {
+            Operation::SendMessages => self.durability().is_persisted(),
+            Operation::StoreConsumerOffset | Operation::DeleteConsumerOffset 
=> {
+                self.consumer_offset_durability().is_persisted()
+            }
+            _ => false,
+        }
+    }
+
+    fn submit_prepare_persistence(&self, prepare: Frozen<4096>, operation: 
Operation) -> bool {
+        let Some(persistence) = &self.persistence else {
+            return self.consensus.replica_count() == 1 || 
!self.requires_persistence(operation);
+        };
+        if let Err(error) = persistence.append(prepare, 
self.requires_persistence(operation)) {

Review Comment:
   critical: a rejected WAL append leaves a gap that makes a later prepare fail 
out of order and fence the replica. submit the missing prefix before newer 
prepares, keeping backpressure until it fits.



##########
core/shard/src/lib.rs:
##########
@@ -4272,6 +4310,14 @@ where
         // and `CommitJournal` is a no-op in both.
         dispatch_partition_journal_actions(consensus, partition, 
&local_actions).await;
         if partition.persist_superblock_if_needed().await {
+            let wire_actions = if partition.requires_state_transfer() {

Review Comment:
   simplification: the same wire-action filter is repeated across dispatch 
sites; extract one private helper, preserving superblock gates, action order, 
and both dispatchers.
   
   also at lines 4398, 4584, 4784, 5663, 7354.



##########
core/server/src/http/reads.rs:
##########
@@ -398,6 +398,79 @@ pub(in crate::http) fn authorize_data_plane(
         .authorize(|permissioner| rule(permissioner, user_id, stream_id, 
topic_id))
 }
 
+static DURABILITY_KEY: std::sync::LazyLock<iggy_common::HeaderKey> =
+    std::sync::LazyLock::new(|| "durability".parse().expect("catalog key is 
valid"));
+
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub(in crate::http) struct TopicDurability {
+    stream_id: usize,
+    topic_id: usize,
+    created_revision: u64,
+    pub durability: iggy_common::Durability,
+}
+
+impl TopicDurability {
+    pub fn confirmed_policy(self, state: &HttpInner) -> 
iggy_common::Durability {
+        let unchanged = state
+            .shard
+            .plane
+            .metadata()
+            .mux_stm
+            .streams()
+            .read(|inner| {
+                inner
+                    .items
+                    .get(self.stream_id)
+                    .and_then(|stream| stream.topics.get(self.topic_id))
+                    .and_then(|topic| topic.partitions.first())
+                    .is_some_and(|partition| partition.created_revision == 
self.created_revision)
+            });
+        self.confirmed_policy_for(unchanged.then_some(self))
+    }
+
+    fn confirmed_policy_for(self, current: Option<Self>) -> 
iggy_common::Durability {

Review Comment:
   simplification: `confirmed_policy_for` only receives `Some(self)` or `None` 
in production. select durability directly from `unchanged`, preserving the 
incarnation check and replicated fallback, and test that lookup.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to