hubcio commented on code in PR #4092:
URL: https://github.com/apache/iggy/pull/4092#discussion_r3983323151
##########
core/partitions/src/state_transfer.rs:
##########
@@ -3194,6 +3465,9 @@ where
minted_next_offset: u64,
staged_was_empty: bool,
) -> Result<(), iggy_common::IggyError> {
+ if let Some(persistence) = &self.persistence {
Review Comment:
removed the unreachable retained-offset branch.
##########
core/partitions/src/offset_storage.rs:
##########
@@ -121,15 +121,43 @@ pub fn decode_offset_record(bytes: &[u8]) -> OffsetRecord
{
///
/// # Errors
/// [`IggyError`] when the directory, file, or write cannot be created or
completed.
-pub async fn persist_offset(path: &str, offset: u64, enforce_fsync: bool) ->
Result<(), IggyError> {
+pub async fn persist_offset(path: &str, offset: u64, persisted: bool) ->
Result<(), IggyError> {
let record = encode_offset_record(offset);
- if enforce_fsync {
+ if persisted {
replace_file(path, record, true, false).await
} else {
write_in_place(path, record).await
}
}
+/// Keep the original writer open so checkpoint observes its writeback errors.
Review Comment:
documented the barriers required before wal reclamation.
##########
core/partitions/src/iggy_partition.rs:
##########
@@ -2139,10 +2669,53 @@ where
Ok(())
}
+ async fn write_consumer_offset(
+ &self,
+ path: &str,
+ offset: u64,
+ persisted: bool,
+ ) -> Result<(), IggyError> {
+ if let Some(persistence) = &self.persistence {
+ let file = crate::offset_storage::persist_offset_retained(
+ path,
+ offset,
+ persistence.take_offset_file(path),
+ )
+ .await
+ .inspect_err(|error| {
+ persistence.fail(std::io::Error::other(error.to_string()));
+ })?;
+ persistence.retain_offset_file(path.to_owned(), file);
+ Ok(())
+ } else {
+ persist_offset(path, offset, persisted).await
+ }
+ }
+
+ async fn write_cold_consumer_offset(
+ &self,
+ path: &str,
+ offset: u64,
+ persisted: bool,
+ ) -> Result<(u64, bool), IggyError> {
+ if self.persistence.is_some() {
+ let result = crate::offset_storage::read_offset_max(path,
offset).await?;
Review Comment:
fixed; covered cold-offset writes are skipped.
##########
core/server/src/http/reads.rs:
##########
@@ -404,6 +404,78 @@ pub(in crate::http) fn authorize_data_plane(
.authorize(|permissioner| rule(permissioner, user_id, stream_id,
topic_id))
}
+static DURABILITY_KEY: std::sync::LazyLock<iggy_common::HeaderKey> =
+ std::sync::LazyLock::new(|| "durability".parse().expect("catalog key is
valid"));
+
+#[derive(Clone, Copy, PartialEq, Eq)]
+pub(in crate::http) struct TopicDurability {
+ stream_id: usize,
+ topic_id: usize,
+ created_revision: u64,
+ pub durability: iggy_common::Durability,
+}
+
+impl TopicDurability {
+ pub fn confirmed_policy(self, state: &HttpInner) ->
iggy_common::Durability {
+ state
+ .shard
+ .plane
+ .metadata()
+ .mux_stm
+ .streams()
+ .read(|inner| self.confirmed_policy_in(inner))
+ }
+
+ fn confirmed_policy_in(
Review Comment:
durability is create-only; the revision check protects replacement.
##########
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 {
Review Comment:
documented why reset can install a new layout size.
--
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]