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


##########
core/partitions/src/state_transfer.rs:
##########
@@ -1532,6 +1553,80 @@ fn segment_dir_entries(partition_dir: &str) -> 
std::io::Result<Vec<PathBuf>> {
         .collect())
 }
 
+/// Remove physical tails outside the logical segment list after draining the 
WAL.
+/// The caller holds the write lock and has published a purge or install 
backup.
+pub(crate) async fn remove_public_segment_files(partition_dir: &str) -> 
std::io::Result<()> {

Review Comment:
   keeping the physical scan; the logical list misses uncommitted tails.



##########
core/partitions/src/persistence.rs:
##########
@@ -0,0 +1,1329 @@
+// 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.
+
+use futures::TryStreamExt;
+use iggy_binary_protocol::PrepareHeader;
+use journal::PartitionPrepareJournal;
+use journal::durable_storage::{DiskStorage, DurableFile, DurableStorage};
+use journal::partition_journal::{PARTITION_WAL_BYTES_MAX, SegmentPosition, 
SegmentReference};
+use server_common::Message;
+use server_common::iobuf::Frozen;
+use smallvec::SmallVec;
+use std::cell::{Cell, RefCell};
+use std::collections::{BTreeSet, HashMap, VecDeque};
+use std::io;
+use std::path::{Path, PathBuf};
+use std::rc::Rc;
+use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
+use std::sync::{Arc, LazyLock, Mutex, Weak};
+use std::time::Duration;
+
+const APPEND_BATCH_BYTES_MAX: u64 = 1024 * 1024;
+const APPEND_BATCH_OPS_MAX: usize = 64;
+const CHECKPOINT_DIRTY_FILES_MAX: usize = 1024;
+const PERSISTENCE_DRAIN_TIMEOUT: Duration = Duration::from_secs(30);
+
+static NEXT_INSTANCE: AtomicU64 = AtomicU64::new(1);
+
+#[derive(Clone, Copy, Debug)]
+pub struct PersistenceCompletion {
+    pub group: u64,
+    pub instance: u64,
+    pub epoch: u64,
+}
+
+#[derive(Default)]
+pub struct PersistenceMetrics {
+    pub disk_bytes: u64,
+    pub queued_bytes: u64,
+    pub in_flight_bytes: u64,
+    pub checkpoints_pending: u64,
+    pub completed_batches: u64,
+    pub batched_prepares: u64,
+    pub completed_checkpoints: u64,
+    pub failed_writes: u64,
+}
+
+pub type PersistenceNotifier = Rc<dyn Fn(PersistenceCompletion)>;
+
+pub struct PartitionPersistence<S: DurableStorage = DiskStorage> {
+    group: u64,
+    instance: u64,
+    lease: Option<Arc<WriterLease>>,
+    epoch: Cell<u64>,
+    journal: RefCell<Option<PartitionPrepareJournal<S>>>,
+    queue: RefCell<VecDeque<Mutation<S>>>,
+    offset_files: RefCell<HashMap<String, S::File>>,
+    retired_offset_files: RefCell<Vec<S::File>>,
+    accepted: RefCell<AcceptedPrepares>,
+    accepted_head: Cell<u64>,
+    durable_head: Cell<u64>,
+    checkpoint: Cell<u64>,
+    checkpoint_checksum: Cell<Option<u128>>,
+    certified_log_view: Cell<Option<u32>>,
+    requested_log_view: Cell<Option<(u32, u64, u128)>>,
+    checkpoint_requested: Cell<u64>,
+    checkpoint_running: Cell<bool>,
+    checkpoint_needed: Cell<bool>,
+    dirty_segments: RefCell<BTreeSet<u64>>,
+    dirty_offsets: [RefCell<BTreeSet<u32>>; 2],
+    purge_generation: Cell<u64>,
+    purge_floor: Cell<u64>,
+    capacity: u64,
+    disk_bytes: Cell<u64>,
+    retained_bytes: Cell<u64>,
+    segment_checkpoint: Cell<Option<SegmentPosition>>,
+    queued_bytes: Cell<u64>,
+    in_flight_bytes: Cell<u64>,
+    waiters: RefCell<Vec<std::task::Waker>>,
+    running: Cell<bool>,
+    writer_active: Cell<bool>,
+    retired: Cell<bool>,
+    failure: RefCell<Option<Arc<io::Error>>>,
+    notifier: RefCell<Option<PersistenceNotifier>>,
+    completed_batches: Cell<u64>,
+    batched_prepares: Cell<u64>,
+    completed_checkpoints: Cell<u64>,
+    failed_writes: Cell<u64>,
+}
+
+struct WriterLease {
+    key: PathBuf,
+    id: u64,
+    retired: AtomicBool,
+    running: AtomicBool,
+    waiters: Mutex<Vec<std::task::Waker>>,
+}
+
+struct WriterRegistration {
+    id: u64,
+    writer: Weak<WriterLease>,
+    interrupted: bool,
+}
+
+static WRITERS: LazyLock<Mutex<HashMap<PathBuf, WriterRegistration>>> =
+    LazyLock::new(|| Mutex::new(HashMap::new()));
+
+impl WriterLease {
+    async fn acquire(key: PathBuf) -> io::Result<Arc<Self>> {
+        loop {
+            let previous = {
+                let mut writers = WRITERS
+                    .lock()
+                    .unwrap_or_else(std::sync::PoisonError::into_inner);
+                if writers
+                    .get(&key)
+                    .is_some_and(|registration| registration.interrupted)
+                {
+                    return Err(io::Error::other(
+                        "interrupted partition WAL writer requires process 
restart",
+                    ));
+                }
+                if let Some(previous) = writers
+                    .get(&key)
+                    .and_then(|registration| registration.writer.upgrade())
+                {
+                    drop(writers);
+                    Some(previous)
+                } else {
+                    let lease = Arc::new(Self {
+                        key: key.clone(),
+                        id: NEXT_INSTANCE.fetch_add(1, Ordering::Relaxed),
+                        retired: AtomicBool::new(false),
+                        running: AtomicBool::new(false),
+                        waiters: Mutex::new(Vec::new()),
+                    });
+                    writers.insert(
+                        key.clone(),
+                        WriterRegistration {
+                            id: lease.id,
+                            writer: Arc::downgrade(&lease),
+                            interrupted: false,
+                        },
+                    );
+                    drop(writers);
+                    return Ok(lease);
+                }
+            };
+            let Some(previous) = previous else { continue };
+            if !previous.retired.load(Ordering::Acquire) {
+                return Err(io::Error::new(
+                    io::ErrorKind::AlreadyExists,
+                    "partition WAL already has an active writer",
+                ));
+            }
+            compio::runtime::time::timeout(
+                PERSISTENCE_DRAIN_TIMEOUT,
+                futures::future::poll_fn(|context| {
+                    let mut waiters = previous
+                        .waiters
+                        .lock()
+                        .unwrap_or_else(std::sync::PoisonError::into_inner);
+                    if !previous.running.load(Ordering::Acquire) {
+                        return std::task::Poll::Ready(());
+                    }
+                    if !waiters
+                        .iter()
+                        .any(|waiter| waiter.will_wake(context.waker()))
+                    {
+                        waiters.push(context.waker().clone());
+                    }
+                    std::task::Poll::Pending
+                }),
+            )
+            .await
+            .map_err(|_| {
+                io::Error::new(
+                    io::ErrorKind::TimedOut,
+                    "retired partition WAL writer did not stop",
+                )
+            })?;
+            let mut writers = WRITERS
+                .lock()
+                .unwrap_or_else(std::sync::PoisonError::into_inner);
+            if writers
+                .get(&key)
+                .is_some_and(|registration| registration.id == previous.id)
+            {
+                writers.remove(&key);
+            }
+        }
+    }
+
+    fn finish(&self) {
+        self.running.store(false, Ordering::Release);
+        for waiter in self
+            .waiters
+            .lock()
+            .unwrap_or_else(std::sync::PoisonError::into_inner)
+            .drain(..)
+        {
+            waiter.wake();
+        }
+    }
+}
+
+impl Drop for WriterLease {
+    fn drop(&mut self) {
+        let mut writers = WRITERS
+            .lock()
+            .unwrap_or_else(std::sync::PoisonError::into_inner);
+        if writers
+            .get(&self.key)
+            .is_some_and(|registration| registration.id == self.id && 
!registration.interrupted)
+        {
+            writers.remove(&self.key);
+        }
+    }
+}
+
+struct WriterTicket<S: DurableStorage> {
+    owner: Rc<PartitionPersistence<S>>,
+    started: bool,
+}
+
+impl<S: DurableStorage> Drop for WriterTicket<S> {
+    fn drop(&mut self) {
+        if self.started {
+            return;
+        }
+        self.owner.fail(io::Error::new(
+            io::ErrorKind::Interrupted,
+            "partition WAL writer was dropped before polling",
+        ));
+        self.owner.running.set(false);
+        if let Some(lease) = &self.owner.lease {
+            lease.finish();
+        }
+        for waiter in self.owner.waiters.borrow_mut().drain(..) {
+            waiter.wake();
+        }
+    }
+}
+
+struct WriterGuard<'a, S: DurableStorage> {
+    owner: &'a PartitionPersistence<S>,
+    journal: Option<PartitionPrepareJournal<S>>,
+    complete: bool,
+}
+
+impl<S: DurableStorage> Drop for WriterGuard<'_, S> {
+    fn drop(&mut self) {
+        *self.owner.journal.borrow_mut() = self.journal.take();
+        self.owner.in_flight_bytes.set(0);
+        self.owner.checkpoint_running.set(false);
+        if !self.complete
+            && let Some(lease) = &self.owner.lease
+        {
+            let mut writers = WRITERS
+                .lock()
+                .unwrap_or_else(std::sync::PoisonError::into_inner);
+            if let Some(registration) = writers.get_mut(&lease.key)
+                && registration.id == lease.id
+            {
+                registration.interrupted = true;
+            }
+        }
+        if !self.complete && self.owner.failure.borrow().is_none() {
+            *self.owner.failure.borrow_mut() = Some(Arc::new(io::Error::new(
+                io::ErrorKind::Interrupted,
+                "partition WAL writer stopped during a mutation",
+            )));
+        }
+        self.owner.writer_active.set(false);
+        self.owner.running.set(false);
+        if let Some(lease) = &self.owner.lease {
+            lease.finish();
+        }
+        for waiter in self.owner.waiters.borrow_mut().drain(..) {
+            waiter.wake();
+        }
+    }
+}
+
+struct AcceptedPrepares {
+    base: u64,
+    checksums: VecDeque<u128>,
+}
+
+impl AcceptedPrepares {
+    fn checksum(&self, op: u64) -> Option<u128> {
+        let index = 
usize::try_from(op.checked_sub(self.base)?.checked_sub(1)?).ok()?;
+        self.checksums.get(index).copied()
+    }
+
+    fn truncate_from(&mut self, op: u64) {
+        let keep =
+            
usize::try_from(op.saturating_sub(self.base).saturating_sub(1)).unwrap_or(usize::MAX);
+        self.checksums.truncate(keep);
+    }
+
+    fn checkpoint(&mut self, op: u64) {
+        let count = usize::try_from(op.saturating_sub(self.base))
+            .unwrap_or(usize::MAX)
+            .min(self.checksums.len());
+        self.checksums.drain(..count);
+        self.base = op;
+    }
+}
+
+enum Mutation<S: DurableStorage> {
+    EnableSegments {
+        epoch: u64,
+        initial: SegmentPosition,
+        max_size: u64,
+    },
+    CertifyView {
+        epoch: u64,
+        view: u32,
+        op: u64,
+        checksum: u128,
+    },
+    Purge {
+        epoch: u64,
+        generation: u64,
+        floor: u64,
+    },
+    Append {
+        epoch: u64,
+        prepare: Frozen<4096>,
+        durable: bool,
+        bytes: u64,
+    },
+    Truncate {
+        epoch: u64,
+        from_op: u64,
+    },
+    Checkpoint {
+        epoch: u64,
+        through_op: u64,
+        files: Vec<PathBuf>,
+        directories: Vec<PathBuf>,
+        offset_files: Vec<S::File>,
+    },
+    Reset {
+        epoch: u64,
+        op: u64,
+        checksum: Option<u128>,
+        prepare: Option<Frozen<4096>>,
+        segments: Option<(SegmentPosition, u64)>,
+    },
+}
+
+impl PartitionPersistence {
+    /// # Errors
+    /// Returns an error if the partition WAL cannot be recovered.
+    pub async fn open(
+        directory: &Path,
+        group: u64,
+        incarnation: u64,
+    ) -> io::Result<(Rc<Self>, Vec<Message<PrepareHeader>>)> {
+        Self::open_with_storage(directory, group, incarnation, 
DiskStorage).await
+    }
+}
+
+impl<S: DurableStorage> PartitionPersistence<S> {
+    /// # Errors
+    /// Returns an error if persistence fails, capacity is exhausted, or 
history is invalid.
+    pub async fn open_with_storage(
+        directory: &Path,
+        group: u64,
+        incarnation: u64,
+        storage: S,
+    ) -> io::Result<(Rc<Self>, Vec<Message<PrepareHeader>>)> {
+        Self::open_with_capacity(
+            directory,
+            group,
+            incarnation,
+            storage,
+            PARTITION_WAL_BYTES_MAX,
+            false,
+        )
+        .await
+    }
+
+    /// # Errors
+    /// Returns an error for invalid capacity or unverifiable durable history.
+    pub async fn open_with_capacity(
+        directory: &Path,
+        group: u64,
+        incarnation: u64,
+        storage: S,
+        capacity: u64,
+        preallocate_segments: bool,
+    ) -> io::Result<(Rc<Self>, Vec<Message<PrepareHeader>>)> {
+        let lease = if let Some(key) = storage.writer_identity(directory)? {
+            Some(WriterLease::acquire(key).await?)
+        } else {
+            None
+        };
+        let mut journal = 
PartitionPrepareJournal::open_with_storage_and_capacity(
+            directory,
+            group,
+            incarnation,
+            storage,
+            capacity,
+            preallocate_segments,
+        )
+        .await?;
+        let prepares = journal.take_recovered_prepares();
+        let mut accepted = AcceptedPrepares {
+            base: journal.checkpoint_op(),
+            checksums: VecDeque::with_capacity(prepares.len()),
+        };
+        for prepare in &prepares {
+            let header = prepare.header();
+            if header.op > journal.checkpoint_op() {
+                accepted.checksums.push_back(header.checksum);
+            }
+        }
+        let persistence = Rc::new(Self {
+            group,
+            instance: NEXT_INSTANCE.fetch_add(1, Ordering::Relaxed),
+            lease,
+            epoch: Cell::new(0),
+            accepted_head: Cell::new(journal.head()),
+            durable_head: Cell::new(journal.head()),
+            checkpoint: Cell::new(journal.checkpoint_op()),
+            checkpoint_checksum: Cell::new(journal.checkpoint_checksum()),
+            certified_log_view: Cell::new(journal.certified_log_view()),
+            requested_log_view: Cell::new(None),
+            checkpoint_requested: Cell::new(journal.checkpoint_op()),
+            checkpoint_running: Cell::new(false),
+            checkpoint_needed: Cell::new(false),
+            dirty_segments: RefCell::new(BTreeSet::new()),
+            dirty_offsets: std::array::from_fn(|_| 
RefCell::new(BTreeSet::new())),
+            purge_generation: Cell::new(journal.purge_marker().0),
+            purge_floor: Cell::new(journal.purge_marker().1),
+            capacity,
+            disk_bytes: Cell::new(journal.size_bytes()),
+            retained_bytes: Cell::new(journal.retained_bytes()),
+            segment_checkpoint: Cell::new(journal.segment_checkpoint()),
+            journal: RefCell::new(Some(journal)),
+            queue: RefCell::new(VecDeque::new()),
+            offset_files: RefCell::new(HashMap::new()),
+            retired_offset_files: RefCell::new(Vec::new()),
+            accepted: RefCell::new(accepted),
+            queued_bytes: Cell::new(0),
+            in_flight_bytes: Cell::new(0),
+            waiters: RefCell::new(Vec::new()),
+            running: Cell::new(false),
+            writer_active: Cell::new(false),
+            retired: Cell::new(false),
+            failure: RefCell::new(None),
+            notifier: RefCell::new(None),
+            completed_batches: Cell::new(0),
+            batched_prepares: Cell::new(0),
+            completed_checkpoints: Cell::new(0),
+            failed_writes: Cell::new(0),
+        });
+        Ok((persistence, prepares))
+    }
+
+    #[cfg(test)]
+    pub(crate) fn exhaust_capacity_for_test(&self) {
+        self.disk_bytes.set(self.capacity);
+        self.retained_bytes.set(self.capacity);
+    }
+
+    #[cfg(test)]
+    pub(crate) fn release_capacity_for_test(&self) {
+        self.disk_bytes.set(0);
+        self.retained_bytes.set(0);
+    }
+
+    pub const fn certified_log_view(&self) -> Option<u32> {
+        self.certified_log_view.get()
+    }
+
+    pub fn certify_log_view(&self, view: u32, op: u64, checksum: u128) -> bool 
{
+        if self.certified_log_view.get() == Some(view) {
+            return true;
+        }
+        if self.retired.get()
+            || self.failure.borrow().is_some()
+            || op > self.head()
+            || !(op == 0 && checksum == 0 || self.checksum(op) == 
Some(checksum))
+        {
+            return false;
+        }
+        let target = (view, op, checksum);
+        if self
+            .requested_log_view
+            .get()
+            .is_none_or(|(pending, _, _)| pending != view)
+        {
+            self.requested_log_view.set(Some(target));
+            self.queue.borrow_mut().push_back(Mutation::CertifyView {
+                epoch: self.epoch.get(),
+                view,
+                op,
+                checksum,
+            });
+        }
+        false
+    }
+
+    pub fn set_notifier(&self, notifier: PersistenceNotifier) {
+        *self.notifier.borrow_mut() = Some(notifier);
+    }
+
+    pub const fn accepts_completion(&self, completion: PersistenceCompletion) 
-> bool {
+        completion.instance == self.instance
+            && completion.epoch == self.epoch.get()
+            && !self.retired.get()
+    }
+
+    pub fn is_durable(&self, header: &PrepareHeader) -> bool {
+        !self.retired.get()
+            && self.failure.borrow().is_none()
+            && header.op <= self.durable_head.get()
+            && self.checksum(header.op) == Some(header.checksum)
+    }
+
+    pub const fn segment_checkpoint(&self) -> Option<SegmentPosition> {
+        self.segment_checkpoint.get()
+    }
+
+    pub fn enable_segment_storage(&self, initial: SegmentPosition, max_size: 
u64) {
+        self.queue.borrow_mut().push_back(Mutation::EnableSegments {
+            epoch: self.epoch.get(),
+            initial,
+            max_size,
+        });
+    }
+
+    /// Verify the physical prefix before making its indexes and logical sizes 
visible.
+    ///
+    /// # Errors
+    /// Returns an error for a failed barrier or a body outside the durable 
segment range.
+    pub async fn validate_segment_prefix(
+        &self,
+        prepares: &[Frozen<4096>],
+        start_offset: u64,
+        mut position: u64,
+    ) -> io::Result<u64> {
+        self.drain_with_timeout().await?;
+        let journal = self.journal.borrow();
+        let journal = journal
+            .as_ref()
+            .ok_or_else(|| io::Error::other("segment writer has no journal"))?;
+        let initial = position;
+        for prepare in prepares {
+            let header = prepare_header(prepare)?;
+            let reference: SegmentReference = journal
+                .segment_reference(header)
+                .ok_or_else(|| io::Error::other("committed prepare has no 
durable segment body"))?;
+            if reference.start_offset != start_offset
+                || reference.position != position
+                || reference.length != (prepare.len() - 
size_of::<PrepareHeader>()) as u64
+            {
+                return Err(io::Error::other(
+                    "committed prepare differs from its segment position",
+                ));
+            }
+            position = position
+                .checked_add(reference.length)
+                .ok_or_else(|| io::Error::other("segment prefix overflow"))?;
+        }
+        Ok(position - initial)
+    }
+
+    pub const fn durable_op(&self) -> u64 {
+        self.durable_head.get()
+    }
+
+    pub fn is_durable_through(&self, op: u64) -> bool {
+        !self.retired.get() && self.failure.borrow().is_none() && op <= 
self.durable_head.get()
+    }
+
+    pub fn has_capacity(&self, frame_bytes: usize) -> bool {
+        let Ok(bytes) = journal::partition_journal::record_length(frame_bytes) 
else {
+            return false;
+        };
+        let bytes = bytes as u64;
+        !self.retired.get()
+            && self.failure.borrow().is_none()
+            && self
+                .retained_bytes
+                .get()
+                .saturating_add(self.queued_bytes.get())
+                .saturating_add(self.in_flight_bytes.get())
+                .saturating_add(bytes)
+                <= self.capacity
+    }
+
+    /// # Errors
+    /// Returns an error if persistence fails, capacity is exhausted, or 
history is invalid.
+    pub fn append(&self, prepare: Frozen<4096>, durable: bool) -> 
io::Result<()> {
+        let header = prepare_header(&prepare)?;
+        let bytes = journal::partition_journal::record_length(prepare.len())? 
as u64;
+        if self.accepted.borrow().checksum(header.op) == Some(header.checksum) 
{
+            return Ok(());
+        }
+        if !self.has_capacity(prepare.len()) {
+            return Err(io::Error::new(
+                io::ErrorKind::WouldBlock,
+                "partition WAL capacity exhausted",
+            ));
+        }
+        if header.op != self.accepted_head.get().saturating_add(1) {
+            return Err(io::Error::new(
+                io::ErrorKind::InvalidData,
+                "partition WAL submission is out of order",
+            ));
+        }
+        self.queued_bytes.set(self.queued_bytes.get() + bytes);
+        self.accepted
+            .borrow_mut()
+            .checksums
+            .push_back(header.checksum);
+        self.accepted_head.set(header.op);
+        self.queue.borrow_mut().push_back(Mutation::Append {
+            epoch: self.epoch.get(),
+            prepare,
+            durable,
+            bytes,
+        });
+        Ok(())
+    }
+
+    pub fn mark_segment_dirty(&self, start_offset: u64) {
+        self.dirty_segments.borrow_mut().insert(start_offset);
+    }
+
+    pub fn take_offset_file(&self, path: &str) -> Option<S::File> {
+        self.offset_files.borrow_mut().remove(path)
+    }
+
+    /// # Panics
+    /// Panics if the previous writer was not taken before replacement.
+    pub fn retain_offset_file(&self, path: String, file: S::File) {
+        assert!(self.offset_files.borrow_mut().insert(path, file).is_none());

Review Comment:
   fixed; displaced writers stay alive through the checkpoint.



##########
core/journal/src/partition_journal/segments.rs:
##########
@@ -0,0 +1,663 @@
+// 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.
+
+use std::collections::{BTreeMap, BTreeSet};
+use std::io;
+use std::path::{Path, PathBuf};
+
+use iggy_binary_protocol::batch::BatchHeader;
+use iggy_binary_protocol::{Operation, PrepareHeader};
+use server_common::iobuf::Frozen;
+
+use super::{JournalState, PartitionPrepareJournal, SegmentReference, 
StoredPrepare, invalid};
+use crate::durable_storage::{DurableFile, DurableStorage, OpenMode};
+
+pub(super) const SEGMENT_STATE_FLAG: usize = 99;
+pub(super) const SEGMENT_STATE_OFFSET: usize = 128;
+pub(super) const SEGMENT_STATE_BYTES: usize = 10 * size_of::<u64>();
+
+/// A whole-batch boundary. `next_offset` is the first offset after this 
prefix.
+/// Physical tails may exceed the checkpoint boundary but cannot be polled yet.
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+pub struct SegmentPosition {
+    pub start_offset: u64,
+    pub length: u64,
+    pub next_offset: u64,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(super) struct SegmentCursor {
+    pub generation: u64,
+    pub position: SegmentPosition,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(super) struct SegmentState {
+    pub max_size: u64,
+    pub next_generation: u64,
+    pub tail: SegmentCursor,
+    pub checkpoint: SegmentCursor,
+}
+
+impl<S: DurableStorage> PartitionPrepareJournal<S> {
+    /// Enable segment body ownership after recovering the legacy committed 
view.
+    /// The initial boundary must describe durable materialized messages only.
+    ///
+    /// # Errors
+    /// Returns an error for inconsistent boundaries or a failed storage 
barrier.
+    pub async fn enable_segment_storage(
+        &mut self,
+        initial: SegmentPosition,
+        max_size: u64,
+    ) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if let Some(segments) = self.state.segment_storage {
+            if segments.max_size != max_size {
+                return Err(invalid("segment size differs from the durable WAL 
layout"));
+            }
+            return Ok(());
+        }
+        if max_size == 0 || !initial.valid() {
+            return Err(invalid("invalid initial segment boundary"));
+        }
+        let generation = self
+            .entries
+            .values()
+            .filter_map(|entry| entry.reference)
+            .map(|reference| reference.generation)
+            .max()
+            .map_or(Some(0), |generation| generation.checked_add(1))
+            .ok_or_else(|| invalid("segment generation exhausted"))?;
+        let cursor = SegmentCursor {
+            generation,
+            position: initial,
+        };
+        let segments = SegmentState {
+            max_size,
+            next_generation: generation
+                .checked_add(1)
+                .ok_or_else(|| invalid("segment generation exhausted"))?,
+            tail: cursor,
+            checkpoint: cursor,
+        };
+        let state = JournalState {
+            segment_references: true,
+            segment_storage: Some(segments),
+            ..self.state
+        };
+        self.poisoned = true;
+        if initial.length > 0 {
+            self.open_segment_file(cursor, 
self.preallocate_segments.then_some(max_size))
+                .await?;
+            self.sync_segment_files().await?;
+        }
+        self.file.sync().await?;
+        // Upgrade before any uncommitted body reaches an offset-named file.
+        self.publish(state).await?;
+        self.state = state;
+        self.durable_head = state.head;
+        self.poisoned = false;
+        self.migrate_segment_prepares().await
+    }
+
+    pub const fn segment_checkpoint(&self) -> Option<SegmentPosition> {
+        match self.state.segment_storage {
+            Some(segments) => Some(segments.checkpoint.position),
+            None => None,
+        }
+    }
+
+    pub fn segment_reference(&self, header: &PrepareHeader) -> 
Option<SegmentReference> {
+        if !self.contains(header) {
+            return None;
+        }
+        self.entries
+            .get(&header.op)
+            .and_then(|entry| entry.reference)
+    }
+
+    /// Install a transferred checkpoint after all replacement segment files 
are durable.
+    ///
+    /// # Errors
+    /// Returns an error if the checkpoint contradicts the installed bytes or 
a barrier fails.
+    #[allow(clippy::too_many_lines)]
+    pub async fn reset_with_segment_checkpoint(
+        &mut self,
+        op: u64,
+        checksum: Option<u128>,
+        prepare: Option<Frozen<4096>>,
+        initial: SegmentPosition,
+        max_size: u64,
+    ) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if !initial.valid() || max_size == 0 {
+            return Err(invalid("invalid installed segment boundary"));
+        }
+        if let Some(prepare) = &prepare {
+            let header = self.validate_checkpoint_prepare(prepare)?;
+            if header.op != op || Some(header.checksum) != checksum {
+                return Err(invalid("installed checkpoint prepare identity 
mismatch"));
+            }
+        }
+        let generation = self
+            .state
+            .segment_storage
+            .map_or(0, |segments| segments.next_generation);
+        let cursor = SegmentCursor {
+            generation,
+            position: initial,
+        };
+        let mut segments = SegmentState {
+            max_size,
+            next_generation: generation
+                .checked_add(1)
+                .ok_or_else(|| invalid("segment generation exhausted"))?,
+            tail: cursor,
+            checkpoint: cursor,
+        };
+        self.poisoned = true;
+        self.segment_files.clear();
+        self.open_segment_file(cursor, 
self.preallocate_segments.then_some(max_size))
+            .await?;
+        if self
+            .segment_files
+            .get(&(generation, initial.start_offset))
+            .ok_or_else(|| invalid("installed segment handle is absent"))?
+            .length()
+            .await?
+            != initial.length
+        {
+            return Err(invalid(
+                "installed segment size differs from its checkpoint",
+            ));
+        }
+        let checkpoint_prepare = if let Some(prepare) = prepare {
+            let header = self.validate_checkpoint_prepare(&prepare)?;
+            let reference = if header.operation == Operation::SendMessages {
+                let batch = decode_batch(prepare.as_slice())?;
+                if initial.length >= batch.batch_length
+                    && batch_next_offset(prepare.as_slice())? == 
initial.next_offset
+                {
+                    let reference = SegmentReference {
+                        generation,
+                        start_offset: initial.start_offset,
+                        position: initial.length - batch.batch_length,
+                        length: batch.batch_length,
+                    };
+                    let file = self
+                        .segment_files
+                        .get(&(generation, initial.start_offset))
+                        .ok_or_else(|| invalid("installed segment handle is 
absent"))?;
+                    if file
+                        .read(
+                            reference.position,
+                            prepare.len() - size_of::<PrepareHeader>(),
+                        )
+                        .await?
+                        != prepare.as_slice()[size_of::<PrepareHeader>()..]
+                    {
+                        return Err(invalid(
+                            "checkpoint prepare differs from installed segment 
bytes",
+                        ));
+                    }
+                    Some(reference)
+                } else if let Some(reference) = self
+                    .entries
+                    .get(&op)
+                    .filter(|entry| entry.checksum == header.checksum)
+                    .and_then(|entry| entry.reference)
+                {
+                    Some(reference)
+                } else {
+                    let retained = segments.allocate(SegmentPosition {
+                        start_offset: batch.base_offset,
+                        length: batch.batch_length,
+                        next_offset: batch_next_offset(prepare.as_slice())?,
+                    })?;
+                    let reference = SegmentReference {
+                        generation: retained.generation,
+                        start_offset: batch.base_offset,
+                        position: 0,
+                        length: batch.batch_length,
+                    };
+                    let mut file = self
+                        .storage
+                        .open(&reference.path(&self.directory), 
OpenMode::Create)
+                        .await?;
+                    file.write_frozen(0, 
prepare.slice(size_of::<PrepareHeader>()..))
+                        .await?;
+                    file.sync().await?;
+                    Some(reference)
+                }
+            } else {
+                None
+            };
+            Some((prepare, reference))
+        } else {
+            None
+        };
+        self.sync_segment_files().await?;
+        self.storage.sync_directory(&self.directory).await?;
+        self.state.segment_references = true;
+        self.state.segment_storage = Some(segments);
+        self.state.checkpoint = op;
+        self.state.anchor_known = checksum.is_some();
+        self.state.certified_log_view = None;
+        self.rewrite(op, checksum.unwrap_or(0), Some(0), checkpoint_prepare)
+            .await
+    }
+
+    pub(super) async fn migrate_segment_prepares(&mut self) -> io::Result<()> {
+        if self.state.segment_storage.is_some()
+            && self
+                .entries
+                .range(self.state.purge_floor.saturating_add(1)..)
+                .any(|(_, entry)| entry.reference.is_none())
+        {
+            self.rewrite(
+                self.state.checkpoint,
+                self.state.checkpoint_checksum,
+                None,
+                None,
+            )
+            .await?;
+        }
+        Ok(())
+    }
+
+    pub(super) async fn write_segment_bodies(
+        &mut self,
+        prepares: &[Frozen<4096>],
+        records: &[(u64, StoredPrepare, usize)],
+    ) -> io::Result<()> {
+        if self.state.segment_storage.is_none() {

Review Comment:
   external refs are supported; the caller writes their bodies first.



##########
core/journal/src/partition_journal/segments.rs:
##########
@@ -0,0 +1,663 @@
+// 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.
+
+use std::collections::{BTreeMap, BTreeSet};
+use std::io;
+use std::path::{Path, PathBuf};
+
+use iggy_binary_protocol::batch::BatchHeader;
+use iggy_binary_protocol::{Operation, PrepareHeader};
+use server_common::iobuf::Frozen;
+
+use super::{JournalState, PartitionPrepareJournal, SegmentReference, 
StoredPrepare, invalid};
+use crate::durable_storage::{DurableFile, DurableStorage, OpenMode};
+
+pub(super) const SEGMENT_STATE_FLAG: usize = 99;
+pub(super) const SEGMENT_STATE_OFFSET: usize = 128;
+pub(super) const SEGMENT_STATE_BYTES: usize = 10 * size_of::<u64>();
+
+/// A whole-batch boundary. `next_offset` is the first offset after this 
prefix.
+/// Physical tails may exceed the checkpoint boundary but cannot be polled yet.
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+pub struct SegmentPosition {
+    pub start_offset: u64,
+    pub length: u64,
+    pub next_offset: u64,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(super) struct SegmentCursor {
+    pub generation: u64,
+    pub position: SegmentPosition,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(super) struct SegmentState {
+    pub max_size: u64,
+    pub next_generation: u64,
+    pub tail: SegmentCursor,
+    pub checkpoint: SegmentCursor,
+}
+
+impl<S: DurableStorage> PartitionPrepareJournal<S> {
+    /// Enable segment body ownership after recovering the legacy committed 
view.
+    /// The initial boundary must describe durable materialized messages only.
+    ///
+    /// # Errors
+    /// Returns an error for inconsistent boundaries or a failed storage 
barrier.
+    pub async fn enable_segment_storage(
+        &mut self,
+        initial: SegmentPosition,
+        max_size: u64,
+    ) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if let Some(segments) = self.state.segment_storage {
+            if segments.max_size != max_size {
+                return Err(invalid("segment size differs from the durable WAL 
layout"));
+            }
+            return Ok(());
+        }
+        if max_size == 0 || !initial.valid() {
+            return Err(invalid("invalid initial segment boundary"));
+        }
+        let generation = self
+            .entries
+            .values()
+            .filter_map(|entry| entry.reference)
+            .map(|reference| reference.generation)
+            .max()
+            .map_or(Some(0), |generation| generation.checked_add(1))
+            .ok_or_else(|| invalid("segment generation exhausted"))?;
+        let cursor = SegmentCursor {
+            generation,
+            position: initial,
+        };
+        let segments = SegmentState {
+            max_size,
+            next_generation: generation
+                .checked_add(1)
+                .ok_or_else(|| invalid("segment generation exhausted"))?,
+            tail: cursor,
+            checkpoint: cursor,
+        };
+        let state = JournalState {
+            segment_references: true,
+            segment_storage: Some(segments),
+            ..self.state
+        };
+        self.poisoned = true;
+        if initial.length > 0 {
+            self.open_segment_file(cursor, 
self.preallocate_segments.then_some(max_size))
+                .await?;
+            self.sync_segment_files().await?;
+        }
+        self.file.sync().await?;
+        // Upgrade before any uncommitted body reaches an offset-named file.
+        self.publish(state).await?;
+        self.state = state;
+        self.durable_head = state.head;
+        self.poisoned = false;
+        self.migrate_segment_prepares().await
+    }
+
+    pub const fn segment_checkpoint(&self) -> Option<SegmentPosition> {
+        match self.state.segment_storage {
+            Some(segments) => Some(segments.checkpoint.position),
+            None => None,
+        }
+    }
+
+    pub fn segment_reference(&self, header: &PrepareHeader) -> 
Option<SegmentReference> {
+        if !self.contains(header) {
+            return None;
+        }
+        self.entries
+            .get(&header.op)
+            .and_then(|entry| entry.reference)
+    }
+
+    /// Install a transferred checkpoint after all replacement segment files 
are durable.
+    ///
+    /// # Errors
+    /// Returns an error if the checkpoint contradicts the installed bytes or 
a barrier fails.
+    #[allow(clippy::too_many_lines)]
+    pub async fn reset_with_segment_checkpoint(
+        &mut self,
+        op: u64,
+        checksum: Option<u128>,
+        prepare: Option<Frozen<4096>>,
+        initial: SegmentPosition,
+        max_size: u64,
+    ) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if !initial.valid() || max_size == 0 {
+            return Err(invalid("invalid installed segment boundary"));
+        }
+        if let Some(prepare) = &prepare {
+            let header = self.validate_checkpoint_prepare(prepare)?;
+            if header.op != op || Some(header.checksum) != checksum {
+                return Err(invalid("installed checkpoint prepare identity 
mismatch"));
+            }
+        }
+        let generation = self
+            .state
+            .segment_storage
+            .map_or(0, |segments| segments.next_generation);
+        let cursor = SegmentCursor {
+            generation,
+            position: initial,
+        };
+        let mut segments = SegmentState {
+            max_size,
+            next_generation: generation
+                .checked_add(1)
+                .ok_or_else(|| invalid("segment generation exhausted"))?,
+            tail: cursor,
+            checkpoint: cursor,
+        };
+        self.poisoned = true;
+        self.segment_files.clear();
+        self.open_segment_file(cursor, 
self.preallocate_segments.then_some(max_size))
+            .await?;
+        if self
+            .segment_files
+            .get(&(generation, initial.start_offset))
+            .ok_or_else(|| invalid("installed segment handle is absent"))?
+            .length()
+            .await?
+            != initial.length
+        {
+            return Err(invalid(
+                "installed segment size differs from its checkpoint",
+            ));
+        }
+        let checkpoint_prepare = if let Some(prepare) = prepare {
+            let header = self.validate_checkpoint_prepare(&prepare)?;
+            let reference = if header.operation == Operation::SendMessages {
+                let batch = decode_batch(prepare.as_slice())?;
+                if initial.length >= batch.batch_length
+                    && batch_next_offset(prepare.as_slice())? == 
initial.next_offset
+                {
+                    let reference = SegmentReference {
+                        generation,
+                        start_offset: initial.start_offset,
+                        position: initial.length - batch.batch_length,
+                        length: batch.batch_length,
+                    };
+                    let file = self
+                        .segment_files
+                        .get(&(generation, initial.start_offset))
+                        .ok_or_else(|| invalid("installed segment handle is 
absent"))?;
+                    if file
+                        .read(
+                            reference.position,
+                            prepare.len() - size_of::<PrepareHeader>(),
+                        )
+                        .await?
+                        != prepare.as_slice()[size_of::<PrepareHeader>()..]
+                    {
+                        return Err(invalid(
+                            "checkpoint prepare differs from installed segment 
bytes",
+                        ));
+                    }
+                    Some(reference)
+                } else if let Some(reference) = self
+                    .entries
+                    .get(&op)
+                    .filter(|entry| entry.checksum == header.checksum)
+                    .and_then(|entry| entry.reference)
+                {
+                    Some(reference)
+                } else {
+                    let retained = segments.allocate(SegmentPosition {
+                        start_offset: batch.base_offset,
+                        length: batch.batch_length,
+                        next_offset: batch_next_offset(prepare.as_slice())?,
+                    })?;
+                    let reference = SegmentReference {
+                        generation: retained.generation,
+                        start_offset: batch.base_offset,
+                        position: 0,
+                        length: batch.batch_length,
+                    };
+                    let mut file = self
+                        .storage
+                        .open(&reference.path(&self.directory), 
OpenMode::Create)
+                        .await?;
+                    file.write_frozen(0, 
prepare.slice(size_of::<PrepareHeader>()..))
+                        .await?;
+                    file.sync().await?;
+                    Some(reference)
+                }
+            } else {
+                None
+            };
+            Some((prepare, reference))
+        } else {
+            None
+        };
+        self.sync_segment_files().await?;
+        self.storage.sync_directory(&self.directory).await?;
+        self.state.segment_references = true;
+        self.state.segment_storage = Some(segments);
+        self.state.checkpoint = op;
+        self.state.anchor_known = checksum.is_some();
+        self.state.certified_log_view = None;
+        self.rewrite(op, checksum.unwrap_or(0), Some(0), checkpoint_prepare)
+            .await
+    }
+
+    pub(super) async fn migrate_segment_prepares(&mut self) -> io::Result<()> {
+        if self.state.segment_storage.is_some()
+            && self
+                .entries
+                .range(self.state.purge_floor.saturating_add(1)..)
+                .any(|(_, entry)| entry.reference.is_none())
+        {
+            self.rewrite(
+                self.state.checkpoint,
+                self.state.checkpoint_checksum,
+                None,
+                None,
+            )
+            .await?;
+        }
+        Ok(())
+    }
+
+    pub(super) async fn write_segment_bodies(
+        &mut self,
+        prepares: &[Frozen<4096>],
+        records: &[(u64, StoredPrepare, usize)],
+    ) -> io::Result<()> {
+        if self.state.segment_storage.is_none() {
+            return Ok(());
+        }
+        for (_, record, index) in records {
+            if let Some(reference) = record.reference {
+                self.write_segment_body(reference, &prepares[*index])
+                    .await?;
+            }
+        }
+        self.sync_segment_files().await
+    }
+
+    pub(super) async fn write_segment_body(
+        &mut self,
+        reference: SegmentReference,
+        prepare: &Frozen<4096>,
+    ) -> io::Result<()> {
+        let cursor = SegmentCursor {
+            generation: reference.generation,
+            position: SegmentPosition {
+                start_offset: reference.start_offset,
+                length: reference.position,
+                next_offset: reference.start_offset,
+            },
+        };
+        let preallocate_size = self
+            .state
+            .segment_storage
+            .filter(|_| self.preallocate_segments)
+            .map(|segments| segments.max_size);
+        self.open_segment_file(cursor, preallocate_size).await?;
+        self.segment_files
+            .get_mut(&(reference.generation, reference.start_offset))
+            .ok_or_else(|| invalid("segment writing handle is absent"))?
+            .write_frozen(
+                reference.position,
+                prepare.slice(size_of::<PrepareHeader>()..),
+            )
+            .await
+    }
+
+    pub(super) async fn sync_segment_files(&self) -> io::Result<()> {
+        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?;
+        }
+        Ok(())
+    }
+
+    pub(super) fn retain_active_segment_file(&mut self) {
+        if let Some(segments) = self.state.segment_storage {
+            let active = (
+                segments.tail.generation,
+                segments.tail.position.start_offset,
+            );
+            self.segment_files.retain(|key, _| *key == active);
+        }
+    }
+
+    pub(super) fn retained_segment_paths(
+        &self,
+        entries: &BTreeMap<u64, StoredPrepare>,
+        segments: Option<SegmentState>,
+    ) -> BTreeSet<PathBuf> {
+        let mut paths = super::referenced_segments(entries, &self.directory);
+        if let Some(segments) = segments {
+            paths.insert(segments.tail.path(&self.directory));
+            paths.insert(segments.checkpoint.path(&self.directory));
+        }
+        paths
+    }
+
+    pub(super) fn segment_boundary(&self, through_op: u64) -> 
Option<SegmentCursor> {
+        let segments = self.state.segment_storage?;
+        self.entries
+            .range(..=through_op)
+            .rev()
+            .find_map(|(&op, entry)| {
+                let reference = entry.reference?;
+                let next_offset = entry.next_offset?;
+                (op > self.state.purge_floor
+                    && op > self.state.checkpoint
+                    && next_offset >= segments.checkpoint.position.next_offset)
+                    .then_some(SegmentCursor {
+                        generation: reference.generation,
+                        position: SegmentPosition {
+                            start_offset: reference.start_offset,
+                            length: reference.position + reference.length,
+                            next_offset,
+                        },
+                    })
+            })
+            .or(Some(segments.checkpoint))
+    }
+
+    pub(super) async fn recover_segment_files(&mut self) -> io::Result<()> {
+        let Some(segments) = self.state.segment_storage else {
+            return Ok(());
+        };
+        self.poisoned = true;
+        let parent = self.segment_directory()?.to_path_buf();
+        let cursor = segments.tail;
+        let public = parent.join(format!("{:020}.log", 
cursor.position.start_offset));
+        // Retention may remove a fully checkpointed sealed segment. Its 
private
+        // checkpoint prepare remains available without restoring polled data.
+        let restore = cursor.position.length < segments.max_size
+            || cursor.position.next_offset > 
segments.checkpoint.position.next_offset
+            || self.storage.exists(&public).await?;
+        if restore {
+            self.open_segment_file(cursor, None).await?;
+            let file = self
+                .segment_files
+                .get(&(cursor.generation, cursor.position.start_offset))
+                .ok_or_else(|| invalid("segment rollback handle is absent"))?;
+            if file.length().await? < cursor.position.length {
+                return Err(invalid("segment lost durably published bytes"));
+            }
+            file.truncate(cursor.position.length).await?;
+            if self.preallocate_segments {
+                file.preallocate(&public, segments.max_size);
+            }
+            file.sync().await?;
+            self.remove_segment_name(&public).await?;
+            self.storage
+                .hard_link(&cursor.path(&self.directory), &public)
+                .await?;
+        }
+        for entry in self.storage.entries(&parent).await? {
+            let Some(name) = entry.name.to_str() else {
+                continue;
+            };
+            let offset = name
+                .strip_suffix(".log")
+                .or_else(|| name.strip_suffix(".index"))
+                .and_then(|offset| offset.parse::<u64>().ok());
+            if !entry.directory
+                && offset.is_some_and(|offset| {
+                    offset > cursor.position.start_offset && offset >= 
cursor.position.next_offset
+                })
+            {
+                self.storage.remove_file(&parent.join(entry.name)).await?;
+            }
+        }
+        self.storage.sync_directory(&parent).await?;
+        self.retain_active_segment_file();
+        self.poisoned = false;
+        Ok(())
+    }
+
+    async fn open_segment_file(
+        &mut self,
+        cursor: SegmentCursor,
+        preallocate_size: Option<u64>,
+    ) -> io::Result<()> {
+        let key = (cursor.generation, cursor.position.start_offset);
+        if self.segment_files.contains_key(&key) {
+            return Ok(());
+        }
+        let parent = self.segment_directory()?.to_path_buf();
+        let retained = cursor.path(&self.directory);
+        let public = parent.join(format!("{:020}.log", 
cursor.position.start_offset));
+        let file = if self.storage.exists(&retained).await? {
+            self.storage.open(&retained, OpenMode::ReadWrite).await?
+        } else if cursor.position.length > 0
+            || (self.storage.exists(&public).await?
+                && self
+                    .storage
+                    .open(&public, OpenMode::Read)
+                    .await?
+                    .length()
+                    .await?
+                    == 0)
+        {
+            self.storage.hard_link(&public, &retained).await?;
+            self.storage.open(&retained, OpenMode::ReadWrite).await?
+        } else {
+            let file = self.storage.open(&retained, OpenMode::Create).await?;
+            // Offset names can still point at a purged inode retained by older
+            // prepares. Unlink before reuse; never truncate that inode in 
place.
+            self.remove_segment_name(&public).await?;
+            self.storage.hard_link(&retained, &public).await?;
+            file
+        };
+        if let Some(size) = preallocate_size {
+            file.preallocate(&public, size);
+        }
+        self.storage.sync_directory(&self.directory).await?;
+        self.storage.sync_directory(&parent).await?;
+        self.segment_files.insert(key, file);
+        Ok(())
+    }
+
+    fn segment_directory(&self) -> io::Result<&Path> {
+        self.directory
+            .parent()
+            .ok_or_else(|| invalid("WAL has no segment directory"))
+    }
+
+    async fn remove_segment_name(&self, path: &Path) -> io::Result<()> {
+        match self.storage.remove_file(path).await {
+            Ok(()) => Ok(()),
+            Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
+            Err(error) => Err(error),
+        }
+    }
+}
+
+impl SegmentPosition {
+    const fn valid(self) -> bool {
+        self.next_offset >= self.start_offset
+            && ((self.length == 0) == (self.next_offset == self.start_offset))
+    }
+}
+
+impl SegmentCursor {
+    fn path(self, directory: &Path) -> PathBuf {
+        SegmentReference {
+            generation: self.generation,
+            start_offset: self.position.start_offset,
+            position: 0,
+            length: self.position.length,
+        }
+        .path(directory)
+    }
+}
+
+impl SegmentState {
+    pub(super) fn reserve(
+        &mut self,
+        header: &PrepareHeader,
+        prepare: &[u8],
+    ) -> io::Result<(Option<SegmentReference>, Option<u64>)> {
+        if header.operation != Operation::SendMessages {
+            return Ok((None, None));
+        }
+        let batch = decode_batch(prepare)?;
+        if batch.base_offset != self.tail.position.next_offset {
+            return Err(invalid(
+                "message offsets do not extend the physical segment tail",
+            ));
+        }
+        if self.tail.position.length >= self.max_size {
+            self.tail = self.allocate(SegmentPosition {
+                start_offset: batch.base_offset,
+                length: 0,
+                next_offset: batch.base_offset,
+            })?;
+        }
+        let reference = SegmentReference {
+            generation: self.tail.generation,
+            start_offset: self.tail.position.start_offset,
+            position: self.tail.position.length,
+            length: batch.batch_length,
+        };
+        self.tail.position.length = reference
+            .position
+            .checked_add(reference.length)
+            .ok_or_else(|| invalid("segment position exhausted"))?;
+        self.tail.position.next_offset = batch_next_offset(prepare)?;
+        Ok((Some(reference), Some(self.tail.position.next_offset)))
+    }
+
+    pub(super) fn reset_position(&mut self, position: SegmentPosition) -> 
io::Result<()> {
+        if !position.valid() {
+            return Err(invalid("invalid replacement segment boundary"));
+        }
+        self.tail = self.allocate(position)?;
+        self.checkpoint = self.tail;
+        Ok(())
+    }
+
+    fn allocate(&mut self, position: SegmentPosition) -> 
io::Result<SegmentCursor> {
+        let generation = self.next_generation;
+        self.next_generation = generation
+            .checked_add(1)
+            .ok_or_else(|| invalid("segment generation exhausted"))?;
+        Ok(SegmentCursor {
+            generation,
+            position,
+        })
+    }
+
+    pub(super) fn encode(self, bytes: &mut [u8]) {
+        let values = [
+            self.max_size,
+            self.next_generation,
+            self.tail.generation,
+            self.tail.position.start_offset,
+            self.tail.position.length,
+            self.tail.position.next_offset,
+            self.checkpoint.generation,
+            self.checkpoint.position.start_offset,
+            self.checkpoint.position.length,
+            self.checkpoint.position.next_offset,
+        ];
+        for (field, value) in bytes
+            .as_chunks_mut::<{ size_of::<u64>() }>()
+            .0
+            .iter_mut()
+            .zip(values)
+        {
+            *field = value.to_le_bytes();
+        }
+    }
+
+    pub(super) fn decode(bytes: &[u8]) -> io::Result<Self> {
+        let mut values = [0; SEGMENT_STATE_BYTES / size_of::<u64>()];
+        for (value, field) in values
+            .iter_mut()
+            .zip(bytes.as_chunks::<{ size_of::<u64>() }>().0)
+        {
+            *value = u64::from_le_bytes(*field);
+        }
+        let [
+            max_size,
+            next_generation,
+            tail_generation,
+            tail_start,
+            tail_length,
+            tail_next,
+            checkpoint_generation,
+            checkpoint_start,
+            checkpoint_length,
+            checkpoint_next,
+        ] = values;
+        let state = Self {
+            max_size,
+            next_generation,
+            tail: SegmentCursor {
+                generation: tail_generation,
+                position: SegmentPosition {
+                    start_offset: tail_start,
+                    length: tail_length,
+                    next_offset: tail_next,
+                },
+            },
+            checkpoint: SegmentCursor {
+                generation: checkpoint_generation,
+                position: SegmentPosition {
+                    start_offset: checkpoint_start,
+                    length: checkpoint_length,
+                    next_offset: checkpoint_next,
+                },
+            },
+        };
+        if max_size == 0
+            || tail_generation >= next_generation
+            || checkpoint_generation >= next_generation
+            || !state.tail.position.valid()
+            || !state.checkpoint.position.valid()

Review Comment:
   overflow already fails closed; a lower cap would just refuse sooner.



##########
core/journal/src/partition_journal/segments.rs:
##########
@@ -0,0 +1,663 @@
+// 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.
+
+use std::collections::{BTreeMap, BTreeSet};
+use std::io;
+use std::path::{Path, PathBuf};
+
+use iggy_binary_protocol::batch::BatchHeader;
+use iggy_binary_protocol::{Operation, PrepareHeader};
+use server_common::iobuf::Frozen;
+
+use super::{JournalState, PartitionPrepareJournal, SegmentReference, 
StoredPrepare, invalid};
+use crate::durable_storage::{DurableFile, DurableStorage, OpenMode};
+
+pub(super) const SEGMENT_STATE_FLAG: usize = 99;
+pub(super) const SEGMENT_STATE_OFFSET: usize = 128;
+pub(super) const SEGMENT_STATE_BYTES: usize = 10 * size_of::<u64>();
+
+/// A whole-batch boundary. `next_offset` is the first offset after this 
prefix.
+/// Physical tails may exceed the checkpoint boundary but cannot be polled yet.
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+pub struct SegmentPosition {
+    pub start_offset: u64,
+    pub length: u64,
+    pub next_offset: u64,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(super) struct SegmentCursor {
+    pub generation: u64,
+    pub position: SegmentPosition,
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(super) struct SegmentState {
+    pub max_size: u64,
+    pub next_generation: u64,
+    pub tail: SegmentCursor,
+    pub checkpoint: SegmentCursor,
+}
+
+impl<S: DurableStorage> PartitionPrepareJournal<S> {
+    /// Enable segment body ownership after recovering the legacy committed 
view.
+    /// The initial boundary must describe durable materialized messages only.
+    ///
+    /// # Errors
+    /// Returns an error for inconsistent boundaries or a failed storage 
barrier.
+    pub async fn enable_segment_storage(
+        &mut self,
+        initial: SegmentPosition,
+        max_size: u64,
+    ) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if let Some(segments) = self.state.segment_storage {
+            if segments.max_size != max_size {
+                return Err(invalid("segment size differs from the durable WAL 
layout"));
+            }
+            return Ok(());
+        }
+        if max_size == 0 || !initial.valid() {
+            return Err(invalid("invalid initial segment boundary"));
+        }
+        let generation = self
+            .entries
+            .values()
+            .filter_map(|entry| entry.reference)
+            .map(|reference| reference.generation)
+            .max()
+            .map_or(Some(0), |generation| generation.checked_add(1))
+            .ok_or_else(|| invalid("segment generation exhausted"))?;
+        let cursor = SegmentCursor {
+            generation,
+            position: initial,
+        };
+        let segments = SegmentState {
+            max_size,
+            next_generation: generation
+                .checked_add(1)
+                .ok_or_else(|| invalid("segment generation exhausted"))?,
+            tail: cursor,
+            checkpoint: cursor,
+        };
+        let state = JournalState {
+            segment_references: true,
+            segment_storage: Some(segments),
+            ..self.state
+        };
+        self.poisoned = true;
+        if initial.length > 0 {
+            self.open_segment_file(cursor, 
self.preallocate_segments.then_some(max_size))
+                .await?;
+            self.sync_segment_files().await?;
+        }
+        self.file.sync().await?;
+        // Upgrade before any uncommitted body reaches an offset-named file.
+        self.publish(state).await?;
+        self.state = state;
+        self.durable_head = state.head;
+        self.poisoned = false;
+        self.migrate_segment_prepares().await
+    }
+
+    pub const fn segment_checkpoint(&self) -> Option<SegmentPosition> {
+        match self.state.segment_storage {
+            Some(segments) => Some(segments.checkpoint.position),
+            None => None,
+        }
+    }
+
+    pub fn segment_reference(&self, header: &PrepareHeader) -> 
Option<SegmentReference> {
+        if !self.contains(header) {
+            return None;
+        }
+        self.entries
+            .get(&header.op)
+            .and_then(|entry| entry.reference)
+    }
+
+    /// Install a transferred checkpoint after all replacement segment files 
are durable.
+    ///
+    /// # Errors
+    /// Returns an error if the checkpoint contradicts the installed bytes or 
a barrier fails.
+    #[allow(clippy::too_many_lines)]
+    pub async fn reset_with_segment_checkpoint(
+        &mut self,
+        op: u64,
+        checksum: Option<u128>,
+        prepare: Option<Frozen<4096>>,
+        initial: SegmentPosition,
+        max_size: u64,
+    ) -> io::Result<()> {
+        self.ensure_healthy()?;
+        if !initial.valid() || max_size == 0 {
+            return Err(invalid("invalid installed segment boundary"));
+        }
+        if let Some(prepare) = &prepare {
+            let header = self.validate_checkpoint_prepare(prepare)?;
+            if header.op != op || Some(header.checksum) != checksum {
+                return Err(invalid("installed checkpoint prepare identity 
mismatch"));
+            }
+        }
+        let generation = self
+            .state
+            .segment_storage
+            .map_or(0, |segments| segments.next_generation);
+        let cursor = SegmentCursor {
+            generation,
+            position: initial,
+        };
+        let mut segments = SegmentState {
+            max_size,
+            next_generation: generation
+                .checked_add(1)
+                .ok_or_else(|| invalid("segment generation exhausted"))?,
+            tail: cursor,
+            checkpoint: cursor,
+        };
+        self.poisoned = true;
+        self.segment_files.clear();
+        self.open_segment_file(cursor, 
self.preallocate_segments.then_some(max_size))
+            .await?;
+        if self
+            .segment_files
+            .get(&(generation, initial.start_offset))
+            .ok_or_else(|| invalid("installed segment handle is absent"))?
+            .length()
+            .await?
+            != initial.length
+        {
+            return Err(invalid(
+                "installed segment size differs from its checkpoint",
+            ));
+        }
+        let checkpoint_prepare = if let Some(prepare) = prepare {
+            let header = self.validate_checkpoint_prepare(&prepare)?;
+            let reference = if header.operation == Operation::SendMessages {
+                let batch = decode_batch(prepare.as_slice())?;
+                if initial.length >= batch.batch_length
+                    && batch_next_offset(prepare.as_slice())? == 
initial.next_offset
+                {
+                    let reference = SegmentReference {
+                        generation,
+                        start_offset: initial.start_offset,
+                        position: initial.length - batch.batch_length,
+                        length: batch.batch_length,
+                    };
+                    let file = self
+                        .segment_files
+                        .get(&(generation, initial.start_offset))
+                        .ok_or_else(|| invalid("installed segment handle is 
absent"))?;
+                    if file
+                        .read(
+                            reference.position,
+                            prepare.len() - size_of::<PrepareHeader>(),
+                        )
+                        .await?
+                        != prepare.as_slice()[size_of::<PrepareHeader>()..]
+                    {
+                        return Err(invalid(
+                            "checkpoint prepare differs from installed segment 
bytes",
+                        ));
+                    }
+                    Some(reference)
+                } else if let Some(reference) = self
+                    .entries
+                    .get(&op)
+                    .filter(|entry| entry.checksum == header.checksum)
+                    .and_then(|entry| entry.reference)
+                {
+                    Some(reference)
+                } else {
+                    let retained = segments.allocate(SegmentPosition {
+                        start_offset: batch.base_offset,
+                        length: batch.batch_length,
+                        next_offset: batch_next_offset(prepare.as_slice())?,
+                    })?;
+                    let reference = SegmentReference {
+                        generation: retained.generation,
+                        start_offset: batch.base_offset,
+                        position: 0,
+                        length: batch.batch_length,
+                    };
+                    let mut file = self
+                        .storage
+                        .open(&reference.path(&self.directory), 
OpenMode::Create)
+                        .await?;
+                    file.write_frozen(0, 
prepare.slice(size_of::<PrepareHeader>()..))
+                        .await?;
+                    file.sync().await?;
+                    Some(reference)
+                }
+            } else {
+                None
+            };
+            Some((prepare, reference))
+        } else {
+            None
+        };
+        self.sync_segment_files().await?;
+        self.storage.sync_directory(&self.directory).await?;
+        self.state.segment_references = true;
+        self.state.segment_storage = Some(segments);
+        self.state.checkpoint = op;
+        self.state.anchor_known = checksum.is_some();
+        self.state.certified_log_view = None;
+        self.rewrite(op, checksum.unwrap_or(0), Some(0), checkpoint_prepare)
+            .await
+    }
+
+    pub(super) async fn migrate_segment_prepares(&mut self) -> io::Result<()> {
+        if self.state.segment_storage.is_some()
+            && self
+                .entries
+                .range(self.state.purge_floor.saturating_add(1)..)
+                .any(|(_, entry)| entry.reference.is_none())
+        {
+            self.rewrite(
+                self.state.checkpoint,
+                self.state.checkpoint_checksum,
+                None,
+                None,
+            )
+            .await?;
+        }
+        Ok(())
+    }
+
+    pub(super) async fn write_segment_bodies(
+        &mut self,
+        prepares: &[Frozen<4096>],
+        records: &[(u64, StoredPrepare, usize)],
+    ) -> io::Result<()> {
+        if self.state.segment_storage.is_none() {
+            return Ok(());
+        }
+        for (_, record, index) in records {
+            if let Some(reference) = record.reference {
+                self.write_segment_body(reference, &prepares[*index])
+                    .await?;
+            }
+        }
+        self.sync_segment_files().await
+    }
+
+    pub(super) async fn write_segment_body(
+        &mut self,
+        reference: SegmentReference,
+        prepare: &Frozen<4096>,
+    ) -> io::Result<()> {
+        let cursor = SegmentCursor {
+            generation: reference.generation,
+            position: SegmentPosition {
+                start_offset: reference.start_offset,
+                length: reference.position,
+                next_offset: reference.start_offset,
+            },
+        };
+        let preallocate_size = self
+            .state
+            .segment_storage
+            .filter(|_| self.preallocate_segments)
+            .map(|segments| segments.max_size);
+        self.open_segment_file(cursor, preallocate_size).await?;
+        self.segment_files
+            .get_mut(&(reference.generation, reference.start_offset))
+            .ok_or_else(|| invalid("segment writing handle is absent"))?
+            .write_frozen(
+                reference.position,
+                prepare.slice(size_of::<PrepareHeader>()..),
+            )
+            .await
+    }
+
+    pub(super) async fn sync_segment_files(&self) -> io::Result<()> {
+        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?;
+        }
+        Ok(())
+    }
+
+    pub(super) fn retain_active_segment_file(&mut self) {
+        if let Some(segments) = self.state.segment_storage {
+            let active = (
+                segments.tail.generation,
+                segments.tail.position.start_offset,
+            );
+            self.segment_files.retain(|key, _| *key == active);
+        }
+    }
+
+    pub(super) fn retained_segment_paths(
+        &self,
+        entries: &BTreeMap<u64, StoredPrepare>,
+        segments: Option<SegmentState>,
+    ) -> BTreeSet<PathBuf> {
+        let mut paths = super::referenced_segments(entries, &self.directory);
+        if let Some(segments) = segments {
+            paths.insert(segments.tail.path(&self.directory));
+            paths.insert(segments.checkpoint.path(&self.directory));
+        }
+        paths
+    }
+
+    pub(super) fn segment_boundary(&self, through_op: u64) -> 
Option<SegmentCursor> {
+        let segments = self.state.segment_storage?;
+        self.entries
+            .range(..=through_op)
+            .rev()
+            .find_map(|(&op, entry)| {
+                let reference = entry.reference?;
+                let next_offset = entry.next_offset?;
+                (op > self.state.purge_floor
+                    && op > self.state.checkpoint
+                    && next_offset >= segments.checkpoint.position.next_offset)
+                    .then_some(SegmentCursor {
+                        generation: reference.generation,
+                        position: SegmentPosition {
+                            start_offset: reference.start_offset,
+                            length: reference.position + reference.length,
+                            next_offset,
+                        },
+                    })
+            })
+            .or(Some(segments.checkpoint))
+    }
+
+    pub(super) async fn recover_segment_files(&mut self) -> io::Result<()> {
+        let Some(segments) = self.state.segment_storage else {
+            return Ok(());
+        };
+        self.poisoned = true;
+        let parent = self.segment_directory()?.to_path_buf();
+        let cursor = segments.tail;
+        let public = parent.join(format!("{:020}.log", 
cursor.position.start_offset));
+        // Retention may remove a fully checkpointed sealed segment. Its 
private
+        // checkpoint prepare remains available without restoring polled data.
+        let restore = cursor.position.length < segments.max_size
+            || cursor.position.next_offset > 
segments.checkpoint.position.next_offset
+            || self.storage.exists(&public).await?;
+        if restore {
+            self.open_segment_file(cursor, None).await?;
+            let file = self
+                .segment_files
+                .get(&(cursor.generation, cursor.position.start_offset))
+                .ok_or_else(|| invalid("segment rollback handle is absent"))?;
+            if file.length().await? < cursor.position.length {
+                return Err(invalid("segment lost durably published bytes"));
+            }
+            file.truncate(cursor.position.length).await?;
+            if self.preallocate_segments {
+                file.preallocate(&public, segments.max_size);
+            }
+            file.sync().await?;
+            self.remove_segment_name(&public).await?;
+            self.storage
+                .hard_link(&cursor.path(&self.directory), &public)
+                .await?;
+        }
+        for entry in self.storage.entries(&parent).await? {
+            let Some(name) = entry.name.to_str() else {
+                continue;
+            };
+            let offset = name
+                .strip_suffix(".log")
+                .or_else(|| name.strip_suffix(".index"))
+                .and_then(|offset| offset.parse::<u64>().ok());
+            if !entry.directory
+                && offset.is_some_and(|offset| {
+                    offset > cursor.position.start_offset && offset >= 
cursor.position.next_offset
+                })
+            {
+                self.storage.remove_file(&parent.join(entry.name)).await?;
+            }
+        }
+        self.storage.sync_directory(&parent).await?;
+        self.retain_active_segment_file();
+        self.poisoned = false;
+        Ok(())
+    }
+
+    async fn open_segment_file(
+        &mut self,
+        cursor: SegmentCursor,
+        preallocate_size: Option<u64>,
+    ) -> io::Result<()> {
+        let key = (cursor.generation, cursor.position.start_offset);
+        if self.segment_files.contains_key(&key) {
+            return Ok(());
+        }
+        let parent = self.segment_directory()?.to_path_buf();
+        let retained = cursor.path(&self.directory);
+        let public = parent.join(format!("{:020}.log", 
cursor.position.start_offset));
+        let file = if self.storage.exists(&retained).await? {
+            self.storage.open(&retained, OpenMode::ReadWrite).await?
+        } else if cursor.position.length > 0
+            || (self.storage.exists(&public).await?
+                && self
+                    .storage
+                    .open(&public, OpenMode::Read)
+                    .await?
+                    .length()
+                    .await?
+                    == 0)
+        {
+            self.storage.hard_link(&public, &retained).await?;
+            self.storage.open(&retained, OpenMode::ReadWrite).await?
+        } else {
+            let file = self.storage.open(&retained, OpenMode::Create).await?;
+            // Offset names can still point at a purged inode retained by older
+            // prepares. Unlink before reuse; never truncate that inode in 
place.
+            self.remove_segment_name(&public).await?;
+            self.storage.hard_link(&retained, &public).await?;
+            file
+        };
+        if let Some(size) = preallocate_size {
+            file.preallocate(&public, size);
+        }
+        self.storage.sync_directory(&self.directory).await?;
+        self.storage.sync_directory(&parent).await?;
+        self.segment_files.insert(key, file);
+        Ok(())
+    }
+
+    fn segment_directory(&self) -> io::Result<&Path> {
+        self.directory
+            .parent()
+            .ok_or_else(|| invalid("WAL has no segment directory"))
+    }
+
+    async fn remove_segment_name(&self, path: &Path) -> io::Result<()> {
+        match self.storage.remove_file(path).await {
+            Ok(()) => Ok(()),
+            Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
+            Err(error) => Err(error),
+        }
+    }
+}
+
+impl SegmentPosition {
+    const fn valid(self) -> bool {
+        self.next_offset >= self.start_offset
+            && ((self.length == 0) == (self.next_offset == self.start_offset))
+    }
+}
+
+impl SegmentCursor {
+    fn path(self, directory: &Path) -> PathBuf {
+        SegmentReference {
+            generation: self.generation,
+            start_offset: self.position.start_offset,
+            position: 0,
+            length: self.position.length,
+        }
+        .path(directory)
+    }
+}
+
+impl SegmentState {
+    pub(super) fn reserve(
+        &mut self,
+        header: &PrepareHeader,
+        prepare: &[u8],
+    ) -> io::Result<(Option<SegmentReference>, Option<u64>)> {
+        if header.operation != Operation::SendMessages {
+            return Ok((None, None));
+        }
+        let batch = decode_batch(prepare)?;
+        if batch.base_offset != self.tail.position.next_offset {
+            return Err(invalid(
+                "message offsets do not extend the physical segment tail",
+            ));
+        }
+        if self.tail.position.length >= self.max_size {
+            self.tail = self.allocate(SegmentPosition {
+                start_offset: batch.base_offset,
+                length: 0,
+                next_offset: batch.base_offset,
+            })?;
+        }
+        let reference = SegmentReference {
+            generation: self.tail.generation,
+            start_offset: self.tail.position.start_offset,
+            position: self.tail.position.length,
+            length: batch.batch_length,
+        };
+        self.tail.position.length = reference
+            .position
+            .checked_add(reference.length)
+            .ok_or_else(|| invalid("segment position exhausted"))?;
+        self.tail.position.next_offset = batch_next_offset(prepare)?;
+        Ok((Some(reference), Some(self.tail.position.next_offset)))
+    }
+
+    pub(super) fn reset_position(&mut self, position: SegmentPosition) -> 
io::Result<()> {
+        if !position.valid() {
+            return Err(invalid("invalid replacement segment boundary"));
+        }
+        self.tail = self.allocate(position)?;
+        self.checkpoint = self.tail;
+        Ok(())
+    }
+
+    fn allocate(&mut self, position: SegmentPosition) -> 
io::Result<SegmentCursor> {
+        let generation = self.next_generation;
+        self.next_generation = generation
+            .checked_add(1)
+            .ok_or_else(|| invalid("segment generation exhausted"))?;
+        Ok(SegmentCursor {
+            generation,
+            position,
+        })
+    }
+
+    pub(super) fn encode(self, bytes: &mut [u8]) {
+        let values = [
+            self.max_size,
+            self.next_generation,
+            self.tail.generation,
+            self.tail.position.start_offset,
+            self.tail.position.length,
+            self.tail.position.next_offset,
+            self.checkpoint.generation,
+            self.checkpoint.position.start_offset,
+            self.checkpoint.position.length,
+            self.checkpoint.position.next_offset,
+        ];
+        for (field, value) in bytes
+            .as_chunks_mut::<{ size_of::<u64>() }>()
+            .0
+            .iter_mut()
+            .zip(values)
+        {
+            *field = value.to_le_bytes();
+        }
+    }
+
+    pub(super) fn decode(bytes: &[u8]) -> io::Result<Self> {
+        let mut values = [0; SEGMENT_STATE_BYTES / size_of::<u64>()];
+        for (value, field) in values
+            .iter_mut()
+            .zip(bytes.as_chunks::<{ size_of::<u64>() }>().0)
+        {
+            *value = u64::from_le_bytes(*field);
+        }
+        let [
+            max_size,
+            next_generation,
+            tail_generation,
+            tail_start,
+            tail_length,
+            tail_next,
+            checkpoint_generation,
+            checkpoint_start,
+            checkpoint_length,
+            checkpoint_next,
+        ] = values;
+        let state = Self {
+            max_size,
+            next_generation,
+            tail: SegmentCursor {
+                generation: tail_generation,
+                position: SegmentPosition {
+                    start_offset: tail_start,
+                    length: tail_length,
+                    next_offset: tail_next,
+                },
+            },
+            checkpoint: SegmentCursor {
+                generation: checkpoint_generation,
+                position: SegmentPosition {
+                    start_offset: checkpoint_start,
+                    length: checkpoint_length,
+                    next_offset: checkpoint_next,
+                },
+            },
+        };

Review Comment:
   added malformed-layout decode and reopen cases with valid checksums.



-- 
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