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


##########
core/bench/src/args/examples.rs:
##########
@@ -15,176 +15,220 @@
 // specific language governing permissions and limitations
 // under the License.
 
-const EXAMPLES: &str = r#"EXAMPLES:
+const EXAMPLES: &str = r"EXAMPLES:
 
-1) Pinned Mode Benchmarking:
+Start iggy-server separately. The benchmark connects to a running server.
+Global options precede the kind, kind options precede the transport, and
+transport options precede the optional output subcommand.
 
-    Run benchmarks with pinned producers and consumers. This mode pins 
specific producers
-    and consumers to specific streams and partitions (one to one):
+Default producer and consumer counts are six. Pinned workloads default to six 
streams.
+
+1) All benchmark kinds and aliases:
+
+    Pinned producer (pp), consumer (pc), and producer/consumer (ppc):
 
     $ cargo r -r --bin iggy-bench -- pinned-producer --streams 10 --producers 
10 tcp
     $ cargo r -r --bin iggy-bench -- pinned-consumer --streams 10 --consumers 
10 tcp
     $ cargo r -r --bin iggy-bench -- pinned-producer-and-consumer --streams 10 
--producers 10 --consumers 10 tcp
-    $ cargo r -r --bin iggy-bench -- -T 10GB pp --producers 5 tcp
-
-2) Balanced Mode Benchmarking:
 
-    Run benchmarks with balanced distribution of producers and consumers. This 
mode
-    automatically balances the load across streams and consumer groups:
+    Balanced producer (bp), consumer group (bcg), and producer/consumer group 
(bpcg):
 
     $ cargo r -r --bin iggy-bench -- balanced-producer --partitions 24 
--producers 6 tcp
     $ cargo r -r --bin iggy-bench -- balanced-consumer-group --consumers 6 tcp
     $ cargo r -r --bin iggy-bench -- balanced-producer-and-consumer-group 
