numinnex commented on code in PR #4036:
URL: https://github.com/apache/iggy/pull/4036#discussion_r3932143500


##########
core/server/src/boot/threads.rs:
##########
@@ -333,6 +320,173 @@ impl Drop for ShutdownOnDrop {
     }
 }
 
+/// The single post-shutdown budget, shared by the main thread's shard
+/// joins and shard 0's peer wait.
+///
+/// Both waits are bounded by `system.sharding.shutdown_join_timeout` and
+/// they NEST: shard 0 cannot start waiting for its peers until its own
+/// drain returned, which is already inside the join budget. Arming one
+/// instant on first use, whichever wait gets there first, keeps the two
+/// inside one deadline instead of stacking two full budgets, so a
+/// correct shutdown cannot report shard 0 as wedged.
+pub(in crate::boot) struct ShutdownDeadline {
+    armed: OnceLock<Instant>,
+    budget: Duration,
+}
+
+impl ShutdownDeadline {
+    pub(in crate::boot) const fn new(budget: Duration) -> Self {
+        Self {
+            armed: OnceLock::new(),
+            budget,
+        }
+    }
+
+    /// Time left in the shared budget, arming it on the first call.
+    /// Callers must only reach this once the shutdown flag is set: a
+    /// running server would otherwise start the clock.
+    fn remaining(&self) -> Duration {
+        self.armed
+            .get_or_init(|| Instant::now() + self.budget)
+            .saturating_duration_since(Instant::now())
+    }
+}
+
+/// Peer shards still running: each [`PeerExitGuard`] counts one out,
+/// [`PeerExitWait`] blocks shard 0 until the count is zero.
+///
+/// Shard 0 owns the metadata state machine's only write handle and every
+/// peer reads through handles that stop working the moment it drops, so
+/// the shard that owns the writer must outlive every reader on every exit
+/// but one: [`PeerExitWait::drop`] skips the wait while shard 0 is
+/// panicking, because parking an unwinding thread on a `Condvar` inside
+/// `runtime.block_on` would stall the `io_uring` driver, and a panic in
+/// shard 0's own pump can be exactly what the peers are blocked on. A
+/// peer mid-read then panics too; the first panic is already recorded and
+/// the process is going down either way.
+pub(in crate::boot) struct PeerExitCountdown {
+    running: Mutex<usize>,
+    all_exited: Condvar,
+}
+
+impl PeerExitCountdown {
+    pub(in crate::boot) const fn new(peers: usize) -> Self {
+        Self {
+            running: Mutex::new(peers),
+            all_exited: Condvar::new(),
+        }
+    }
+
+    /// Block until every peer has counted itself out or `timeout` elapses.
+    /// `Err` carries the number of peers still running at the deadline.
+    fn wait(&self, timeout: Duration) -> Result<(), usize> {
+        let (guard, _) = self
+            .all_exited
+            .wait_timeout_while(
+                self.running.lock().unwrap_or_else(PoisonError::into_inner),
+                timeout,
+                |running| *running > 0,
+            )
+            .unwrap_or_else(PoisonError::into_inner);
+        let running = *guard;
+        drop(guard);
+        if running == 0 { Ok(()) } else { Err(running) }
+    }
+
+    fn peer_exited(&self) {
+        let mut running = 
self.running.lock().unwrap_or_else(PoisonError::into_inner);
+        *running = running.saturating_sub(1);
+        if *running == 0 {
+            self.all_exited.notify_all();
+        }
+    }
+
+    /// Count out peers that never spawned. The countdown is sized before the
+    /// spawn loop, so a failed `thread::Builder::spawn` leaves peers whose
+    /// [`PeerExitGuard`] will never exist: without this shard 0 waits out its
+    /// whole join budget for threads that were never there.
+    pub(in crate::boot) fn peers_never_spawned(&self, count: usize) {
+        for _ in 0..count {
+            self.peer_exited();
+        }
+    }
+}
+
+/// Counts one peer shard out of the [`PeerExitCountdown`] on drop.
+///
+/// Held by `run_shard_thread` from before the runtime exists, so it drops
+/// after the runtime and its tasks are gone: past that point nothing on
+/// the peer's thread can still read shard 0's metadata. Drop runs on every
+/// exit path, the error `?` returns and panic unwinds included.
+struct PeerExitGuard {
+    countdown: Arc<PeerExitCountdown>,
+}
+
+impl PeerExitGuard {
+    const fn new(countdown: Arc<PeerExitCountdown>) -> Self {
+        Self { countdown }
+    }
+}
+
+impl Drop for PeerExitGuard {
+    fn drop(&mut self) {
+        self.countdown.peer_exited();
+    }
+}
+
+/// Shard 0's side of the [`PeerExitCountdown`]: blocks on drop until every
+/// peer has exited, so whatever is declared before it outlives every
+/// peer's reads.
+///
+/// Flips the shutdown flag before waiting: a peer parked on its bus token
+/// only starts its drain once the flag is set, and the thread-level
+/// `ShutdownOnDrop` flips it only after `block_on` returns, which is after
+/// this wait. Bounded by what is left of the [`ShutdownDeadline`] that
+/// [`ShardHandles::join_all`] shares, with the same abandon-and-log
+/// policy, so a wedged peer cannot hold shard 0 past the deadline the
+/// main thread is itself counting down.
+pub(in crate::boot) struct PeerExitWait {
+    countdown: Arc<PeerExitCountdown>,
+    shutdown_flag: Arc<AtomicBool>,
+    deadline: Arc<ShutdownDeadline>,
+}
+
+impl PeerExitWait {
+    pub(in crate::boot) const fn new(
+        countdown: Arc<PeerExitCountdown>,
+        shutdown_flag: Arc<AtomicBool>,
+        deadline: Arc<ShutdownDeadline>,
+    ) -> Self {
+        Self {
+            countdown,
+            shutdown_flag,
+            deadline,
+        }
+    }
+}
+
+impl Drop for PeerExitWait {
+    fn drop(&mut self) {
+        // The panic is the fault to report and `ShutdownOnDrop` still
+        // flips the flag for the peers; blocking an unwinding thread here
+        // would only delay it.
+        if thread::panicking() {
+            return;
+        }
+        self.shutdown_flag.store(true, Ordering::Relaxed);
+        let remaining = self.deadline.remaining();

Review Comment:
   **blocker** — this closes the nested-budget problem by making the peer-exit 
wait a debtor of the join budget, which reopens the writer-drop race the commit 
exists to close.
   
   `remaining()` is `saturating_duration_since` (`:349-353`), so it floors at 
`Duration::ZERO`. `join_all` arms the shared instant at flag-flip (`:229`), and 
shard 0 cannot reach this drop until `bus.shutdown()` and `await_pump_drain` 
have finished. When the join leg has spent the budget, `countdown.wait(ZERO)` 
re-checks the predicate and returns `Err(running)` **without blocking** 
(`:381-395`), shard 0 releases, the last `Rc<ServerMuxStateMachine>` drops, and 
any peer inside `LeftRight::read` panics on `.expect("read handle should be 
accessible")` (`core/metadata/src/stm/mod.rs:129`).
   
   Reachability: `join == drain` is legal and explicitly blessed by 
`core/configs/src/server_config/sharding.rs:378` 
(`join_equal_to_drain_is_accepted`), and at that setting one full-length drain 
exhausts the budget on its own. At the shipped default (join 30s / drain 10s, 
`core/server/config.toml:888`) zero is not reachable, but the margin still 
silently drops from a guaranteed 30s to a worst-case ~10s.
   
   The distinction that makes this worth blocking on: the join budget bounds a 
LIVENESS concern (do not hang process exit on a wedged shard), while this wait 
enforces a SAFETY invariant (the writer must outlive every reader). Funding the 
second out of the first lets an exit-latency timeout cancel a correctness 
fence. The invariant doc at `:359-365` names exactly one exception (panic); 
deadline exhaustion is now a second one and is not named.
   
   Note the original finding was about `join_all`'s **verdict** — shard 0 being 
reported `Wedged` on a correct shutdown — not about the peer wait's length. 
Fix: floor it (`self.deadline.remaining().max(drain_timeout)`), or give the 
wait its own budget and let `join_all` accept shard 0 returning within `budget 
+ peer_budget` without a `Wedged` row.



##########
core/server/src/boot/threads.rs:
##########
@@ -649,6 +827,188 @@ mod tests {
         );
     }
 
+    #[test]
+    fn peer_exit_countdown_releases_the_waiter_once_every_peer_is_out() {
+        let countdown = Arc::new(PeerExitCountdown::new(2));
+        let guards: [PeerExitGuard; 2] =
+            std::array::from_fn(|_| 
PeerExitGuard::new(Arc::clone(&countdown)));
+        assert_eq!(
+            countdown.wait(Duration::from_millis(10)),
+            Err(2),
+            "two live peers must hold the waiter past a short deadline"
+        );
+        let peers: Vec<_> = guards
+            .into_iter()
+            .map(|guard| {
+                thread::spawn(move || {
+                    thread::sleep(Duration::from_millis(20));
+                    drop(guard);
+                })
+            })
+            .collect();
+        assert_eq!(
+            countdown.wait(Duration::from_secs(30)),
+            Ok(()),
+            "the last guard drop must release the waiter"
+        );
+        for peer in peers {
+            peer.join()
+                .expect("peer thread dropped its guard without panicking");
+        }
+    }
+
+    #[test]
+    fn peer_exit_guard_counts_out_during_unwind() {
+        let countdown = Arc::new(PeerExitCountdown::new(1));
+        let guard = PeerExitGuard::new(Arc::clone(&countdown));
+        // `resume_unwind` skips the panic hook, so the unwind is silent.
+        let unwound = panic::catch_unwind(panic::AssertUnwindSafe(|| {
+            let _guard = guard;
+            panic::resume_unwind(Box::new("peer shard body panicked"));
+        }));
+        assert!(unwound.is_err());
+        assert_eq!(
+            countdown.wait(Duration::ZERO),
+            Ok(()),
+            "a guard dropped by a panic unwind must still count its peer out"
+        );
+    }
+
+    #[test]
+    fn peer_exit_wait_flips_the_flag_before_waiting() {
+        // Stands in for a peer parked on its bus token: it exits only once
+        // the shutdown flag is set, so a waiter that set the flag after
+        // waiting would sit out the whole budget.
+        let countdown = Arc::new(PeerExitCountdown::new(1));
+        let shutdown_flag = Arc::new(AtomicBool::new(false));
+        let peer = thread::spawn({
+            let guard = PeerExitGuard::new(Arc::clone(&countdown));
+            let shutdown_flag = Arc::clone(&shutdown_flag);
+            move || {
+                while !shutdown_flag.load(Ordering::Relaxed) {
+                    thread::sleep(Duration::from_millis(1));
+                }
+                drop(guard);
+            }
+        });
+        let started = Instant::now();
+        drop(PeerExitWait::new(
+            countdown,
+            Arc::clone(&shutdown_flag),
+            Arc::new(ShutdownDeadline::new(Duration::from_secs(30))),
+        ));
+        assert!(
+            started.elapsed() < Duration::from_secs(5),
+            "the waiter must not sit out its budget on a peer that waits for 
the flag"
+        );
+        assert!(shutdown_flag.load(Ordering::Relaxed));
+        peer.join()
+            .expect("peer thread dropped its guard without panicking");
+    }
+
+    #[test]
+    fn peer_exit_wait_gives_up_at_the_deadline() {
+        let countdown = Arc::new(PeerExitCountdown::new(1));
+        let _wedged_peer = PeerExitGuard::new(Arc::clone(&countdown));
+        let timeout = Duration::from_millis(50);
+        let started = Instant::now();
+        drop(PeerExitWait::new(
+            countdown,
+            Arc::new(AtomicBool::new(false)),
+            Arc::new(ShutdownDeadline::new(timeout)),
+        ));
+        let waited = started.elapsed();
+        assert!(
+            waited >= timeout && waited < Duration::from_secs(5),
+            "a wedged peer must be abandoned at the deadline, waited 
{waited:?}"
+        );
+    }
+
+    #[test]
+    fn peer_wait_and_shard_join_share_one_shutdown_budget() {
+        // Regression: the two waits nest (shard 0 cannot start waiting for
+        // its peers until its own drain returned, already inside the join
+        // budget). With a budget each, a clean shutdown reported shard 0 as
+        // wedged and exited non-zero.
+        let deadline = 
Arc::new(ShutdownDeadline::new(Duration::from_millis(200)));
+        let shutdown_flag = AtomicBool::new(true);
+        // Never finishes: stands in for the slow drain that arms and then
+        // spends the shared budget. The thread leaks into the test process.
+        let wedged_shard = thread::spawn(|| -> Result<(), ServerError> {
+            loop {
+                thread::sleep(Duration::from_secs(1));
+            }
+        });
+        assert!(
+            join_until_shutdown_deadline(wedged_shard, &shutdown_flag, 
&deadline).is_none(),
+            "the join must spend the budget it armed"
+        );
+
+        let countdown = Arc::new(PeerExitCountdown::new(1));
+        let _wedged_peer = PeerExitGuard::new(Arc::clone(&countdown));
+        let started = Instant::now();
+        drop(PeerExitWait::new(
+            countdown,
+            Arc::new(AtomicBool::new(false)),
+            Arc::clone(&deadline),
+        ));
+        assert!(
+            started.elapsed() < Duration::from_millis(100),

Review Comment:
   **blocker (same defect, second half)** — this assertion pins the unsafe 
release as intended behaviour, which is why the item above cannot be resolved 
by a comment.
   
   The test builds `PeerExitCountdown::new(1)` and holds `_wedged_peer`, so the 
peer is **never counted out**. It then spends the entire 200ms budget in 
`join_until_shutdown_deadline` and asserts that dropping `PeerExitWait` returns 
in under 100ms. That is an assertion that shard 0 drops the metadata writer 
while a peer is still running — a green test encoding exactly the race 
`PeerExitWait` was added to prevent.
   
   The stated intent ("inherit what is left of the budget, not arm a second 
one") is reasonable; it is the live peer guard that turns it into a 
specification of the starvation. A future reader will treat this as settled and 
no test will catch a regression.
   
   Fix: keep a case that pins budget inheritance with the peer already counted 
out, and add one that asserts the wait still HONOURS a live peer (blocks for 
the floor, or for its own budget).



##########
core/server/src/dispatch/mod.rs:
##########
@@ -351,208 +308,106 @@ fn pop_next_client_request(
     active: &ActiveClientRequests,
     client_id: u128,
 ) -> Option<Message<GenericHeader>> {
-    let mut queues = queues.borrow_mut();
-    let Some(queue) = queues.get_mut(&client_id) else {
-        active.borrow_mut().remove(&client_id);
-        return None;
-    };
-    let message = queue.pop_front();
-    if queue.is_empty() {
-        queues.remove(&client_id);
-    }
+    // The entry survives draining to empty: removing it would cost a

Review Comment:
   **warning** — moving the free to the connection-lost hook makes the hook the 
sole owner, but the hook is not the last toucher of this map, so an entry can 
now leak for process lifetime.
   
   The hook fires from the TRANSPORT task's scopeguard 
(`core/message_bus/src/installer/tcp.rs:118-130`), while the separate DISPATCH 
task still holds buffered frames in `in_rx`. `ctx.in_tx` drops when 
`conn.run(ctx)` returns, and `async_channel` delivers every buffered item 
before `recv()` reports `Closed` — so `client_dispatch_loop` (`tcp.rs:259-289`) 
still calls `on_request` after the hook ran, and `enqueue_client_request` 
re-creates the entry via `.entry(client_id).or_default()` (`:239-243`). Nothing 
frees it again: the hook is once-only (`message_bus/src/lib.rs:983` guards on 
`remove(..).is_some()`) and `client_id` is monotonic `(shard << 112) | seq`, 
never reused. The two guards that could stop the loop — `aborted` 
(`tcp.rs:275`) and `conn_token` — are both insert-race-only, so this is not a 
rare interleaving; it happens whenever a frame was buffered at close.
   
   Pre-PR there was no leak, because this branch freed the entry every drain.
   
   Fix: restore the removal here and keep the hook's removal for the mid-flight 
case. That also resolves the sticky-capacity issue below, so no `shrink_to` is 
needed.
   
   Separately, on the stated motive: `VecDeque::pop_front` never releases 
capacity, so as written a client's PEAK depth stays resident (no 
`shrink_to`/`shrink_to_fit` exists anywhere in `core/server` or `core/shard`). 
For the in-tree SDKs that is small — they are lockstep, so depth stays at the 
4-slot `RawVec` minimum, ~96 B per connection — but it interacts badly with the 
uncapped-queue item you deferred: an unbounded transient spike becomes a 
permanent high-water mark. Restoring the removal makes both moot.



##########
core/server/src/dispatch/partition.rs:
##########
@@ -725,13 +781,15 @@ pub(in crate::dispatch) async fn 
handle_get_consumer_offset<B, MJ, S, SB>(
     SB: SuperblockStore + 'static,
 {
     let Ok(wire) = 
GetConsumerOffsetRequest::decode_from(request_body(request)) else {
-        // Undecodable: an empty body decodes as None (no offset) on the SDK.
-        send_non_replicated_bytes(
+        // Same rule as the poll above: an empty body is the legitimate
+        // "no offset stored" answer, so a malformed request must not send
+        // one. Byte-identical frames on the same channel with the same
+        // context would also be indistinguishable in the send-failure log.
+        send_non_replicated_deny(
             shard,
             request,
             transport_client_id,
-            Bytes::new(),
-            "get_consumer_offset",
+            IggyError::InvalidCommand.as_code(),

Review Comment:
   **warning** — the undecodable-body fail-open is closed here, but the 
not-found one nine lines below is still open, and it is the only finding in 
this round reachable from a **well-formed request through a public SDK method**.
   
   `:827` denies typed only on `IggyError::PartitionNotFound`; 
`StreamIdNotFound` and `TopicIdNotFound` from `resolve_partition_namespace` 
(`core/server/src/responses.rs:466-476`) fall through to the generic `Err` arm 
at `:836` and answer `Bytes::new()` — which every SDK decodes as "no offset 
stored":
   
   - Rust `get_consumer_offset(..) -> Result<Option<ConsumerOffsetInfo>, 
IggyError>` returns `Ok(None)`, minted at 
`core/common/src/traits/binary_impls/consumer_offsets.rs:85`
   - Go `GetConsumerOffset(..)` returns `(nil, nil)` 
(`binary_serialization/binary_response_deserializer.go:41-43`)
   - C# `GetOffsetAsync(..)` returns `null` 
(`Iggy_SDK/IggyClient/Implementations/TcpMessageStream.cs:386-389`)
   - Java `Optional.empty()`, Node `null`
   
   So `get_consumer_offset` against a typo'd or deleted stream is 
indistinguishable from a fresh consumer: it resumes from its configured default 
and silently reprocesses, with no error anywhere. Two contrasts make it hard to 
defend: the SAME typo through `poll_messages` returns `Err(StreamIdNotFound)` 
(`:654-659` denies all three variants), and the SAME read over REST 404s 
(`core/server/src/http/handlers.rs:1312`). One lookup, three answers.
   
   Fix: extend the `:827` match to `StreamIdNotFound | TopicIdNotFound`, 
mirroring `:654-659`. Error-path only, and cheaper than today since it drops 
the `Bytes::new()` + `send_non_replicated_bytes` for a header-only deny.



##########
core/server/src/dispatch/failure.rs:
##########
@@ -0,0 +1,677 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! The wire failure channels and the one send exit for host-built frames.
+//!
+//! Every frame the dispatch host builds, success or rejection, leaves through
+//! [`send_host_frame`], so the send-failure log has one shape. Which channel a
+//! failure rides is a wire contract with the SDK:
+//!
+//! | channel | carrier | when |
+//! |---|---|---|
+//! | [`FrameChannel::TypedDeny`] | Reply, nonzero status + empty body, or a 
result-framed rejection body | rejections that must unblock the SDK's lockstep 
request slot: checksum, authz, pre-consensus rewrite, unknown or unsupported 
non-replicated code, unbound non-PING read, transient replay hints |
+//! | [`FrameChannel::Eviction`] | session-terminal Eviction frame with a 
typed reason | the client must register again: `NoSession`, `MalformedLogin`, 
heartbeat and login evictions. The reason rides the channel label, since one 
`context` covers four of them |
+//! | [`FrameChannel::ResyncSentinel`] | status-0 poll reply, body carries 
`RESYNC_REQUIRED_PARTITION_SENTINEL` | a fenced consumer-group poll: the 
consumer must re-sync its assignment; HTTP mirrors it as 
`resync_required_polled_messages` in `crate::http::wire` |
+//! | [`FrameChannel::EmptyFrame`] | status-0 fail-fast body, empty or the 
16-byte empty poll | the partition cannot answer yet; the SDK fails fast (empty 
poll) and retries. A permanent client error never rides this channel: an 
undecodable body denies typed, because there is nothing to retry |

Review Comment:
   **warning** — "A permanent client error never rides this channel" is false 
as written, and it is a rule this commit authored.
   
   A consumer-group poll that omits the partition id decodes cleanly 
(`core/binary_protocol/src/requests/messages/poll_messages.rs:95` maps 
`partition_flag != 1` to `None`), then `resolve_poll_request` returns 
`IggyError::InvalidIdentifier` (`partition.rs:946`). That variant is absent 
from the deny set at `partition.rs:654-659`, so it falls to 
`empty_poll_fallback(0)` (`:754`) and rides `EmptyFrame` as a status-0 16-byte 
empty poll — a permanent client error answered as a successful 0-message read, 
which is the exact symptom your commit message describes.
   
   Reachability is foreign/hand-rolled clients only (every in-tree SDK resolves 
the partition client-side, e.g. 
`core/common/src/traits/binary_impls/messages.rs:317-320`) — but that is the 
same reachability class as the three fail-opens this commit closed, so it 
cannot be the reason to leave the sentence absolute.
   
   Fix: add `InvalidIdentifier` to the `:654-659` match (one pattern, error 
path, strictly cheaper than the current `Bytes` allocation), or drop the 
absolute claim from this row.



##########
core/server/src/dispatch/partition.rs:
##########
@@ -581,8 +583,10 @@ async fn relay_partition_reply<B, MJ, S, SB>(
 /// the owning shard ([`shard::IggyShard::partition_read`]), and re-encode
 /// the stored batches into the legacy wire `PolledMessages` body.
 ///
-/// Failures reply with an empty body so the SDK fails fast on decode
-/// instead of hanging until its read timeout.
+/// A partition that cannot answer yet replies with the 16-byte empty poll,
+/// which the SDK reads as 0 messages and retries. Permanent client errors
+/// (undecodable body, authz, unresolved target) deny with a nonzero status

Review Comment:
   **warning** — this sentence, added by this commit, says permanent client 
errors including an "unresolved target" deny with a nonzero status. 
`IggyError::InvalidIdentifier` IS an unresolved target and does not: it is 
missing from the deny match at `:654-659` and falls to 
`empty_poll_fallback(0)`. Same defect as the `failure.rs:29` row. Fix the match 
or narrow the sentence.



##########
core/server/src/dispatch/partition.rs:
##########
@@ -695,18 +677,92 @@ pub(in crate::dispatch) async fn handle_poll_messages<B, 
MJ, S, SB>(
                 "poll_messages request rejected; replying empty poll"
             );
             let partition_id = if matches!(error, 
IggyError::ConsumerGroupPartitionNotOwned(..)) {
-                iggy_common::RESYNC_REQUIRED_PARTITION_SENTINEL
+                RESYNC_REQUIRED_PARTITION_SENTINEL
             } else {
                 0
             };
-            empty_polled_messages_body(partition_id)
+            empty_poll_fallback(partition_id)
         }
     };
-    send_non_replicated_bytes(shard, request, transport_client_id, body, 
"poll_messages").await;
+    send_non_replicated_bytes(
+        shard,
+        request,
+        transport_client_id,
+        body,
+        channel,
+        "poll_messages",
+    )
+    .await;
 }
 
-/// Serve `get_consumer_offset`. An empty body decodes as `None` on the SDK
-/// side (no offset stored / partition unknown).
+/// Run the resolved poll on the owning shard and re-encode the stored
+/// batches into the wire `PolledMessages` reply. A failed read or re-encode
+/// hands back the fail-fast empty poll for the partition instead.
+#[allow(clippy::future_not_send)]
+async fn read_polled_messages<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    transport_client_id: u128,
+    request: &Message<RoutedRequestHeader>,
+    (namespace, partition_id, consumer, args): DecodedPollRequest,
+) -> Result<BusMessage, (Bytes, FrameChannel)>
+where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    match shard
+        .partition_read(namespace, PartitionRead::Poll { consumer, args })
+        .await
+    {
+        Some(PartitionReadReply::Poll {
+            fragments,
+            current_offset,
+        }) => build_polled_messages_reply(
+            request.header(),
+            current_metadata_commit(shard),
+            partition_id,
+            current_offset,
+            fragments,
+            shard.plane.partitions().config().encryptor.as_deref(),
+        )
+        .map_err(|error| {
+            warn!(
+                transport_client_id,
+                error = %error,
+                "failed to re-encode polled batches; replying empty poll"
+            );
+            empty_poll_fallback(partition_id)
+        }),
+        other => {
+            warn!(
+                transport_client_id,
+                namespace = namespace.inner(),
+                reply_was_none = other.is_none(),
+                "partition read failed; replying empty poll"
+            );
+            Err(empty_poll_fallback(partition_id))
+        }
+    }
+}
+
+/// The fail-fast poll reply for a partition that could not answer: the
+/// 16-byte empty poll for `partition_id`, riding the re-sync sentinel
+/// channel when the id is the sentinel and the empty-frame channel
+/// otherwise.
+fn empty_poll_fallback(partition_id: u32) -> (Bytes, FrameChannel) {
+    let channel = if partition_id == RESYNC_REQUIRED_PARTITION_SENTINEL {
+        FrameChannel::ResyncSentinel
+    } else {
+        FrameChannel::EmptyFrame
+    };
+    (empty_polled_messages_body(partition_id), channel)
+}
+
+/// Serve `get_consumer_offset`. An empty status-0 body decodes as `None` on
+/// the SDK side (no offset stored / partition unknown), so it is reserved for

Review Comment:
   **warning** — this docstring, rewritten by this commit, is contradicted by 
its own function body. It reserves the empty status-0 body for "no offset 
stored / partition unknown", but `:827` denies `PartitionNotFound` typed, so 
partition-unknown never takes that path. Meanwhile the case that DOES take it — 
an unresolved stream or topic — is not mentioned.
   
   Fix the `:827` match (see the comment there) and the sentence becomes true 
as written. If the match is deferred instead, this needs to say plainly that an 
unresolved stream or topic reads as "no offset stored", with an issue link — a 
documented gap is defensible, this sentence is not.



##########
core/server/src/dispatch/mod.rs:
##########
@@ -188,50 +137,46 @@ where
 {
     let shard_handle = Rc::clone(shard_handle);
     let sessions = Rc::clone(sessions);
-    let queues: ClientRequestQueues = Rc::new(RefCell::new(HashMap::new()));
-    let active: ActiveClientRequests = Rc::new(RefCell::new(HashSet::new()));
+    let queues: ClientRequestQueues = Rc::new(RefCell::new(AHashMap::new()));
+    let active: ActiveClientRequests = Rc::new(RefCell::new(AHashSet::new()));
+    let queues_for_disconnect = Rc::clone(&queues);
+    let active_for_disconnect = Rc::clone(&active);
     let sessions_for_disconnect = Rc::clone(&sessions);
     let shard_handle_for_disconnect = Rc::clone(&shard_handle);
     let bus_for_spawn = (*bus).clone();
     bus.set_client_connection_lost_fn(Rc::new(move |client_id| {
-        if let Some((vsr_client_id, session)) = sessions_for_disconnect
-            .borrow_mut()
-            .remove_connection(client_id)
-            && let Some(shard) = 
upgrade_shard_handle(&shard_handle_for_disconnect)
+        // The socket is gone: nothing will drain what is still queued and no
+        // later frame will release the active slot, so both entries go here.
+        // This is also the only recovery from a panic inside the drain task
+        // (compio catches it), which would otherwise leave the slot taken and
+        // every later frame for the client queued forever.
+        queues_for_disconnect.borrow_mut().remove(&client_id);
+        active_for_disconnect.borrow_mut().remove(&client_id);

Review Comment:
   **nit** — clearing `active` from the hook breaks the one-drain-per-client 
invariant this handler's doc promises at `:118-123`. The hook can fire while 
drain #1 is suspended at an `.await` inside `handle_client_request`, so the 
post-hook `enqueue_client_request` (see the `:311` comment for why one can 
still arrive) wins `active.insert` at `:244` and spawns drain #2 concurrently. 
Two drains then pop the same `VecDeque`.
   
   Impact is limited and I want to be accurate about it: because this closure 
is synchronous with no await between here and `remove_connection` at `:161`, 
the session is already stripped, and the hook fires FROM `remove_client_meta`, 
so `ensure_transport_connection` finds nothing and `classify(header, false)` 
routes every replicated op to `UnboundReplicated`. Drain #2 therefore commits 
nothing — it only emits `Eviction(NoSession)` frames at a dead socket. No op 
loss, no reordering.
   
   Fix: release the slot from a scopeguard inside the spawned drain task 
instead. That still covers the panic case this line was written for, without 
racing the hook.



##########
core/server/src/dispatch/mod.rs:
##########
@@ -681,353 +532,193 @@ async fn handle_client_request<B, MJ, S, SB>(
                 IggyError::Unauthenticated.as_code(),
             )
             .await;
-            return;
         }
-        handle_non_replicated_request(shard, sessions, system_config, 
transport_client_id, request)
-            .await;
-        return;
-    }
-
-    if header.operation == Operation::Register && header.session == 0 && 
header.request == 0 {
-        handle_login_register_request(shard, sessions, transport_client_id, 
request).await;
-        return;
-    }
-
-    if header.operation == Operation::Logout {
-        handle_logout_request(shard, sessions, transport_client_id, 
request).await;
-        return;
-    }
-
-    let bound = sessions.borrow().get_session(transport_client_id);
-    if bound.is_none() {
-        // Replicated request on an unbound transport. Without this short-
-        // circuit, the rewrite below overwrites `header.client` with
-        // `transport_client_id` and dispatches; the request_preflight then
-        // rejects with `NoSession`/`Fenced` and the failure disappears
-        // silently, wedging the SDK until the socket timeout. A typed
-        // `Eviction(NoSession)` is right here, unlike the pre-auth read
-        // guard above: a replicated request implies the client believes it
-        // has a session, and that session is gone, so it must register
-        // again. An empty status-0 Reply is not safe here, because
-        // SendMessages is the one replicated operation without a result
-        // section, and its decoder would read the empty body as a
-        // successful send.
-        warn!(
-            transport_client_id,
-            operation = ?header.operation,
-            "rejecting replicated request from unbound transport with 
Eviction(NoSession)"
-        );
-        send_unauthenticated_eviction(shard, transport_client_id).await;
-        return;
-    }
-
-    // DeleteSegments is neither a partition nor a metadata consensus op: the
-    // owning shard resolves the requested count to a concrete offset, then a
-    // `TruncatePartition` is replicated through metadata (Option A). Each
-    // replica's reconciler trims to the committed watermark. Handle it here,
-    // ahead of the partition/metadata routing below.
-    if header.operation == Operation::DeleteSegments {
-        handle_delete_segments_request(shard, transport_client_id, bound, 
&request).await;
-        return;
-    }
-
-    if header.operation.is_partition() {
-        // `bound` is Some here (unbound transports returned above).
-        let (vsr_client_id, bound_session) = bound.unwrap_or((0, 0));
-        // `get_session` discards the acting user id the partition gate needs;
-        // resolve it from the same bound connection. A bound transport always
-        // has one, but the gate fails closed on `None` rather than trust that.
-        let acting_user_id = 
sessions.borrow().get_user_id(transport_client_id);
-        dispatch_partition_request(
-            shard,
-            request,
-            vsr_client_id,
-            bound_session,
-            transport_client_id,
-            acting_user_id,
-        )
-        .await;
-        return;
-    }
-
-    let request = request.transmute_header(|header, new_header: &mut 
RoutedRequestHeader| {
-        *new_header = header;
-        // Metadata-plane ops route by operation: stamp the sentinel group.
-        new_header.group = server_common::sharding::METADATA_GROUP;
-        // `bound` is always Some here (unbound transports early-return above);
-        // this sets the consensus client id + session for the replicated op.
-        if let Some((bound_client_id, bound_session)) = bound {
-            new_header.client = bound_client_id;
-            new_header.session = bound_session;
-        }
-    });
-    let (request, raw_pat_token) = match maybe_rewrite_pat_request(
-        sessions,
-        transport_client_id,
-        max_tokens_per_user,
-        |user_id| {
-            shard
-                .plane
-                .metadata()
-                .mux_stm
-                .users()
-                .read(|users| users.pat_count_of(user_id))
-        },
-        request,
-    ) {
-        Ok(rewritten) => rewritten,
-        Err(error) => {
-            // Token cap reached, malformed body, or a lost session binding.
-            send_pre_consensus_deny(
+        RequestClass::NonReplicatedRead => {
+            // The auth-bypass guard is `classify`'s `UnauthenticatedRead` 
class:
+            // `PING`, the liveness probe, is the only pre-auth code, on every
+            // roster shape. `GET_CLUSTER_METADATA` describes the private 
replica
+            // network and is not something an unauthenticated caller gets to
+            // read; a client that dialed a backup no longer needs it to find 
the
+            // leader, because the backup authenticates the login locally and
+            // forwards only the consensus proposal
+            // (`submit_register_local_or_forward`). Every other non-replicated
+            // code MUST go through Register first, which binds the acting user
+            // the per-op authz gates resolve.
+            handle_non_replicated_request(
                 shard,
-                &header,
+                sessions,
+                system_config,
                 transport_client_id,
-                &error,
-                "personal-access-token",
+                request,
+                (user_id, client_address),
             )
             .await;
-            return;
         }
-    };
-    // Hash raw passwords and, for ChangePassword, verify the current password
-    // on the primary before replication; see `crate::users`. Replicas store 
the
-    // hash directly. A wrong current password is not denied here: it rides
-    // consensus and applies as a committed InvalidCredentials no-op, so the 
only
-    // Err returned is a malformed body.
-    let request = match maybe_rewrite_user_password_request(shard, request) {
-        Ok(rewritten) => rewritten,
-        Err(error) => {
-            // Malformed body: deny fast with InvalidCommand.
-            send_pre_consensus_deny(shard, &header, transport_client_id, 
&error, "user-password")
-                .await;
-            return;
+        RequestClass::LoginRegister => {
+            handle_login_register_request(shard, sessions, 
transport_client_id, request).await;
         }
-    };
-    // Static bounds run pre-consensus so a rejected request burns no
-    // replicated log entry; HTTP covers the same bounds via
-    // `command.validate()`. A body that fails to decode denies typed too
-    // (`InvalidCommand`), instead of riding consensus just to fail there.
-    let bounds = match header.operation {
-        Operation::CreateTopic => 
CreateTopicRequest::decode_from(request_body(&request))
-            .map_err(|_| IggyError::InvalidCommand)
-            .and_then(|create_topic| {
-                // `parse` doubles as the catalog gate: an unknown key or a
-                // malformed value denies typed here, pre-consensus.
-                let options = 
TopicCreateOptions::parse(&create_topic.options)?;
-                if let Some(segment_size) = options.segment_size {
-                    validate_topic_segment_size(
-                        segment_size.as_bytes_u64(),
-                        iggy_common::MAX_TOPIC_SEGMENT_SIZE,
-                    )?;
-                }
-                let segment_size = options.segment_size.map_or_else(
-                    || iggy_common::DEFAULT_SEGMENT_SIZE,
-                    |segment_size| segment_size.as_bytes_u64(),
-                );
-                if options
-                    .preallocate_segments
-                    .unwrap_or(iggy_common::DEFAULT_PREALLOCATE_SEGMENTS)
-                {
-                    validate_preallocated_topic_bytes(segment_size, 
create_topic.partitions_count)?;
-                }
-                let max_topic_size = options
-                    .max_topic_size
-                    .unwrap_or(MaxTopicSize::ServerDefault);
-                validate_topic_bounds(create_topic.partitions_count, 
max_topic_size, segment_size)?;
-                warn_unenforceable_topic_size(
-                    max_topic_size,
-                    segment_size,
-                    shard.bus_max_message_size(),
-                    create_topic.partitions_count,
-                );
-                Ok(())
-            }),
-        Operation::CreatePartitions => 
CreatePartitionsRequest::decode_from(request_body(&request))
-            .map_err(|_| IggyError::InvalidCommand)
-            .and_then(|create_partitions| {
-                
validate_partitions_change_count(create_partitions.partitions_count)?;
-                let metadata = shard.plane.metadata();
-                warn_unenforceable_topic_size_on_partition_add(
-                    metadata.mux_stm.streams(),
-                    &create_partitions.stream_id,
-                    &create_partitions.topic_id,
-                    shard.bus_max_message_size(),
-                    create_partitions.partitions_count,
-                );
-                Ok(())
-            }),
-        Operation::DeletePartitions => 
DeletePartitionsRequest::decode_from(request_body(&request))
-            .map_err(|_| IggyError::InvalidCommand)
-            .and_then(|delete_partitions| {
-                
validate_partitions_change_count(delete_partitions.partitions_count)
-            }),
-        // Only the updatable subset: the create-time knobs are pushed to
-        // partitions when the topic is built and nothing re-pushes them, so
-        // accepting one here would store a value no partition ever sees.
-        Operation::UpdateTopic => 
UpdateTopicRequest::decode_from(request_body(&request))
-            .map_err(|_| IggyError::InvalidCommand)
-            .and_then(|update_topic| {
-                validate_option_keys(&update_topic.options, 
UPDATABLE_TOPIC_OPTION_KEYS)?;
-                let options = 
TopicCreateOptions::parse(&update_topic.options)?;
-                let Some(max_topic_size) = options.max_topic_size else {
-                    return Ok(());
-                };
-                // An update can lower the cap below one segment just as a
-                // create can, and the stored map would then report a size the
-                // topic can never enforce. The floor is this topic's own
-                // segment size, since that key is create-only.
-                let metadata = shard.plane.metadata();
-                let streams = metadata.mux_stm.streams();
-                let segment_size = streams
-                    .topic_segment_size(&update_topic.stream_id, 
&update_topic.topic_id)
-                    .map_or_else(
-                        || iggy_common::DEFAULT_SEGMENT_SIZE,
-                        |segment_size| segment_size.as_bytes_u64(),
-                    );
-                validate_topic_size_floor(max_topic_size, segment_size)?;
-                let partitions_count = streams
-                    .topic_partitions_count(&update_topic.stream_id, 
&update_topic.topic_id)
-                    .unwrap_or(0);
-                warn_unenforceable_topic_size(
-                    max_topic_size,
-                    segment_size,
-                    shard.bus_max_message_size(),
-                    u32::try_from(partitions_count).unwrap_or(u32::MAX),
-                );
-                Ok(())
-            }),
-        Operation::UpdateStream => 
UpdateStreamRequest::decode_from(request_body(&request))
-            .map_err(|_| IggyError::InvalidCommand)
-            .and_then(|update_stream| {
-                validate_option_keys(&update_stream.options, 
UPDATABLE_STREAM_OPTION_KEYS)
-            }),
-        Operation::UpdateUser => 
UpdateUserRequest::decode_from(request_body(&request))
-            .map_err(|_| IggyError::InvalidCommand)
-            .and_then(|update_user| {
-                validate_option_keys(&update_user.options, 
UPDATABLE_USER_OPTION_KEYS)
-            }),
-        Operation::CreateStream => 
CreateStreamRequest::decode_from(request_body(&request))
-            .map_err(|_| IggyError::InvalidCommand)
-            .and_then(|create_stream| 
validate_option_keys(&create_stream.options, &[])),
-        Operation::CreateUser => 
CreateUserRequest::decode_from(request_body(&request))
-            .map_err(|_| IggyError::InvalidCommand)
-            .and_then(|create_user| validate_option_keys(&create_user.options, 
&[])),
-        _ => Ok(()),
-    };
-    if let Err(error) = bounds {
-        send_pre_consensus_deny(shard, &header, transport_client_id, &error, 
"static-bounds").await;
-        return;
-    }
-    // Enrich consumer-group Join/Leave with the client's VSR id (+ topic
-    // partition count for Join) before replication; see 
`crate::consumer_group`.
-    let request = match maybe_rewrite_consumer_group_request(shard, 
request).await {
-        Ok(rewritten) => rewritten,
-        Err(error) => {
+        RequestClass::Logout => {
+            handle_logout_request(shard, sessions, transport_client_id, 
request).await;
+        }
+        RequestClass::UnboundReplicated => {
+            // Replicated request on an unbound transport. Without this short-
+            // circuit, the rewrite below overwrites `header.client` with
+            // `transport_client_id` and dispatches; the request_preflight then
+            // rejects with `NoSession`/`Fenced` and the failure disappears
+            // silently, wedging the SDK until the socket timeout. A typed
+            // `Eviction(NoSession)` is right here, unlike the plain deny of
+            // `UnauthenticatedRead`: a replicated request implies the client
+            // believes it has a session, and that session is gone, so it must
+            // register again. An empty status-0 Reply is not safe here, 
because
+            // SendMessages is the one replicated operation without a result
+            // section, and its decoder would read the empty body as a
+            // successful send.
             warn!(
                 transport_client_id,
-                error = %error,
                 operation = ?header.operation,
-                "dropping consumer-group request with invalid payload"
+                "rejecting replicated request from unbound transport with 
Eviction(NoSession)"
             );
-            return;
+            // The eviction context is best-effort off the metadata consensus
+            // (peer shards have none; zeroes are cosmetic -- the SDK only
+            // reads the reason), and the evicted id is the transport id: no
+            // VSR session exists to name.
+            send_eviction(
+                shard,
+                transport_client_id,
+                transport_client_id,
+                EvictionReason::NoSession,
+                "unbound replicated request",
+            )
+            .await;
         }
-    };
-    let request_header = *request.header();
-    // Replicated request: run consensus on the metadata owner (shard 0) and
-    // bring the committed reply back here. This shard owns the connection,
-    // so it writes the reply to the socket via the transport client id --
-    // shard 0 can't route by the consensus client id (no home-shard bits).
-    match submit_client_request_on_owner(shard, request).await {
-        Some(reply) => {
-            // The raw PAT token never enters consensus (it is 
non-deterministic
-            // and secret), so the committed reply body is empty. Substitute 
the
-            // raw-token response here, on the minting client's home shard, 
using
-            // the confirmed commit position from the committed reply.
-            let reply = match build_raw_pat_reply(&request_header, reply, 
raw_pat_token) {
-                Ok(reply) => reply,
+        RequestClass::DeleteSegments => {
+            // DeleteSegments is neither a partition nor a metadata consensus 
op: the
+            // owning shard resolves the requested count to a concrete offset, 
then a
+            // `TruncatePartition` is replicated through metadata (Option A). 
Each
+            // replica's reconciler trims to the committed watermark. 
`classify`
+            // names it ahead of the partition and metadata classes.
+            handle_delete_segments_request(shard, transport_client_id, bound, 
&request).await;
+        }
+        RequestClass::Partition => {
+            // `bound` is Some here: `classify` sends unbound transports to
+            // `UnboundReplicated`.
+            let (vsr_client_id, bound_session) = bound.unwrap_or((0, 0));
+            // The acting user comes from the prologue's lookup. A bound
+            // transport always has one, but the gate below fails closed on
+            // `None` rather than trust that.
+            dispatch_partition_request(
+                shard,
+                request,
+                vsr_client_id,
+                bound_session,
+                transport_client_id,
+                user_id,
+            )
+            .await;
+        }
+        RequestClass::ReplicatedMetadata => {
+            // The one arm that needs the by-value copy: the rewrite below
+            // stamps the consensus client / session / group over the header,
+            // and a pre-consensus deny still has to echo what the client sent.
+            let header = *header;
+            let request =
+                request.transmute_header(|header, new_header: &mut 
RoutedRequestHeader| {
+                    *new_header = header;
+                    // Metadata-plane ops route by operation: stamp the 
sentinel group.
+                    new_header.group = server_common::sharding::METADATA_GROUP;
+                    // `bound` is always Some here (`classify` sends unbound 
transports to
+                    // `UnboundReplicated`); this sets the consensus client id 
+ session
+                    // for the replicated op.
+                    if let Some((bound_client_id, bound_session)) = bound {
+                        new_header.client = bound_client_id;
+                        new_header.session = bound_session;
+                    }
+                });
+            let (request, raw_pat_token) = match tcp_chain(
+                shard,
+                sessions,
+                transport_client_id,
+                max_tokens_per_user,
+                request,
+            ) {
+                Ok(rewritten) => rewritten,
+                Err(RewriteDeny { stage, error }) => {
+                    send_pre_consensus_deny(shard, transport_client_id, 
&header, &error, stage)
+                        .await;
+                    return;
+                }
+            };
+            // Enrich consumer-group Join/Leave with the client's VSR id (+ 
topic
+            // partition count for Join) before replication; see 
`crate::consumer_group`.
+            let request = match maybe_rewrite_consumer_group_request(shard, 
request).await {
+                Ok(rewritten) => rewritten,
                 Err(error) => {
-                    warn!(
+                    // Both of the rewrite's own failures are `InvalidCommand`
+                    // decode errors, so a replay cannot help: deny typed
+                    // instead of leaving the lockstep connection to its read
+                    // timeout. (Its third error path needs a body past
+                    // `u32::MAX` against a 64 MiB message cap, so no client
+                    // frame reaches it; the deny is correct there too.)
+                    send_pre_consensus_deny(
+                        shard,
                         transport_client_id,
-                        error = %error,
-                        "failed to build raw PAT reply"
-                    );
+                        &header,
+                        &error,
+                        RewriteStage::ConsumerGroup,
+                    )
+                    .await;
                     return;
                 }
             };
-            if let Err(error) = shard
-                .bus
-                .send_to_client(transport_client_id, reply.into_frozen())
-                .await
-            {
-                warn!(
-                    transport_client_id,
-                    error = %error,
-                    operation = ?header.operation,
-                    "failed to deliver committed reply to client"
-                );
+            let request_header = *request.header();
+            // Replicated request: run consensus on the metadata owner (shard 
0) and
+            // bring the committed reply back here. This shard owns the 
connection,
+            // so it writes the reply to the socket via the transport client 
id --
+            // shard 0 can't route by the consensus client id (no home-shard 
bits).
+            match submit_client_request_on_owner(shard, request).await {
+                Some(reply) => {
+                    // The raw PAT token never enters consensus (it is 
non-deterministic
+                    // and secret), so the committed reply body is empty. 
Substitute the
+                    // raw-token response here, on the minting client's home 
shard, using
+                    // the confirmed commit position from the committed reply.
+                    let reply = match build_raw_pat_reply(&request_header, 
reply, raw_pat_token) {
+                        Ok(reply) => reply,
+                        Err(error) => {
+                            warn!(
+                                transport_client_id,
+                                error = %error,
+                                "failed to build raw PAT reply"
+                            );
+                            // The op COMMITTED; only the reply could not be
+                            // rendered. A typed deny is still the right frame:
+                            // silence wedges the lockstep connection on a
+                            // request that succeeded server-side, and the
+                            // client can read the minted token back.

Review Comment:
   **nit** — "the client can read the minted token back" is false, so the 
comment justifies the right frame with a recovery path that does not exist. 
`PersonalAccessTokenResponse` carries `name` + `expiry_at` only 
(`core/binary_protocol/src/responses/personal_access_tokens/get_personal_access_tokens.rs:30-33`,
 `core/common/src/types/permissions/personal_access_token.rs:35-40`); the raw 
secret exists in that one reply and nowhere else. A retry hits `ClientTable` 
dedup and replays the empty committed body, so it is unrecoverable and the 
caller must delete and re-create.
   
   The deny is still correct — silence on a committed op is worse. Just drop 
the last clause.



##########
core/server/src/dispatch/partition.rs:
##########
@@ -1176,8 +1186,7 @@ pub(in crate::dispatch) async fn 
handle_delete_segments_request<B, MJ, S, SB>(
 /// truncate commits under. A resolvable namespace with nothing sealed to 
delete
 /// still yields a `TruncatePartition(up_to_offset = 0)` so the metadata 
request
 /// sequence stays contiguous. `Err` on a malformed body or an unresolved
-/// namespace: the TCP caller drops it to a silent replay, the HTTP caller 
renders
-/// the error.
+/// namespace: the TCP caller denies typed, the HTTP caller renders the error.

Review Comment:
   **nit** — two stale rustdoc sentences around the lines this commit edited.
   
   Here: "`Err` on a malformed body or an unresolved namespace" — an unresolved 
namespace returns 
`Ok(build_truncate_partition_client_message_with_identifiers(..))` at 
`:1223-1230`. The only `Err` values are `InvalidCommand` and 
`TransientNotAccepted`.
   
   And at `:1075-1076`: "Only a malformed / unresolvable request is acked empty 
without a commit" — now false on both halves, since malformed denies typed at 
`:1143` and unresolvable returns `Ok`.



##########
core/server/src/dispatch/mod.rs:
##########
@@ -188,50 +137,46 @@ where
 {
     let shard_handle = Rc::clone(shard_handle);
     let sessions = Rc::clone(sessions);
-    let queues: ClientRequestQueues = Rc::new(RefCell::new(HashMap::new()));
-    let active: ActiveClientRequests = Rc::new(RefCell::new(HashSet::new()));
+    let queues: ClientRequestQueues = Rc::new(RefCell::new(AHashMap::new()));
+    let active: ActiveClientRequests = Rc::new(RefCell::new(AHashSet::new()));
+    let queues_for_disconnect = Rc::clone(&queues);
+    let active_for_disconnect = Rc::clone(&active);
     let sessions_for_disconnect = Rc::clone(&sessions);
     let shard_handle_for_disconnect = Rc::clone(&shard_handle);
     let bus_for_spawn = (*bus).clone();
     bus.set_client_connection_lost_fn(Rc::new(move |client_id| {
-        if let Some((vsr_client_id, session)) = sessions_for_disconnect
-            .borrow_mut()
-            .remove_connection(client_id)
-            && let Some(shard) = 
upgrade_shard_handle(&shard_handle_for_disconnect)
+        // The socket is gone: nothing will drain what is still queued and no
+        // later frame will release the active slot, so both entries go here.
+        // This is also the only recovery from a panic inside the drain task
+        // (compio catches it), which would otherwise leave the slot taken and
+        // every later frame for the client queued forever.
+        queues_for_disconnect.borrow_mut().remove(&client_id);
+        active_for_disconnect.borrow_mut().remove(&client_id);
+        // Upgrade FIRST: `remove_connection` strips the `SessionManager`
+        // entry, so running it ahead of a failed upgrade would drop the
+        // binding without ever submitting the replicated `Logout`, leaking
+        // the `ClientTable` entry and its consumer-group memberships. The
+        // window is pre-build / post-runtime-drop only.
+        if let Some(shard) = upgrade_shard_handle(&shard_handle_for_disconnect)
+            && let Some((vsr_client_id, session)) = sessions_for_disconnect

Review Comment:
   **nit** — upgrading before `remove_connection` is the right trade (a leaked 
local `SessionManager` entry beats a leaked replicated `ClientTable` entry), 
but the failure is now silent: if the upgrade fails, `connections` and 
`client_to_connection` keep their entries and nothing says so. Worth noting the 
heartbeat verifier does not always reap it — it is gated on 
`config.heartbeat.enabled` (`boot/mod.rs:806`) and `collect_stale` only 
considers `Bound`/`Authenticated` (`session_manager.rs:196-199`), so a 
`Connected` leak is never reaped. Window is pre-build / post-runtime-drop only, 
so this is about observability, not incidence. Fix: `error!` on the 
upgrade-failed branch.



##########
core/server/src/boot/threads.rs:
##########
@@ -406,6 +560,22 @@ pub(in crate::boot) fn validate_sharding_runtime_knobs(
             drain: drain_timeout,
         });
     }
+    let join_timeout = sharding.shutdown_join_timeout.get_duration();

Review Comment:
   **nit** — the restored `join >= drain` and `join <= MAX` checks match 
`ShardingConfig::validate` exactly 
(`core/configs/src/server_config/sharding.rs:251-267`), but the doc at 
`:519-522` claims this fn "Mirrors `ShardingConfig::validate`" and one clause 
is still missing: `reconcile_periodic_interval`, which `validate` rejects at 
zero and above `RECONCILE_PERIODIC_INTERVAL_MAX` (`sharding.rs:269-287`).
   
   That one is not cosmetic for a caller that bypasses `validate`: zero feeds 
`run_reconciler`'s `ctx.shard.bus.sleep(periodic)` inside an unconditional loop 
(`core/server/src/partition_reconciler.rs:338-355`), so the sleep arm is ready 
every iteration and `reconcile_once` runs back-to-back, starving the pump on 
that core. Add the two checks, or narrow the doc to the knobs it actually 
mirrors.



##########
core/server/src/dispatch/authz.rs:
##########
@@ -328,46 +261,19 @@ where
             |request| (&request.stream_id, &request.topic_id),
             Permissioner::get_consumer_groups,
         ),
-        _ => Ok(()),
-    }
-}
-
-/// Reply to a denied non-replicated read with the request's reply frame: empty
-/// body + nonzero `status`. The SDK peeks the status before body decode and
-/// surfaces the typed error, so a poll denial never reaches the empty-poll
-/// "0 messages" body path.
-#[allow(clippy::future_not_send)]
-pub(in crate::dispatch) async fn send_non_replicated_deny<B, MJ, S, SB>(
-    shard: &Rc<ShellShard<B, MJ, S, SB>>,
-    request: &Message<RoutedRequestHeader>,
-    transport_client_id: u128,
-    status: u32,
-) where
-    B: ShellBus,
-    MJ: JournalHandle + 'static,
-    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
-    S: 'static,
-    SB: SuperblockStore + 'static,
-{
-    let commit = current_metadata_commit(shard);
-    let reply = build_deny_reply(
-        request.header(),
-        request.header().client,
-        request.header().session,
-        commit,
-        status,
-    );
-    if let Err(error) = shard
-        .bus
-        .send_to_client(transport_client_id, 
reply.into_generic().into_frozen())
-        .await
-    {
-        warn!(
-            transport_client_id,
-            status,
-            error = %error,
-            "failed to surface non-replicated authz denial"
-        );
+        // No on-demand flush primitive exists. This arm is the only thing
+        // answering `FeatureUnavailable` for flush now: the HTTP path reaches
+        // the builder through `read_local`, whose call sites pass eleven fixed

Review Comment:
   **nit** — the corrected rationale is true today but rests on a cross-module 
count nothing enforces: "whose call sites pass eleven fixed read codes, none of 
them flush". I counted eleven `read_local` sites too 
(`core/server/src/http/handlers.rs:286,310,344,378,415,443,481,522,562,590,1718`),
 so it is accurate — but a twelfth site silently falsifies it, and the reader 
has no link to check. Prefer stating the local requirement (flush has no HTTP 
route, so this arm is the only source of `FeatureUnavailable` for it) without 
the number.



##########
core/server/src/dispatch/failure.rs:
##########
@@ -0,0 +1,677 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! The wire failure channels and the one send exit for host-built frames.
+//!
+//! Every frame the dispatch host builds, success or rejection, leaves through
+//! [`send_host_frame`], so the send-failure log has one shape. Which channel a
+//! failure rides is a wire contract with the SDK:
+//!
+//! | channel | carrier | when |
+//! |---|---|---|
+//! | [`FrameChannel::TypedDeny`] | Reply, nonzero status + empty body, or a 
result-framed rejection body | rejections that must unblock the SDK's lockstep 
request slot: checksum, authz, pre-consensus rewrite, unknown or unsupported 
non-replicated code, unbound non-PING read, transient replay hints |
+//! | [`FrameChannel::Eviction`] | session-terminal Eviction frame with a 
typed reason | the client must register again: `NoSession`, `MalformedLogin`, 
heartbeat and login evictions. The reason rides the channel label, since one 
`context` covers four of them |
+//! | [`FrameChannel::ResyncSentinel`] | status-0 poll reply, body carries 
`RESYNC_REQUIRED_PARTITION_SENTINEL` | a fenced consumer-group poll: the 
consumer must re-sync its assignment; HTTP mirrors it as 
`resync_required_polled_messages` in `crate::http::wire` |
+//! | [`FrameChannel::EmptyFrame`] | status-0 fail-fast body, empty or the 
16-byte empty poll | the partition cannot answer yet; the SDK fails fast (empty 
poll) and retries. A permanent client error never rides this channel: an 
undecodable body denies typed, because there is nothing to retry |
+//! | [`FrameChannel::Reply`] | status-0 success frame | host-built success 
replies: login/register, ping, logout, non-replicated read bodies, committed 
metadata replies |
+//! | silent drop | no frame | one deliberate case, a transient consensus 
submit failure: the SDK read-timeout replays the same request id, and a 
synthesized failure could contradict a write that commits moments later. A 
header `RequestHeader::validate` rejected also drops, but that one is a GAP, 
not a contract - the fields decode, so a deny could be echoed under the 
transport id, and the client instead waits out its read timeout |
+//! | HTTP status | HTTP status code | the HTTP spine maps the same rejections 
in `crate::http::error`; it never rides these frames |
+//!
+//! The last two send nothing, so [`FrameChannel`] has no variant for them.
+//!
+//! Scope: the table covers `crate::dispatch` only. Two neighbours answer on
+//! their own paths by design - the partitions engine builds and sends
+//! produce/poll replies, and the shard crate builds client-shaped denies of
+//! its own (`IggyShard::deny_partition_request_transient` and
+//! `stage_transient_deny`, both `TypedDeny`-shaped, the latter shedding the
+//! frame outright when its lifecycle queue is full).
+
+use crate::responses::{
+    NonReplicatedResponse, build_deny_reply, build_empty_reply, 
current_metadata_commit,
+};
+use crate::rewrite::RewriteStage;
+use crate::shell::{ShellBus, ShellShard};
+use bytes::Bytes;
+use consensus::{
+    EvictionContext, MetadataHandle, build_eviction_message,
+    build_incompatible_protocol_eviction_message, build_result_rejection_reply,
+};
+use iggy_binary_protocol::{EvictionReason, PrepareHeader, RoutedRequestHeader};
+use iggy_common::IggyError;
+use journal::superblock::SuperblockStore;
+use journal::{Journal, JournalHandle};
+use message_bus::BusMessage;
+use server_common::Message;
+use std::rc::Rc;
+use tracing::warn;
+
+/// Labels the channel a host-built frame rides, for the send-failure log.
+/// The taxonomy, including the two channels that never construct a frame,
+/// is on the module doc.
+#[derive(Clone, Copy, Debug)]
+pub(in crate::dispatch) enum FrameChannel {
+    TypedDeny,
+    /// The reason travels with the channel: five call sites share the
+    /// `"login_rejection"` context across four distinct reasons, so
+    /// `context` alone cannot tell a `MalformedLogin` send failure from an
+    /// `InvalidCredentials` one.
+    Eviction(EvictionReason),
+    ResyncSentinel,
+    EmptyFrame,
+    Reply,
+}
+
+impl std::fmt::Display for FrameChannel {
+    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result 
{
+        match self {
+            Self::TypedDeny => formatter.write_str("typed_deny"),
+            Self::Eviction(reason) => write!(formatter, 
"eviction({reason:?})"),
+            Self::ResyncSentinel => formatter.write_str("resync_sentinel"),
+            Self::EmptyFrame => formatter.write_str("empty_frame"),
+            Self::Reply => formatter.write_str("reply"),
+        }
+    }
+}
+
+/// The one send exit for host-built client frames. Best-effort: a failed
+/// send means the connection is gone (or its queue is full), and there is
+/// nothing left to reply on, so the error is logged and dropped. `frame` is
+/// any bus message: a contiguous frozen frame or the vectored poll reply.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_host_frame<B: ShellBus>(
+    bus: &B,
+    transport_client_id: u128,
+    frame: impl Into<BusMessage>,
+    channel: FrameChannel,
+    context: &'static str,
+) {
+    if let Err(send_error) = bus.send_to_client(transport_client_id, 
frame).await {
+        warn!(
+            transport_client_id,
+            error = %send_error,
+            channel = %channel,
+            context,
+            "failed to send host frame to client"
+        );
+    }
+}
+
+/// Reply to a request rejected before it reached its plane with the request's
+/// own frame: empty body + nonzero `status`. The nonzero status is the whole
+/// point: the SDK peeks it and surfaces the typed error, whereas a status-0
+/// frame reads as a committed ack for work that never happened. Silence is no
+/// better, the connection decodes replies in lockstep and would wedge on every
+/// later request.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_deny_reply<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    transport_client_id: u128,
+    request_header: &RoutedRequestHeader,
+    status: u32,
+) where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    let commit = current_metadata_commit(shard);
+    let reply = build_deny_reply(request_header, transport_client_id, 0, 
commit, status);
+    send_host_frame(
+        &shard.bus,
+        transport_client_id,
+        reply.into_generic().into_frozen(),
+        FrameChannel::TypedDeny,
+        "request_denial",
+    )
+    .await;
+}
+
+/// Deny a request from an unbound transport without disclosing the metadata
+/// commit frontier. The status is the only field a pre-authenticated caller
+/// needs, while the live commit would expose cluster write activity.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_unbound_deny_reply<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    transport_client_id: u128,
+    request_header: &RoutedRequestHeader,
+    status: u32,
+) where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    let reply = build_deny_reply(request_header, transport_client_id, 0, 0, 
status);
+    send_host_frame(
+        &shard.bus,
+        transport_client_id,
+        reply.into_generic().into_frozen(),
+        FrameChannel::TypedDeny,
+        "unbound_request_denial",
+    )
+    .await;
+}
+
+/// Reply to a denied non-replicated read with the request's reply frame: empty
+/// body + nonzero `status`. The SDK peeks the status before body decode and
+/// surfaces the typed error, so a poll denial never reaches the empty-poll
+/// "0 messages" body path.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_non_replicated_deny<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    request: &Message<RoutedRequestHeader>,
+    transport_client_id: u128,
+    status: u32,
+) where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    let commit = current_metadata_commit(shard);
+    let reply = build_deny_reply(
+        request.header(),
+        request.header().client,
+        request.header().session,
+        commit,
+        status,
+    );
+    send_host_frame(
+        &shard.bus,
+        transport_client_id,
+        reply.into_generic().into_frozen(),
+        FrameChannel::TypedDeny,
+        "non_replicated_denial",
+    )
+    .await;
+}
+
+/// Reject a request before it reaches consensus: warn, then send the typed
+/// deny reply. A silent drop would wedge every later request on the
+/// connection until the socket read timeout. `stage` names the chain step
+/// for both log lines, so the set stays enumerable.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_pre_consensus_deny<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    transport_client_id: u128,
+    request_header: &RoutedRequestHeader,
+    error: &IggyError,
+    stage: RewriteStage,
+) where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    let context = stage.as_str();
+    warn!(
+        transport_client_id,
+        error = %error,
+        operation = ?request_header.operation,
+        context,
+        "denying request pre-consensus"
+    );
+    let commit = current_metadata_commit(shard);
+    let reply = build_deny_reply(
+        request_header,
+        transport_client_id,
+        0,
+        commit,
+        error.as_code(),
+    );
+    send_host_frame(
+        &shard.bus,
+        transport_client_id,
+        reply.into_generic().into_frozen(),
+        FrameChannel::TypedDeny,
+        context,
+    )
+    .await;
+}
+
+/// Result-framed rejection Reply: status 0, body `[count=1][index][code]`.
+/// The SDK decodes the nonzero result code, so a transient code makes it
+/// replay the same request at once instead of waiting out its read timeout.
+/// Replying empty instead would surface as a hard `InvalidFormat` decode
+/// failure and break the replay.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_result_rejection<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    transport_client_id: u128,
+    request_header: &RoutedRequestHeader,
+    error: &IggyError,
+    context: &'static str,
+) where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    let commit = current_metadata_commit(shard);
+    let reply = build_result_rejection_reply(request_header, commit, 
error.as_code());
+    send_host_frame(
+        &shard.bus,
+        transport_client_id,
+        reply.into_generic().into_frozen(),
+        FrameChannel::TypedDeny,
+        context,
+    )
+    .await;
+}
+
+/// Best-effort session-terminal `Eviction` frame: the client's session is
+/// gone (or was never granted), so it must register again. Every frame
+/// transport decodes `Command::Eviction` and maps the typed reason
+/// (`NoSession` -> `Unauthenticated`, ...), so clients fail fast with the
+/// real cause instead of a body-decode failure or a timeout. Consensus
+/// context (cluster/view/replica) is stamped on the metadata shard and
+/// zeroed elsewhere; the SDK only reads the reason, plus the protocol
+/// window on `IncompatibleProtocol`.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_eviction<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    transport_client_id: u128,
+    vsr_client_id: u128,
+    reason: EvictionReason,
+    context: &'static str,
+) where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    let ctx = shard.plane.metadata().consensus.as_ref().map_or(
+        EvictionContext {
+            cluster: 0,
+            view: 0,
+            replica: 0,
+        },
+        EvictionContext::from_consensus,
+    );
+    let eviction = match reason {
+        EvictionReason::IncompatibleProtocol => {
+            build_incompatible_protocol_eviction_message(ctx, vsr_client_id)
+        }
+        _ => build_eviction_message(ctx, vsr_client_id, reason),
+    };
+    send_host_frame(
+        &shard.bus,
+        transport_client_id,
+        eviction.into_generic().into_frozen(),
+        FrameChannel::Eviction(reason),
+        context,
+    )
+    .await;
+}
+
+/// Send a non-replicated reply body to a client, stamping the current
+/// metadata commit. Shared by the non-replicated read arms; `channel`
+/// labels the body shape (a real answer, a fail-fast empty poll, or the
+/// re-sync sentinel) for the send-failure log.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_non_replicated_bytes<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    request: &Message<RoutedRequestHeader>,
+    transport_client_id: u128,
+    bytes: Bytes,
+    channel: FrameChannel,
+    context: &'static str,
+) where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    let commit = current_metadata_commit(shard);
+    let reply = NonReplicatedResponse::Bytes(bytes).into_reply(
+        request.header(),
+        request.header().client,
+        request.header().session,
+        commit,
+    );
+    send_host_frame(
+        &shard.bus,
+        transport_client_id,
+        reply.into_generic().into_frozen(),
+        channel,
+        context,
+    )
+    .await;
+}
+
+/// Ack a consumer-offset op whose body could not be rewritten for the
+/// partition plane with an empty Reply. The SDK connection processes replies
+/// in lockstep, so a silent drop wedges every subsequent request on that
+/// connection.
+#[allow(clippy::future_not_send)]
+pub(in crate::dispatch) async fn send_empty_partition_reply<B, MJ, S, SB>(

Review Comment:
   **nit** — `send_empty_partition_reply` appears to have no reachable 
producer. Its only caller is the `Err` arm of 
`maybe_rewrite_consumer_offset_request` (`partition.rs:472-483`), but 
`resolve_partition_request_namespace` already decodes the same bytes with the 
same `decode_from` earlier at `partition.rs:377` and denies typed at `:396-405` 
(`core/server/src/responses.rs:361-384`), so an undecodable body never reaches 
the rewrite. The second decode is on unmutated bytes, and 
`rewrite_request_body`'s only other error needs a body past `u32::MAX` against 
a 64 MiB cap.
   
   If that holds, the `EmptyFrame` row's "empty" carrier documents a shape 
nothing produces. Worth confirming and then either deleting this helper or 
noting why it is kept.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to