This is an automated email from the ASF dual-hosted git repository.
numinnex pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new 565e1a7b2 refactor(shard): drop journal generic, metrics scrape, poll
cursor (#3914)
565e1a7b2 is described below
commit 565e1a7b20ff1a1c58a9ac4eddead3c708cc0c04
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Thu Aug 20 13:08:58 2026 +0200
refactor(shard): drop journal generic, metrics scrape, poll cursor (#3914)
---
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 | 205 ++++++++++---------------------
core/shard/src/metrics.rs | 157 ++++++++++++++++-------
core/shard/src/router.rs | 16 +--
core/simulator/src/deps.rs | 3 +-
21 files changed, 556 insertions(+), 338 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 5f9073ed5..35c1b672b 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,
{
@@ -755,7 +755,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,
{
@@ -926,7 +926,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,
{
@@ -955,7 +955,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,
{
@@ -1106,7 +1106,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,
{
@@ -1127,7 +1127,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,
{
@@ -1171,7 +1171,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,
{
@@ -1208,7 +1208,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,
{
@@ -1241,7 +1241,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,
{
@@ -1640,7 +1640,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..eda20d4dd 100644
--- a/core/shard/src/lib.rs
+++ b/core/shard/src/lib.rs
@@ -1808,9 +1808,11 @@ where
shards_table: T,
partition_consensus: PartitionConsensusConfig<B>,
) -> Self {
- // TODO(hubcio): crossfire's Flavor trait blocks unbounded channels
- // with the current type setup; revisit when crossfire grows an
- // unbounded variant or we replace it.
+ // Placeholder inbox: the simulator delivers frames straight to
+ // `on_message` (see the `shard_count` note below), so nothing ever
+ // sends here and capacity 1 exists only to satisfy the field. The
+ // real inbox is bounded on purpose (`inbox_capacity` is the shard's
+ // backpressure), so no unbounded variant is wanted here either.
let (_tx, inbox) = channel(1);
let nonce_seed = forward_nonce_seed(metadata.consensus.as_ref());
let plane = MuxPlane::new(variadic!(metadata, partitions));
@@ -2450,11 +2452,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 +3307,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 +3326,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 +3345,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 +3383,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 +3553,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 +3600,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 +3687,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 +3892,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 +3993,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 +4042,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 +4278,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 +4449,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 +4817,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 +4987,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 +5105,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 +5383,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 +5520,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 +5962,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 +6013,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 +6232,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 +6261,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 +6302,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 +8158,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 +8285,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 +8379,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 +8578,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 +8879,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 697559246..0432ecc22 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 {