--partitions 24 --producers 6 --consumers 6 tcp
-    $ cargo r -r --bin iggy-bench -- -T 10GB bpc tcp
-
-    Durability-matched run, where every produce ack waits for an fsync
-    (--partitions 1 routes every producer straight to that partition):
-
-    $ cargo r -r --bin iggy-bench -- --enforce-fsync 
--messages-required-to-save 1 \
-        balanced-producer --partitions 1 --producers 8 tcp
-
-3) End-to-End Benchmarking:
-
-    Run end-to-end benchmarks that measure performance for a producer that is 
also a consumer:
-
-    $ cargo r -r --bin iggy-bench -- end-to-end-producing-consumer --producers 
12 --streams 12 tcp
-    $ cargo r -r --bin iggy-bench -- end-to-end-producing-consumer-group 
--partitions 24 --producers 6 tcp
-
-4) Advanced Configuration:
-
-    You can customize various parameters for any benchmark mode:
-
-    Global options (before the benchmark command):
-    --messages-per-batch (-P): Number of messages per batch [default: 1000]
-                               For random batch sizes, use range format: 
"100..1000"
-    --message-batches (-b): Total number of batches [default: 1000]
-    --total-messages-size (-T): Total size of messages to send (e.g., "1GB", 
"500MB")
-                                Mutually exclusive with --message-batches
-    --message-size (-m): Message size in bytes [default: 1000]
-                        For random sizes, use range format: "100..1000"
-    --rate-limit (-r): Optional throughput limit per producer (e.g., "50KB", 
"10MB")
-    --warmup-time (-w): Warmup duration [default: 0s]
-    --sampling-time (-t): Metrics sampling interval [default: 10ms]
-    --moving-average-window (-W): Window size for moving average [default: 20]
-    --username (-u): Username for server authentication [default: iggy]
-    --password (-p): Password for server authentication [default: iggy]
-    --reuse-streams: Reuse existing bench streams instead of deleting them
-
-    Benchmark-specific options (after the benchmark command):
-    --streams (-s): Number of streams
-    --partitions (-a): Number of partitions
-    --producers (-c): Number of producers
-    --consumers (-c): Number of consumers
-    --max-topic-size (-T): Max topic size (e.g., "1GiB")
-    --message-expiry (-e): Topic message expiry time (e.g., "1s", "5min", "1h")
-
-    Examples with detailed configuration:
-
-    # Fixed message and batch sizes:
-    $ cargo r -r --bin iggy-bench -- \
-        --message-size 1000 \
-        --messages-per-batch 100 \
-        --message-batches 1000 \
-        --rate-limit "100MB" \
-        balanced-producer \
-        --streams 5 \
-        --producers 5 \
-        tcp
-
-    # Random message sizes (100-1000 bytes):
-    $ cargo r -r --bin iggy-bench -- \
-        --message-size "100..1000" \
-        --messages-per-batch 100 \
-        --total-messages-size "1GB" \
-        balanced-producer \
-        --streams 5 \
-        --producers 5 \
-        tcp
-
-    # Random batch sizes (10-100 messages per batch):
-    $ cargo r -r --bin iggy-bench -- \
-        --message-size 1000 \
-        --messages-per-batch "10..100" \
-        --total-messages-size "500MB" \
-        balanced-producer \
-        --streams 5 \
-        --producers 5 \
-        tcp
-
-    # Random message and batch sizes with rate limiting:
-    $ cargo r -r --bin iggy-bench -- \
-        --message-size "500..2000" \
-        --messages-per-batch "50..200" \
-        --total-messages-size "2GB" \
-        --rate-limit "50MB" \
-        balanced-producer \
-        --streams 5 \
-        --producers 5 \
-        tcp
-
-5) Remote Server Benchmarking:
-
-    To benchmark a remote server, specify the server address in the transport 
subcommand.
-    Both IP addresses and hostnames are supported:
-
-    $ cargo r -r --bin iggy-bench -- pinned-producer \
-        --streams 5 --producers 5 \
-        tcp --server-address 192.168.1.100:8090
-    $ cargo r -r --bin iggy-bench -- pinned-producer \
-        --streams 5 --producers 5 \
-        tcp --server-address localhost:8090
-
-    With custom credentials:
-
-    $ cargo r -r --bin iggy-bench -- \
-        --username admin --password secret \
-        pinned-producer --streams 5 --producers 5 \
-        tcp --server-address 192.168.1.100:8090
-
-6) Output Data and Results:
-
-    The benchmark tool can store detailed results for analysis and comparison:
-
-    # Basic result storage (results will be stored in ./performance_results):
-    $ cargo r -r --bin iggy-bench -- pinned-producer --streams 10 --producers 
10 tcp output
-
-
-    # Organized benchmarking with metadata:
-    $ cargo r -r --bin iggy-bench -- balanced-producer --partitions 24 
--producers 6 tcp \
-        output \
-        --identifier "prod-test-$(date +%Y%m%d)" \
-        --remark "production-config" \
-        --gitref "$(git rev-parse --short HEAD)" \
-        --gitref-date "$(git show -s --format=%cI HEAD)"
-
-    # Quick result visualization:
-    $ cargo r -r --bin iggy-bench -- end-to-end-producing-consumer --producers 
12 --streams 12 tcp \
-        output --open-charts
-
-    Output configuration options:
-    --open-charts (-c)  : Open charts after the benchmark
-    --output-dir (-o)   : Directory for storing results [default: 
performance_results]
-    --identifier        : Benchmark run ID (if not provided defaults to 
hostname)
-    --remark            : Additional context (e.g., "production-config")
-    --extra-info        : Custom metadata for future analysis, currently unused
-
-7) Help and Documentation:
-
-    For more details on available options:
-
-    # General help
-    $ cargo r -r --bin iggy-bench -- --help
+    $ cargo r -r --bin iggy-bench -- --total-data 10GiB bpcg tcp
 
-    # Specific benchmark help
-    $ cargo r -r --bin iggy-bench -- pinned-producer --help
+    End-to-end producing consumer (e2e) and producing consumer group (e2ecg):
 
