hubcio commented on code in PR #4092:
URL: https://github.com/apache/iggy/pull/4092#discussion_r3983320200
##########
core/simulator/src/replica.rs:
##########
@@ -379,8 +379,7 @@ pub fn new_shard(
let partitions_config = PartitionsConfig {
messages_required_to_save: 1000,
size_of_messages_required_to_save: IggyByteSize::from(4 * 1024 * 1024),
- enforce_fsync: false, //Disable fsync for simulation
- consumer_offset_enforce_fsync: false,
+
validate_checksum: true,
segment_size: IggyByteSize::from(1024 * 1024 * 1024),
Review Comment:
reused the shared segment-size default in the simulator.
##########
core/shard/src/lib.rs:
##########
@@ -7553,6 +7606,14 @@ where
walk_cursor.get_or_insert(namespace);
}
}
+ if partition.needs_persistence_checkpoint() {
Review Comment:
fixed; checkpoint failures are returned in the same sweep.
##########
core/shard/src/lib.rs:
##########
@@ -7381,6 +7416,24 @@ 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() {
Review Comment:
fixed; fatal partitions can't checkpoint or send pending acks.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -711,8 +727,510 @@ 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> {
+ self.open_persistence_with_recovered(capacity, None).await
+ }
+
+ /// # Errors
+ /// Returns an error if durable history cannot be opened, migrated, or
replayed.
+ #[allow(clippy::too_many_lines)]
+ pub async fn open_persistence_with_recovered(
+ &mut self,
+ capacity: u64,
+ recovered: Option<(Rc<PartitionPersistence>,
Vec<Message<PrepareHeader>>)>,
+ ) -> Result<(), IggyError> {
+ if self.consensus.replica_count() > 1
+ && let Some(directory) = &self.partition_dir
+ {
+ self.materialization_missing =
+ crate::state_transfer::materialization_is_missing(directory,
self.created_revision)
+ .await
+ .map_err(|_| IggyError::CannotReadFile)?;
+ self.ensure_materialization_recovery();
+ }
+ if self.consensus.replica_count() == 1
+ || !(self.durability().is_persisted()
+ || self.consumer_offset_durability().is_persisted())
+ {
+ return Ok(());
+ }
+ let directory = self
+ .partition_dir
+ .as_ref()
+ .ok_or(IggyError::CannotReadFile)?;
+ let directory =
+ std::path::Path::new(directory).join(format!("prepares-{}",
self.created_revision));
+ let (persistence, prepares) = if let Some(recovered) = recovered {
+ recovered
+ } else {
+ PartitionPersistence::open_with_capacity(
+ &directory,
+ self.namespace().inner(),
+ self.created_revision,
+ journal::durable_storage::DiskStorage,
+ capacity,
+ self.runtime_options
+ .preallocate_segments
+ .unwrap_or(iggy_common::DEFAULT_PREALLOCATE_SEGMENTS),
+ )
+ .await
+ .map_err(|error| {
+ warn!(%error, "cannot open partition prepare WAL");
+ IggyError::CannotReadFile
+ })?
+ };
+ if self.durability().is_persisted() && !self.materialization_missing {
+ let segment = self.log.active_segment();
+ let length = segment.size.as_bytes_u64();
+ let initial = journal::partition_journal::SegmentPosition {
+ start_offset: segment.start_offset,
+ length,
+ next_offset: if length == 0 {
+ segment.start_offset
+ } else {
+ segment
+ .end_offset
+ .checked_add(1)
+ .ok_or(IggyError::CannotReadFile)?
+ },
+ };
+ persistence.enable_segment_storage(initial,
segment.max_size.as_bytes_u64());
+ if persistence.start() {
+ self.consensus
+ .message_bus()
+ .spawn(Rc::clone(&persistence).run());
+ }
+ persistence.drain_with_timeout().await.map_err(|error| {
+ warn!(%error, "cannot enable partition segment persistence");
+ IggyError::CannotSyncFile
+ })?;
+ }
+ self.restore_certified_log_view(&persistence).await?;
+ if self.materialization_missing {
+ self.persistence = Some(persistence);
+ return Ok(());
+ }
+ let (purge_generation, purge_floor) = persistence.purge_marker();
+ if purge_generation <= self.applied_purge_generation {
+ self.purge_floor_op = self.purge_floor_op.max(purge_floor);
+ }
+ let checkpoint = persistence.checkpoint_op();
+ let head = persistence.head();
+ let mut commit = checkpoint;
+ for message in prepares {
+ let header = *message.header();
+ if header.op == checkpoint {
+ self.log
+ .journal()
+ .inner
+ .restore_checkpoint_prepare(checkpoint,
message.into_frozen());
+ continue;
+ }
+ commit = commit.max(header.commit.min(head));
+ if header.operation == Operation::SendMessages {
+ self.append_repaired_send_messages(message).await?;
+ } else {
+ self.apply_replicated_operation(message).await?;
+ }
+ }
+ if head > 0 {
+ self.consensus.sequencer().set_sequence(head);
+ if let Some(checksum) = persistence.checksum(head) {
+ self.consensus.set_last_prepare_checksum(checksum);
+ }
+ self.consensus.restore_commit_state(checkpoint, commit);
+ }
+ // Recovery can reuse completed writes whose last barrier was
interrupted.
+ // Include them once before reclaiming any recovered WAL history.
+ for segment in self.log.segments() {
+ persistence.mark_segment_dirty(segment.start_offset);
+ }
+ for kind in [ConsumerKind::Consumer, ConsumerKind::ConsumerGroup] {
+ self.durable_consumer_offsets.with_entries(kind, |entries| {
+ for consumer_id in entries.keys() {
+ persistence.mark_offset_dirty(
+ crate::state_transfer::consumer_kind_index(kind),
+ *consumer_id,
+ true,
+ );
+ }
+ });
+ }
+ self.persistence = Some(persistence);
+ Ok(())
+ }
+
+ async fn restore_certified_log_view(
+ &mut self,
+ persistence: &Rc<PartitionPersistence>,
+ ) -> Result<(), IggyError> {
+ if !self.materialization_missing
+ && self.recovered_log_view.is_none()
+ && persistence.head() == 0
+ {
+ persistence.certify_log_view(self.consensus.log_view(), 0, 0);
+ if persistence.start() {
+ self.consensus
+ .message_bus()
+ .spawn(Rc::clone(persistence).run());
+ }
+ persistence
+ .drain_with_timeout()
+ .await
+ .map_err(|_| IggyError::CannotSyncFile)?;
+ }
+ if !self.materialization_missing {
+ match persistence.certified_log_view() {
+ Some(view) if view >= self.consensus.log_view() => {
+ if view > self.consensus.view() {
+ self.consensus.set_view(view);
+ }
+ self.consensus.set_log_view(view);
+ }
+ _ => {
+ self.materialization_missing = true;
+ self.ensure_materialization_recovery();
+ }
+ }
+ }
+ Ok(())
+ }
+
+ pub const fn requires_state_transfer(&self) -> bool {
+ self.materialization_missing
+ }
+
+ pub fn ensure_materialization_recovery(&self) {
+ if self.materialization_missing
+ && (self.consensus.state_transfer_stage() ==
consensus::StateTransferStage::Idle
+ || self.consensus.status() == consensus::Status::ViewChange)
+ {
+ self.consensus.begin_view_probe();
+ if self.consensus.state_transfer_stage() ==
consensus::StateTransferStage::Idle {
+ self.consensus.begin_state_transfer_await();
+ }
+ }
+ }
+
+ pub async fn on_persistence_completed(&mut self, completion:
PersistenceCompletion) {
+ if !self
+ .persistence
+ .as_ref()
+ .is_some_and(|persistence|
persistence.accepts_completion(completion))
+ {
+ return;
+ }
+ self.drive_persistence().await;
+ }
+
+ pub fn needs_persistence_checkpoint(&self) -> bool {
+ self.persistence.as_ref().is_some_and(|persistence| {
+ persistence.needs_checkpoint()
+ && self.consensus.commit_min().min(persistence.head()) >
persistence.checkpoint_op()
+ })
+ }
+
+ pub async fn checkpoint_persistence(&mut self, config: &PartitionsConfig) {
+ let Some(persistence) = self
+ .persistence
+ .as_ref()
+ .filter(|persistence| persistence.needs_checkpoint())
+ .cloned()
+ else {
+ return;
+ };
+ let through_op = self.consensus.commit_min().min(persistence.head());
+ if through_op <= persistence.checkpoint_op() {
+ return;
+ }
+ if let Err(error) = self.commit_messages_inner(config, true,
through_op).await {
+ error!(%error, namespace_raw = self.namespace().inner(),
"partition checkpoint failed");
+ self.fatal = Some(FatalCommit {
+ namespace_raw: self.namespace().inner(),
+ op: through_op,
+ operation: Operation::SendMessages,
+ });
+ return;
+ }
+ let (files, directories) = self.persistence_checkpoint_files(config);
+ persistence.checkpoint_files(through_op, files, directories);
+ self.start_persistence();
+ }
+
+ fn persistence_checkpoint_files(
+ &self,
+ config: &PartitionsConfig,
+ ) -> (Vec<std::path::PathBuf>, Vec<std::path::PathBuf>) {
+ let namespace = self.namespace();
+ let Some(persistence) = &self.persistence else {
+ return (Vec::new(), Vec::new());
+ };
+ let (segments, offsets) = persistence.take_dirty_files();
+ let mut paths = Vec::with_capacity(
+ segments.len() * 2
+ + offsets
+ .iter()
+ .map(std::collections::BTreeSet::len)
+ .sum::<usize>(),
+ );
+ for start_offset in segments {
+ // Retention may already have removed a dirty sealed segment.
+ if self
+ .log
+ .segments()
+ .binary_search_by_key(&start_offset, |segment|
segment.start_offset)
+ .is_err()
+ {
+ continue;
+ }
+ paths.push(config.get_messages_path(
+ namespace.stream_id(),
+ namespace.topic_id(),
+ namespace.partition_id(),
+ start_offset,
+ ));
+ paths.push(config.get_index_path(
+ namespace.stream_id(),
+ namespace.topic_id(),
+ namespace.partition_id(),
+ start_offset,
+ ));
+ }
+ for (kind, consumers) in [ConsumerKind::Consumer,
ConsumerKind::ConsumerGroup]
+ .into_iter()
+ .zip(offsets)
+ {
+ for consumer_id in consumers {
+ if self.durable_consumer_offsets.contains(kind, consumer_id)
+ && let Some(path) = self.persisted_offset_path(kind,
consumer_id)
+ {
+ paths.push(path);
+ }
+ }
+ }
+ let mut directories = Vec::with_capacity(4);
+ for directory in [
+ &self.consumer_offsets_path,
+ &self.consumer_group_offsets_path,
+ ]
+ .into_iter()
+ .flatten()
+ {
+ let path = std::path::PathBuf::from(directory);
+ directories.push(path.clone());
+ if let Some(parent) = path.parent()
+ && !directories.iter().any(|existing| existing == parent)
+ {
+ directories.push(parent.to_path_buf());
+ }
+ }
+ if let Some(directory) = &self.partition_dir {
+ directories.push(std::path::PathBuf::from(directory));
+ }
+ (
+ paths.into_iter().map(std::path::PathBuf::from).collect(),
+ directories,
+ )
+ }
+
+ pub fn take_persistence_metrics(&self) ->
Option<crate::persistence::PersistenceMetrics> {
+ self.persistence
+ .as_ref()
+ .map(|persistence| persistence.take_metrics())
+ }
+
+ fn persistence_checkpoint_pending(&self) -> bool {
+ self.persistence
+ .as_ref()
+ .is_some_and(|persistence| persistence.checkpoint_pending())
+ }
+
+ pub async fn drive_persistence(&mut self) {
+ let Some(persistence) = self.persistence.as_ref() else {
+ return;
+ };
+ if let Some(error) = persistence.failure() {
+ error!(%error, namespace_raw = self.namespace().inner(),
"partition prepare persistence failed");
+ if self.fatal.is_none() {
+ self.fatal = Some(FatalCommit {
+ namespace_raw: self.namespace().inner(),
+ op: persistence.head(),
+ operation: Operation::SendMessages,
+ });
+ }
+ return;
+ }
+ if self.materialization_missing || !self.ensure_wal_view() {
+ return;
+ }
+ let durable_op = persistence.durable_op();
+ loop {
+ let ready = self
+ .pending_persisted_acks
+ .borrow()
+ .first_key_value()
+ .is_some_and(|(op, _)| *op <= durable_op);
+ if !ready {
+ break;
+ }
+ let Some((_, header)) =
self.pending_persisted_acks.borrow_mut().pop_first() else {
+ break;
+ };
+ if persistence.checksum(header.op) == Some(header.checksum)
Review Comment:
fixed; pending acks survive until they're actually queued.
##########
core/shard/src/lib.rs:
##########
@@ -8501,11 +8562,10 @@ where
/// disk as they complete, so this bounds corruption, not memory.
const PARTITION_TRANSFER_TOTAL_LEN_MAX: u64 = 1 << 40;
- /// Alloc cap for the `CONSUMER_OFFSETS` artifact, which accumulates whole
- /// in `ArtifactProgress::buf` before decode can reject it. Its decoder
- /// ceilings imply ~24 MiB (two sections of 2^20 12-byte entries); this
- /// leaves headroom without letting a hostile manifest stage gigabytes.
- const CONSUMER_OFFSETS_ARTIFACT_LEN_MAX: u64 = 32 << 20;
+ /// Bound the buffered offset and dedup state plus one maximum-sized
Review Comment:
kept the cap for full checkpoint prepares; corrected the size docs.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -711,8 +727,510 @@ 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> {
+ self.open_persistence_with_recovered(capacity, None).await
+ }
+
+ /// # Errors
+ /// Returns an error if durable history cannot be opened, migrated, or
replayed.
+ #[allow(clippy::too_many_lines)]
+ pub async fn open_persistence_with_recovered(
+ &mut self,
+ capacity: u64,
+ recovered: Option<(Rc<PartitionPersistence>,
Vec<Message<PrepareHeader>>)>,
+ ) -> Result<(), IggyError> {
+ if self.consensus.replica_count() > 1
+ && let Some(directory) = &self.partition_dir
+ {
+ self.materialization_missing =
+ crate::state_transfer::materialization_is_missing(directory,
self.created_revision)
+ .await
+ .map_err(|_| IggyError::CannotReadFile)?;
+ self.ensure_materialization_recovery();
+ }
+ if self.consensus.replica_count() == 1
+ || !(self.durability().is_persisted()
+ || self.consumer_offset_durability().is_persisted())
+ {
+ return Ok(());
+ }
+ let directory = self
+ .partition_dir
+ .as_ref()
+ .ok_or(IggyError::CannotReadFile)?;
+ let directory =
+ std::path::Path::new(directory).join(format!("prepares-{}",
self.created_revision));
+ let (persistence, prepares) = if let Some(recovered) = recovered {
+ recovered
+ } else {
+ PartitionPersistence::open_with_capacity(
+ &directory,
+ self.namespace().inner(),
+ self.created_revision,
+ journal::durable_storage::DiskStorage,
+ capacity,
+ self.runtime_options
+ .preallocate_segments
+ .unwrap_or(iggy_common::DEFAULT_PREALLOCATE_SEGMENTS),
+ )
+ .await
+ .map_err(|error| {
+ warn!(%error, "cannot open partition prepare WAL");
+ IggyError::CannotReadFile
+ })?
+ };
+ if self.durability().is_persisted() && !self.materialization_missing {
+ let segment = self.log.active_segment();
+ let length = segment.size.as_bytes_u64();
+ let initial = journal::partition_journal::SegmentPosition {
+ start_offset: segment.start_offset,
+ length,
+ next_offset: if length == 0 {
+ segment.start_offset
+ } else {
+ segment
+ .end_offset
+ .checked_add(1)
+ .ok_or(IggyError::CannotReadFile)?
+ },
+ };
+ persistence.enable_segment_storage(initial,
segment.max_size.as_bytes_u64());
+ if persistence.start() {
+ self.consensus
+ .message_bus()
+ .spawn(Rc::clone(&persistence).run());
+ }
+ persistence.drain_with_timeout().await.map_err(|error| {
+ warn!(%error, "cannot enable partition segment persistence");
+ IggyError::CannotSyncFile
+ })?;
+ }
+ self.restore_certified_log_view(&persistence).await?;
+ if self.materialization_missing {
+ self.persistence = Some(persistence);
+ return Ok(());
+ }
+ let (purge_generation, purge_floor) = persistence.purge_marker();
+ if purge_generation <= self.applied_purge_generation {
+ self.purge_floor_op = self.purge_floor_op.max(purge_floor);
+ }
+ let checkpoint = persistence.checkpoint_op();
+ let head = persistence.head();
+ let mut commit = checkpoint;
+ for message in prepares {
+ let header = *message.header();
+ if header.op == checkpoint {
+ self.log
+ .journal()
+ .inner
+ .restore_checkpoint_prepare(checkpoint,
message.into_frozen());
+ continue;
+ }
+ commit = commit.max(header.commit.min(head));
+ if header.operation == Operation::SendMessages {
+ self.append_repaired_send_messages(message).await?;
+ } else {
+ self.apply_replicated_operation(message).await?;
+ }
+ }
+ if head > 0 {
+ self.consensus.sequencer().set_sequence(head);
+ if let Some(checksum) = persistence.checksum(head) {
+ self.consensus.set_last_prepare_checksum(checksum);
+ }
+ self.consensus.restore_commit_state(checkpoint, commit);
+ }
+ // Recovery can reuse completed writes whose last barrier was
interrupted.
+ // Include them once before reclaiming any recovered WAL history.
+ for segment in self.log.segments() {
+ persistence.mark_segment_dirty(segment.start_offset);
+ }
+ for kind in [ConsumerKind::Consumer, ConsumerKind::ConsumerGroup] {
+ self.durable_consumer_offsets.with_entries(kind, |entries| {
+ for consumer_id in entries.keys() {
+ persistence.mark_offset_dirty(
+ crate::state_transfer::consumer_kind_index(kind),
+ *consumer_id,
+ true,
+ );
+ }
+ });
+ }
+ self.persistence = Some(persistence);
+ Ok(())
+ }
+
+ async fn restore_certified_log_view(
+ &mut self,
+ persistence: &Rc<PartitionPersistence>,
+ ) -> Result<(), IggyError> {
+ if !self.materialization_missing
+ && self.recovered_log_view.is_none()
+ && persistence.head() == 0
+ {
+ persistence.certify_log_view(self.consensus.log_view(), 0, 0);
+ if persistence.start() {
+ self.consensus
+ .message_bus()
+ .spawn(Rc::clone(persistence).run());
+ }
+ persistence
+ .drain_with_timeout()
+ .await
+ .map_err(|_| IggyError::CannotSyncFile)?;
+ }
+ if !self.materialization_missing {
+ match persistence.certified_log_view() {
+ Some(view) if view >= self.consensus.log_view() => {
+ if view > self.consensus.view() {
+ self.consensus.set_view(view);
+ }
+ self.consensus.set_log_view(view);
+ }
+ _ => {
+ self.materialization_missing = true;
+ self.ensure_materialization_recovery();
+ }
+ }
+ }
+ Ok(())
+ }
+
+ pub const fn requires_state_transfer(&self) -> bool {
+ self.materialization_missing
+ }
+
+ pub fn ensure_materialization_recovery(&self) {
+ if self.materialization_missing
+ && (self.consensus.state_transfer_stage() ==
consensus::StateTransferStage::Idle
+ || self.consensus.status() == consensus::Status::ViewChange)
+ {
+ self.consensus.begin_view_probe();
+ if self.consensus.state_transfer_stage() ==
consensus::StateTransferStage::Idle {
+ self.consensus.begin_state_transfer_await();
+ }
+ }
+ }
+
+ pub async fn on_persistence_completed(&mut self, completion:
PersistenceCompletion) {
+ if !self
+ .persistence
+ .as_ref()
+ .is_some_and(|persistence|
persistence.accepts_completion(completion))
+ {
+ return;
+ }
+ self.drive_persistence().await;
+ }
+
+ pub fn needs_persistence_checkpoint(&self) -> bool {
+ self.persistence.as_ref().is_some_and(|persistence| {
+ persistence.needs_checkpoint()
+ && self.consensus.commit_min().min(persistence.head()) >
persistence.checkpoint_op()
+ })
+ }
+
+ pub async fn checkpoint_persistence(&mut self, config: &PartitionsConfig) {
+ let Some(persistence) = self
+ .persistence
+ .as_ref()
+ .filter(|persistence| persistence.needs_checkpoint())
+ .cloned()
+ else {
+ return;
+ };
+ let through_op = self.consensus.commit_min().min(persistence.head());
+ if through_op <= persistence.checkpoint_op() {
+ return;
+ }
+ if let Err(error) = self.commit_messages_inner(config, true,
through_op).await {
+ error!(%error, namespace_raw = self.namespace().inner(),
"partition checkpoint failed");
+ self.fatal = Some(FatalCommit {
Review Comment:
covered by the fatal guard; the first fault stays intact.
--
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]