hubcio commented on code in PR #4092:
URL: https://github.com/apache/iggy/pull/4092#discussion_r3957118691
##########
core/server/src/partition_helpers.rs:
##########
@@ -1399,6 +1376,10 @@ pub async fn build_partition_fresh(
});
}
+ partition
+
.open_persistence_with_capacity(config.partition.wal_bytes_max.as_bytes_u64())
Review Comment:
critical: quarantine rebuild restores the WAL checkpoint over empty
replacement segments, letting the replica rejoin with committed messages
missing. require full state transfer before accepting that checkpoint as
applied.
##########
core/common/src/types/options/mod.rs:
##########
@@ -800,8 +823,11 @@ impl TopicCreateOptions {
let size = parse_byte_size(entry, key)?;
parsed.segment_size = (size !=
0).then_some(IggyByteSize::from(size));
}
- topic_option_keys::ENFORCE_FSYNC => {
- parsed.enforce_fsync = Some(parse_bool(entry, key)?);
+ topic_option_keys::DURABILITY => {
Review Comment:
critical: existing `enforce_fsync=true` topic metadata silently becomes
`Durability::Replicated` after config migration, disabling fsync on reopened
writers. explicitly migrate the stored policy or reject it before opening
writers.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -711,6 +720,292 @@ where
partition
}
+ pub fn set_persistence_notifier(&self, notifier: PersistenceNotifier) {
+ if let Some(persistence) = &self.persistence {
+ persistence.set_notifier(notifier);
+ }
+ }
+
+ /// # Errors
+ /// Returns an error if durable prepare history cannot be opened or
replayed.
+ pub async fn open_persistence(&mut self) -> Result<(), IggyError> {
+
self.open_persistence_with_capacity(journal::partition_journal::PARTITION_WAL_BYTES_MAX)
+ .await
+ }
+
+ /// # Errors
+ /// Returns an error if durable prepare history cannot be opened or
replayed.
+ pub async fn open_persistence_with_capacity(&mut self, capacity: u64) ->
Result<(), IggyError> {
+ if self.consensus.replica_count() == 1
+ || !(self.durability().is_persisted()
+ || self.consumer_offset_durability().is_persisted())
+ {
+ return Ok(());
+ }
+ let directory = self
+ .partition_dir
+ .as_ref()
+ .ok_or(IggyError::CannotReadFile)?;
+ let directory =
+ std::path::Path::new(directory).join(format!("prepares-{}",
self.created_revision));
+ let (persistence, prepares) = PartitionPersistence::open_with_capacity(
+ &directory,
+ self.namespace().inner(),
+ self.created_revision,
+ journal::durable_storage::DiskStorage,
+ capacity,
+ )
+ .await
+ .map_err(|error| {
+ warn!(%error, "cannot open partition prepare WAL");
+ IggyError::CannotReadFile
+ })?;
+ let (purge_generation, purge_floor) = persistence.purge_marker();
+ if purge_generation <= self.applied_purge_generation {
+ self.purge_floor_op = self.purge_floor_op.max(purge_floor);
+ }
+ let checkpoint = persistence.checkpoint_op();
+ let head = persistence.head();
+ let mut commit = checkpoint;
+ for prepare in prepares {
+ let mut owned =
server_common::iobuf::Owned::<4096>::zeroed(prepare.len());
+ owned.as_mut_slice().copy_from_slice(prepare.as_slice());
+ let message =
+ Message::<PrepareHeader>::try_from(owned).map_err(|_|
IggyError::InvalidCommand)?;
+ let header = *message.header();
+ commit = commit.max(header.commit.min(head));
+ if header.operation == Operation::SendMessages {
+ self.append_repaired_send_messages(message).await?;
+ } else {
+ self.apply_replicated_operation(message).await?;
+ }
+ }
+ if head > 0 {
+ self.consensus.sequencer().set_sequence(head);
+ if let Some(checksum) = persistence.checksum(head) {
+ self.consensus.set_last_prepare_checksum(checksum);
+ }
+ self.consensus.restore_commit_state(checkpoint, commit);
+ }
+ // Recovery can reuse completed writes whose last barrier was
interrupted.
+ // Include them once before reclaiming any recovered WAL history.
+ for segment in self.log.segments() {
+ persistence.mark_segment_dirty(segment.start_offset);
+ }
+ for kind in [ConsumerKind::Consumer, ConsumerKind::ConsumerGroup] {
+ self.durable_consumer_offsets.with_entries(kind, |entries| {
+ for consumer_id in entries.keys() {
+ persistence.mark_offset_dirty(
+ crate::state_transfer::consumer_kind_index(kind),
+ *consumer_id,
+ true,
+ );
+ }
+ });
+ }
+ self.persistence = Some(persistence);
+ Ok(())
+ }
+
+ pub async fn on_persistence_completed(&mut self, completion:
PersistenceCompletion) {
+ if !self
+ .persistence
+ .as_ref()
+ .is_some_and(|persistence|
persistence.accepts_completion(completion))
+ {
+ return;
+ }
+ self.drive_persistence().await;
+ }
+
+ pub async fn checkpoint_persistence(&mut self, config: &PartitionsConfig) {
+ let Some(persistence) = self
+ .persistence
+ .as_ref()
+ .filter(|persistence| persistence.needs_checkpoint())
+ .cloned()
+ else {
+ return;
+ };
+ let through_op = self.consensus.commit_min();
+ if through_op <= persistence.checkpoint_op() {
+ return;
+ }
+ if let Err(error) = self.flush_committed_messages(config).await {
+ error!(%error, namespace_raw = self.namespace().inner(),
"partition checkpoint failed");
+ self.fatal = Some(FatalCommit {
+ namespace_raw: self.namespace().inner(),
+ op: through_op,
+ operation: Operation::SendMessages,
+ });
+ return;
+ }
+ let (files, directories) = self.persistence_checkpoint_files(config);
+ persistence.checkpoint_files(through_op, files, directories);
+ self.start_persistence();
+ }
+
+ fn persistence_checkpoint_files(
+ &self,
+ config: &PartitionsConfig,
+ ) -> (Vec<std::path::PathBuf>, Vec<std::path::PathBuf>) {
+ let namespace = self.namespace();
+ let Some(persistence) = &self.persistence else {
+ return (Vec::new(), Vec::new());
+ };
+ let (segments, offsets) = persistence.take_dirty_files();
+ let mut paths = Vec::with_capacity(
+ segments.len() * 2
+ + offsets
+ .iter()
+ .map(std::collections::BTreeSet::len)
+ .sum::<usize>(),
+ );
+ for start_offset in segments {
+ // Retention may already have removed a dirty sealed segment.
+ if self
+ .log
+ .segments()
+ .binary_search_by_key(&start_offset, |segment|
segment.start_offset)
+ .is_err()
+ {
+ continue;
+ }
+ paths.push(config.get_messages_path(
+ namespace.stream_id(),
+ namespace.topic_id(),
+ namespace.partition_id(),
+ start_offset,
+ ));
+ paths.push(config.get_index_path(
+ namespace.stream_id(),
+ namespace.topic_id(),
+ namespace.partition_id(),
+ start_offset,
+ ));
+ }
+ for (kind, consumers) in [ConsumerKind::Consumer,
ConsumerKind::ConsumerGroup]
+ .into_iter()
+ .zip(offsets)
+ {
+ for consumer_id in consumers {
+ if self.durable_consumer_offsets.contains(kind, consumer_id)
+ && let Some(path) = self.persisted_offset_path(kind,
consumer_id)
+ {
+ paths.push(path);
+ }
+ }
+ }
+ let directories = [
Review Comment:
critical: checkpointing skips syncing `offsets/`, so power loss can remove
consumer or group directories after their WAL history is reclaimed. sync the
parent directory before publishing the checkpoint.
##########
helm/charts/iggy/templates/_helpers.tpl:
##########
@@ -290,9 +290,9 @@ Secret owns it.
{{- $server := .Values.server }}
{{- $generated := include "iggy.secretName" . }}
{{- if $server.encryption.enabled }}
-- name: IGGY_SYSTEM_ENCRYPTION_ENABLED
+- name: IGGY_ENCRYPTION_ENABLED
Review Comment:
critical: the default `apache/iggy:0.9.0-edge.6` image ignores these renamed
encryption variables, so enabling encryption leaves messages unencrypted. pin a
compatible image or render variables supported by the selected image.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -7034,6 +7429,18 @@ where
if !self.persist_superblock_if_needed().await {
return;
}
+ if self.consensus.replica_count() > 1
Review Comment:
critical: `StartView` sends `PrepareOk` before the backup WAL is durable,
allowing persisted success without durable quorum copies. route every partition
acknowledgment through the WAL durability gate.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -711,6 +720,292 @@ where
partition
}
+ pub fn set_persistence_notifier(&self, notifier: PersistenceNotifier) {
+ if let Some(persistence) = &self.persistence {
+ persistence.set_notifier(notifier);
+ }
+ }
+
+ /// # Errors
+ /// Returns an error if durable prepare history cannot be opened or
replayed.
+ pub async fn open_persistence(&mut self) -> Result<(), IggyError> {
+
self.open_persistence_with_capacity(journal::partition_journal::PARTITION_WAL_BYTES_MAX)
+ .await
+ }
+
+ /// # Errors
+ /// Returns an error if durable prepare history cannot be opened or
replayed.
+ pub async fn open_persistence_with_capacity(&mut self, capacity: u64) ->
Result<(), IggyError> {
+ if self.consensus.replica_count() == 1
+ || !(self.durability().is_persisted()
+ || self.consumer_offset_durability().is_persisted())
+ {
+ return Ok(());
+ }
+ let directory = self
+ .partition_dir
+ .as_ref()
+ .ok_or(IggyError::CannotReadFile)?;
+ let directory =
+ std::path::Path::new(directory).join(format!("prepares-{}",
self.created_revision));
+ let (persistence, prepares) = PartitionPersistence::open_with_capacity(
+ &directory,
+ self.namespace().inner(),
+ self.created_revision,
+ journal::durable_storage::DiskStorage,
+ capacity,
+ )
+ .await
+ .map_err(|error| {
+ warn!(%error, "cannot open partition prepare WAL");
+ IggyError::CannotReadFile
+ })?;
+ let (purge_generation, purge_floor) = persistence.purge_marker();
+ if purge_generation <= self.applied_purge_generation {
+ self.purge_floor_op = self.purge_floor_op.max(purge_floor);
+ }
+ let checkpoint = persistence.checkpoint_op();
+ let head = persistence.head();
+ let mut commit = checkpoint;
+ for prepare in prepares {
+ let mut owned =
server_common::iobuf::Owned::<4096>::zeroed(prepare.len());
+ owned.as_mut_slice().copy_from_slice(prepare.as_slice());
+ let message =
+ Message::<PrepareHeader>::try_from(owned).map_err(|_|
IggyError::InvalidCommand)?;
+ let header = *message.header();
+ commit = commit.max(header.commit.min(head));
+ if header.operation == Operation::SendMessages {
+ self.append_repaired_send_messages(message).await?;
+ } else {
+ self.apply_replicated_operation(message).await?;
+ }
+ }
+ if head > 0 {
+ self.consensus.sequencer().set_sequence(head);
+ if let Some(checksum) = persistence.checksum(head) {
+ self.consensus.set_last_prepare_checksum(checksum);
+ }
+ self.consensus.restore_commit_state(checkpoint, commit);
+ }
+ // Recovery can reuse completed writes whose last barrier was
interrupted.
+ // Include them once before reclaiming any recovered WAL history.
+ for segment in self.log.segments() {
+ persistence.mark_segment_dirty(segment.start_offset);
+ }
+ for kind in [ConsumerKind::Consumer, ConsumerKind::ConsumerGroup] {
+ self.durable_consumer_offsets.with_entries(kind, |entries| {
+ for consumer_id in entries.keys() {
+ persistence.mark_offset_dirty(
+ crate::state_transfer::consumer_kind_index(kind),
+ *consumer_id,
+ true,
+ );
+ }
+ });
+ }
+ self.persistence = Some(persistence);
+ Ok(())
+ }
+
+ pub async fn on_persistence_completed(&mut self, completion:
PersistenceCompletion) {
+ if !self
+ .persistence
+ .as_ref()
+ .is_some_and(|persistence|
persistence.accepts_completion(completion))
+ {
+ return;
+ }
+ self.drive_persistence().await;
+ }
+
+ pub async fn checkpoint_persistence(&mut self, config: &PartitionsConfig) {
+ let Some(persistence) = self
+ .persistence
+ .as_ref()
+ .filter(|persistence| persistence.needs_checkpoint())
+ .cloned()
+ else {
+ return;
+ };
+ let through_op = self.consensus.commit_min();
+ if through_op <= persistence.checkpoint_op() {
+ return;
+ }
+ if let Err(error) = self.flush_committed_messages(config).await {
+ error!(%error, namespace_raw = self.namespace().inner(),
"partition checkpoint failed");
+ self.fatal = Some(FatalCommit {
+ namespace_raw: self.namespace().inner(),
+ op: through_op,
+ operation: Operation::SendMessages,
+ });
+ return;
+ }
+ let (files, directories) = self.persistence_checkpoint_files(config);
+ persistence.checkpoint_files(through_op, files, directories);
+ self.start_persistence();
+ }
+
+ fn persistence_checkpoint_files(
+ &self,
+ config: &PartitionsConfig,
+ ) -> (Vec<std::path::PathBuf>, Vec<std::path::PathBuf>) {
+ let namespace = self.namespace();
+ let Some(persistence) = &self.persistence else {
+ return (Vec::new(), Vec::new());
+ };
+ let (segments, offsets) = persistence.take_dirty_files();
+ let mut paths = Vec::with_capacity(
+ segments.len() * 2
+ + offsets
+ .iter()
+ .map(std::collections::BTreeSet::len)
+ .sum::<usize>(),
+ );
+ for start_offset in segments {
+ // Retention may already have removed a dirty sealed segment.
+ if self
+ .log
+ .segments()
+ .binary_search_by_key(&start_offset, |segment|
segment.start_offset)
+ .is_err()
+ {
+ continue;
+ }
+ paths.push(config.get_messages_path(
+ namespace.stream_id(),
+ namespace.topic_id(),
+ namespace.partition_id(),
+ start_offset,
+ ));
+ paths.push(config.get_index_path(
+ namespace.stream_id(),
+ namespace.topic_id(),
+ namespace.partition_id(),
+ start_offset,
+ ));
+ }
+ for (kind, consumers) in [ConsumerKind::Consumer,
ConsumerKind::ConsumerGroup]
+ .into_iter()
+ .zip(offsets)
+ {
+ for consumer_id in consumers {
+ if self.durable_consumer_offsets.contains(kind, consumer_id)
+ && let Some(path) = self.persisted_offset_path(kind,
consumer_id)
+ {
+ paths.push(path);
+ }
+ }
+ }
+ let directories = [
+ &self.partition_dir,
+ &self.consumer_offsets_path,
+ &self.consumer_group_offsets_path,
+ ]
+ .into_iter()
+ .flatten()
+ .map(std::path::PathBuf::from)
+ .collect();
+ (
+ paths.into_iter().map(std::path::PathBuf::from).collect(),
+ directories,
+ )
+ }
+
+ pub fn take_persistence_metrics(&self) ->
Option<crate::persistence::PersistenceMetrics> {
+ self.persistence
+ .as_ref()
+ .map(|persistence| persistence.take_metrics())
+ }
+
+ fn persistence_checkpoint_pending(&self) -> bool {
+ self.persistence
+ .as_ref()
+ .is_some_and(|persistence| persistence.checkpoint_pending())
+ }
+
+ pub async fn drive_persistence(&mut self) {
+ let Some(persistence) = self.persistence.as_ref() else {
+ return;
+ };
+ if let Some(error) = persistence.failure() {
+ error!(%error, namespace_raw = self.namespace().inner(),
"partition prepare persistence failed");
+ if self.fatal.is_none() {
+ self.fatal = Some(FatalCommit {
+ namespace_raw: self.namespace().inner(),
+ op: persistence.head(),
+ operation: Operation::SendMessages,
+ });
+ }
+ return;
+ }
+ let pending = std::mem::take(&mut
*self.pending_persisted_acks.borrow_mut());
+ for header in pending.into_values() {
+ if self
+ .log
+ .journal()
+ .inner
+ .header_by_op(header.op)
Review Comment:
warning: every tick rescans headers for all pending acknowledgments and
rebuilds the waiting map, even without durable progress. process only newly
durable entries with indexed lookup and leave the rest queued.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -711,6 +720,292 @@ where
partition
}
+ pub fn set_persistence_notifier(&self, notifier: PersistenceNotifier) {
+ if let Some(persistence) = &self.persistence {
+ persistence.set_notifier(notifier);
+ }
+ }
+
+ /// # Errors
+ /// Returns an error if durable prepare history cannot be opened or
replayed.
+ pub async fn open_persistence(&mut self) -> Result<(), IggyError> {
+
self.open_persistence_with_capacity(journal::partition_journal::PARTITION_WAL_BYTES_MAX)
+ .await
+ }
+
+ /// # Errors
+ /// Returns an error if durable prepare history cannot be opened or
replayed.
+ pub async fn open_persistence_with_capacity(&mut self, capacity: u64) ->
Result<(), IggyError> {
+ if self.consensus.replica_count() == 1
+ || !(self.durability().is_persisted()
+ || self.consumer_offset_durability().is_persisted())
+ {
+ return Ok(());
+ }
+ let directory = self
+ .partition_dir
+ .as_ref()
+ .ok_or(IggyError::CannotReadFile)?;
+ let directory =
+ std::path::Path::new(directory).join(format!("prepares-{}",
self.created_revision));
+ let (persistence, prepares) = PartitionPersistence::open_with_capacity(
+ &directory,
+ self.namespace().inner(),
+ self.created_revision,
+ journal::durable_storage::DiskStorage,
+ capacity,
+ )
+ .await
+ .map_err(|error| {
+ warn!(%error, "cannot open partition prepare WAL");
+ IggyError::CannotReadFile
+ })?;
+ let (purge_generation, purge_floor) = persistence.purge_marker();
+ if purge_generation <= self.applied_purge_generation {
+ self.purge_floor_op = self.purge_floor_op.max(purge_floor);
+ }
+ let checkpoint = persistence.checkpoint_op();
+ let head = persistence.head();
+ let mut commit = checkpoint;
+ for prepare in prepares {
+ let mut owned =
server_common::iobuf::Owned::<4096>::zeroed(prepare.len());
+ owned.as_mut_slice().copy_from_slice(prepare.as_slice());
+ let message =
+ Message::<PrepareHeader>::try_from(owned).map_err(|_|
IggyError::InvalidCommand)?;
+ let header = *message.header();
+ commit = commit.max(header.commit.min(head));
+ if header.operation == Operation::SendMessages {
+ self.append_repaired_send_messages(message).await?;
+ } else {
+ self.apply_replicated_operation(message).await?;
+ }
+ }
+ if head > 0 {
+ self.consensus.sequencer().set_sequence(head);
+ if let Some(checksum) = persistence.checksum(head) {
+ self.consensus.set_last_prepare_checksum(checksum);
+ }
+ self.consensus.restore_commit_state(checkpoint, commit);
+ }
+ // Recovery can reuse completed writes whose last barrier was
interrupted.
+ // Include them once before reclaiming any recovered WAL history.
+ for segment in self.log.segments() {
+ persistence.mark_segment_dirty(segment.start_offset);
+ }
+ for kind in [ConsumerKind::Consumer, ConsumerKind::ConsumerGroup] {
+ self.durable_consumer_offsets.with_entries(kind, |entries| {
+ for consumer_id in entries.keys() {
+ persistence.mark_offset_dirty(
+ crate::state_transfer::consumer_kind_index(kind),
+ *consumer_id,
+ true,
+ );
+ }
+ });
+ }
+ self.persistence = Some(persistence);
+ Ok(())
+ }
+
+ pub async fn on_persistence_completed(&mut self, completion:
PersistenceCompletion) {
+ if !self
+ .persistence
+ .as_ref()
+ .is_some_and(|persistence|
persistence.accepts_completion(completion))
+ {
+ return;
+ }
+ self.drive_persistence().await;
+ }
+
+ pub async fn checkpoint_persistence(&mut self, config: &PartitionsConfig) {
+ let Some(persistence) = self
+ .persistence
+ .as_ref()
+ .filter(|persistence| persistence.needs_checkpoint())
+ .cloned()
+ else {
+ return;
+ };
+ let through_op = self.consensus.commit_min();
+ if through_op <= persistence.checkpoint_op() {
+ return;
+ }
+ if let Err(error) = self.flush_committed_messages(config).await {
+ error!(%error, namespace_raw = self.namespace().inner(),
"partition checkpoint failed");
+ self.fatal = Some(FatalCommit {
+ namespace_raw: self.namespace().inner(),
+ op: through_op,
+ operation: Operation::SendMessages,
+ });
+ return;
+ }
+ let (files, directories) = self.persistence_checkpoint_files(config);
+ persistence.checkpoint_files(through_op, files, directories);
Review Comment:
warning: repair can commit entries that the WAL rejected for capacity, so
the next checkpoint fails with "unknown WAL checkpoint" and fences the node.
keep the checkpoint frontier within accepted WAL history.
##########
core/journal/src/partition_journal.rs:
##########
@@ -0,0 +1,980 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#![allow(clippy::future_not_send)]
+
+use crate::durable_storage::{DiskStorage, DurableFile, DurableStorage,
OpenMode};
+use iggy_binary_protocol::{Command, PrepareHeader};
+use server_common::{
+ Message,
+ iobuf::{Frozen, Owned},
+};
+use std::collections::BTreeMap;
+use std::io;
+use std::path::{Path, PathBuf};
+use twox_hash::XxHash3_64;
+
+pub const PARTITION_WAL_BLOCK_SIZE: usize = 4096;
+pub const PARTITION_WAL_BYTES_MAX: u64 = 256 * 1024 * 1024;
+pub const PARTITION_WAL_CAPACITY_MIN: u64 = 64 * 1024 * 1024 + 4096;
+pub const PARTITION_WAL_CAPACITY_MAX: u64 = 4 * 1024 * 1024 * 1024;
+const RECORD_PREFIX: usize = 32;
+const PREPARE_BYTES_MAX: usize = 64 * 1024 * 1024;
+const STATE_MAGIC: &[u8; 8] = b"IGGYWAL1";
+
+pub trait DurableAppend {
+ /// # Errors
+ /// Returns an error if persistence fails or the prepare does not extend
the journal.
+ fn append(&mut self, prepare: Frozen<4096>) -> impl Future<Output =
io::Result<()>>;
+}
+
+pub struct PartitionPrepareJournal<S: DurableStorage = DiskStorage> {
+ directory: PathBuf,
+ file: S::File,
+ storage: S,
+ capacity: u64,
+ state: JournalState,
+ entries: BTreeMap<u64, StoredPrepare>,
+ poisoned: bool,
+ durable_head: u64,
+}
+
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+struct JournalState {
+ group: u64,
+ incarnation: u64,
+ generation: u64,
+ length: u64,
+ checkpoint: u64,
+ checkpoint_checksum: u128,
+ head: u64,
+ head_checksum: u128,
+ anchor_known: bool,
+ purge_generation: u64,
+ purge_floor: u64,
+}
+
+#[derive(Clone, Copy)]
+struct StoredPrepare {
+ position: u64,
+ length: usize,
+ checksum: u128,
+}
+
+impl PartitionPrepareJournal {
+ /// # Errors
+ /// Returns an error on I/O failure or invalid durable history.
+ pub async fn open(directory: &Path, group: u64, incarnation: u64) ->
io::Result<Self> {
+ Self::open_with_storage(directory, group, incarnation,
DiskStorage).await
+ }
+}
+
+impl<S: DurableStorage> PartitionPrepareJournal<S> {
+ /// Open and verify the durably published partition history.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn open_with_storage(
+ directory: &Path,
+ group: u64,
+ incarnation: u64,
+ storage: S,
+ ) -> io::Result<Self> {
+ Self::open_with_storage_and_capacity(
+ directory,
+ group,
+ incarnation,
+ storage,
+ PARTITION_WAL_BYTES_MAX,
+ )
+ .await
+ }
+
+ /// Open history independently of the current admission capacity.
+ ///
+ /// # Errors
+ /// Returns an error for invalid capacity or unverifiable durable history.
+ pub async fn open_with_storage_and_capacity(
+ directory: &Path,
+ group: u64,
+ incarnation: u64,
+ storage: S,
+ capacity: u64,
+ ) -> io::Result<Self> {
+ if
!(PARTITION_WAL_CAPACITY_MIN..=PARTITION_WAL_CAPACITY_MAX).contains(&capacity)
+ || !capacity.is_multiple_of(PARTITION_WAL_BLOCK_SIZE as u64)
+ {
+ return Err(invalid(
+ "partition WAL capacity is out of bounds or unaligned",
+ ));
+ }
+ storage.create_directories(directory).await?;
+ let mut ancestor = directory;
+ while let Some(parent) = ancestor.parent() {
+ if parent.as_os_str().is_empty() {
+ break;
+ }
+ storage.sync_directory(parent).await?;
+ ancestor = parent;
+ }
+ let state_path = directory.join("frontier");
+ let existing = match storage.open(&state_path, OpenMode::Read).await {
+ Ok(file) => {
+ let bytes = file.read(0, PARTITION_WAL_BLOCK_SIZE).await?;
+ let state = JournalState::decode(&bytes)?;
+ if state.group != group || state.incarnation != incarnation {
+ return Err(invalid("partition WAL identity mismatch"));
+ }
+ Some(state)
+ }
+ Err(error) if error.kind() == io::ErrorKind::NotFound => None,
+ Err(error) => return Err(error),
+ };
+ if existing.is_none() {
+ Self::validate_unpublished_history(&storage, directory).await?;
+ }
+ let state = existing.unwrap_or_else(|| JournalState {
+ group,
+ incarnation,
+ ..JournalState::default()
+ });
+ let mode = if existing.is_none() {
+ OpenMode::Create
+ } else {
+ OpenMode::ReadWrite
+ };
+ let file = storage
+ .open(&data_path(directory, state.generation), mode)
+ .await?;
+ if file.length().await? < state.length {
+ return Err(invalid("partition WAL lost acknowledged bytes"));
+ }
+ let mut journal = Self {
+ directory: directory.to_path_buf(),
+ file,
+ storage,
+ capacity,
+ state,
+ entries: BTreeMap::new(),
+ poisoned: false,
+ durable_head: state.head,
+ };
+ let mut position = 0;
+ let mut previous = state.checkpoint;
+ let mut checksum = state.checkpoint_checksum;
+ while position < state.length {
+ let (header, length, _) = journal.read_record(position).await?;
+ if header.op
+ != previous
+ .checked_add(1)
+ .ok_or_else(|| invalid("WAL op overflow"))?
+ || header.parent != checksum
+ {
+ return Err(invalid("partition WAL prepare chain is broken"));
+ }
+ journal.entries.insert(
+ header.op,
+ StoredPrepare {
+ position,
+ length,
+ checksum: header.checksum,
+ },
+ );
+ previous = header.op;
+ checksum = header.checksum;
+ position += length as u64;
+ }
+ if position != state.length || previous != state.head || checksum !=
state.head_checksum {
+ return Err(invalid("partition WAL frontier disagrees with its
data"));
+ }
+ // Only bytes covered by the durable frontier could have released an
ack.
+ journal.file.truncate(state.length).await?;
+ journal.file.sync().await?;
+ journal.storage.sync_directory(directory).await?;
+ if existing.is_none() {
+ journal.publish(state).await?;
+ }
+ Ok(journal)
+ }
+
+ #[must_use]
+ pub const fn durable_op(&self) -> u64 {
+ self.durable_head
+ }
+
+ #[must_use]
+ pub const fn head(&self) -> u64 {
+ self.state.head
+ }
+
+ #[must_use]
+ pub const fn checkpoint_op(&self) -> u64 {
+ self.state.checkpoint
+ }
+
+ #[must_use]
+ pub const fn checkpoint_checksum(&self) -> Option<u128> {
+ if self.state.anchor_known {
+ Some(self.state.checkpoint_checksum)
+ } else {
+ None
+ }
+ }
+
+ #[must_use]
+ pub const fn generation(&self) -> u64 {
+ self.state.generation
+ }
+
+ #[must_use]
+ pub const fn size_bytes(&self) -> u64 {
+ self.state.length
+ }
+
+ #[must_use]
+ pub fn contains(&self, header: &PrepareHeader) -> bool {
+ if header.op > self.durable_head {
+ return false;
+ }
+ self.entries
+ .get(&header.op)
+ .is_some_and(|entry| entry.checksum == header.checksum)
+ }
+
+ /// Read the retained prepares in operation order.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn prepares(&self) -> io::Result<Vec<Frozen<4096>>> {
+ let mut prepares = Vec::with_capacity(self.entries.len());
+ for entry in self.entries.values() {
+ let (_, length, prepare) = self.read_record(entry.position).await?;
+ if length != entry.length {
+ return Err(invalid("partition WAL index length mismatch"));
+ }
+ prepares.push(prepare);
+ }
+ Ok(prepares)
+ }
+
+ /// Durably replace the uncommitted suffix.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn truncate_from(&mut self, from_op: u64) -> io::Result<()> {
+ self.ensure_healthy()?;
+ if from_op <= self.state.checkpoint {
+ return Err(invalid("cannot truncate checkpointed partition
operations"));
+ }
+ if from_op > self.state.head {
+ return Ok(());
+ }
+ self.rewrite(
+ self.state.checkpoint,
+ self.state.checkpoint_checksum,
+ Some(from_op),
+ )
+ .await
+ }
+
+ /// Synchronize materialized files before removing their WAL coverage.
+ ///
+ /// # Errors
+ /// Returns an error if any file or directory barrier fails.
+ pub async fn checkpoint_files(
+ &mut self,
+ through_op: u64,
+ files: &[std::path::PathBuf],
+ directories: &[std::path::PathBuf],
+ ) -> io::Result<()> {
+ for path in files {
+ self.storage
+ .open(path, OpenMode::Read)
+ .await?
+ .sync()
+ .await?;
+ }
+ for path in directories {
+ self.storage.sync_directory(path).await?;
+ }
+ self.checkpoint(through_op).await
+ }
+
+ /// The caller has durably materialized every operation through this point.
+ /// Reclaim a prefix already materialized durably by the caller.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn checkpoint(&mut self, through_op: u64) -> io::Result<()> {
+ self.ensure_healthy()?;
+ if through_op <= self.state.checkpoint {
+ return Ok(());
+ }
+ let checksum = self
+ .entries
+ .get(&through_op)
+ .ok_or_else(|| invalid("unknown WAL checkpoint"))?
+ .checksum;
+ self.state.anchor_known = true;
+ self.rewrite(through_op, checksum, None).await
+ }
+
+ /// Install an already durable replacement state, such as a completed
transfer.
+ /// Replace the journal with an already durable state-transfer checkpoint.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn reset(&mut self, op: u64, checksum: Option<u128>) ->
io::Result<()> {
+ self.ensure_healthy()?;
+ self.state.anchor_known = checksum.is_some();
+ self.rewrite(op, checksum.unwrap_or(0), Some(0)).await
+ }
+
+ #[must_use]
+ pub const fn purge_marker(&self) -> (u64, u64) {
+ (self.state.purge_generation, self.state.purge_floor)
+ }
+
+ /// # Errors
+ /// Returns an error if the purge marker cannot be published durably.
+ pub async fn mark_purge(&mut self, generation: u64, floor: u64) ->
io::Result<()> {
+ self.ensure_healthy()?;
+ if generation <= self.state.purge_generation {
+ return Ok(());
+ }
+ if floor > self.state.head {
+ return Err(invalid("purge exceeds WAL history"));
+ }
+ self.poisoned = true;
+ self.file.sync().await?;
+ let state = JournalState {
+ purge_generation: generation,
+ purge_floor: floor,
+ ..self.state
+ };
+ self.publish(state).await?;
+ self.state = state;
+ self.durable_head = state.head;
+ self.poisoned = false;
+ Ok(())
+ }
+
+ /// Make every buffered predecessor recoverable with one frontier
publication.
+ ///
+ /// # Errors
+ /// Returns an error unless the buffered prefix and its frontier are
durable.
+ pub async fn sync(&mut self) -> io::Result<()> {
+ self.ensure_healthy()?;
+ if self.durable_head == self.state.head {
+ return Ok(());
+ }
+ self.poisoned = true;
+ self.file.sync().await?;
+ self.publish(self.state).await?;
+ self.durable_head = self.state.head;
+ self.poisoned = false;
+ Ok(())
+ }
+
+ /// Append a predecessor without releasing a durable acknowledgment.
+ ///
+ /// # Errors
+ /// Returns an error on write failure, capacity exhaustion, or a history
conflict.
+ pub async fn append_buffered(&mut self, prepare: Frozen<4096>) ->
io::Result<()> {
+ self.ensure_healthy()?;
+ let header = bytemuck::checked::try_from_bytes::<PrepareHeader>(
+ prepare
+ .as_slice()
+ .get(..size_of::<PrepareHeader>())
+ .ok_or_else(|| invalid("short WAL prepare"))?,
+ )
+ .map_err(|_| invalid("invalid WAL prepare alignment"))?;
+ if self
+ .entries
+ .get(&header.op)
+ .is_some_and(|entry| entry.checksum == header.checksum)
+ {
+ return Ok(());
+ }
+ if header.op
+ != self
+ .state
+ .head
+ .checked_add(1)
+ .ok_or_else(|| invalid("WAL op exhausted"))?
+ || (self.state.anchor_known && header.parent !=
self.state.head_checksum)
+ || header.group != self.state.group
+ {
+ return Err(invalid("partition WAL append does not extend its
history"));
+ }
+ let encoded = encode_record(prepare.as_slice(),
self.state.generation)?;
+ let length = encoded.len();
+ if self.state.length + length as u64 > self.capacity {
+ return Err(io::Error::new(
+ io::ErrorKind::WouldBlock,
+ "partition WAL requires checkpoint",
+ ));
+ }
+ self.poisoned = true;
+ self.file.write(self.state.length, encoded).await?;
+ let state = JournalState {
+ length: self.state.length + length as u64,
+ head: header.op,
+ head_checksum: header.checksum,
+ checkpoint_checksum: if self.state.anchor_known {
+ self.state.checkpoint_checksum
+ } else {
+ header.parent
+ },
+ anchor_known: true,
+ ..self.state
+ };
+ self.entries.insert(
+ header.op,
+ StoredPrepare {
+ position: self.state.length,
+ length,
+ checksum: header.checksum,
+ },
+ );
+ self.state = state;
+ self.poisoned = false;
+ Ok(())
+ }
+
+ async fn validate_unpublished_history(storage: &S, directory: &Path) ->
io::Result<()> {
+ for entry in storage.entries(directory)? {
+ if entry.directory || entry.name == "frontier.tmp" {
+ continue;
+ }
+ // An interrupted first open can leave only its empty
generation-zero file.
+ // Any other history without its frontier may contain acknowledged
data.
+ if entry.name != "prepares-0.wal"
+ || storage
+ .open(&directory.join(&entry.name), OpenMode::Read)
+ .await?
+ .length()
+ .await?
+ != 0
+ {
+ return Err(invalid(
+ "partition WAL history exists without its durable
frontier",
+ ));
+ }
+ }
+ Ok(())
+ }
+
+ fn ensure_healthy(&self) -> io::Result<()> {
+ if self.poisoned {
+ Err(invalid(
+ "partition WAL requires recovery after failed mutation",
+ ))
+ } else {
+ Ok(())
+ }
+ }
+
+ async fn read_record(&self, position: u64) -> io::Result<(PrepareHeader,
usize, Frozen<4096>)> {
+ let prefix = self.file.read(position, RECORD_PREFIX).await?;
+ let frame_length = u32::from_le_bytes(
+ prefix[..4]
+ .try_into()
+ .map_err(|_| invalid("invalid WAL prefix"))?,
+ ) as usize;
+ let length = record_length(frame_length)?;
+ if position
+ .checked_add(length as u64)
+ .is_none_or(|end| end > self.state.length)
+ {
+ return Err(invalid("partition WAL record crosses durable
frontier"));
+ }
+ let mut bytes = self.file.read(position, length).await?;
+ let stored_hash = u64::from_le_bytes(
+ bytes[16..24]
+ .try_into()
+ .map_err(|_| invalid("invalid WAL checksum"))?,
+ );
+ bytes[16..24].fill(0);
+ if XxHash3_64::oneshot(&bytes) != stored_hash {
+ return Err(invalid("partition WAL record checksum mismatch"));
+ }
+ let generation = u64::from_le_bytes(
+ bytes[8..16]
+ .try_into()
+ .map_err(|_| invalid("invalid WAL generation"))?,
+ );
+ if generation != self.state.generation {
+ return Err(invalid("partition WAL stale record generation"));
+ }
+ let mut owned = Owned::<4096>::zeroed(frame_length);
+ owned
+ .as_mut_slice()
+ .copy_from_slice(&bytes[RECORD_PREFIX..RECORD_PREFIX +
frame_length]);
+ let message = Message::<PrepareHeader>::try_from(owned)
+ .map_err(|_| invalid("invalid partition WAL prepare"))?;
+ let header = *message.header();
+ if header.command != Command::Prepare
+ || header.group != self.state.group
+ || header.size as usize != frame_length
+ {
+ return Err(invalid("partition WAL prepare identity mismatch"));
+ }
+ Ok((header, length, message.into_frozen()))
+ }
+
+ async fn publish(&self, state: JournalState) -> io::Result<()> {
+ let temporary = self.directory.join("frontier.tmp");
+ let mut file = self.storage.open(&temporary, OpenMode::Create).await?;
+ file.write(0, state.encode()).await?;
+ file.sync().await?;
+ self.storage
+ .rename(&temporary, &self.directory.join("frontier"))
+ .await?;
+ self.storage.sync_directory(&self.directory).await
+ }
+
+ async fn rewrite(
+ &mut self,
+ checkpoint: u64,
+ checksum: u128,
+ truncate: Option<u64>,
+ ) -> io::Result<()> {
+ // A failed publication can leave a newer frontier visible on disk.
+ // Poison until reopen instead of overwriting that possibly durable
state.
+ self.poisoned = true;
+ let generation = self
+ .state
+ .generation
+ .checked_add(1)
+ .ok_or_else(|| invalid("WAL generation exhausted"))?;
+ let mut file = self
+ .storage
+ .open(&data_path(&self.directory, generation), OpenMode::Create)
+ .await?;
+ let mut entries = BTreeMap::new();
+ let mut state = JournalState {
+ generation,
+ checkpoint,
+ checkpoint_checksum: checksum,
+ head: checkpoint,
+ head_checksum: checksum,
+ length: 0,
+ ..self.state
+ };
+ for (&op, entry) in &self.entries {
+ if op <= checkpoint || truncate.is_some_and(|from| op >= from) {
Review Comment:
critical: if all replicas checkpoint the same head and restart, the
discarded prepare leaves elections without the required commit-point header and
body. retain and recover that prepare across checkpoints.
##########
core/shard/src/lib.rs:
##########
@@ -7103,6 +7134,25 @@ where
// spreads over every group instead of replaying the same prefix.
rotate_sweep_to_cursor(namespace_scratch,
self.partition_walk_cursor.get());
+ let mut persistence_metrics =
partitions::PersistenceMetrics::default();
+ for namespace in namespace_scratch.iter() {
+ if let Some(partition) = partitions.get_mut_by_ns(namespace) {
+ partition.drive_persistence().await;
+ partition.checkpoint_persistence(partitions.config()).await;
Review Comment:
warning: this pass serially awaits every eligible partition's checkpoint
flush before consensus ticks, delaying unrelated heartbeats and requests during
checkpoint bursts. bound checkpoint work using the existing rotating partition
budget.
##########
core/journal/src/partition_journal.rs:
##########
@@ -0,0 +1,980 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#![allow(clippy::future_not_send)]
+
+use crate::durable_storage::{DiskStorage, DurableFile, DurableStorage,
OpenMode};
+use iggy_binary_protocol::{Command, PrepareHeader};
+use server_common::{
+ Message,
+ iobuf::{Frozen, Owned},
+};
+use std::collections::BTreeMap;
+use std::io;
+use std::path::{Path, PathBuf};
+use twox_hash::XxHash3_64;
+
+pub const PARTITION_WAL_BLOCK_SIZE: usize = 4096;
+pub const PARTITION_WAL_BYTES_MAX: u64 = 256 * 1024 * 1024;
+pub const PARTITION_WAL_CAPACITY_MIN: u64 = 64 * 1024 * 1024 + 4096;
+pub const PARTITION_WAL_CAPACITY_MAX: u64 = 4 * 1024 * 1024 * 1024;
+const RECORD_PREFIX: usize = 32;
+const PREPARE_BYTES_MAX: usize = 64 * 1024 * 1024;
+const STATE_MAGIC: &[u8; 8] = b"IGGYWAL1";
+
+pub trait DurableAppend {
+ /// # Errors
+ /// Returns an error if persistence fails or the prepare does not extend
the journal.
+ fn append(&mut self, prepare: Frozen<4096>) -> impl Future<Output =
io::Result<()>>;
+}
+
+pub struct PartitionPrepareJournal<S: DurableStorage = DiskStorage> {
+ directory: PathBuf,
+ file: S::File,
+ storage: S,
+ capacity: u64,
+ state: JournalState,
+ entries: BTreeMap<u64, StoredPrepare>,
+ poisoned: bool,
+ durable_head: u64,
+}
+
+#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
+struct JournalState {
+ group: u64,
+ incarnation: u64,
+ generation: u64,
+ length: u64,
+ checkpoint: u64,
+ checkpoint_checksum: u128,
+ head: u64,
+ head_checksum: u128,
+ anchor_known: bool,
+ purge_generation: u64,
+ purge_floor: u64,
+}
+
+#[derive(Clone, Copy)]
+struct StoredPrepare {
+ position: u64,
+ length: usize,
+ checksum: u128,
+}
+
+impl PartitionPrepareJournal {
+ /// # Errors
+ /// Returns an error on I/O failure or invalid durable history.
+ pub async fn open(directory: &Path, group: u64, incarnation: u64) ->
io::Result<Self> {
+ Self::open_with_storage(directory, group, incarnation,
DiskStorage).await
+ }
+}
+
+impl<S: DurableStorage> PartitionPrepareJournal<S> {
+ /// Open and verify the durably published partition history.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn open_with_storage(
+ directory: &Path,
+ group: u64,
+ incarnation: u64,
+ storage: S,
+ ) -> io::Result<Self> {
+ Self::open_with_storage_and_capacity(
+ directory,
+ group,
+ incarnation,
+ storage,
+ PARTITION_WAL_BYTES_MAX,
+ )
+ .await
+ }
+
+ /// Open history independently of the current admission capacity.
+ ///
+ /// # Errors
+ /// Returns an error for invalid capacity or unverifiable durable history.
+ pub async fn open_with_storage_and_capacity(
+ directory: &Path,
+ group: u64,
+ incarnation: u64,
+ storage: S,
+ capacity: u64,
+ ) -> io::Result<Self> {
+ if
!(PARTITION_WAL_CAPACITY_MIN..=PARTITION_WAL_CAPACITY_MAX).contains(&capacity)
+ || !capacity.is_multiple_of(PARTITION_WAL_BLOCK_SIZE as u64)
+ {
+ return Err(invalid(
+ "partition WAL capacity is out of bounds or unaligned",
+ ));
+ }
+ storage.create_directories(directory).await?;
+ let mut ancestor = directory;
+ while let Some(parent) = ancestor.parent() {
+ if parent.as_os_str().is_empty() {
+ break;
+ }
+ storage.sync_directory(parent).await?;
+ ancestor = parent;
+ }
+ let state_path = directory.join("frontier");
+ let existing = match storage.open(&state_path, OpenMode::Read).await {
+ Ok(file) => {
+ let bytes = file.read(0, PARTITION_WAL_BLOCK_SIZE).await?;
+ let state = JournalState::decode(&bytes)?;
+ if state.group != group || state.incarnation != incarnation {
+ return Err(invalid("partition WAL identity mismatch"));
+ }
+ Some(state)
+ }
+ Err(error) if error.kind() == io::ErrorKind::NotFound => None,
+ Err(error) => return Err(error),
+ };
+ if existing.is_none() {
+ Self::validate_unpublished_history(&storage, directory).await?;
+ }
+ let state = existing.unwrap_or_else(|| JournalState {
+ group,
+ incarnation,
+ ..JournalState::default()
+ });
+ let mode = if existing.is_none() {
+ OpenMode::Create
+ } else {
+ OpenMode::ReadWrite
+ };
+ let file = storage
+ .open(&data_path(directory, state.generation), mode)
+ .await?;
+ if file.length().await? < state.length {
+ return Err(invalid("partition WAL lost acknowledged bytes"));
+ }
+ let mut journal = Self {
+ directory: directory.to_path_buf(),
+ file,
+ storage,
+ capacity,
+ state,
+ entries: BTreeMap::new(),
+ poisoned: false,
+ durable_head: state.head,
+ };
+ let mut position = 0;
+ let mut previous = state.checkpoint;
+ let mut checksum = state.checkpoint_checksum;
+ while position < state.length {
+ let (header, length, _) = journal.read_record(position).await?;
+ if header.op
+ != previous
+ .checked_add(1)
+ .ok_or_else(|| invalid("WAL op overflow"))?
+ || header.parent != checksum
+ {
+ return Err(invalid("partition WAL prepare chain is broken"));
+ }
+ journal.entries.insert(
+ header.op,
+ StoredPrepare {
+ position,
+ length,
+ checksum: header.checksum,
+ },
+ );
+ previous = header.op;
+ checksum = header.checksum;
+ position += length as u64;
+ }
+ if position != state.length || previous != state.head || checksum !=
state.head_checksum {
+ return Err(invalid("partition WAL frontier disagrees with its
data"));
+ }
+ // Only bytes covered by the durable frontier could have released an
ack.
+ journal.file.truncate(state.length).await?;
+ journal.file.sync().await?;
+ journal.storage.sync_directory(directory).await?;
+ if existing.is_none() {
+ journal.publish(state).await?;
+ }
+ Ok(journal)
+ }
+
+ #[must_use]
+ pub const fn durable_op(&self) -> u64 {
+ self.durable_head
+ }
+
+ #[must_use]
+ pub const fn head(&self) -> u64 {
+ self.state.head
+ }
+
+ #[must_use]
+ pub const fn checkpoint_op(&self) -> u64 {
+ self.state.checkpoint
+ }
+
+ #[must_use]
+ pub const fn checkpoint_checksum(&self) -> Option<u128> {
+ if self.state.anchor_known {
+ Some(self.state.checkpoint_checksum)
+ } else {
+ None
+ }
+ }
+
+ #[must_use]
+ pub const fn generation(&self) -> u64 {
+ self.state.generation
+ }
+
+ #[must_use]
+ pub const fn size_bytes(&self) -> u64 {
+ self.state.length
+ }
+
+ #[must_use]
+ pub fn contains(&self, header: &PrepareHeader) -> bool {
+ if header.op > self.durable_head {
+ return false;
+ }
+ self.entries
+ .get(&header.op)
+ .is_some_and(|entry| entry.checksum == header.checksum)
+ }
+
+ /// Read the retained prepares in operation order.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn prepares(&self) -> io::Result<Vec<Frozen<4096>>> {
+ let mut prepares = Vec::with_capacity(self.entries.len());
+ for entry in self.entries.values() {
+ let (_, length, prepare) = self.read_record(entry.position).await?;
+ if length != entry.length {
+ return Err(invalid("partition WAL index length mismatch"));
+ }
+ prepares.push(prepare);
+ }
+ Ok(prepares)
+ }
+
+ /// Durably replace the uncommitted suffix.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn truncate_from(&mut self, from_op: u64) -> io::Result<()> {
+ self.ensure_healthy()?;
+ if from_op <= self.state.checkpoint {
+ return Err(invalid("cannot truncate checkpointed partition
operations"));
+ }
+ if from_op > self.state.head {
+ return Ok(());
+ }
+ self.rewrite(
+ self.state.checkpoint,
+ self.state.checkpoint_checksum,
+ Some(from_op),
+ )
+ .await
+ }
+
+ /// Synchronize materialized files before removing their WAL coverage.
+ ///
+ /// # Errors
+ /// Returns an error if any file or directory barrier fails.
+ pub async fn checkpoint_files(
+ &mut self,
+ through_op: u64,
+ files: &[std::path::PathBuf],
+ directories: &[std::path::PathBuf],
+ ) -> io::Result<()> {
+ for path in files {
+ self.storage
+ .open(path, OpenMode::Read)
+ .await?
+ .sync()
+ .await?;
+ }
+ for path in directories {
+ self.storage.sync_directory(path).await?;
+ }
+ self.checkpoint(through_op).await
+ }
+
+ /// The caller has durably materialized every operation through this point.
+ /// Reclaim a prefix already materialized durably by the caller.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn checkpoint(&mut self, through_op: u64) -> io::Result<()> {
+ self.ensure_healthy()?;
+ if through_op <= self.state.checkpoint {
+ return Ok(());
+ }
+ let checksum = self
+ .entries
+ .get(&through_op)
+ .ok_or_else(|| invalid("unknown WAL checkpoint"))?
+ .checksum;
+ self.state.anchor_known = true;
+ self.rewrite(through_op, checksum, None).await
+ }
+
+ /// Install an already durable replacement state, such as a completed
transfer.
+ /// Replace the journal with an already durable state-transfer checkpoint.
+ ///
+ /// # Errors
+ /// Returns an error on I/O failure, invalid history, or a poisoned
journal.
+ pub async fn reset(&mut self, op: u64, checksum: Option<u128>) ->
io::Result<()> {
+ self.ensure_healthy()?;
+ self.state.anchor_known = checksum.is_some();
+ self.rewrite(op, checksum.unwrap_or(0), Some(0)).await
+ }
+
+ #[must_use]
+ pub const fn purge_marker(&self) -> (u64, u64) {
+ (self.state.purge_generation, self.state.purge_floor)
+ }
+
+ /// # Errors
+ /// Returns an error if the purge marker cannot be published durably.
+ pub async fn mark_purge(&mut self, generation: u64, floor: u64) ->
io::Result<()> {
+ self.ensure_healthy()?;
+ if generation <= self.state.purge_generation {
+ return Ok(());
+ }
+ if floor > self.state.head {
+ return Err(invalid("purge exceeds WAL history"));
+ }
+ self.poisoned = true;
+ self.file.sync().await?;
+ let state = JournalState {
+ purge_generation: generation,
+ purge_floor: floor,
+ ..self.state
+ };
+ self.publish(state).await?;
+ self.state = state;
+ self.durable_head = state.head;
+ self.poisoned = false;
+ Ok(())
+ }
+
+ /// Make every buffered predecessor recoverable with one frontier
publication.
+ ///
+ /// # Errors
+ /// Returns an error unless the buffered prefix and its frontier are
durable.
+ pub async fn sync(&mut self) -> io::Result<()> {
+ self.ensure_healthy()?;
+ if self.durable_head == self.state.head {
+ return Ok(());
+ }
+ self.poisoned = true;
+ self.file.sync().await?;
+ self.publish(self.state).await?;
+ self.durable_head = self.state.head;
+ self.poisoned = false;
+ Ok(())
+ }
+
+ /// Append a predecessor without releasing a durable acknowledgment.
+ ///
+ /// # Errors
+ /// Returns an error on write failure, capacity exhaustion, or a history
conflict.
+ pub async fn append_buffered(&mut self, prepare: Frozen<4096>) ->
io::Result<()> {
+ self.ensure_healthy()?;
+ let header = bytemuck::checked::try_from_bytes::<PrepareHeader>(
+ prepare
+ .as_slice()
+ .get(..size_of::<PrepareHeader>())
+ .ok_or_else(|| invalid("short WAL prepare"))?,
+ )
+ .map_err(|_| invalid("invalid WAL prepare alignment"))?;
+ if self
+ .entries
+ .get(&header.op)
+ .is_some_and(|entry| entry.checksum == header.checksum)
+ {
+ return Ok(());
+ }
+ if header.op
+ != self
+ .state
+ .head
+ .checked_add(1)
+ .ok_or_else(|| invalid("WAL op exhausted"))?
+ || (self.state.anchor_known && header.parent !=
self.state.head_checksum)
+ || header.group != self.state.group
+ {
+ return Err(invalid("partition WAL append does not extend its
history"));
+ }
+ let encoded = encode_record(prepare.as_slice(),
self.state.generation)?;
+ let length = encoded.len();
+ if self.state.length + length as u64 > self.capacity {
+ return Err(io::Error::new(
+ io::ErrorKind::WouldBlock,
+ "partition WAL requires checkpoint",
+ ));
+ }
+ self.poisoned = true;
+ self.file.write(self.state.length, encoded).await?;
+ let state = JournalState {
+ length: self.state.length + length as u64,
+ head: header.op,
+ head_checksum: header.checksum,
+ checkpoint_checksum: if self.state.anchor_known {
+ self.state.checkpoint_checksum
+ } else {
+ header.parent
+ },
+ anchor_known: true,
+ ..self.state
+ };
+ self.entries.insert(
+ header.op,
+ StoredPrepare {
+ position: self.state.length,
+ length,
+ checksum: header.checksum,
+ },
+ );
+ self.state = state;
+ self.poisoned = false;
+ Ok(())
+ }
+
+ async fn validate_unpublished_history(storage: &S, directory: &Path) ->
io::Result<()> {
+ for entry in storage.entries(directory)? {
+ if entry.directory || entry.name == "frontier.tmp" {
+ continue;
+ }
+ // An interrupted first open can leave only its empty
generation-zero file.
+ // Any other history without its frontier may contain acknowledged
data.
+ if entry.name != "prepares-0.wal"
+ || storage
+ .open(&directory.join(&entry.name), OpenMode::Read)
+ .await?
+ .length()
+ .await?
+ != 0
+ {
+ return Err(invalid(
+ "partition WAL history exists without its durable
frontier",
+ ));
+ }
+ }
+ Ok(())
+ }
+
+ fn ensure_healthy(&self) -> io::Result<()> {
+ if self.poisoned {
+ Err(invalid(
+ "partition WAL requires recovery after failed mutation",
+ ))
+ } else {
+ Ok(())
+ }
+ }
+
+ async fn read_record(&self, position: u64) -> io::Result<(PrepareHeader,
usize, Frozen<4096>)> {
+ let prefix = self.file.read(position, RECORD_PREFIX).await?;
+ let frame_length = u32::from_le_bytes(
+ prefix[..4]
+ .try_into()
+ .map_err(|_| invalid("invalid WAL prefix"))?,
+ ) as usize;
+ let length = record_length(frame_length)?;
+ if position
+ .checked_add(length as u64)
+ .is_none_or(|end| end > self.state.length)
+ {
+ return Err(invalid("partition WAL record crosses durable
frontier"));
+ }
+ let mut bytes = self.file.read(position, length).await?;
+ let stored_hash = u64::from_le_bytes(
+ bytes[16..24]
+ .try_into()
+ .map_err(|_| invalid("invalid WAL checksum"))?,
+ );
+ bytes[16..24].fill(0);
+ if XxHash3_64::oneshot(&bytes) != stored_hash {
+ return Err(invalid("partition WAL record checksum mismatch"));
+ }
+ let generation = u64::from_le_bytes(
+ bytes[8..16]
+ .try_into()
+ .map_err(|_| invalid("invalid WAL generation"))?,
+ );
+ if generation != self.state.generation {
+ return Err(invalid("partition WAL stale record generation"));
+ }
+ let mut owned = Owned::<4096>::zeroed(frame_length);
+ owned
+ .as_mut_slice()
+ .copy_from_slice(&bytes[RECORD_PREFIX..RECORD_PREFIX +
frame_length]);
+ let message = Message::<PrepareHeader>::try_from(owned)
+ .map_err(|_| invalid("invalid partition WAL prepare"))?;
+ let header = *message.header();
+ if header.command != Command::Prepare
+ || header.group != self.state.group
+ || header.size as usize != frame_length
+ {
+ return Err(invalid("partition WAL prepare identity mismatch"));
+ }
+ Ok((header, length, message.into_frozen()))
+ }
+
+ async fn publish(&self, state: JournalState) -> io::Result<()> {
+ let temporary = self.directory.join("frontier.tmp");
+ let mut file = self.storage.open(&temporary, OpenMode::Create).await?;
+ file.write(0, state.encode()).await?;
+ file.sync().await?;
+ self.storage
+ .rename(&temporary, &self.directory.join("frontier"))
+ .await?;
+ self.storage.sync_directory(&self.directory).await
+ }
+
+ async fn rewrite(
+ &mut self,
+ checkpoint: u64,
+ checksum: u128,
+ truncate: Option<u64>,
+ ) -> io::Result<()> {
+ // A failed publication can leave a newer frontier visible on disk.
+ // Poison until reopen instead of overwriting that possibly durable
state.
+ self.poisoned = true;
+ let generation = self
+ .state
+ .generation
+ .checked_add(1)
+ .ok_or_else(|| invalid("WAL generation exhausted"))?;
+ let mut file = self
+ .storage
+ .open(&data_path(&self.directory, generation), OpenMode::Create)
+ .await?;
+ let mut entries = BTreeMap::new();
+ let mut state = JournalState {
+ generation,
+ checkpoint,
+ checkpoint_checksum: checksum,
+ head: checkpoint,
+ head_checksum: checksum,
+ length: 0,
+ ..self.state
+ };
+ for (&op, entry) in &self.entries {
+ if op <= checkpoint || truncate.is_some_and(|from| op >= from) {
+ continue;
+ }
+ let prepare = self.read_record(entry.position).await?.2;
+ let encoded = encode_record(prepare.as_slice(), generation)?;
+ let length = encoded.len();
+ file.write(state.length, encoded).await?;
+ entries.insert(
+ op,
+ StoredPrepare {
+ position: state.length,
+ length,
+ checksum: entry.checksum,
+ },
+ );
+ state.length += length as u64;
+ state.head = op;
+ state.head_checksum = entry.checksum;
+ }
+ file.sync().await?;
+ self.storage.sync_directory(&self.directory).await?;
+ self.publish(state).await?;
+ let obsolete = data_path(&self.directory, self.state.generation);
+ self.file = file;
+ self.state = state;
+ self.durable_head = state.head;
+ self.entries = entries;
+ self.poisoned = false;
+ if let Err(error) = self.storage.remove_file(&obsolete).await {
Review Comment:
warning: a crash after publication or failed unlink leaves obsolete WAL
generations permanently, allowing repeated failures to fill the disk. reclaim
obsolete generations on reopen and retry failed cleanup.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4143,6 +4476,9 @@ where
}
};
+ if !self.submit_prepare_persistence(&frozen_for_forward,
header.operation) {
Review Comment:
warning: WAL backpressure leaves the memory journal ahead of the sequencer;
a successful retry acknowledges without advancing it, so the next operation
fails the gap check. keep both in sync across failed admission and retries.
##########
core/integration/tests/data_integrity/storage_compat.rs:
##########
@@ -1012,7 +1012,7 @@ fn data_topic_options() -> TopicCreateOptions {
MAX_TOPIC_SIZE_BYTES,
))),
segment_size: Some(IggyByteSize::from(SEGMENT_SIZE_BYTES)),
- enforce_fsync: Some(true),
+ durability: iggy_common::Durability::Persisted,
Review Comment:
warning: the fixture sends unsupported options and config to the baseline,
then reuses its incompatible config for the replacement; upgrade assertions
never run. seed with the baseline schema and translate config when swapping
binaries.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4631,7 +4983,11 @@ where
if std::mem::take(&mut self.injected_commit_failure) {
return Err(IggyError::CannotSaveMessagesToSegment);
}
- self.commit_messages_inner(config, false).await
+ self.commit_messages_inner(
+ config,
+ self.consensus.replica_count() == 1 &&
self.durability().is_persisted(),
Review Comment:
critical: this pre-existing gap lets singleton persisted replies precede
durable publication of rotated segment names, so power loss can erase
acknowledged messages. sync the partition directory before replying.
##########
core/bench/src/args/common.rs:
##########
@@ -101,15 +101,15 @@ pub struct IggyBenchArgs {
#[arg(long, default_value_t = false)]
pub reuse_streams: bool,
- /// Fsync each journal flush on the benchmark topic. Flush timing stays
- /// governed by `--messages-required-to-save` (server default: 1024), so
- /// acks are durability-gated only with `--messages-required-to-save 1`.
- /// Topic option at creation, so it has no effect with `--reuse-streams`.
- #[arg(long, default_value_t = false)]
- pub enforce_fsync: bool,
+ /// Message completion policy for newly created benchmark topics.
Review Comment:
warning: the performance suite helper still emits the removed fsync flag, so
those benchmark variants exit during argument parsing. update
`construct_bench_command` to set `durability` to `persisted`.
##########
core/common/src/types/options/mod.rs:
##########
@@ -800,8 +823,11 @@ impl TopicCreateOptions {
let size = parse_byte_size(entry, key)?;
parsed.segment_size = (size !=
0).then_some(IggyByteSize::from(size));
}
- topic_option_keys::ENFORCE_FSYNC => {
- parsed.enforce_fsync = Some(parse_bool(entry, key)?);
+ topic_option_keys::DURABILITY => {
Review Comment:
critical: existing `enforce_fsync=true` topic metadata silently becomes
`Durability::Replicated` after config migration, disabling fsync on reopened
writers. explicitly migrate the stored policy or reject it before opening
writers.
--
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]