-    # Protocol help
-    $ cargo r -r --bin iggy-bench -- pinned-producer tcp --help
+    $ cargo r -r --bin iggy-bench -- end-to-end-producing-consumer 
--producing-consumers 12 --streams 12 tcp
+    $ cargo r -r --bin iggy-bench -- end-to-end-producing-consumer-group 
--partitions 24 --producers 6 --consumers 6 tcp
+
+2) All transports:
+
+    $ cargo r -r --bin iggy-bench -- pinned-producer tcp --server-address 
127.0.0.1:8090
+    $ cargo r -r --bin iggy-bench -- pinned-producer quic --server-address 
127.0.0.1:8080
+    $ cargo r -r --bin iggy-bench -- pinned-producer http --server-address 
127.0.0.1:3000
+    $ cargo r -r --bin iggy-bench -- pinned-producer websocket 
--server-address 127.0.0.1:8092
+
+3) Topic durability:
+
+    --durability controls message completion.
+    --consumer-offset-durability controls explicit offset-store/delete 
completion.
+    Both independently default to replicated. Neither inherits the other.
+    Both policies normally write to disk. Persisted additionally waits for
+    recoverable stable storage on the required VSR quorum, or the single 
replica.
+
+    Both replicated (the default):
+    $ cargo r -r --bin iggy-bench -- balanced-producer-and-consumer-group tcp
+
+    Persisted messages, replicated offsets:
+    $ cargo r -r --bin iggy-bench -- --durability persisted 
balanced-producer-and-consumer-group tcp
+
+    Replicated messages, persisted offsets:
+    $ cargo r -r --bin iggy-bench -- --consumer-offset-durability persisted 
balanced-producer-and-consumer-group tcp
+
+    Both persisted:
+    $ cargo r -r --bin iggy-bench -- --durability persisted 
--consumer-offset-durability persisted balanced-producer-and-consumer-group tcp
+
+    These options apply when topics are created. They do not modify topics
+    with --reuse-streams. Consumer polling with auto-commit remains 
asynchronous.
+    Its poll latency is not an acknowledged offset-store latency measurement.
+    In replicated groups, either persisted policy enables a WAL that also 
retains
+    message predecessors, even when message durability is replicated.
+
+4) Topic flush cadence and retention:
+
+    --messages-required-to-save is a global create-time topic option.
+    It controls segment flush cadence, not the acknowledgment guarantee.
+    --max-topic-size and --message-expiry are kind-specific topic options.
+    These flags do not change existing topics with --reuse-streams.
+
+    Persisted messages without forcing a segment flush for every message:
+    $ cargo r -r --bin iggy-bench -- --durability persisted 
--messages-required-to-save 1024 balanced-producer --partitions 1 --producers 8 
tcp
 
