hubcio commented on code in PR #4092: URL: https://github.com/apache/iggy/pull/4092#discussion_r3983315497
########## 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 { Review Comment: the suggested exemption can't cover the next-op case. ########## core/journal/src/partition_journal.rs: ########## @@ -0,0 +1,2517 @@ +// 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 std::collections::{BTreeMap, BTreeSet, VecDeque}; +use std::io; +use std::path::{Path, PathBuf}; + +use futures::TryStreamExt; +use iggy_binary_protocol::batch::BATCH_HEADER_SIZE; +use iggy_binary_protocol::{Command, ConsensusHeader, Operation, PrepareHeader}; +use server_common::{ + Message, + iobuf::{Frozen, Owned}, +}; +use twox_hash::XxHash3_64; + +use crate::durable_storage::{DiskStorage, DurableFile, DurableStorage, OpenMode}; + +mod segments; +pub use segments::SegmentPosition; +use segments::SegmentState; + +pub const PARTITION_WAL_BLOCK_SIZE: usize = 4096; +pub const PARTITION_WAL_BYTES_MAX: u64 = 256 * 1024 * 1024; +pub const PARTITION_WAL_CAPACITY_MIN: u64 = 2 * (64 * 1024 * 1024 + 4096); +pub const PARTITION_WAL_CAPACITY_MAX: u64 = 4 * 1024 * 1024 * 1024; +const RECORD_PREFIX: usize = 32; +pub const PREPARE_BYTES_MAX: usize = 64 * 1024 * 1024; +const STATE_MAGIC: &[u8; 8] = b"IGGYWAL1"; +const REFERENCE_STATE_MAGIC: &[u8; 8] = b"IGGYWAL2"; +const INLINE_RECORD: u32 = 0; +const SEGMENT_RECORD: u32 = 1; +const RECORD_KIND_OFFSET: usize = 4; +const SEGMENT_REFERENCE_BYTES: usize = 4 * size_of::<u64>(); +const REFERENCED_PREPARE_BYTES: usize = size_of::<PrepareHeader>() + SEGMENT_REFERENCE_BYTES; + +/// A prepare body in a segment inode retained by the WAL. +/// +/// The caller assigns a fresh generation whenever an offset-named segment is +/// replaced, and synchronizes the body through its writing handle before +/// submitting this reference. Referenced bytes must remain unchanged until the +/// WAL durably removes their operations. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct SegmentReference { + pub generation: u64, + pub start_offset: u64, + pub position: u64, + pub length: u64, +} + +pub trait DurableAppend { + /// # Errors + /// Returns an error if persistence fails or the prepare does not extend the journal. + fn append(&mut self, prepare: Frozen<4096>) -> impl Future<Output = io::Result<()>>; +} + +pub struct PartitionPrepareJournal<S: DurableStorage = DiskStorage> { + directory: PathBuf, + file: S::File, + storage: S, + capacity: u64, + state: JournalState, + entries: BTreeMap<u64, StoredPrepare>, + poisoned: bool, + durable_head: u64, + obsolete: VecDeque<PathBuf>, + cleanup_directory_dirty: bool, + recovered_prepares: Vec<Message<PrepareHeader>>, + segment_files: BTreeMap<(u64, u64), S::File>, + preallocate_segments: bool, + retained_bytes: 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, + checkpoint_prepare: bool, + certified_log_view: Option<u32>, + purge_generation: u64, + purge_floor: u64, + segment_references: bool, + segment_storage: Option<SegmentState>, +} + +#[derive(Clone, Copy)] +struct StoredPrepare { + position: u64, + length: usize, + checksum: u128, + reference: Option<SegmentReference>, + next_offset: Option<u64>, + retained_bytes: u64, +} + +impl PartitionPrepareJournal { + /// # Errors + /// Returns an error on I/O failure or invalid durable history. + pub async fn open(directory: &Path, group: u64, incarnation: u64) -> io::Result<Self> { + Self::open_with_storage(directory, group, incarnation, DiskStorage).await + } +} + +impl<S: DurableStorage> PartitionPrepareJournal<S> { + /// Open and verify the durably published partition history. + /// The caller must first durably materialize the parent directory. + /// + /// # Errors + /// Returns an error on I/O failure, invalid history, or a poisoned journal. + pub async fn open_with_storage( + directory: &Path, + group: u64, + incarnation: u64, + storage: S, + ) -> io::Result<Self> { + Self::open_with_storage_and_capacity( + directory, + group, + incarnation, + storage, + PARTITION_WAL_BYTES_MAX, + false, + ) + .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, + preallocate_segments: bool, + ) -> io::Result<Self> { + if !(PARTITION_WAL_CAPACITY_MIN..=PARTITION_WAL_CAPACITY_MAX).contains(&capacity) + || !capacity.is_multiple_of(PARTITION_WAL_BLOCK_SIZE as u64) + { + return Err(invalid( + "partition WAL capacity is out of bounds or unaligned", + )); + } + let parent = directory + .parent() + .filter(|parent| !parent.as_os_str().is_empty()) + .unwrap_or_else(|| Path::new(".")); + if !storage.exists(parent).await? { + return Err(invalid( + "partition WAL parent must already be durably materialized", + )); + } + storage.create_directories(directory).await?; + storage.sync_directory(parent).await?; + let state_path = directory.join("frontier"); + let existing = match storage.open(&state_path, OpenMode::Read).await { + Ok(file) => { + let bytes = file.read(0, PARTITION_WAL_BLOCK_SIZE).await?; + let state = JournalState::decode(&bytes)?; + if state.group != group || state.incarnation != incarnation { + return Err(invalid("partition WAL identity mismatch")); + } + Some(state) + } + Err(error) if error.kind() == io::ErrorKind::NotFound => None, + Err(error) => return Err(error), + }; + if existing.is_none() { + Self::validate_unpublished_history(&storage, directory).await?; + } + let state = existing.unwrap_or_else(|| JournalState { + group, + incarnation, + certified_log_view: Some(0), + ..JournalState::default() + }); + let mode = if existing.is_none() { + OpenMode::Create + } else { + OpenMode::ReadWrite + }; + let file = storage + .open(&data_path(directory, state.generation), mode) + .await?; + if file.length().await? < state.length { + return Err(invalid("partition WAL lost acknowledged bytes")); + } + let mut journal = Self { + directory: directory.to_path_buf(), + file, + storage, + capacity, + state, + entries: BTreeMap::new(), + poisoned: false, + durable_head: state.head, + obsolete: VecDeque::new(), + cleanup_directory_dirty: false, + recovered_prepares: Vec::new(), + segment_files: BTreeMap::new(), + preallocate_segments, + retained_bytes: 0, + }; + journal.recover_entries().await?; + // Only bytes covered by the durable frontier could have released an ack. + journal.file.truncate(state.length).await?; + journal.file.sync().await?; + journal.storage.sync_directory(directory).await?; + if existing.is_none() { + journal.publish(state).await?; + } + journal.discover_obsolete().await?; + loop { + let remaining = journal.obsolete.len(); + journal.cleanup_obsolete().await; + if journal.obsolete.is_empty() || journal.obsolete.len() == remaining { + break; + } + } + journal.recover_segment_files().await?; + journal.migrate_segment_prepares().await?; + Ok(journal) + } + + pub const fn certified_log_view(&self) -> Option<u32> { + self.state.certified_log_view + } + + /// Publish a complete canonical view only after its required head is durable. + /// + /// # Errors + /// Returns an error if the expected history is absent or persistence fails. + pub async fn certify_log_view(&mut self, view: u32, op: u64, checksum: u128) -> io::Result<()> { + self.ensure_healthy()?; + let matches = if op == self.state.checkpoint { + (op == 0 && checksum == 0) || self.state.checkpoint_checksum == checksum + } else { + self.entries + .get(&op) + .is_some_and(|entry| entry.checksum == checksum) + }; + if op > self.state.head || !matches { + return Err(invalid("view certificate does not match WAL history")); + } + self.poisoned = true; + self.file.sync().await?; + let state = JournalState { + certified_log_view: Some(view), + ..self.state + }; + self.publish(state).await?; + self.state = state; + self.durable_head = state.head; + self.poisoned = false; + Ok(()) + } + + #[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 + } + + /// Admission includes retained body bytes, even when the WAL stores references. + pub const fn retained_bytes(&self) -> u64 { + self.retained_bytes + } + + #[must_use] + pub fn contains(&self, header: &PrepareHeader) -> bool { + if header.op > self.durable_head { + return false; + } + self.entries + .get(&header.op) + .is_some_and(|entry| entry.checksum == header.checksum) + } + + pub fn take_recovered_prepares(&mut self) -> Vec<Message<PrepareHeader>> { + std::mem::take(&mut self.recovered_prepares) + } + + /// Read the retained prepares in operation order. + /// + /// # Errors + /// Returns an error on I/O failure, invalid history, or a poisoned journal. + pub async fn prepares(&self) -> io::Result<Vec<Message<PrepareHeader>>> { + let mut prepares = Vec::with_capacity(self.entries.len()); + for entry in self.entries.values() { + let (_, length, prepare, _) = self.read_record(entry.position).await?; + if length != entry.length { + return Err(invalid("partition WAL index length mismatch")); + } + prepares.push(prepare); + } + Ok(prepares) + } + + /// Durably replace the uncommitted suffix. + /// + /// # Errors + /// Returns an error on I/O failure, invalid history, or a poisoned journal. + pub async fn truncate_from(&mut self, from_op: u64) -> io::Result<()> { + self.ensure_healthy()?; + if from_op <= self.state.checkpoint { + return Err(invalid("cannot truncate checkpointed partition operations")); + } + if from_op > self.state.head { + return Ok(()); + } + self.state.certified_log_view = None; + self.rewrite( + self.state.checkpoint, + self.state.checkpoint_checksum, + Some(from_op), + None, + ) + .await?; + self.recover_segment_files().await + } + + /// Synchronize required materialized files before removing their WAL coverage. + /// Authorized deletions are excluded by the caller. A missing listed path + /// does not prove deletion was authorized and cannot permit WAL reclamation. + /// + /// # Errors + /// Returns an error if any file or directory barrier fails. + pub async fn checkpoint_files( + &mut self, + through_op: u64, + files: &[std::path::PathBuf], + directories: &[std::path::PathBuf], + ) -> io::Result<()> { + futures::stream::iter(files.iter().map(Ok::<_, io::Error>)) + .try_for_each_concurrent(16, |path| async { + self.storage.open(path, OpenMode::Read).await?.sync().await + }) + .await?; + for path in directories { + self.storage.sync_directory(path).await?; + } + self.checkpoint(through_op).await + } + + /// The caller has durably materialized every operation through this point. + /// Reclaim a prefix already materialized durably by the caller. + /// + /// # Errors + /// Returns an error on I/O failure, invalid history, or a poisoned journal. + pub async fn checkpoint(&mut self, through_op: u64) -> io::Result<()> { + self.ensure_healthy()?; + if through_op <= self.state.checkpoint { + return Ok(()); + } + let checksum = self + .entries + .get(&through_op) + .ok_or_else(|| invalid("unknown WAL checkpoint"))? + .checksum; + self.state.anchor_known = true; + self.rewrite(through_op, checksum, None, 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()?; + if self.state.segment_storage.is_some() { + return Err(invalid( + "segment reset requires the installed body boundary", + )); + } + self.state.anchor_known = checksum.is_some(); + self.state.certified_log_view = None; + self.rewrite(op, checksum.unwrap_or(0), Some(0), None).await + } + + /// Install the committed prepare together with its materialized checkpoint. + /// + /// # Errors + /// Returns an error for an invalid checkpoint prepare or a storage failure. + pub async fn reset_with_prepare(&mut self, prepare: Frozen<4096>) -> io::Result<()> { + self.ensure_healthy()?; + if self.state.segment_storage.is_some() { + return Err(invalid( + "segment reset requires the installed body boundary", + )); + } + let header = self.validate_checkpoint_prepare(&prepare)?; + self.state.anchor_known = true; + self.state.certified_log_view = None; + self.rewrite(header.op, header.checksum, Some(0), Some((prepare, None))) + .await + } + + fn validate_checkpoint_prepare(&self, prepare: &Frozen<4096>) -> io::Result<PrepareHeader> { + let header = *bytemuck::checked::try_from_bytes::<PrepareHeader>( + prepare + .as_slice() + .get(..size_of::<PrepareHeader>()) + .ok_or_else(|| invalid("short checkpoint prepare"))?, + ) + .map_err(|_| invalid("invalid checkpoint prepare"))?; + if header.command != Command::Prepare + || header.group != self.state.group + || (header.checksum != 0 && header.identity_checksum() != header.checksum) + || header.size as usize != prepare.len() + || (header.checksum_body != 0 + && header.checksum_body + != u128::from(XxHash3_64::oneshot( + &prepare.as_slice()[size_of::<PrepareHeader>()..], + ))) + { + return Err(invalid("invalid checkpoint prepare identity or checksum")); + } + Ok(header) + } + + #[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 + || (generation == self.state.purge_generation && floor <= self.state.purge_floor) + { + return Ok(()); + } + if floor > self.state.head { + return Err(invalid("purge exceeds WAL history")); + } + self.poisoned = true; + self.file.sync().await?; + let mut state = JournalState { + purge_generation: generation, + purge_floor: floor, + ..self.state + }; + if let Some(segments) = &mut state.segment_storage { + segments.reset_position(SegmentPosition::default())?; + } + self.publish(state).await?; + self.state = state; + self.durable_head = state.head; + self.segment_files.clear(); + 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.append_batch_buffered(std::slice::from_ref(&prepare)) + .await + } + + /// Append a contiguous extent, validating every operation before allocation or I/O. + /// + /// # Errors + /// Returns an error on invalid history, capacity exhaustion or failed write. + pub async fn append_batch_buffered(&mut self, prepares: &[Frozen<4096>]) -> io::Result<()> { + self.append_batch_inner(prepares, None).await + } + + /// Retain durable segment bodies and append only their headers and locations. + /// Non-message operations keep their inline payloads. The caller must first + /// synchronize every referenced body through the handle that wrote it. + /// + /// Segment names in the parent directory remain offset-based. Hard links + /// owned by this WAL keep referenced inodes alive across retention and purge. + /// + /// # Errors + /// Returns an error for invalid references, history, capacity, or storage. + pub async fn append_batch_referenced_buffered( + &mut self, + prepares: &[Frozen<4096>], + references: &[Option<SegmentReference>], + ) -> io::Result<()> { + if references.len() != prepares.len() { + return Err(invalid("partition WAL reference count mismatch")); + } + self.append_batch_inner(prepares, Some(references)).await + } + + #[allow(clippy::too_many_lines)] + async fn append_batch_inner( + &mut self, + prepares: &[Frozen<4096>], + references: Option<&[Option<SegmentReference>]>, + ) -> io::Result<()> { + self.ensure_healthy()?; + self.cleanup_obsolete().await; + self.recovered_prepares.clear(); + let mut state = self.state; + let mut retained_bytes = self.retained_bytes; + let mut records: Vec<(u64, StoredPrepare, usize)> = Vec::with_capacity(prepares.len()); + for (index, prepare) in prepares.iter().enumerate() { + 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) + || records.last().is_some_and(|(op, entry, _)| { + *op == header.op && entry.checksum == header.checksum + }) + { + continue; + } + if header.op + != state + .head + .checked_add(1) + .ok_or_else(|| invalid("WAL op exhausted"))? + || (state.anchor_known && header.parent != state.head_checksum) + || header.group != state.group + { + return Err(invalid("partition WAL append does not extend its history")); + } + let (reference, next_offset) = if let Some(segments) = &mut state.segment_storage { + if references.is_some() { + return Err(invalid( + "externally assigned reference in owned segment storage", + )); + } + segments.reserve(header, prepare.as_slice())? + } else { + (references.and_then(|references| references[index]), None) + }; + if let Some(reference) = reference { + reference.validate(header, prepare.len())?; + } + let length = + record_length(reference.map_or(prepare.len(), |_| REFERENCED_PREPARE_BYTES))?; + let body_bytes = record_length(prepare.len())? as u64; + if retained_bytes.saturating_add(body_bytes) > self.capacity { + return Err(io::Error::new( + io::ErrorKind::WouldBlock, + "partition WAL requires checkpoint", + )); + } + records.push(( + header.op, + StoredPrepare { + position: state.length, + length, + checksum: header.checksum, + reference, + next_offset, + retained_bytes: body_bytes, + }, + index, + )); + state.length += length as u64; + retained_bytes += body_bytes; + state.head = header.op; + state.head_checksum = header.checksum; + state.segment_references |= reference.is_some(); + if state + .certified_log_view + .is_some_and(|view| header.view > view) + { + state.certified_log_view = None; + } + if !state.anchor_known { + state.checkpoint_checksum = header.parent; + state.anchor_known = true; + } + } + if records.is_empty() { + return Ok(()); + } + let mut extent = Owned::zeroed( Review Comment: let's keep packing separate; it needs crash proof and benchmarks. ########## 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; Review Comment: keeping the body-byte cap; refs alone don't bound the body work. ########## core/journal/src/partition_journal.rs: ########## @@ -0,0 +1,2517 @@ +// 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 std::collections::{BTreeMap, BTreeSet, VecDeque}; +use std::io; +use std::path::{Path, PathBuf}; + +use futures::TryStreamExt; +use iggy_binary_protocol::batch::BATCH_HEADER_SIZE; +use iggy_binary_protocol::{Command, ConsensusHeader, Operation, PrepareHeader}; +use server_common::{ + Message, + iobuf::{Frozen, Owned}, +}; +use twox_hash::XxHash3_64; + +use crate::durable_storage::{DiskStorage, DurableFile, DurableStorage, OpenMode}; + +mod segments; +pub use segments::SegmentPosition; +use segments::SegmentState; + +pub const PARTITION_WAL_BLOCK_SIZE: usize = 4096; +pub const PARTITION_WAL_BYTES_MAX: u64 = 256 * 1024 * 1024; +pub const PARTITION_WAL_CAPACITY_MIN: u64 = 2 * (64 * 1024 * 1024 + 4096); +pub const PARTITION_WAL_CAPACITY_MAX: u64 = 4 * 1024 * 1024 * 1024; +const RECORD_PREFIX: usize = 32; +pub const PREPARE_BYTES_MAX: usize = 64 * 1024 * 1024; +const STATE_MAGIC: &[u8; 8] = b"IGGYWAL1"; +const REFERENCE_STATE_MAGIC: &[u8; 8] = b"IGGYWAL2"; +const INLINE_RECORD: u32 = 0; +const SEGMENT_RECORD: u32 = 1; +const RECORD_KIND_OFFSET: usize = 4; +const SEGMENT_REFERENCE_BYTES: usize = 4 * size_of::<u64>(); +const REFERENCED_PREPARE_BYTES: usize = size_of::<PrepareHeader>() + SEGMENT_REFERENCE_BYTES; + +/// A prepare body in a segment inode retained by the WAL. +/// +/// The caller assigns a fresh generation whenever an offset-named segment is +/// replaced, and synchronizes the body through its writing handle before +/// submitting this reference. Referenced bytes must remain unchanged until the +/// WAL durably removes their operations. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct SegmentReference { + pub generation: u64, + pub start_offset: u64, + pub position: u64, + pub length: u64, +} + +pub trait DurableAppend { + /// # Errors + /// Returns an error if persistence fails or the prepare does not extend the journal. + fn append(&mut self, prepare: Frozen<4096>) -> impl Future<Output = io::Result<()>>; +} + +pub struct PartitionPrepareJournal<S: DurableStorage = DiskStorage> { + directory: PathBuf, + file: S::File, + storage: S, + capacity: u64, + state: JournalState, + entries: BTreeMap<u64, StoredPrepare>, + poisoned: bool, + durable_head: u64, + obsolete: VecDeque<PathBuf>, + cleanup_directory_dirty: bool, + recovered_prepares: Vec<Message<PrepareHeader>>, + segment_files: BTreeMap<(u64, u64), S::File>, + preallocate_segments: bool, + retained_bytes: 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, + checkpoint_prepare: bool, + certified_log_view: Option<u32>, + purge_generation: u64, + purge_floor: u64, + segment_references: bool, + segment_storage: Option<SegmentState>, +} + +#[derive(Clone, Copy)] +struct StoredPrepare { + position: u64, + length: usize, + checksum: u128, + reference: Option<SegmentReference>, + next_offset: Option<u64>, + retained_bytes: u64, +} + +impl PartitionPrepareJournal { + /// # Errors + /// Returns an error on I/O failure or invalid durable history. + pub async fn open(directory: &Path, group: u64, incarnation: u64) -> io::Result<Self> { + Self::open_with_storage(directory, group, incarnation, DiskStorage).await + } +} + +impl<S: DurableStorage> PartitionPrepareJournal<S> { + /// Open and verify the durably published partition history. + /// The caller must first durably materialize the parent directory. + /// + /// # Errors + /// Returns an error on I/O failure, invalid history, or a poisoned journal. + pub async fn open_with_storage( + directory: &Path, + group: u64, + incarnation: u64, + storage: S, + ) -> io::Result<Self> { + Self::open_with_storage_and_capacity( + directory, + group, + incarnation, + storage, + PARTITION_WAL_BYTES_MAX, + false, + ) + .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, + preallocate_segments: bool, + ) -> io::Result<Self> { + if !(PARTITION_WAL_CAPACITY_MIN..=PARTITION_WAL_CAPACITY_MAX).contains(&capacity) + || !capacity.is_multiple_of(PARTITION_WAL_BLOCK_SIZE as u64) + { + return Err(invalid( + "partition WAL capacity is out of bounds or unaligned", + )); + } + let parent = directory + .parent() + .filter(|parent| !parent.as_os_str().is_empty()) + .unwrap_or_else(|| Path::new(".")); + if !storage.exists(parent).await? { + return Err(invalid( + "partition WAL parent must already be durably materialized", + )); + } + storage.create_directories(directory).await?; + storage.sync_directory(parent).await?; + let state_path = directory.join("frontier"); + let existing = match storage.open(&state_path, OpenMode::Read).await { + Ok(file) => { + let bytes = file.read(0, PARTITION_WAL_BLOCK_SIZE).await?; + let state = JournalState::decode(&bytes)?; + if state.group != group || state.incarnation != incarnation { + return Err(invalid("partition WAL identity mismatch")); + } + Some(state) + } + Err(error) if error.kind() == io::ErrorKind::NotFound => None, + Err(error) => return Err(error), + }; + if existing.is_none() { + Self::validate_unpublished_history(&storage, directory).await?; + } + let state = existing.unwrap_or_else(|| JournalState { + group, + incarnation, + certified_log_view: Some(0), + ..JournalState::default() + }); + let mode = if existing.is_none() { + OpenMode::Create + } else { + OpenMode::ReadWrite + }; + let file = storage + .open(&data_path(directory, state.generation), mode) + .await?; + if file.length().await? < state.length { + return Err(invalid("partition WAL lost acknowledged bytes")); + } + let mut journal = Self { + directory: directory.to_path_buf(), + file, + storage, + capacity, + state, + entries: BTreeMap::new(), + poisoned: false, + durable_head: state.head, + obsolete: VecDeque::new(), + cleanup_directory_dirty: false, + recovered_prepares: Vec::new(), + segment_files: BTreeMap::new(), + preallocate_segments, + retained_bytes: 0, + }; + journal.recover_entries().await?; + // Only bytes covered by the durable frontier could have released an ack. + journal.file.truncate(state.length).await?; + journal.file.sync().await?; + journal.storage.sync_directory(directory).await?; + if existing.is_none() { + journal.publish(state).await?; + } + journal.discover_obsolete().await?; + loop { + let remaining = journal.obsolete.len(); + journal.cleanup_obsolete().await; + if journal.obsolete.is_empty() || journal.obsolete.len() == remaining { + break; + } + } + journal.recover_segment_files().await?; + journal.migrate_segment_prepares().await?; + Ok(journal) + } + + pub const fn certified_log_view(&self) -> Option<u32> { + self.state.certified_log_view + } + + /// Publish a complete canonical view only after its required head is durable. + /// + /// # Errors + /// Returns an error if the expected history is absent or persistence fails. + pub async fn certify_log_view(&mut self, view: u32, op: u64, checksum: u128) -> io::Result<()> { + self.ensure_healthy()?; + let matches = if op == self.state.checkpoint { + (op == 0 && checksum == 0) || self.state.checkpoint_checksum == checksum + } else { + self.entries + .get(&op) + .is_some_and(|entry| entry.checksum == checksum) + }; + if op > self.state.head || !matches { + return Err(invalid("view certificate does not match WAL history")); + } + self.poisoned = true; + self.file.sync().await?; + let state = JournalState { + certified_log_view: Some(view), + ..self.state + }; + self.publish(state).await?; + self.state = state; + self.durable_head = state.head; + self.poisoned = false; + Ok(()) + } + + #[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 + } + + /// Admission includes retained body bytes, even when the WAL stores references. + pub const fn retained_bytes(&self) -> u64 { + self.retained_bytes + } + + #[must_use] + pub fn contains(&self, header: &PrepareHeader) -> bool { + if header.op > self.durable_head { + return false; + } + self.entries + .get(&header.op) + .is_some_and(|entry| entry.checksum == header.checksum) + } + + pub fn take_recovered_prepares(&mut self) -> Vec<Message<PrepareHeader>> { + std::mem::take(&mut self.recovered_prepares) + } + + /// Read the retained prepares in operation order. + /// + /// # Errors + /// Returns an error on I/O failure, invalid history, or a poisoned journal. + pub async fn prepares(&self) -> io::Result<Vec<Message<PrepareHeader>>> { + let mut prepares = Vec::with_capacity(self.entries.len()); + for entry in self.entries.values() { + let (_, length, prepare, _) = self.read_record(entry.position).await?; + if length != entry.length { + return Err(invalid("partition WAL index length mismatch")); + } + prepares.push(prepare); + } + Ok(prepares) + } + + /// Durably replace the uncommitted suffix. + /// + /// # Errors + /// Returns an error on I/O failure, invalid history, or a poisoned journal. + pub async fn truncate_from(&mut self, from_op: u64) -> io::Result<()> { + self.ensure_healthy()?; + if from_op <= self.state.checkpoint { + return Err(invalid("cannot truncate checkpointed partition operations")); + } + if from_op > self.state.head { + return Ok(()); + } + self.state.certified_log_view = None; + self.rewrite( + self.state.checkpoint, + self.state.checkpoint_checksum, + Some(from_op), + None, + ) + .await?; + self.recover_segment_files().await + } + + /// Synchronize required materialized files before removing their WAL coverage. + /// Authorized deletions are excluded by the caller. A missing listed path + /// does not prove deletion was authorized and cannot permit WAL reclamation. + /// + /// # Errors + /// Returns an error if any file or directory barrier fails. + pub async fn checkpoint_files( + &mut self, + through_op: u64, + files: &[std::path::PathBuf], + directories: &[std::path::PathBuf], + ) -> io::Result<()> { + futures::stream::iter(files.iter().map(Ok::<_, io::Error>)) + .try_for_each_concurrent(16, |path| async { + self.storage.open(path, OpenMode::Read).await?.sync().await + }) + .await?; + for path in directories { + self.storage.sync_directory(path).await?; + } + self.checkpoint(through_op).await + } + + /// The caller has durably materialized every operation through this point. + /// Reclaim a prefix already materialized durably by the caller. + /// + /// # Errors + /// Returns an error on I/O failure, invalid history, or a poisoned journal. + pub async fn checkpoint(&mut self, through_op: u64) -> io::Result<()> { + self.ensure_healthy()?; + if through_op <= self.state.checkpoint { + return Ok(()); + } + let checksum = self + .entries + .get(&through_op) + .ok_or_else(|| invalid("unknown WAL checkpoint"))? + .checksum; + self.state.anchor_known = true; + self.rewrite(through_op, checksum, None, 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()?; + if self.state.segment_storage.is_some() { + return Err(invalid( + "segment reset requires the installed body boundary", + )); + } + self.state.anchor_known = checksum.is_some(); + self.state.certified_log_view = None; + self.rewrite(op, checksum.unwrap_or(0), Some(0), None).await + } + + /// Install the committed prepare together with its materialized checkpoint. + /// + /// # Errors + /// Returns an error for an invalid checkpoint prepare or a storage failure. + pub async fn reset_with_prepare(&mut self, prepare: Frozen<4096>) -> io::Result<()> { + self.ensure_healthy()?; + if self.state.segment_storage.is_some() { + return Err(invalid( + "segment reset requires the installed body boundary", + )); + } + let header = self.validate_checkpoint_prepare(&prepare)?; + self.state.anchor_known = true; + self.state.certified_log_view = None; + self.rewrite(header.op, header.checksum, Some(0), Some((prepare, None))) + .await + } + + fn validate_checkpoint_prepare(&self, prepare: &Frozen<4096>) -> io::Result<PrepareHeader> { + let header = *bytemuck::checked::try_from_bytes::<PrepareHeader>( + prepare + .as_slice() + .get(..size_of::<PrepareHeader>()) + .ok_or_else(|| invalid("short checkpoint prepare"))?, + ) + .map_err(|_| invalid("invalid checkpoint prepare"))?; + if header.command != Command::Prepare + || header.group != self.state.group + || (header.checksum != 0 && header.identity_checksum() != header.checksum) + || header.size as usize != prepare.len() + || (header.checksum_body != 0 + && header.checksum_body + != u128::from(XxHash3_64::oneshot( + &prepare.as_slice()[size_of::<PrepareHeader>()..], + ))) + { + return Err(invalid("invalid checkpoint prepare identity or checksum")); + } + Ok(header) + } + + #[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 + || (generation == self.state.purge_generation && floor <= self.state.purge_floor) + { + return Ok(()); + } + if floor > self.state.head { + return Err(invalid("purge exceeds WAL history")); + } + self.poisoned = true; + self.file.sync().await?; + let mut state = JournalState { + purge_generation: generation, + purge_floor: floor, + ..self.state + }; + if let Some(segments) = &mut state.segment_storage { + segments.reset_position(SegmentPosition::default())?; + } + self.publish(state).await?; + self.state = state; + self.durable_head = state.head; + self.segment_files.clear(); + 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.append_batch_buffered(std::slice::from_ref(&prepare)) + .await + } + + /// Append a contiguous extent, validating every operation before allocation or I/O. + /// + /// # Errors + /// Returns an error on invalid history, capacity exhaustion or failed write. + pub async fn append_batch_buffered(&mut self, prepares: &[Frozen<4096>]) -> io::Result<()> { + self.append_batch_inner(prepares, None).await + } + + /// Retain durable segment bodies and append only their headers and locations. + /// Non-message operations keep their inline payloads. The caller must first + /// synchronize every referenced body through the handle that wrote it. + /// + /// Segment names in the parent directory remain offset-based. Hard links + /// owned by this WAL keep referenced inodes alive across retention and purge. + /// + /// # Errors + /// Returns an error for invalid references, history, capacity, or storage. + pub async fn append_batch_referenced_buffered( + &mut self, + prepares: &[Frozen<4096>], + references: &[Option<SegmentReference>], + ) -> io::Result<()> { + if references.len() != prepares.len() { + return Err(invalid("partition WAL reference count mismatch")); + } + self.append_batch_inner(prepares, Some(references)).await + } + + #[allow(clippy::too_many_lines)] + async fn append_batch_inner( + &mut self, + prepares: &[Frozen<4096>], + references: Option<&[Option<SegmentReference>]>, + ) -> io::Result<()> { + self.ensure_healthy()?; + self.cleanup_obsolete().await; + self.recovered_prepares.clear(); + let mut state = self.state; + let mut retained_bytes = self.retained_bytes; + let mut records: Vec<(u64, StoredPrepare, usize)> = Vec::with_capacity(prepares.len()); + for (index, prepare) in prepares.iter().enumerate() { + 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) + || records.last().is_some_and(|(op, entry, _)| { + *op == header.op && entry.checksum == header.checksum + }) + { + continue; + } + if header.op + != state + .head + .checked_add(1) + .ok_or_else(|| invalid("WAL op exhausted"))? + || (state.anchor_known && header.parent != state.head_checksum) + || header.group != state.group + { + return Err(invalid("partition WAL append does not extend its history")); + } + let (reference, next_offset) = if let Some(segments) = &mut state.segment_storage { + if references.is_some() { + return Err(invalid( + "externally assigned reference in owned segment storage", + )); + } + segments.reserve(header, prepare.as_slice())? + } else { + (references.and_then(|references| references[index]), None) + }; + if let Some(reference) = reference { + reference.validate(header, prepare.len())?; + } + let length = + record_length(reference.map_or(prepare.len(), |_| REFERENCED_PREPARE_BYTES))?; + let body_bytes = record_length(prepare.len())? as u64; + if retained_bytes.saturating_add(body_bytes) > self.capacity { + return Err(io::Error::new( + io::ErrorKind::WouldBlock, + "partition WAL requires checkpoint", + )); + } + records.push(( + header.op, + StoredPrepare { + position: state.length, + length, + checksum: header.checksum, + reference, + next_offset, + retained_bytes: body_bytes, + }, + index, + )); + state.length += length as u64; + retained_bytes += body_bytes; + state.head = header.op; + state.head_checksum = header.checksum; + state.segment_references |= reference.is_some(); + if state + .certified_log_view + .is_some_and(|view| header.view > view) + { + state.certified_log_view = None; + } + if !state.anchor_known { + state.checkpoint_checksum = header.parent; + state.anchor_known = true; + } + } + if records.is_empty() { + return Ok(()); + } + let mut extent = Owned::zeroed( + usize::try_from(state.length - self.state.length) + .map_err(|_| invalid("WAL extent overflow"))?, + ); + for (_, record, index) in &records { + let position = usize::try_from(record.position - self.state.length) + .map_err(|_| invalid("WAL extent overflow"))?; + encode_record_into( + prepares[*index].as_slice(), + record.reference, + state.generation, + &mut extent.as_mut_slice()[position..position + record.length], + )?; + } + self.poisoned = true; + self.write_segment_bodies(prepares, &records).await?; + self.retain_segment_inodes(&records).await?; + self.file.write_aligned(self.state.length, extent).await?; + for (op, record, _) in records { + self.entries.insert(op, record); + } + self.state = state; + self.retained_bytes = retained_bytes; + self.retain_active_segment_file(); + self.poisoned = false; + Ok(()) + } + + async fn retain_segment_inodes( + &self, + records: &[(u64, StoredPrepare, usize)], + ) -> io::Result<()> { + let mut linked = false; + for (_, record, _) in records { + if let Some(reference) = record.reference { + let retained = reference.path(&self.directory); + if !self.storage.exists(&retained).await? { + let parent = self + .directory + .parent() + .ok_or_else(|| invalid("referenced segment has no partition directory"))?; + let source = parent.join(format!("{:020}.log", reference.start_offset)); + self.storage.hard_link(&source, &retained).await?; + linked = true; + } + } + } + if linked { + // A visible frontier must never name an inode whose last durable + // directory entry a concurrent retention or purge can remove. + self.storage.sync_directory(&self.directory).await?; + } + Ok(()) + } + + async fn recover_entries(&mut self) -> io::Result<()> { + let state = self.state; + let mut position = 0; + let mut previous = state.checkpoint; + let mut checksum = state.checkpoint_checksum; + while position < state.length { + let (header, length, prepare, reference) = self.read_record(position).await?; + let next_offset = if state.segment_storage.is_some() && reference.is_some() { + Some(segments::batch_next_offset(prepare.as_slice())?) + } else { + None + }; + self.recovered_prepares.push(prepare); + let checkpoint_prepare = position == 0 && state.checkpoint_prepare; + if checkpoint_prepare { + if header.op != state.checkpoint || header.checksum != state.checkpoint_checksum { + return Err(invalid("partition WAL checkpoint prepare mismatch")); + } + } else 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")); + } + self.entries.insert( + header.op, + StoredPrepare { + position, + length, + checksum: header.checksum, + reference, + next_offset, + retained_bytes: record_length(header.size as usize)? as u64, + }, + ); + 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")); + } + self.retained_bytes = self + .entries + .values() + .map(|entry| entry.retained_bytes) + .sum(); + Ok(()) + } + + async fn discover_obsolete(&mut self) -> io::Result<()> { + let retained = self.retained_segment_paths(&self.entries, self.state.segment_storage); + for entry in self.storage.entries(&self.directory).await? { + let Some(name) = entry.name.to_str() else { + continue; + }; + let generation = name + .strip_prefix("prepares-") + .and_then(|name| name.strip_suffix(".wal")) + .and_then(|value| value.parse::<u64>().ok()); + if !entry.directory + && (generation.is_some_and(|generation| generation != self.state.generation) + || name == "frontier.tmp" + || (is_retained_segment_name(name) + && !retained.contains(&self.directory.join(&entry.name)))) + { + self.obsolete.push_back(self.directory.join(entry.name)); + } + } + Ok(()) + } + + async fn cleanup_obsolete(&mut self) { + let count = self.obsolete.len().min(16); + for _ in 0..count { + let Some(path) = self.obsolete.pop_front() else { + break; + }; + match self.storage.remove_file(&path).await { + Ok(()) => self.cleanup_directory_dirty = true, + Err(error) if error.kind() == io::ErrorKind::NotFound => { + self.cleanup_directory_dirty = true; + } + Err(error) => { + tracing::warn!(%error, path = %path.display(), "cannot remove obsolete partition WAL generation"); + self.obsolete.push_back(path); + } + } + } + if self.cleanup_directory_dirty { + match self.storage.sync_directory(&self.directory).await { + Ok(()) => self.cleanup_directory_dirty = false, + Err(error) => tracing::warn!(%error, "cannot synchronize partition WAL cleanup"), + } + } + } + + async fn validate_unpublished_history(storage: &S, directory: &Path) -> io::Result<()> { + for entry in storage.entries(directory).await? { + if entry.directory || entry.name == "frontier.tmp" { + 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, + Message<PrepareHeader>, + Option<SegmentReference>, + )> { + let (mut buffer, frame_length, reference) = self.read_encoded_record(position).await?; + let length = buffer.as_slice().len(); + if let Some(reference) = reference { + let payload = &buffer.as_slice()[RECORD_PREFIX..RECORD_PREFIX + frame_length]; + let header = bytemuck::checked::try_from_bytes::<PrepareHeader>( + &payload[..size_of::<PrepareHeader>()], + ) + .map_err(|_| invalid("invalid referenced prepare header"))?; + let mut prepare = Owned::zeroed(header.size as usize); + prepare.as_mut_slice()[..size_of::<PrepareHeader>()] + .copy_from_slice(&payload[..size_of::<PrepareHeader>()]); + buffer = self + .storage + .open(&reference.path(&self.directory), OpenMode::Read) + .await? + .read_aligned_tail(reference.position, prepare, size_of::<PrepareHeader>()) + .await?; + } else { + buffer + .as_mut_slice() + .copy_within(RECORD_PREFIX..RECORD_PREFIX + frame_length, 0); + buffer.truncate(frame_length); + } + let message = Message::<PrepareHeader>::try_from(buffer) + .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 != message.as_slice().len() + { + return Err(invalid("partition WAL prepare identity mismatch")); + } + if reference.is_some() + && ((header.checksum != 0 && header.identity_checksum() != header.checksum) + || (header.checksum_body != 0 + && header.checksum_body + != u128::from(XxHash3_64::oneshot( + &message.as_slice()[size_of::<PrepareHeader>()..], + )))) + { + return Err(invalid("referenced segment prepare checksum mismatch")); + } + Ok((header, length, message, reference)) + } + + async fn read_encoded_record( + &self, + position: u64, + ) -> io::Result<(Owned<4096>, usize, Option<SegmentReference>)> { + let prefix = self + .file + .read_aligned(position, PARTITION_WAL_BLOCK_SIZE) + .await?; + let frame_length = u32::from_le_bytes( + prefix.as_slice()[..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 buffer = if length == PARTITION_WAL_BLOCK_SIZE { + prefix + } else { + let mut bytes = Owned::zeroed(length); + bytes.as_mut_slice()[..PARTITION_WAL_BLOCK_SIZE].copy_from_slice(prefix.as_slice()); + self.file + .read_aligned_tail( + position + PARTITION_WAL_BLOCK_SIZE as u64, + bytes, + PARTITION_WAL_BLOCK_SIZE, + ) + .await? + }; + let bytes = buffer.as_mut_slice(); + 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[..RECORD_PREFIX + frame_length]) != 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 kind = u32::from_le_bytes( + bytes[RECORD_KIND_OFFSET..RECORD_KIND_OFFSET + size_of::<u32>()] + .try_into() + .map_err(|_| invalid("invalid WAL record kind"))?, + ); + if kind != INLINE_RECORD && (kind != SEGMENT_RECORD || !self.state.segment_references) { + return Err(invalid("unknown partition WAL record kind")); + } + let header = bytemuck::checked::try_from_bytes::<PrepareHeader>( + &bytes[RECORD_PREFIX..RECORD_PREFIX + size_of::<PrepareHeader>()], + ) + .map_err(|_| invalid("invalid partition WAL prepare"))?; + header + .validate() + .map_err(|_| invalid("invalid partition WAL prepare"))?; + if header.group != self.state.group + || (kind == INLINE_RECORD && header.size as usize != frame_length) + { + return Err(invalid("partition WAL prepare identity mismatch")); + } + let reference = if kind == SEGMENT_RECORD { + if frame_length != REFERENCED_PREPARE_BYTES { + return Err(invalid("invalid segment reference record size")); + } + let reference = SegmentReference::decode( + &bytes[RECORD_PREFIX + size_of::<PrepareHeader>()..RECORD_PREFIX + frame_length], + )?; + reference.validate(header, header.size as usize)?; Review Comment: reference frames and rebuilt messages differ in size; both are checked. ########## 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 +# ("unlimited" when unset) +# message_expiry - delete sealed segments older than this ("none" when unset) +# Both policies can be active at once. The active segment is never touched. +# Call GET /options/topic (or the SDK's describe_options) for the full catalog +# with this server's defaults. + +# Topic durability and flush scheduling are set at creation and returned by GetTopic. +# durability controls message completion. consumer_offset_durability controls explicit offset stores and deletes. +# Both independently default to "replicated". Neither inherits the other. +# In replicated groups, either persisted policy enables the prepare WAL for full message bodies. Review Comment: fixed the docs; both persisted policies share segment bodies. ########## core/server/src/server_error.rs: ########## @@ -572,7 +572,7 @@ impl std::fmt::Display for PartitionRecoveryRefusal { durable_position, } => write!( f, - "segment {start_offset} runs under enforce_fsync with {entry_count} sparse \ + "segment {start_offset} runs under persisted durability with {entry_count} sparse \ index entries, and the byte-0 rebuild proved its log only through byte \ {walked_position}, short of byte {durable_position} which the last \ entry's own fdatasync ordering proves the log had already made durable; \ Review Comment: fixed the error to report the validated prefix accurately. -- 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]
