atharvalade commented on code in PR #3564:
URL: https://github.com/apache/iggy/pull/3564#discussion_r3478301520
##########
core/metadata/src/stm/stream.rs:
##########
Review Comment:
Isn't this a snapshot backwards-incompatibility? The comment says "a shorter
array" fills trailing defaults, but on master `round_robin_counter: usize` sat
at position 9 in the positional msgpack array. Old snapshots still carry that
integer at position 9, and `rmp_serde` will fail trying to deserialize it as
`Vec<(u64, ConsumerGroupSnapshot)>`. If server-ng has never cut a real snapshot
this is fine, but the comment's compatibility claim doesn't hold for any
existing snapshot file.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -2344,6 +2355,341 @@ where
Ok(())
}
+ /// Minimum committed offset across all consumers and consumer groups, with
+ /// the holder's identity. `None` when nothing has been committed, in which
+ /// case there is no deletion barrier.
+ fn min_committed_offset(&self) -> Option<(u64, ConsumerKind, u32)> {
+ let consumer_guard = self.consumer_offsets.pin();
+ let group_guard = self.consumer_group_offsets.pin();
+ let consumers = consumer_guard.iter().map(|(_, offset)| {
+ (
+ offset.offset.load(Ordering::Relaxed),
+ offset.kind,
+ offset.consumer_id,
+ )
+ });
+ let groups = group_guard.iter().map(|(_, offset)| {
+ (
+ offset.offset.load(Ordering::Relaxed),
+ offset.kind,
+ offset.consumer_id,
+ )
+ });
+ consumers.chain(groups).min_by_key(|(offset, _, _)| *offset)
+ }
+
+ /// Time-expiry plus size-retention in one pass: remove the leading sealed
+ /// segments that have expired or that push the partition past `max_bytes`.
+ /// Returns the `(segments, messages)` removed.
+ pub async fn clean_expired_segments(
+ &mut self,
+ now: IggyTimestamp,
+ message_expiry: IggyExpiry,
+ max_bytes: Option<u64>,
+ ) -> (u64, u64) {
+ let expired = leading_expired_end(self.log.segments(), now,
message_expiry);
+ let oversized =
+ max_bytes.and_then(|max_bytes|
leading_oversized_end(self.log.segments(), max_bytes));
+ let Some(up_to) = expired.into_iter().chain(oversized).max() else {
+ return (0, 0);
+ };
+ self.remove_sealed_segments_up_to(up_to).await
+ }
+
+ /// Remove the oldest sealed segments whose `end_offset <= up_to_offset`,
+ /// never the active segment and never past the consumer barrier (the
+ /// minimum committed consumer/group offset). Unlinks the messages and
+ /// index files and decrements partition stats. Idempotent: an offset below
+ /// the oldest sealed segment removes nothing. Returns the
+ /// `(segments, messages)` removed.
+ ///
+ /// Holds `write_lock` to serialize against the commit/rotate path, which
+ /// runs on the separate consensus-tick loop.
+ pub async fn remove_sealed_segments_up_to(&mut self, up_to_offset: u64) ->
(u64, u64) {
+ let write_lock = self.write_lock.clone();
+ let _guard = write_lock.lock().await;
+
+ let barrier = self.min_committed_offset();
+ let namespace = self.namespace();
+ let removable = {
+ let segments = self.log.segments();
+ let last_idx = segments.len().saturating_sub(1);
+ let mut removable = 0usize;
+ for (idx, segment) in segments.iter().enumerate() {
+ if idx == last_idx || !segment.sealed || segment.end_offset >
up_to_offset {
+ break;
+ }
+ if let Some((barrier_offset, kind, consumer_id)) = barrier
+ && segment.end_offset > barrier_offset
+ {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ start_offset = segment.start_offset,
+ end_offset = segment.end_offset,
+ barrier = barrier_offset,
+ %kind,
+ consumer_id,
+ "segment retained: blocked by committed consumer
offset"
+ );
+ break;
+ }
+ removable += 1;
+ }
+ removable
+ };
+
+ let mut deleted_segments = 0u64;
+ let mut deleted_messages = 0u64;
+ for _ in 0..removable {
+ // The removable run is always a prefix (oldest first), so the next
+ // victim is index 0 once the previous one is gone.
+ let segment = self.log.segments_mut().remove(0);
+ let mut storage = self.log.storages_mut().remove(0);
+ self.log.indexes_mut().remove(0);
+ self.log.messages_writers_mut().remove(0);
+ self.log.index_writers_mut().remove(0);
+
+ let (messages_path, index_path) =
storage.segment_and_index_paths();
+ let _ = storage.shutdown();
+ drop(storage);
+
+ for path in messages_path.into_iter().chain(index_path) {
+ match compio::fs::remove_file(&path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound
=> {}
+ Err(error) => {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ path = %path,
+ %error,
+ "failed to unlink segment file during cleanup"
+ );
+ }
+ }
+ }
+
+ let segment_size = segment.size.as_bytes_u64();
+ let messages_in_segment = if segment.start_offset ==
segment.end_offset {
+ 0
+ } else {
+ segment.end_offset - segment.start_offset + 1
+ };
Review Comment:
This undercounts by 1 for sealed segments that contain exactly one message.
A sealed segment with one message at offset N has `start_offset == end_offset
== N`, but this branch treats it as 0 messages. Since the removal loop already
guarantees `segment.sealed == true` (line 2419), a sealed segment can't be
empty, so the formula should always be `end_offset - start_offset + 1`.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -2344,6 +2355,341 @@ where
Ok(())
}
+ /// Minimum committed offset across all consumers and consumer groups, with
+ /// the holder's identity. `None` when nothing has been committed, in which
+ /// case there is no deletion barrier.
+ fn min_committed_offset(&self) -> Option<(u64, ConsumerKind, u32)> {
+ let consumer_guard = self.consumer_offsets.pin();
+ let group_guard = self.consumer_group_offsets.pin();
+ let consumers = consumer_guard.iter().map(|(_, offset)| {
+ (
+ offset.offset.load(Ordering::Relaxed),
+ offset.kind,
+ offset.consumer_id,
+ )
+ });
+ let groups = group_guard.iter().map(|(_, offset)| {
+ (
+ offset.offset.load(Ordering::Relaxed),
+ offset.kind,
+ offset.consumer_id,
+ )
+ });
+ consumers.chain(groups).min_by_key(|(offset, _, _)| *offset)
+ }
+
+ /// Time-expiry plus size-retention in one pass: remove the leading sealed
+ /// segments that have expired or that push the partition past `max_bytes`.
+ /// Returns the `(segments, messages)` removed.
+ pub async fn clean_expired_segments(
+ &mut self,
+ now: IggyTimestamp,
+ message_expiry: IggyExpiry,
+ max_bytes: Option<u64>,
+ ) -> (u64, u64) {
+ let expired = leading_expired_end(self.log.segments(), now,
message_expiry);
+ let oversized =
+ max_bytes.and_then(|max_bytes|
leading_oversized_end(self.log.segments(), max_bytes));
+ let Some(up_to) = expired.into_iter().chain(oversized).max() else {
+ return (0, 0);
+ };
+ self.remove_sealed_segments_up_to(up_to).await
+ }
+
+ /// Remove the oldest sealed segments whose `end_offset <= up_to_offset`,
+ /// never the active segment and never past the consumer barrier (the
+ /// minimum committed consumer/group offset). Unlinks the messages and
+ /// index files and decrements partition stats. Idempotent: an offset below
+ /// the oldest sealed segment removes nothing. Returns the
+ /// `(segments, messages)` removed.
+ ///
+ /// Holds `write_lock` to serialize against the commit/rotate path, which
+ /// runs on the separate consensus-tick loop.
+ pub async fn remove_sealed_segments_up_to(&mut self, up_to_offset: u64) ->
(u64, u64) {
+ let write_lock = self.write_lock.clone();
+ let _guard = write_lock.lock().await;
+
+ let barrier = self.min_committed_offset();
+ let namespace = self.namespace();
+ let removable = {
+ let segments = self.log.segments();
+ let last_idx = segments.len().saturating_sub(1);
+ let mut removable = 0usize;
+ for (idx, segment) in segments.iter().enumerate() {
+ if idx == last_idx || !segment.sealed || segment.end_offset >
up_to_offset {
+ break;
+ }
+ if let Some((barrier_offset, kind, consumer_id)) = barrier
+ && segment.end_offset > barrier_offset
+ {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ start_offset = segment.start_offset,
+ end_offset = segment.end_offset,
+ barrier = barrier_offset,
+ %kind,
+ consumer_id,
+ "segment retained: blocked by committed consumer
offset"
+ );
+ break;
+ }
+ removable += 1;
+ }
+ removable
+ };
+
+ let mut deleted_segments = 0u64;
+ let mut deleted_messages = 0u64;
+ for _ in 0..removable {
+ // The removable run is always a prefix (oldest first), so the next
+ // victim is index 0 once the previous one is gone.
+ let segment = self.log.segments_mut().remove(0);
+ let mut storage = self.log.storages_mut().remove(0);
+ self.log.indexes_mut().remove(0);
+ self.log.messages_writers_mut().remove(0);
+ self.log.index_writers_mut().remove(0);
+
+ let (messages_path, index_path) =
storage.segment_and_index_paths();
+ let _ = storage.shutdown();
+ drop(storage);
+
+ for path in messages_path.into_iter().chain(index_path) {
+ match compio::fs::remove_file(&path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound
=> {}
+ Err(error) => {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ path = %path,
+ %error,
+ "failed to unlink segment file during cleanup"
+ );
+ }
+ }
+ }
+
+ let segment_size = segment.size.as_bytes_u64();
+ let messages_in_segment = if segment.start_offset ==
segment.end_offset {
+ 0
+ } else {
+ segment.end_offset - segment.start_offset + 1
+ };
+ self.stats.decrement_size_bytes(segment_size);
+ self.stats.decrement_segments_count(1);
+ self.stats.decrement_messages_count(messages_in_segment);
+
+ deleted_segments += 1;
+ deleted_messages += messages_in_segment;
+
+ debug!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ start_offset = segment.start_offset,
+ end_offset = segment.end_offset,
+ "deleted sealed segment during cleanup"
+ );
+ }
+
+ (deleted_segments, deleted_messages)
+ }
+
+ /// Build and install a fresh empty segment starting at `start_offset` with
+ /// real on-disk writers. Paths are derived from the partition directory
+ /// (see `rotate_segment`); falls back to the config-derived path for
+ /// in-memory partitions with no directory.
+ ///
+ /// # Errors
+ /// If the segment's log / index file cannot be created.
+ async fn install_empty_segment(
+ &mut self,
+ config: &PartitionsConfig,
+ start_offset: u64,
+ ) -> Result<(), IggyError> {
+ let namespace = self.namespace();
+ let (messages_path, index_path) = self.partition_dir().map_or_else(
+ || {
+ (
+ config.get_messages_path(
+ namespace.stream_id(),
+ namespace.topic_id(),
+ namespace.partition_id(),
+ start_offset,
+ ),
+ config.get_index_path(
+ namespace.stream_id(),
+ namespace.topic_id(),
+ namespace.partition_id(),
+ start_offset,
+ ),
+ )
+ },
+ |dir| {
+ (
+ format!("{dir}/{start_offset:0>20}.log"),
+ format!("{dir}/{start_offset:0>20}.index"),
+ )
+ },
+ );
+ let segment = Segment::new(start_offset, config.segment_size);
+ let storage = SegmentStorage::new(
+ &messages_path,
+ &index_path,
+ 0,
+ 0,
+ config.enforce_fsync,
+ config.enforce_fsync,
+ false,
+ )
+ .await
+ .map_err(|_|
IggyError::CannotCreateSegmentLogFile(messages_path.clone()))?;
+ let messages_size_bytes = storage
+ .messages_writer
+ .as_ref()
+ .ok_or_else(||
IggyError::CannotCreateSegmentLogFile(messages_path.clone()))?
+ .size_counter();
+ let messages_writer = Rc::new(
+ MessagesWriter::new(
+ &messages_path,
+ messages_size_bytes,
+ config.enforce_fsync,
+ false,
+ )
+ .await
+ .map_err(|_|
IggyError::CannotCreateSegmentLogFile(messages_path.clone()))?,
+ );
+ let index_writer = Rc::new(
+ IggyIndexWriter::new(
+ &index_path,
+ Rc::new(std::sync::atomic::AtomicU64::new(0)),
+ config.enforce_fsync,
+ false,
+ )
+ .await
+ .map_err(|_|
IggyError::CannotCreateSegmentIndexFile(index_path.clone()))?,
+ );
+ self.log
+ .add_persisted_segment(segment, storage, Some(messages_writer),
Some(index_writer));
+ Ok(())
+ }
+
+ /// Reset the partition to a single empty segment at offset 0 and clear all
+ /// consumer / consumer-group offsets (memory + disk). This is the local
+ /// effect of a committed `PurgeTopic`: it wipes message data and offsets
but
+ /// preserves the partition and its consumer-group membership. Mirrors the
+ /// legacy server's `purge_all_segments` + offset-file deletion.
+ ///
+ /// Records `generation` as the applied purge generation so the reconciler
+ /// does not re-wipe a partition already purged at this generation (a later
+ /// `PurgeTopic` advances the committed generation and triggers a fresh
pass).
+ ///
+ /// # Errors
+ /// If the replacement segment's log / index file cannot be created.
+ pub async fn purge(
+ &mut self,
+ config: &PartitionsConfig,
+ generation: u64,
+ ) -> Result<(), IggyError> {
+ let write_lock = self.write_lock.clone();
+ let _guard = write_lock.lock().await;
+
+ let namespace = self.namespace();
+
+ // Drain every segment (including the active one) and unlink its files.
+ let segment_count = self.log.segments().len();
+ for _ in 0..segment_count {
+ self.log.segments_mut().remove(0);
+ let mut storage = self.log.storages_mut().remove(0);
+ self.log.indexes_mut().remove(0);
+ self.log.messages_writers_mut().remove(0);
+ self.log.index_writers_mut().remove(0);
+
+ let (messages_path, index_path) =
storage.segment_and_index_paths();
+ let _ = storage.shutdown();
+ drop(storage);
+
+ for path in messages_path.into_iter().chain(index_path) {
+ match compio::fs::remove_file(&path).await {
+ Ok(()) => {}
+ Err(error) if error.kind() == std::io::ErrorKind::NotFound
=> {}
+ Err(error) => {
+ warn!(
+ target: "iggy.partitions.diag",
+ plane = "partitions",
+ namespace_raw = namespace.inner(),
+ path = %path,
+ %error,
+ "failed to unlink segment file during purge"
+ );
+ }
+ }
+ }
+ }
+
+ // Recreate a fresh empty segment at offset 0 with real writers.
+ let start_offset = 0u64;
+ self.install_empty_segment(config, start_offset).await?;
+
+ // Reset the offset counters so new messages start at offset 0.
+ self.offset.store(start_offset, Ordering::Release);
+ self.dirty_offset.store(start_offset, Ordering::Relaxed);
+ self.should_increment_offset = false;
+
+ // Clear consumer + consumer-group offsets (memory + disk). Collect the
+ // file paths before deleting so the map guard is not held across an
+ // await.
+ let consumer_paths: Vec<String> = {
+ let guard = self.consumer_offsets.pin();
+ let paths = guard
+ .iter()
+ .filter_map(|(key, _)| {
+ u32::try_from(*key)
+ .ok()
+ .and_then(|id|
self.persisted_offset_path(ConsumerKind::Consumer, id))
+ })
+ .collect();
+ guard.clear();
+ paths
+ };
+ let group_paths: Vec<String> = {
+ let guard = self.consumer_group_offsets.pin();
+ let paths = guard
+ .iter()
+ .filter_map(|(key, _)| {
+ u32::try_from(key.0)
+ .ok()
+ .and_then(|id|
self.persisted_offset_path(ConsumerKind::ConsumerGroup, id))
+ })
+ .collect();
+ guard.clear();
+ paths
Review Comment:
`last_polled_offsets` is not cleared during purge. After purge resets
everything to offset 0, stale `last_polled` entries survive. The reconciler's
completable check (`committed >= last_polled`) then sees `None >=
Some(old_high_offset)` and can't complete pending revocations until timeout. I
think you need `self.last_polled_offsets.pin().clear();` alongside the other
two clears.
--
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]