-    # Output help
+    Replicated acknowledgments with eager segment flushing:
+    $ cargo r -r --bin iggy-bench -- --messages-required-to-save 1 
balanced-producer tcp
+
+    Retention can delete data during a long run. Use matching policies when 
comparing:
+    $ cargo r -r --bin iggy-bench -- balanced-producer --max-topic-size 10GiB 
--message-expiry 1h tcp
+
+5) Workload configuration:
+
+    --messages-per-batch (-P): Messages per batch, or a range such as 
100..1000.
+    --message-batches (-b): Batches per actor, mutually exclusive with 
--total-data.
+    --total-data (-T): Total message bytes across actors, such as 10GiB.
+    --message-size (-m): Message bytes, or a range such as 100..1000.
+    --rate-limit (-r): Aggregate throughput limit across actors, such as 100MB.
+    --warmup-time (-w): Warmup duration, such as 10s.
+    --sampling-time (-t): Metrics sampling interval.
+    --moving-average-window (-W): Moving-average window size.
+
+    Fixed message and batch sizes:
+    $ cargo r -r --bin iggy-bench -- --message-size 1000 --messages-per-batch 
100 --message-batches 1000 --rate-limit 100MB balanced-producer --streams 5 
--producers 5 tcp
+
+    Random message sizes:
+    $ cargo r -r --bin iggy-bench -- --message-size 100..1000 
--messages-per-batch 100 --total-data 1GiB balanced-producer --streams 5 
--producers 5 tcp
+
+    Random batch sizes:
+    $ cargo r -r --bin iggy-bench -- --message-size 1000 --messages-per-batch 
10..100 --total-data 500MiB balanced-producer tcp
+
+    Random message and batch sizes with a warmup and aggregate rate limit:
+    $ cargo r -r --bin iggy-bench -- --message-size 500..2000 
--messages-per-batch 50..200 --total-data 2GiB --warmup-time 10s --rate-limit 
50MB balanced-producer tcp
+
+6) Remote server and output:
+
+    $ cargo r -r --bin iggy-bench -- pinned-producer --streams 5 --producers 5 
tcp --server-address localhost:8090
+    $ cargo r -r --bin iggy-bench -- --username admin --password secret 
pinned-producer tcp --server-address 192.168.1.100:8090
+    $ cargo r -r --bin iggy-bench -- pinned-producer tcp output
+    $ cargo r -r --bin iggy-bench -- --durability persisted balanced-producer 
tcp output --identifier dedicated-host --remark persisted-messages --gitref 
abc123
+    $ cargo r -r --bin iggy-bench -- end-to-end-producing-consumer tcp output 
--open-charts
+
+    Output options include --output-dir (-o), --identifier, --remark, --gitref,
+    --gitref-date, --extra-info, and --open-charts (-c).
+    Record both durability policies, CPU allocation, host tuning, and storage.
+    See core/bench/README.md and 
https://iggy.apache.org/docs/server/linux-tuning.
+
+7) Help:
+
+    $ cargo r -r --bin iggy-bench -- --help
+    $ cargo r -r --bin iggy-bench -- pinned-producer --help
+    $ cargo r -r --bin iggy-bench -- pinned-producer tcp --help
     $ cargo r -r --bin iggy-bench -- pinned-producer tcp output --help
-"#;
+";
 
 pub fn print_examples() {
     println!("{EXAMPLES}");
 }
+
+#[cfg(test)]
+mod tests {
+    use super::EXAMPLES;
+    use crate::args::common::IggyBenchArgs;
+    use clap::{CommandFactory, Parser, error::ErrorKind};
+    use iggy::prelude::Durability;
+    use std::collections::BTreeSet;
+
+    #[test]
+    fn published_examples_parse_and_cover_every_kind_and_transport() {

Review Comment:
   added semantic validation for every published example.



##########
core/server/config.toml:
##########
@@ -1005,6 +844,68 @@ clients_table_max = 8192
 # plane), a pipeline exists per partition, so raising this multiplies pinned
 # request-buffer memory by the partition count. Keep it modest.
 [partition]
+# Retention is per topic, set at CreateTopic and readable on GetTopic:
+#   max_topic_size   - delete oldest sealed segments past this size

Review Comment:
   documented retention waiting for a requested wal checkpoint.



##########
core/metadata/src/stm/stream.rs:
##########
@@ -2820,11 +2819,21 @@ mod tests {
         assert_eq!(topic.compression_algorithm, CompressionAlgorithm::None);
 
         // partitions_count is create-consumed, never persisted.
-        assert_eq!(topic.options.len(), 3);
+        assert_eq!(topic.options.len(), 5);
         let expiry_key = 
HeaderKey::from_str(topic_option_keys::MESSAGE_EXPIRY).unwrap();
         let expiry = topic.options.get(&expiry_key).unwrap();
         assert!(expiry.explicit, "client-sent key keeps its provenance");
         assert_eq!(expiry.value.kind(), HeaderKind::Uint64);

Review Comment:
   the test now checks derived defaults through explicit-wire encoding.



##########
core/bench/src/args/kinds/pinned/consumer.rs:
##########
@@ -16,7 +16,7 @@
 // under the License.

Review Comment:
   fixed the messages to match the inequalities we actually enforce.



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

Review Comment:
   documented the existing-layout segment size mismatch error.



##########
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"));

Review Comment:
   fixed; damaged wal state is fenced and recovered from peers.



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