This is an automated email from the ASF dual-hosted git repository. numinnex pushed a commit to branch fix_cleanup_issues_nextnextnext in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 258c662e076c5f7f7647dbe41b164d9e67ebbb2c Author: Grzegorz Koszyk <[email protected]> AuthorDate: Tue Aug 18 13:49:57 2026 +0200 refactor(shard): drop journal generic, metrics scrape, poll cursor --- core/journal/src/lib.rs | 9 +- core/journal/src/prepare_journal.rs | 3 +- core/metadata/src/impls/metadata.rs | 8 +- core/partitions/src/iggy_index_reader.rs | 82 ++++++++++++- core/partitions/src/iggy_partition.rs | 2 +- core/partitions/src/journal.rs | 2 +- core/partitions/src/log.rs | 37 ++---- core/partitions/src/poll_plan.rs | 138 +++++++++++++++++++++- core/server/src/auth.rs | 12 +- core/server/src/bootstrap.rs | 19 ++- core/server/src/consumer_group.rs | 8 +- core/server/src/dispatch.rs | 86 +++++++------- core/server/src/dispatch/authz.rs | 24 ++-- core/server/src/http.rs | 3 +- core/server/src/http/metrics.rs | 36 +++++- core/server/src/responses.rs | 42 +++---- core/server/src/users.rs | 2 +- core/shard/src/lib.rs | 197 +++++++++---------------------- core/shard/src/metrics.rs | 157 +++++++++++++++++------- core/shard/src/router.rs | 16 +-- core/simulator/src/deps.rs | 3 +- 21 files changed, 551 insertions(+), 335 deletions(-) diff --git a/core/journal/src/lib.rs b/core/journal/src/lib.rs index f97dc35f4..cca455a18 100644 --- a/core/journal/src/lib.rs +++ b/core/journal/src/lib.rs @@ -26,10 +26,7 @@ pub mod local_gate; pub mod prepare_journal; pub mod superblock; -pub trait Journal<S> -where - S: Storage, -{ +pub trait Journal { type Header; type Entry; type HeaderRef<'a>: Deref<Target = Self::Header> @@ -101,8 +98,7 @@ where } pub trait JournalHandle { - type Storage: Storage; - type Target: Journal<Self::Storage>; + type Target: Journal; fn handle(&self) -> &Self::Target; } @@ -112,7 +108,6 @@ pub trait JournalHandle { /// the metadata WAL across a replica restart: the bytes and index survive the /// shard being dropped and rebuilt. impl<T: JournalHandle> JournalHandle for Rc<T> { - type Storage = T::Storage; type Target = T::Target; fn handle(&self) -> &Self::Target { diff --git a/core/journal/src/prepare_journal.rs b/core/journal/src/prepare_journal.rs index dc15e0217..a5f1fde47 100644 --- a/core/journal/src/prepare_journal.rs +++ b/core/journal/src/prepare_journal.rs @@ -741,7 +741,7 @@ impl PrepareJournal { clippy::cast_sign_loss, clippy::future_not_send )] -impl Journal<FileStorage> for PrepareJournal { +impl Journal for PrepareJournal { fn last_op(&self) -> Option<u64> { self.last_op.get() } @@ -1162,7 +1162,6 @@ impl Journal<FileStorage> for PrepareJournal { } impl JournalHandle for PrepareJournal { - type Storage = FileStorage; type Target = Self; fn handle(&self) -> &Self::Target { diff --git a/core/metadata/src/impls/metadata.rs b/core/metadata/src/impls/metadata.rs index 1fda606d1..86170553d 100644 --- a/core/metadata/src/impls/metadata.rs +++ b/core/metadata/src/impls/metadata.rs @@ -945,7 +945,7 @@ where B: MessageBus, SB: SuperblockStore, J: JournalHandle, - J::Target: Journal<J::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + J::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StreamsFrontend + StateMachine< Input = Message<PrepareHeader>, @@ -1313,7 +1313,7 @@ where B: MessageBus, P: Pipeline<Entry = PipelineEntry>, J: JournalHandle, - J::Target: Journal<J::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + J::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StateMachine<Input = Message<PrepareHeader>>, { fn is_applicable<H>(&self, message: &<VsrConsensus<B, P> as Consensus>::Message<H>) -> bool @@ -1456,7 +1456,7 @@ where B: MessageBus, SB: SuperblockStore, J: JournalHandle, - J::Target: Journal<J::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + J::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StreamsFrontend + StateMachine< Input = Message<PrepareHeader>, @@ -3090,7 +3090,7 @@ where P: Pipeline<Entry = PipelineEntry>, SB: SuperblockStore, J: JournalHandle, - J::Target: Journal<J::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + J::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StreamsFrontend + StateMachine< Input = Message<PrepareHeader>, diff --git a/core/partitions/src/iggy_index_reader.rs b/core/partitions/src/iggy_index_reader.rs index bc53a0276..a5dacf142 100644 --- a/core/partitions/src/iggy_index_reader.rs +++ b/core/partitions/src/iggy_index_reader.rs @@ -170,8 +170,42 @@ impl IggyIndexReader { entry_count: u64, offset: u64, ) -> Result<Option<IggyIndex>, IggyError> { - self.lower_bound_by(entry_count, |entry| entry.offset, offset) - .await + Ok(self + .lower_bound_by(entry_count, |entry| entry.offset, offset) + .await? + .map(|(entry, _)| entry)) + } + + /// [`Self::offset_lower_bound`] that also reports the successor entry's + /// offset (`None` when the match is the last entry), bounding the offset + /// interval the match resolves for. Costs one extra entry read; callers + /// memoizing the resolution use the bound to answer later in-interval + /// queries without reopening the index. + /// + /// # Errors + /// + /// Returns an error if any probed entry cannot be read. + pub async fn offset_lower_bound_with_successor( + &self, + entry_count: u64, + offset: u64, + ) -> Result<Option<(IggyIndex, Option<u64>)>, IggyError> { + let Some((entry, successor_index)) = self + .lower_bound_by(entry_count, |entry| entry.offset, offset) + .await? + else { + return Ok(None); + }; + let successor_offset = if successor_index < entry_count { + Some( + self.read_entry_at(successor_index * IGGY_INDEX_SIZE as u64) + .await? + .offset, + ) + } else { + None + }; + Ok(Some((entry, successor_offset))) } /// Last entry with `timestamp` at or below the target; `None` semantics @@ -185,16 +219,21 @@ impl IggyIndexReader { entry_count: u64, timestamp: u64, ) -> Result<Option<IggyIndex>, IggyError> { - self.lower_bound_by(entry_count, |entry| entry.timestamp, timestamp) - .await + Ok(self + .lower_bound_by(entry_count, |entry| entry.timestamp, timestamp) + .await? + .map(|(entry, _)| entry)) } + /// Binary search for the last entry with `key` at or below `target`, + /// returning it with its successor's entry index (the standard + /// lower-bound exit: `low` lands on the first entry above the target). async fn lower_bound_by( &self, entry_count: u64, key: impl Fn(&IggyIndex) -> u64, target: u64, - ) -> Result<Option<IggyIndex>, IggyError> { + ) -> Result<Option<(IggyIndex, u64)>, IggyError> { let mut low = 0u64; let mut high = entry_count; let mut result = None; @@ -208,7 +247,7 @@ impl IggyIndexReader { high = middle; } } - Ok(result) + Ok(result.map(|entry| (entry, low))) } } @@ -271,6 +310,37 @@ mod tests { let _ = std::fs::remove_dir_all(&dir); } + #[compio::test] + async fn offset_lower_bound_with_successor_reports_the_interval_ceiling() { + let (dir, path) = write_index_file(&entries()).await; + let reader = IggyIndexReader::new(&path).await.expect("open index"); + let count = reader.entry_count().await.expect("entry count"); + + let mid = reader + .offset_lower_bound_with_successor(count, 25) + .await + .expect("lookup") + .expect("in range"); + assert_eq!((mid.0.offset, mid.1), (20, Some(30))); + let last = reader + .offset_lower_bound_with_successor(count, 35) + .await + .expect("lookup") + .expect("in range"); + assert_eq!( + (last.0.offset, last.1), + (30, None), + "the last entry has no successor to bound its interval", + ); + let below_range = reader + .offset_lower_bound_with_successor(count, 5) + .await + .expect("lookup"); + assert!(below_range.is_none()); + + let _ = std::fs::remove_dir_all(&dir); + } + #[compio::test] async fn timestamp_lower_bound_on_file_returns_predecessor() { let (dir, path) = write_index_file(&entries()).await; diff --git a/core/partitions/src/iggy_partition.rs b/core/partitions/src/iggy_partition.rs index 369db555a..23a9f1b93 100644 --- a/core/partitions/src/iggy_partition.rs +++ b/core/partitions/src/iggy_partition.rs @@ -96,7 +96,7 @@ where B: MessageBus, { consensus: VsrConsensus<B>, - pub log: SegmentedLog<PartitionJournal<PartitionJournalMemStorage>, PartitionJournalMemStorage>, + pub log: SegmentedLog<PartitionJournal<PartitionJournalMemStorage>>, /// Highest durably persisted offset. pub offset: Arc<AtomicU64>, /// Highest offset assigned to prepares that may still only live in the in-memory journal. diff --git a/core/partitions/src/journal.rs b/core/partitions/src/journal.rs index 905ae6119..ac5f1909f 100644 --- a/core/partitions/src/journal.rs +++ b/core/partitions/src/journal.rs @@ -929,7 +929,7 @@ where } } -impl Journal<PartitionJournalMemStorage> for PartitionJournal<PartitionJournalMemStorage> { +impl Journal for PartitionJournal<PartitionJournalMemStorage> { type Header = PrepareHeader; type Entry = JournalBuffer; #[rustfmt::skip] diff --git a/core/partitions/src/log.rs b/core/partitions/src/log.rs index fc3a7744b..6c43176f3 100644 --- a/core/partitions/src/log.rs +++ b/core/partitions/src/log.rs @@ -21,7 +21,7 @@ use crate::messages_writer::MessagesWriter; use crate::poll_plan::SealedSegmentHandle; use crate::segment::Segment; use iggy_common::IggyByteSize; -use journal::{Journal, Storage}; +use journal::Journal; use ringbuffer::AllocRingBuffer; use server_common::SegmentStorage; use std::collections::VecDeque; @@ -65,10 +65,9 @@ pub struct JournalState<J> { pub info: JournalInfo, } -impl<J, S> Journal<S> for JournalState<J> +impl<J> Journal for JournalState<J> where - S: Storage, - J: Journal<S>, + J: Journal, { type Header = J::Header; type Entry = J::Entry; @@ -135,18 +134,12 @@ impl<J: Default> Default for JournalState<J> { } } -// TODO: Structure this better, the segmented log does not need to be generic over S, the Journal needs. - -// This struct aliases in terms of the code contained the `SegmentedLog` from `core/server/src/streaming/partitions/log.rs`. -// The only difference is the `Journal` generic, we use different trait. #[derive(Debug)] -pub struct SegmentedLog<J, S> +pub struct SegmentedLog<J> where - S: Storage, - J: Debug + Journal<S>, + J: Debug + Journal, { journal: JournalState<J>, - _pd: std::marker::PhantomData<S>, // Ring buffer tracking recently accessed segment indices for cleanup optimization. // A background task uses this to identify and close file descriptors for unused segments. _access_map: AllocRingBuffer<usize>, @@ -168,15 +161,13 @@ where sealed_lru: VecDeque<u64>, } -impl<J, S> Default for SegmentedLog<J, S> +impl<J> Default for SegmentedLog<J> where - S: Storage, - J: Debug + Default + Journal<S>, + J: Debug + Default + Journal, { fn default() -> Self { Self { journal: JournalState::default(), - _pd: std::marker::PhantomData, _access_map: AllocRingBuffer::with_capacity_power_of_2(ACCESS_MAP_CAPACITY), _cache: (), segments: Vec::with_capacity(SEGMENTS_CAPACITY), @@ -190,10 +181,9 @@ where } } -impl<J, S> SegmentedLog<J, S> +impl<J> SegmentedLog<J> where - S: Storage, - J: Debug + Journal<S>, + J: Debug + Journal, { pub const fn has_segments(&self) -> bool { !self.segments.is_empty() @@ -277,6 +267,7 @@ where handle.tracked.set(false); *handle.fd.borrow_mut() = None; *handle.index.borrow_mut() = None; + handle.offset_cursor.set(None); } self.sealed_lru.clear(); } @@ -431,10 +422,9 @@ where } } -impl<J, S> SegmentedLog<J, S> +impl<J> SegmentedLog<J> where - S: Storage, - J: Debug + Journal<S>, + J: Debug + Journal, { pub const fn journal_mut(&mut self) -> &mut JournalState<J> { &mut self.journal @@ -450,8 +440,7 @@ mod tests { use super::*; use crate::journal::{PartitionJournal, PartitionJournalMemStorage}; - type TestLog = - SegmentedLog<PartitionJournal<PartitionJournalMemStorage>, PartitionJournalMemStorage>; + type TestLog = SegmentedLog<PartitionJournal<PartitionJournalMemStorage>>; /// Push a sealed segment with a resident (index-filled) read handle and /// return a clone of that handle, standing in for an in-flight poll's clone. diff --git a/core/partitions/src/poll_plan.rs b/core/partitions/src/poll_plan.rs index 3b19f0dbb..03a74e2da 100644 --- a/core/partitions/src/poll_plan.rs +++ b/core/partitions/src/poll_plan.rs @@ -104,6 +104,27 @@ pub struct SealedSegmentReadState { /// `SEALED_READ_STATE_CAP` budget never counts. Set on touch, cleared on /// evict; plain `Cell`, all access is same-thread (`Rc` handle). pub(crate) tracked: Cell<bool>, + /// One-slot memo of the last file-backed offset resolution, so a + /// sequentially advancing consumer re-polling the same segment resolves + /// inside `[offset, valid_until)` with zero index-file reads. Only the + /// too-large-to-materialize index path consults it (a resident + /// [`Self::index`] already resolves in memory). Sealed segments are + /// immutable, so the memo cannot go stale; the one exception is a purge + /// recreating the same paths, which wipes this slot with the others. + /// Timestamp polls bypass it. + pub(crate) offset_cursor: Cell<Option<SealedOffsetCursor>>, +} + +/// See [`SealedSegmentReadState::offset_cursor`]. +#[derive(Debug, Clone, Copy)] +pub struct SealedOffsetCursor { + /// Offset of the resolved index entry (interval floor, inclusive). + pub(crate) offset: u64, + /// The successor index entry's offset (interval ceiling, exclusive); + /// `u64::MAX` when the resolved entry is the segment's last. + pub(crate) valid_until: u64, + /// Start byte the whole interval resolves to. + pub(crate) position: u64, } pub type SealedSegmentHandle = Rc<SealedSegmentReadState>; @@ -640,9 +661,6 @@ impl DiskReadPlan { query: MessageLookup, partition_dir: &str, ) -> Option<u64> { - // TODO: a per-consumer cursor hint (the previous sealed poll's resolved - // position) could seed this so a sequentially advancing consumer skips - // the sparse-index lookup on repeated polls of the same segment. let handle = segment.read_state.as_ref()?; // Cache hit: resolve under a short borrow, never across the await. let cached = handle @@ -653,6 +671,16 @@ impl DiskReadPlan { if let Some(resolved) = cached { return resolved; } + // No resident index, so this is the file-backed path a sequential + // consumer would otherwise binary-search on disk every poll: answer + // from the memoized interval when the query lands inside it. + if let (MessageLookup::Offset { offset, .. }, Some(cursor)) = + (query, handle.offset_cursor.get()) + && offset >= cursor.offset + && offset < cursor.valid_until + { + return Some(cursor.position); + } let path = format!("{partition_dir}/{:0>20}.index", segment.start_offset); let reader = match IggyIndexReader::new(&path).await { Ok(reader) => reader, @@ -682,7 +710,20 @@ impl DiskReadPlan { } let looked_up = match query { MessageLookup::Offset { offset, .. } => { - reader.offset_lower_bound(entry_count, offset).await + match reader + .offset_lower_bound_with_successor(entry_count, offset) + .await + { + Ok(resolved) => Ok(resolved.map(|(entry, successor_offset)| { + handle.offset_cursor.set(Some(SealedOffsetCursor { + offset: entry.offset, + valid_until: successor_offset.unwrap_or(u64::MAX), + position: entry.position, + })); + entry + })), + Err(error) => Err(error), + } } MessageLookup::Timestamp { timestamp, .. } => { reader.timestamp_lower_bound(entry_count, timestamp).await @@ -980,8 +1021,97 @@ struct ChunkWalk { #[cfg(test)] mod tests { use super::*; + use crate::iggy_index::IggyIndex; + use compio::io::AsyncWriteAtExt; use server_common::iobuf::Owned; + /// Write a sealed-segment index file too large to materialize + /// (`entry_count * IGGY_INDEX_SIZE > SEALED_INDEX_RESIDENT_MAX_BYTES`), so + /// `resolve_sealed_start` takes the file-backed lookup path the offset + /// cursor memoizes. Entry `i` maps offset `i * 10` to position `i * 100`. + async fn write_oversized_index(dir: &std::path::Path, start_offset: u64) -> u64 { + let entry_count = + SEALED_INDEX_RESIDENT_MAX_BYTES / crate::iggy_index::IGGY_INDEX_SIZE as u64 + 1; + let mut bytes = Vec::with_capacity( + usize::try_from(entry_count).unwrap() * crate::iggy_index::IGGY_INDEX_SIZE, + ); + for i in 0..entry_count { + bytes.extend_from_slice(&crate::iggy_index::IggyIndexCache::serialize( + &IggyIndex::new(i * 10, i + 1, i * 100), + )); + } + let path = format!("{}/{:0>20}.index", dir.display(), start_offset); + let mut file = compio::fs::File::create(&path).await.expect("create index"); + let (written, _) = file.write_all_at(bytes, 0).await.into(); + written.expect("write index"); + file.sync_all().await.expect("sync index"); + entry_count + } + + fn offset_query(offset: u64) -> MessageLookup { + MessageLookup::Offset { + offset, + count: 1, + ceiling: u64::MAX, + } + } + + #[compio::test] + async fn sealed_offset_cursor_answers_in_interval_polls_without_the_index_file() { + let dir = std::env::temp_dir().join(format!( + "iggy-poll-cursor-{}-{}", + std::process::id(), + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .expect("system clock after epoch") + .as_nanos(), + )); + compio::fs::create_dir_all(&dir).await.expect("create dir"); + write_oversized_index(&dir, 0).await; + + let handle: SealedSegmentHandle = Rc::new(SealedSegmentReadState::default()); + let segment = DiskSegment { + start_offset: 0, + persisted: u64::MAX, + read_state: Some(Rc::clone(&handle)), + }; + let plan = DiskReadPlan { + partition_dir: PartitionDirResolution::Resolved(dir.display().to_string()), + segments: Vec::new(), + start_position: 0, + namespace_raw: 0, + validate_checksum: false, + }; + let partition_dir = dir.display().to_string(); + + // First poll pays the on-file lookup and memoizes entry 2's interval + // [20, 30): offset 25 resolves to entry 2 (offset 20 -> position 200). + let first = plan + .resolve_sealed_start(&segment, offset_query(25), &partition_dir) + .await; + assert_eq!(first, Some(200)); + let cursor = handle.offset_cursor.get().expect("cursor memoized"); + assert_eq!( + (cursor.offset, cursor.valid_until, cursor.position), + (20, 30, 200), + ); + + // Delete the index file: an in-interval re-poll must still resolve + // (proof the cursor answered with zero index-file reads)... + std::fs::remove_dir_all(&dir).expect("remove dir"); + let in_interval = plan + .resolve_sealed_start(&segment, offset_query(29), &partition_dir) + .await; + assert_eq!(in_interval, Some(200)); + + // ...while an offset past the interval misses the cursor, reaches for + // the (now gone) file, and falls back to the byte-0 scan. + let past_interval = plan + .resolve_sealed_start(&segment, offset_query(30), &partition_dir) + .await; + assert_eq!(past_interval, None); + } + fn non_empty_fragments() -> PollFragments<4096> { let mut fragments = PollFragments::new(); fragments.push(crate::types::Fragment::whole( diff --git a/core/server/src/auth.rs b/core/server/src/auth.rs index f752df655..589214852 100644 --- a/core/server/src/auth.rs +++ b/core/server/src/auth.rs @@ -70,7 +70,7 @@ pub(crate) fn verify_login_credentials<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -116,7 +116,7 @@ pub(crate) fn verify_pat_credentials<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -134,7 +134,7 @@ pub(crate) fn verify_pat_credentials_with_expiry<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -191,7 +191,7 @@ pub(crate) async fn complete_login_register<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -302,7 +302,7 @@ pub(crate) async fn surface_login_failure<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -340,7 +340,7 @@ async fn send_login_transient_reply<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { diff --git a/core/server/src/bootstrap.rs b/core/server/src/bootstrap.rs index f3cf877ce..f2c5a9d11 100644 --- a/core/server/src/bootstrap.rs +++ b/core/server/src/bootstrap.rs @@ -204,7 +204,7 @@ pub fn wire_shell_handlers<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -790,6 +790,13 @@ pub fn bootstrap( // Shared metadata-group view: written by shard 0's publisher task, read by // every shard's cluster-metadata roster so leader marking works off-shard. let metadata_view = Arc::new(AtomicU64::new(crate::cluster_meta::METADATA_VIEW_UNKNOWN)); + // Every shard's metric handles, minted before the threads spawn: each + // shard bumps its own entry, and shard 0's HTTP scrape endpoint registers + // the whole set (counters are Arc-backed, so cross-thread reads see the + // owning shard's bumps). + let shard_metrics_all: Vec<ShardMetrics> = (0..shards_count) + .map(|_| ShardMetrics::for_shard()) + .collect(); for (idx, assignment) in assignments.into_iter().enumerate() { #[allow(clippy::cast_possible_truncation)] let shard_id = idx as u16; @@ -820,6 +827,7 @@ pub fn bootstrap( }; let metadata_view_for_shard = Arc::clone(&metadata_view); + let shard_metrics_for_shard = shard_metrics_all.clone(); let handle = match thread::Builder::new() .name(format!("shard-{shard_id}")) .spawn(move || -> Result<(), ServerError> { @@ -836,6 +844,7 @@ pub fn bootstrap( barrier_for_shard, owner_table_for_shard, metadata_view_for_shard, + shard_metrics_for_shard, ) }) { Ok(handle) => handle, @@ -899,6 +908,7 @@ fn run_shard_thread( barrier: BootstrapBarrier, owner_table: Arc<ReplicaOwnerTable>, metadata_view: Arc<AtomicU64>, + shard_metrics_all: Vec<ShardMetrics>, ) -> Result<(), ServerError> { // Armed for the whole thread body: a post-spawn error `?` or a panic // unwind here must flip `shutdown_flag` so sibling watchdogs drive @@ -940,6 +950,7 @@ fn run_shard_thread( barrier, owner_table, metadata_view, + shard_metrics_all, )) .await }); @@ -967,6 +978,7 @@ async fn shard_main( barrier: BootstrapBarrier, owner_table: Arc<ReplicaOwnerTable>, metadata_view: Arc<AtomicU64>, + shard_metrics_all: Vec<ShardMetrics>, ) -> Result<(), ServerError> { let topology = resolve_tcp_topology(config, replica_id)?; let bus = Rc::new(IggyMessageBus::with_config_and_owner_table( @@ -1142,7 +1154,7 @@ async fn shard_main( // margin (config validation keeps journal_slots >= 4x this). metadata.set_checkpoint_margin(config.metadata.checkpoint_margin()); - let shard_metrics = ShardMetrics::for_shard(); + let shard_metrics = shard_metrics_all[usize::from(shard_id)].clone(); // Notifier install deferred until after tick handler wires below. let senders_for_notifier = senders.clone(); let metrics_for_notifier = shard_metrics.clone(); @@ -1405,6 +1417,7 @@ async fn shard_main( accepted_replica, dialed_replica, accepted_client, + &shard_metrics_all, ) .await { @@ -2970,6 +2983,7 @@ async fn start_tcp_runtime( accepted_replica: AcceptedReplicaFn, dialed_replica: DialedReplicaFn, accepted_clients: LocalClientAcceptFns, + shard_metrics_all: &[ShardMetrics], ) -> Result<(), ServerError> { if config.tcp.enabled && !config.tcp.tls.enabled { start_via_replica_io( @@ -3015,6 +3029,7 @@ async fn start_tcp_runtime( &config.cluster, Arc::clone(&config.system), self_ports, + shard_metrics_all, ) .await?; } diff --git a/core/server/src/consumer_group.rs b/core/server/src/consumer_group.rs index 5eebb3eff..00a6df558 100644 --- a/core/server/src/consumer_group.rs +++ b/core/server/src/consumer_group.rs @@ -67,7 +67,7 @@ pub(crate) async fn maybe_rewrite_consumer_group_request<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -122,7 +122,7 @@ async fn gather_in_flight<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -212,7 +212,7 @@ pub(crate) fn maybe_rewrite_consumer_offset_request<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -266,7 +266,7 @@ fn resolve_group_offset_id<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { diff --git a/core/server/src/dispatch.rs b/core/server/src/dispatch.rs index d3821b15b..429c46fd7 100644 --- a/core/server/src/dispatch.rs +++ b/core/server/src/dispatch.rs @@ -141,7 +141,7 @@ pub(crate) fn make_client_request_handler<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -203,7 +203,7 @@ pub(crate) fn make_partition_read_handler<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -295,7 +295,7 @@ fn spawn_poll_io<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -351,7 +351,7 @@ fn submit_auto_commit<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -456,7 +456,7 @@ pub(crate) fn make_deferred_replica_message_handler<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -478,7 +478,7 @@ pub(crate) fn make_deferred_client_request_handler<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -543,7 +543,7 @@ pub(crate) fn make_metadata_submit_handler<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -700,7 +700,7 @@ fn enqueue_client_request<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -740,7 +740,7 @@ async fn drain_client_requests<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -944,7 +944,7 @@ async fn send_pre_consensus_deny<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -982,7 +982,7 @@ async fn handle_client_request<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1407,7 +1407,7 @@ async fn handle_get_personal_access_tokens<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1434,7 +1434,7 @@ async fn handle_get_me<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1477,7 +1477,7 @@ pub(crate) async fn dispatch_partition_request<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1615,7 +1615,7 @@ async fn handle_non_replicated_request<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1772,7 +1772,7 @@ async fn handle_default_non_replicated<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1850,7 +1850,7 @@ async fn handle_get_snapshot<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1932,7 +1932,7 @@ async fn send_non_replicated_bytes<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1969,7 +1969,7 @@ async fn send_unauthenticated_eviction<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2013,7 +2013,7 @@ pub(crate) async fn run_heartbeat_verifier<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2076,7 +2076,7 @@ async fn evict_stale_client<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2130,7 +2130,7 @@ async fn handle_poll_messages<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2251,7 +2251,7 @@ async fn handle_get_consumer_offset<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2337,7 +2337,7 @@ async fn handle_sync_consumer_group<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2391,7 +2391,7 @@ async fn send_empty_partition_reply<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2439,7 +2439,7 @@ async fn wait_for_partition_routable<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2487,7 +2487,7 @@ pub(crate) fn resolve_poll_request<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2553,7 +2553,7 @@ pub(crate) fn resolve_consumer_offset_request<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2638,7 +2638,7 @@ async fn answer_forwarded_register<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2730,7 +2730,7 @@ async fn submit_register_local_or_forward<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2892,7 +2892,7 @@ async fn answer_forwarded_logout<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -2938,7 +2938,7 @@ async fn submit_logout_local_or_forward<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3088,7 +3088,7 @@ pub(crate) async fn submit_register_on_owner<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3119,7 +3119,7 @@ pub(crate) async fn submit_logout_on_owner<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3157,7 +3157,7 @@ async fn handle_delete_segments_request<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3291,7 +3291,7 @@ pub(crate) async fn resolve_delete_segments_truncate<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3411,7 +3411,7 @@ fn submit_disconnect_logout<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3453,7 +3453,7 @@ pub(crate) async fn submit_client_request_on_owner<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3482,7 +3482,7 @@ async fn handle_logout_request<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3577,7 +3577,7 @@ fn ensure_transport_connection<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3598,7 +3598,7 @@ async fn handle_login_register_request<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3768,7 +3768,7 @@ pub(crate) async fn send_login_eviction<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -3806,7 +3806,7 @@ pub(crate) fn upgrade_shard_handle<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { diff --git a/core/server/src/dispatch/authz.rs b/core/server/src/dispatch/authz.rs index 73cfc09f8..5e8790e7a 100644 --- a/core/server/src/dispatch/authz.rs +++ b/core/server/src/dispatch/authz.rs @@ -70,7 +70,7 @@ pub(super) fn authorize_partition_op<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -147,7 +147,7 @@ pub(super) async fn send_deny_reply<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -180,7 +180,7 @@ pub(super) async fn send_unbound_deny_reply<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -210,7 +210,7 @@ pub(super) fn authorize_uid<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -237,7 +237,7 @@ pub(super) fn authorize_partition_read<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -272,7 +272,7 @@ pub(super) fn authorize_default_read<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -345,7 +345,7 @@ pub(super) async fn send_non_replicated_deny<B, MJ, S, SB>( ) where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -386,7 +386,7 @@ fn gate_user_scoped<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -422,7 +422,7 @@ fn gate_stream_scoped<T: WireDecode, B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -451,7 +451,7 @@ fn gate_topic_scoped<T: WireDecode, B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -476,7 +476,7 @@ fn resolve_stream_scope<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -498,7 +498,7 @@ fn resolve_topic_scope<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { diff --git a/core/server/src/http.rs b/core/server/src/http.rs index 28c6ebeb0..8f49f3023 100644 --- a/core/server/src/http.rs +++ b/core/server/src/http.rs @@ -104,6 +104,7 @@ pub async fn start( cluster: &ClusterConfig, system_config: Arc<ServerSystemConfig>, self_ports: TransportPorts, + shard_metrics_all: &[shard::metrics::ShardMetrics], ) -> Result<(), ServerError> { // In cluster mode with no configured JWT secret the signing key derives // from the cluster PSK, so a bearer minted on any node verifies on every @@ -165,7 +166,7 @@ pub async fn start( max_tokens_per_user, in_flight_writes: Cell::new(0), forward, - metrics: metrics::HttpMetrics::init(), + metrics: metrics::HttpMetrics::init(shard_metrics_all), })); let router = router( state, diff --git a/core/server/src/http/metrics.rs b/core/server/src/http/metrics.rs index 8261d6dc6..3fdb1b423 100644 --- a/core/server/src/http/metrics.rs +++ b/core/server/src/http/metrics.rs @@ -53,7 +53,7 @@ pub(in crate::http) struct HttpMetrics { } impl HttpMetrics { - pub(in crate::http) fn init() -> Self { + pub(in crate::http) fn init(shard_metrics_all: &[shard::metrics::ShardMetrics]) -> Self { let mut registry = Registry::default(); let http_requests = Counter::default(); let streams = Gauge::default(); @@ -79,6 +79,19 @@ impl HttpMetrics { registry.register("messages", "total count of messages", messages.clone()); registry.register("users", "total count of users", users.clone()); registry.register("clients", "total count of clients", clients.clone()); + // Every shard's drop / reconcile / partition counters, one + // `shard`-labelled sub-registry per shard so series stay per-shard + // without a `shard_id` label in the counter label sets (see + // `shard::metrics::FrameDropLabel`). The counters are Arc-backed: + // each shard bumps its own handle on its own thread, and the scrape + // on shard 0 reads the shared atomics. + for (shard_id, shard_metrics) in shard_metrics_all.iter().enumerate() { + let sub_registry = registry.sub_registry_with_label(( + std::borrow::Cow::Borrowed("shard"), + std::borrow::Cow::Owned(shard_id.to_string()), + )); + shard_metrics.register(sub_registry); + } Self { registry, http_requests, @@ -209,6 +222,23 @@ fn gauge_value(count: u64) -> i64 { #[cfg(test)] mod tests { use super::*; + use shard::metrics::{ShardMetrics, frame_drop_reason, frame_drop_variant}; + + #[test] + fn shard_metrics_land_in_the_exposition_under_a_shard_label() { + let shard_metrics = ShardMetrics::for_shard(); + shard_metrics.record_frame_drop(frame_drop_variant::CONSENSUS, frame_drop_reason::FULL); + let metrics = HttpMetrics::init(&[shard_metrics]); + let output = metrics.formatted_output(); + assert!( + output.contains("frame_drops_total"), + "shard frame-drop counter missing from exposition:\n{output}" + ); + assert!( + output.contains(r#"shard="0""#), + "shard label missing from exposition:\n{output}" + ); + } const PARITY_METRIC_NAMES: [&str; 8] = [ "http_requests", @@ -230,7 +260,7 @@ mod tests { #[test] fn formatted_output_exposes_every_parity_metric() { - let metrics = HttpMetrics::init(); + let metrics = HttpMetrics::init(&[]); let output = metrics.formatted_output(); for name in PARITY_METRIC_NAMES { assert!( @@ -246,7 +276,7 @@ mod tests { #[test] fn scraped_values_land_in_the_exposition() { - let metrics = HttpMetrics::init(); + let metrics = HttpMetrics::init(&[]); metrics.streams.set(1); metrics.topics.set(2); metrics.partitions.set(3); diff --git a/core/server/src/responses.rs b/core/server/src/responses.rs index 6b7fc4ec5..86dbce69f 100644 --- a/core/server/src/responses.rs +++ b/core/server/src/responses.rs @@ -114,7 +114,7 @@ pub(crate) fn build_get_personal_access_tokens_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -151,7 +151,7 @@ pub(crate) fn build_get_me_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -225,7 +225,7 @@ pub(crate) fn connected_client_to_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -267,7 +267,7 @@ fn fence_group_offset<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -311,7 +311,7 @@ fn fence_and_resolve_offset_namespace<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -335,7 +335,7 @@ pub(crate) fn resolve_partition_request_namespace<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -402,7 +402,7 @@ fn resolve_send_messages_namespace<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -440,7 +440,7 @@ pub(crate) fn resolve_partition_namespace<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -506,7 +506,7 @@ pub(crate) fn build_non_replicated_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -607,7 +607,7 @@ fn build_consumer_group_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -632,7 +632,7 @@ fn build_consumer_groups_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -662,7 +662,7 @@ fn build_cluster_metadata_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -710,7 +710,7 @@ fn build_stats_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -910,7 +910,7 @@ fn build_get_stream_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -939,7 +939,7 @@ fn build_get_streams_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1090,7 +1090,7 @@ fn build_get_users_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1111,7 +1111,7 @@ fn build_get_user_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1155,7 +1155,7 @@ fn build_get_topic_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1192,7 +1192,7 @@ fn build_get_topics_response<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1225,7 +1225,7 @@ fn ensure_topic_exists<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { @@ -1603,7 +1603,7 @@ pub(crate) fn current_metadata_commit<B, MJ, S, SB>(shard: &Rc<ShellShard<B, MJ, where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { diff --git a/core/server/src/users.rs b/core/server/src/users.rs index 802d092e3..63368cd5b 100644 --- a/core/server/src/users.rs +++ b/core/server/src/users.rs @@ -65,7 +65,7 @@ pub(crate) fn maybe_rewrite_user_password_request<B, MJ, S, SB>( where B: ShellBus, MJ: JournalHandle + 'static, - MJ::Target: Journal<MJ::Storage, Entry = Message<PrepareHeader>, Header = PrepareHeader>, + MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, S: 'static, SB: SuperblockStore + 'static, { diff --git a/core/shard/src/lib.rs b/core/shard/src/lib.rs index a51dfa25f..001a49bcd 100644 --- a/core/shard/src/lib.rs +++ b/core/shard/src/lib.rs @@ -2450,11 +2450,8 @@ where where B: MessageBus + 'static, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StateMachine< Input = Message<PrepareHeader>, Output = metadata::stm::result::ApplyReply, @@ -3308,11 +3305,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StateMachine< Input = Message<PrepareHeader>, Output = metadata::stm::result::ApplyReply, @@ -3330,11 +3324,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StateMachine< Input = Message<PrepareHeader>, Output = metadata::stm::result::ApplyReply, @@ -3352,11 +3343,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StateMachine< Input = Message<PrepareHeader>, Output = metadata::stm::result::ApplyReply, @@ -3393,11 +3381,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StateMachine< Input = Message<PrepareHeader>, Output = metadata::stm::result::ApplyReply, @@ -3566,11 +3551,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { let header = *msg.header(); let planes = self.plane.inner(); @@ -3616,11 +3598,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: MetadataStm, { let header = *msg.header(); @@ -3706,11 +3685,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: MetadataStm, { let header = *msg.header(); @@ -3914,11 +3890,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: MetadataStm, { let header = *msg.header(); @@ -4018,11 +3991,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { let header = *msg.header(); let planes = self.plane.inner(); @@ -4070,11 +4040,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StreamsFrontend, { let header = *msg.header(); @@ -4309,11 +4276,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { tracing::debug!( shard = self.id, @@ -4483,11 +4447,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: MetadataStm, { let header = *msg.header(); @@ -4854,11 +4815,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { let partitions = self.plane.partitions(); let started = { @@ -5027,11 +4985,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: MetadataStm, { let metadata = self.plane.metadata(); @@ -5148,11 +5103,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: MetadataStm, { let metadata = self.plane.metadata(); @@ -5429,11 +5381,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: MetadataStm, { let header = *msg.header(); @@ -5569,11 +5518,8 @@ where B: MessageBus + 'static, T: ShardsTable, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: RestorableMetadataStm, { /// Alloc cap per artifact: a corrupt length field must not OOM the @@ -6014,11 +5960,8 @@ where where B: MessageBus + 'static, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: RestorableMetadataStm, T: ShardsTable, { @@ -6068,11 +6011,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: RestorableMetadataStm, { let planes = self.plane.inner(); @@ -6290,27 +6230,25 @@ where /// Tick partition consensuses. Loop partitions. No partitions-plane journal. #[allow(clippy::future_not_send)] #[allow(clippy::too_many_lines)] - pub async fn tick_partitions(&self) + pub async fn tick_partitions(&self, namespace_scratch: &mut Vec<IggyNamespace>) where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { + debug_assert!( + namespace_scratch.is_empty(), + "namespace_scratch must be empty on entry", + ); let partitions = self.plane.partitions(); let repair_retry_ticks = self.repair_retry_ticks.get(); // Fan out over every group (each partition's heartbeat/retransmit timer // must advance), so the keyed single-namespace lookup the control-frame // handlers use does not apply here. The namespaces are snapshotted into - // an owned Vec so no partitions-plane borrow is held across the tick - // `.await`. - // TODO(hubcio): reuse the pump's `namespace_scratch` (as - // `process_loopback` does) to drop this per-tick alloc; a quiet cluster - // still pays one Vec per heartbeat. - let namespaces: Vec<_> = partitions.namespaces().copied().collect(); + // the pump's owned scratch (as `process_loopback` does) so no + // partitions-plane borrow is held across the tick `.await`. + namespace_scratch.extend(partitions.namespaces().copied()); // Pre-pass: issue every group's pending superblock persist // CONCURRENTLY. A cluster-wide view change makes every group on @@ -6321,7 +6259,7 @@ where // its store, lock, and failure bookkeeping, all behind `&self`), // and the per-group loop below re-checks the gate on its lock-free // fast path, so gating semantics are unchanged. - let pending_persists: Vec<_> = namespaces + let pending_persists: Vec<_> = namespace_scratch .iter() .copied() .filter(|namespace| { @@ -6362,7 +6300,7 @@ where // already accepts. let mut transfers_inflight: Option<usize> = None; - for namespace in namespaces { + for namespace in namespace_scratch.drain(..) { let Some(partition) = partitions.get_by_ns(&namespace) else { continue; }; @@ -8218,11 +8156,8 @@ where where B: MessageBus, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: StateMachine< Input = Message<PrepareHeader>, Output = metadata::stm::result::ApplyReply, @@ -8348,11 +8283,7 @@ where B: MessageBus, P: Pipeline<Entry = consensus::PipelineEntry>, J: JournalHandle, - <J as JournalHandle>::Target: Journal< - <J as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <J as JournalHandle>::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { tracing::info!( view = consensus.view(), @@ -8446,11 +8377,7 @@ where B: MessageBus, P: Pipeline<Entry = consensus::PipelineEntry>, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { if !consensus.local_dvc_suffix_stale() { return; @@ -8649,11 +8576,7 @@ fn build_metadata_dvc_suffix<J>( ) -> DvcSuffix where J: JournalHandle, - <J as JournalHandle>::Target: Journal< - <J as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <J as JournalHandle>::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { let Some(journal) = journal else { return DvcSuffix::empty(); @@ -8954,11 +8877,7 @@ async fn dispatch_vsr_actions<B, P, J>( B: MessageBus, P: Pipeline<Entry = consensus::PipelineEntry>, J: JournalHandle, - <J as JournalHandle>::Target: Journal< - <J as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <J as JournalHandle>::Target: Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, { use std::mem::size_of; diff --git a/core/shard/src/metrics.rs b/core/shard/src/metrics.rs index 2df15c357..99f70a224 100644 --- a/core/shard/src/metrics.rs +++ b/core/shard/src/metrics.rs @@ -31,18 +31,16 @@ //! this shard's own inbox was refused. //! //! The counter uses atomic interior mutability, safe to bump from `!Send` -//! compio reactor contexts. Each shard owns its own instance. It is not -//! yet exposed through a scrape endpoint; every drop site also logs via -//! `tracing`, so the counter is a structured complement to those logs -//! until a per-shard exporter lands. -//! -//! TODO(hubcio): register `frame_drops_total` with a per-shard prometheus -//! exporter so the counter is scrape-able; until then the `tracing` -//! drop-site logs are the only alertable signal. +//! compio reactor contexts. Each shard owns its own instance, and the server +//! exposes every shard's instance through the `[http.metrics]` scrape +//! endpoint via [`ShardMetrics::register`] (one `shard`-labelled +//! sub-registry per shard); every drop site also logs via `tracing`. use prometheus_client::encoding::EncodeLabelSet; use prometheus_client::metrics::counter::Counter; use prometheus_client::metrics::family::Family; +use prometheus_client::registry::Registry; +use std::sync::{Arc, OnceLock}; /// Label for `frame_drops_total`. /// @@ -71,9 +69,9 @@ pub struct FrameDropLabel { /// closure fails. Unlike `CONSENSUS` drops (which VSR retransmit /// recovers), a `FORWARD_CLIENT_SEND` drop is terminal: the client never /// receives the reply and request / response semantics break above the -/// bus. Operators should alert on the drop-site `tracing` logs (the -/// counter is not scrape-able yet, see the module doc) and size -/// `inbox_capacity` for the worst-case cross-shard reply burst. +/// bus. Operators should alert on this pair in the scrape (backed by the +/// drop-site `tracing` logs) and size `inbox_capacity` for the worst-case +/// cross-shard reply burst. /// `FORWARD_REPLICA_SEND` is the symmetric variant for replica forwards; /// VSR retransmit covers its loss so it stays informational. /// @@ -131,10 +129,10 @@ pub mod frame_drop_reason { pub const PARK_DROPPED: &str = "park_dropped"; } -// Minted in full, so 7 x 7 includes pairs no drop site produces: -// `park_overflow` and `park_dropped` pair only with `PARTITION`, leaving 12 -// unreachable. Free while nothing scrapes these (module `TODO(hubcio)`); mint -// per drop site once a registry lands, so the scrape carries no permanent zeroes. +// The tables only index the lazy fast-path cache below; a `{variant, reason}` +// pair enters the `Family` (and therefore the scrape) the first time a drop +// site actually produces it, so the unreachable corners of the 7 x 7 cross +// product never appear as permanent zero-valued series. const VARIANT_COUNT: usize = 7; const REASON_COUNT: usize = 7; @@ -169,11 +167,14 @@ fn reason_index(s: &str) -> Option<usize> { /// Per-shard metric handles. /// /// Cheap to clone (`Arc` of a `Family` under the hood). Each shard owns -/// one instance produced by [`ShardMetrics::for_shard`]. The -/// `VARIANT_COUNT * REASON_COUNT` cross product of `Counter`s is minted -/// at construction so the drop-site hot path never re-enters -/// `Family::get_or_create` (which acquires a `RwLock` read guard per -/// drop and stalls under VSR retransmit / drop-burst storms). +/// one instance produced by [`ShardMetrics::for_shard`]. A known +/// `{variant, reason}` pair's `Counter` is minted through +/// `Family::get_or_create` (a `RwLock` read guard, too dear per drop +/// under VSR retransmit / drop-burst storms) exactly once, on the pair's +/// first drop, and cached; later drops are an array index + atomic +/// increment. Lazy rather than pre-minted so a pair no drop site produces +/// never enters the family, keeping the scrape free of permanent +/// zero-valued series. /// /// `partitions_materialised_total` / `partitions_removed_total` / /// `partitions_reconcile_failures_total` are simple unlabelled counters @@ -182,7 +183,7 @@ fn reason_index(s: &str) -> Option<usize> { #[derive(Clone)] pub struct ShardMetrics { frame_drops_total: Family<FrameDropLabel, Counter>, - cached_counters: [[Counter; REASON_COUNT]; VARIANT_COUNT], + cached_counters: Arc<[[OnceLock<Counter>; REASON_COUNT]; VARIANT_COUNT]>, partitions_materialised_total: Counter, partitions_removed_total: Counter, partitions_reconcile_failures_total: Counter, @@ -198,23 +199,12 @@ impl ShardMetrics { /// Create a metrics handle for a shard. The handle is per-shard by /// virtue of being constructed once per shard; the shard id does not /// appear in the label set (see [`FrameDropLabel`] doc). - /// - /// All `VARIANT_COUNT * REASON_COUNT` counters are pre-registered - /// with the underlying [`Family`] so the drop-site hot path is a - /// constant-time array index + atomic increment. #[must_use] pub fn for_shard() -> Self { let frame_drops_total: Family<FrameDropLabel, Counter> = Family::default(); - let cached_counters = std::array::from_fn(|v_idx| { - std::array::from_fn(|r_idx| { - frame_drops_total - .get_or_create(&FrameDropLabel { - variant: VARIANTS[v_idx], - reason: REASONS[r_idx], - }) - .clone() - }) - }); + let cached_counters = Arc::new(std::array::from_fn(|_| { + std::array::from_fn(|_| OnceLock::new()) + })); Self { frame_drops_total, cached_counters, @@ -239,7 +229,13 @@ impl ShardMetrics { /// to extend the const tables above. pub fn record_frame_drop(&self, variant: &'static str, reason: &'static str) { if let (Some(v_idx), Some(r_idx)) = (variant_index(variant), reason_index(reason)) { - self.cached_counters[v_idx][r_idx].inc(); + self.cached_counters[v_idx][r_idx] + .get_or_init(|| { + self.frame_drops_total + .get_or_create(&FrameDropLabel { variant, reason }) + .clone() + }) + .inc(); } else { self.frame_drops_total .get_or_create(&FrameDropLabel { variant, reason }) @@ -304,10 +300,6 @@ impl ShardMetrics { /// (delete + recreate recycled the namespace's slab keys). Serving it would /// have written a dead topic's op into the topic that replaced it, so a /// non-zero value is a caught correctness anomaly, not routine churn. - /// - /// Like every counter in this module it is not scrape-able yet (see the - /// module-level `TODO(hubcio)`); the `warn!` at the reject site is what an - /// operator can actually alert on today. pub fn record_partition_frame_rejected_stale(&self) { self.partition_frames_rejected_stale_total.inc(); } @@ -334,6 +326,7 @@ impl ShardMetrics { self.cached_counters .iter() .flatten() + .filter_map(OnceLock::get) .map(prometheus_client::metrics::counter::Counter::get) .sum() } @@ -425,10 +418,70 @@ impl ShardMetrics { #[must_use] pub fn frame_drop_count(&self, variant: &'static str, reason: &'static str) -> u64 { match (variant_index(variant), reason_index(reason)) { - (Some(v_idx), Some(r_idx)) => self.cached_counters[v_idx][r_idx].get(), + (Some(v_idx), Some(r_idx)) => self.cached_counters[v_idx][r_idx] + .get() + .map_or(0, prometheus_client::metrics::counter::Counter::get), _ => 0, } } + + /// Register every metric this handle owns with `registry`, which the + /// server scopes per shard (a `shard`-labelled sub-registry) before the + /// `[http.metrics]` scrape encodes it. Names are registered without the + /// `_total` suffix; the prometheus text exposition appends it for + /// counters. + pub fn register(&self, registry: &mut Registry) { + registry.register( + "frame_drops", + "frames shed instead of delivered, by frame class and refusal reason", + self.frame_drops_total.clone(), + ); + registry.register( + "partitions_materialised", + "partitions materialised by the reconciliation loop", + self.partitions_materialised_total.clone(), + ); + registry.register( + "partitions_removed", + "partitions dropped after their namespace left the committed metadata", + self.partitions_removed_total.clone(), + ); + registry.register( + "partitions_reconcile_failures", + "partition build or delete attempts the reconciler will retry", + self.partitions_reconcile_failures_total.clone(), + ); + registry.register( + "partitions_duplicate_builds_discarded", + "duplicate partition builds discarded by the pump; non-zero is a caught anomaly", + self.partitions_duplicate_builds_discarded_total.clone(), + ); + registry.register( + "partition_transfer_refusals", + "partition state transfers refused by the serving peer", + self.partition_transfer_refusals_total.clone(), + ); + registry.register( + "partition_frames_rejected_stale", + "parked frames rejected for a dead incarnation; non-zero is a caught anomaly", + self.partition_frames_rejected_stale_total.clone(), + ); + registry.register( + "partition_frames_rejected_ahead", + "parked frames rejected for an epoch ahead of the materialised one", + self.partition_frames_rejected_ahead_total.clone(), + ); + registry.register( + "partition_requests_denied_transient", + "partition requests answered with a retriable transient denial", + self.partition_requests_denied_transient_total.clone(), + ); + registry.register( + "partition_repair_serves_deferred_purge", + "partition repair serves or completions deferred until a committed purge applies", + self.partition_repair_serves_deferred_purge_total.clone(), + ); + } } #[cfg(test)] @@ -467,6 +520,28 @@ mod tests { ); } + #[test] + fn unproduced_pairs_never_enter_the_scrape() { + // The lazy fast-path cache must not mint the full variant x reason + // cross product: a pair no drop site produced would otherwise sit in + // every scrape as a permanent zero-valued series. + let metrics = ShardMetrics::for_shard(); + metrics.record_frame_drop(frame_drop_variant::CONSENSUS, frame_drop_reason::FULL); + let mut registry = Registry::default(); + metrics.register(&mut registry); + let mut buffer = String::new(); + prometheus_client::encoding::text::encode(&mut buffer, ®istry) + .expect("scrape encoding succeeds"); + assert!( + buffer.contains(frame_drop_variant::CONSENSUS), + "the produced pair must appear in the scrape", + ); + assert!( + !buffer.contains(frame_drop_reason::PARK_OVERFLOW), + "a pair no drop site produced must not appear in the scrape", + ); + } + #[test] fn cached_counter_aliases_family_entry() { // Cached fast-path counters must point at the same underlying diff --git a/core/shard/src/router.rs b/core/shard/src/router.rs index 0acb78a13..2b82390f2 100644 --- a/core/shard/src/router.rs +++ b/core/shard/src/router.rs @@ -230,11 +230,8 @@ where where B: MessageBus + 'static, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: RestorableMetadataStm, { // Reused across every pump iteration; pre-size to skip the @@ -272,7 +269,7 @@ where // decoupled from the pump again without reintroducing the // partition-ref-across-`.await` UB this fold closed. self.tick_metadata().await; - self.tick_partitions().await; + self.tick_partitions(&mut namespace_scratch).await; // Runs here, not inside `tick_metadata`: that early-returns // on shards without metadata consensus, and partition-plane // offers live on every shard that hosts a serving group -- @@ -367,11 +364,8 @@ where where B: MessageBus + 'static, MJ: JournalHandle, - <MJ as JournalHandle>::Target: Journal< - <MJ as JournalHandle>::Storage, - Entry = Message<PrepareHeader>, - Header = PrepareHeader, - >, + <MJ as JournalHandle>::Target: + Journal<Entry = Message<PrepareHeader>, Header = PrepareHeader>, M: RestorableMetadataStm, { match frame { diff --git a/core/simulator/src/deps.rs b/core/simulator/src/deps.rs index f44ab3ce9..9b7f89b79 100644 --- a/core/simulator/src/deps.rs +++ b/core/simulator/src/deps.rs @@ -168,7 +168,7 @@ impl Drop for JournalAccessGuard<'_> { } #[allow(clippy::future_not_send)] -impl<S: Storage<Buffer = Vec<u8>>> Journal<S> for SimJournal<S> { +impl<S: Storage<Buffer = Vec<u8>>> Journal for SimJournal<S> { type Header = PrepareHeader; type Entry = Message<PrepareHeader>; type HeaderRef<'a> @@ -283,7 +283,6 @@ impl<S: Storage<Buffer = Vec<u8>>> Journal<S> for SimJournal<S> { } impl JournalHandle for SimJournal<MemStorage> { - type Storage = MemStorage; type Target = Self; fn handle(&self) -> &Self::Target {
