numinnex commented on code in PR #4036:
URL: https://github.com/apache/iggy/pull/4036#discussion_r3926679128
##########
core/server/src/boot/mod.rs:
##########
@@ -627,6 +645,18 @@ async fn shard_main(
))
.await?;
+ // The shard above owns the metadata state machine's only write handle,
+ // and the peers read through it until their runtimes are gone. Declared
Review Comment:
**warning** — this comment claims more than the code delivers: "every exit
from here on, clean or `?`, waits for the peers" is false for the exits that
matter.
`build_shard_for_thread` takes `metadata: ServerMetadata` **by value**
(`boot/recovery.rs:75`), so the `?` at `recovery.rs:112`, `:201` and `:237`
drop shard 0's only metadata `WriteCell` *inside the callee*, and a failed
`.build()` at `:265` drops it inside the builder. All four happen before
`_peer_exit_wait` exists. Peers are reading concurrently — their own build's
first statement is `metadata.mux_stm.streams().read(..)` (`recovery.rs:84`,
again `:144`), then the reconciler's first pass (`partition_reconciler.rs:336`
-> `:1115`) — and hit `.expect("read handle should be accessible")` at
`metadata/src/stm/mod.rs:129`. Nothing serializes the two builds:
`BootstrapBarrier` is signalled after both.
Same shape one level up: `broadcast_metadata_bundle(..).await?`
(`boot/mod.rs:525-533`) can `Err` after serving *some* peers
(`boot/handoff.rs:153`, `:157`), and `recovered.mux_stm` is a local of that
match arm, so no outer binding can drop after it.
Drop-order cannot reach a callee-owned drop, so moving this declaration up
does not fix it. Fix: arm the wait right after `broadcast_metadata_bundle`, and
keep the writer out of a fallible call that drops it (`Rc` the mux state
machine, clone held before the guard). Note returning `metadata` in the `Err`
variant does **not** work — the receiving binding is declared after the guard
and drops first.
The by-value signature predates this PR; the invariant claim is new, which
is what makes it worth fixing here. Consequence today is misattribution, not
data loss: panics unwind (no profile sets `panic = "abort"`), so `join_all`
still reports shard 0's real `Err` first, but `first_panic`
(`threads.rs:205-207`) records a peer's read-handle panic instead.
##########
core/server/src/boot/mod.rs:
##########
@@ -627,6 +645,18 @@ async fn shard_main(
))
.await?;
+ // The shard above owns the metadata state machine's only write handle,
+ // and the peers read through it until their runtimes are gone. Declared
+ // after the shard so it drops first: every exit from here on, clean or
+ // `?`, waits for the peers before the write side goes with the shard.
+ let _peer_exit_wait = (shard_id == 0).then(|| {
+ PeerExitWait::new(
+ peer_exit,
+ Arc::clone(&shutdown_flag_for_handoff),
Review Comment:
**warning** — the peer wait and the outer thread join share
`shutdown_join_timeout` and **nest**, so shard 0 can be reported as the wedged
shard while it is legitimately waiting.
`join_all` arms ONE shared `deadline = now + join_timeout` on the first 25ms
poll that sees the shutdown flag (`boot/threads.rs:110-113`, `:236-238`). On
graceful Ctrl-C the signal handler sets that flag (`threads.rs:71`) at T0, so
the outer clock arms at T0+25ms independent of shard 0. Shard 0 cannot *start*
this wait until after `await_pump_drain` (<= `shutdown_drain_timeout`) plus the
watchdog's own `bus.shutdown(drain_timeout)` (`threads.rs:729`) — a second
drain budget. So shard 0 returns at up to `2*drain + join` = 50s at defaults,
against a 30s deadline.
Result: `"shard thread still running at the shutdown join deadline;
abandoning it"`, a `Wedged` row, `JoinHandle` dropped while the OS thread still
holds the metadata writer, and a non-zero exit on a correct shutdown. Because
the deadline is shared and already blown, every shard after shard 0 in the loop
is abandoned too.
Fix: bound this wait by the *remaining* join budget (thread a deadline
through rather than passing the raw knob). Related:
`validate_sharding_runtime_knobs` (`boot/threads.rs:496-536`) re-checks `poll
<= drain` but not `join >= drain` (or `join <= MAX`), which
`ShardingConfig::validate` does enforce
(`configs/src/server_config/sharding.rs:251-268`) — that gap reaches the same
abandon-mid-wait outcome by a second route, and `PeerExitWait` is a new
consumer of `join`.
##########
core/server/src/boot/threads.rs:
##########
@@ -333,6 +333,132 @@ impl Drop for ShutdownOnDrop {
}
}
+/// 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.
Review Comment:
**nit (doc)** — stated unconditionally, but `PeerExitWait::drop` skips the
wait entirely on `thread::panicking()` (`:447`), and `install_panic_hook` only
sets the flag without aborting (`:196-210`). So on a shard-0 panic the writer
drops with peers still reading and their pumps panic at
`metadata/src/stm/mod.rs:129`.
The early return is the **right** behaviour and should stay — blocking an
unwinding thread on a `Condvar` for up to 30s inside `runtime.block_on` would
stall the io_uring driver during a panic, and if the panic came from shard 0's
pump the peers may be waiting on shard 0, so the wait would burn its whole
budget or deadlock. Only the comment needs narrowing: say the invariant holds
on every exit **except** the panic path, and why.
##########
core/server/src/dispatch/failure.rs:
##########
@@ -0,0 +1,652 @@
+// 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 |
+//! | [`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 |
+//! | [`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 | deliberate only where a reply would be wrong:
an undecodable header (nothing to echo), a transient consensus submit failure
(the SDK read-timeout replays) |
Review Comment:
**warning** — this row is not exhaustive, and its justification is false for
one of the classes it covers.
"an undecodable header (nothing to echo)" —
`try_into_typed::<RequestHeader>` (`dispatch/mod.rs:415`) runs `verify_frame()`
then `validate()` (`server_common/src/consensus_message.rs:310-311`), and
`RequestHeader::FRAME_SEALED = false`
(`binary_protocol/src/consensus/header.rs:501`) so `verify_frame` is a no-op —
`validate` is the whole check. `validate_request_fields` (`header.rs:442-490`)
rejects `client == 0`, `Operation::Reserved`, `Register` with `session !=
0`/`request != 0`, and non-register non-`NonReplicated` with `session ==
0`/`request == 0`. Every one of those **decodes cleanly**, so
`build_deny_reply` has all the fields it needs — there is plenty to echo. Those
clients get no frame at all and wedge to their read timeout, and the victim of
the `session == 0` case is exactly the lost-session client
`RequestClass::UnboundReplicated` was built to fail fast.
`dispatch/partition.rs:483`'s own `request.max(1)` comment ("older and internal
callers may still send zero")
names callers this silently drops, and that normalization sits *after*
`validate`.
Two more routes answer nothing and are not listed: `dispatch/mod.rs:663`
(`build_raw_pat_reply` `Err` returns with no frame *after* the op committed —
both `Err` arms look practically dead, but the row is still missing) and
`shard/src/lib.rs:3400-3402` (`stage_transient_deny` sheds on a `try_send`
failure and counts a frame drop).
Fix: enumerate the real set, or scope the table explicitly to
`crate::dispatch`. Certifying an echoable drop as deliberate is the part that
should not ship — it is what stops the next reader from checking.
##########
core/server/src/dispatch/failure.rs:
##########
@@ -0,0 +1,652 @@
+// 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 |
+//! | [`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 |
Review Comment:
**warning** — "the partition cannot answer yet ... and retries" does not
describe two of the cases that ride this channel. An **undecodable client
body** is a permanent client error, not a transient partition condition, and
there is nothing to retry: `dispatch/partition.rs:514` (poll) and `:698`
(consumer-offset) both answer a malformed request with a success-shaped
status-0 frame.
The taxonomy has no row for "undecodable client body answered
success-shaped". Add one, or give those two sites a distinct `context` — today
`:698` and the legitimate no-stored-offset success at `:757` send
byte-identical `Bytes::new()` on the same channel with the same
`"get_consumer_offset"` label, so nothing in the log tells them apart.
##########
core/server/src/dispatch/failure.rs:
##########
@@ -0,0 +1,652 @@
+// 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 |
+//! | [`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 |
+//! | [`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 | deliberate only where a reply would be wrong:
an undecodable header (nothing to echo), a transient consensus submit failure
(the SDK read-timeout replays) |
+//! | 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.
+//! This exit covers HOST-built frames only: the partitions engine builds and
Review Comment:
**warning** — the exemption is scoped to "produce/poll replies", but two
`TypedDeny`-shaped client frames are built and sent from outside this module,
so "the one send exit every host-built frame takes" holds only for frames built
in this crate.
`IggyShard::deny_partition_request_transient` (`shard/src/lib.rs:3346-3358`)
and `stage_transient_deny` (`:3382-3390`) both build
`build_deny_reply_from_request_header(request_header,
IggyError::TransientNotAccepted.as_code())` — a nonzero-status client Reply,
exactly the `TypedDeny` shape at `:26` — and deliver it via
`bus.send_to_client` and `LifecycleFrame::ForwardClientSend` respectively. A
transient deny is neither a produce nor a poll reply.
Fix: widen the exemption to name the shard-crate partition denies, or route
them through `send_host_frame`.
##########
core/server/src/dispatch/failure.rs:
##########
@@ -0,0 +1,652 @@
+// 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 |
+//! | [`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 |
+//! | [`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 | deliberate only where a reply would be wrong:
an undecodable header (nothing to echo), a transient consensus submit failure
(the SDK read-timeout replays) |
+//! | 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.
+//! This exit covers HOST-built frames only: the partitions engine builds and
+//! sends produce/poll replies on its own path by design.
+
+use crate::responses::{
+ NonReplicatedResponse, build_deny_reply, build_empty_reply,
current_metadata_commit,
+};
+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,
+ Eviction,
+ ResyncSentinel,
+ EmptyFrame,
+ 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. `context` labels the rejection
+/// site in both log lines.
+#[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,
+ context: &'static str,
+) where
+ B: ShellBus,
+ MJ: JournalHandle + 'static,
+ MJ::Target: Journal<Entry = Message<PrepareHeader>, Header =
PrepareHeader>,
+ S: 'static,
+ SB: SuperblockStore + 'static,
+{
+ 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,
+ 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>(
+ shard: &Rc<ShellShard<B, MJ, S, SB>>,
+ transport_client_id: u128,
+ request_header: &RoutedRequestHeader,
+) 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_empty_reply(request_header, transport_client_id, 0,
commit);
+ send_host_frame(
+ &shard.bus,
+ transport_client_id,
+ reply.into_generic().into_frozen(),
+ FrameChannel::EmptyFrame,
+ "empty partition reply",
+ )
+ .await;
+}
+
+// Byte snapshots pinning each channel's frame to the pre-refactor inline
+// construction. DELIBERATELY temporary: they freeze the refactor, not the
Review Comment:
**warning** — repo rule 8 says comments must never reference the current or
a future task/PR, and "a later PR removes them" is an unenforceable removal
trigger: it leaves ~250 lines of snapshot test with nothing a future reader can
evaluate. State the invariant these snapshots pin instead, or file the removal
and drop the note.
##########
core/server/src/dispatch/failure.rs:
##########
@@ -0,0 +1,652 @@
+// 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 |
+//! | [`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 |
+//! | [`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 | deliberate only where a reply would be wrong:
an undecodable header (nothing to echo), a transient consensus submit failure
(the SDK read-timeout replays) |
+//! | 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.
+//! This exit covers HOST-built frames only: the partitions engine builds and
+//! sends produce/poll replies on its own path by design.
+
+use crate::responses::{
+ NonReplicatedResponse, build_deny_reply, build_empty_reply,
current_metadata_commit,
+};
+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,
+ Eviction,
+ ResyncSentinel,
+ EmptyFrame,
+ 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,
Review Comment:
**warning** — this PR deleted a diagnostic it had in hand. The
`send_login_eviction` helper this replaced logged `reason = ?reason` on a send
failure; `send_host_frame` logs only `transport_client_id`, `error`, `channel`
and `context`. `send_eviction` (`:267-302`) receives `reason` and never logs
it, and five call sites in `dispatch/session_ops.rs` (`:339`, `:1186`, `:1204`,
`:1302`, `:1317`) all pass the same `"login rejection"` label across four
distinct `EvictionReason`s — so `MalformedLogin`, `IncompatibleProtocol` and
`InvalidCredentials` are indistinguishable when a send fails. Add `reason` to
the eviction send-failure line.
##########
core/server/src/dispatch/authz.rs:
##########
@@ -328,46 +261,17 @@ 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, so the code stays the builder's
+ // `FeatureUnavailable` (its arm still serves the HTTP caller).
Review Comment:
**warning** — "(its arm still serves the HTTP caller)" is false.
`read_local` (`http/reads.rs:104`) is the only HTTP entry to
`build_non_replicated_response`, and its eleven call sites pass fixed constants
(`http/handlers.rs:290`, `314`, `348`, `382`, `419`, `447`, `485`, `526`,
`566`, `594`, `1722`) — `DESCRIBE_OPTIONS`, `GET_STREAM(S)`, `GET_TOPIC(S)`,
`GET_USER(S)`, `GET_CONSUMER_GROUP(S)`, `GET_STATS`,
`GET_PERSONAL_ACCESS_TOKENS`. No `FLUSH_UNSAVED_BUFFER`, and there are zero
flush references anywhere under `core/server/src/http/` — which is why
`flush_vsr.rs` is declared `test_client_transport = [Tcp]`.
So this gate arm is now the only thing producing `FeatureUnavailable` for
flush (keep it), and `responses.rs:590` is dead. Drop the false clause from the
comment.
##########
core/server/src/dispatch/authz.rs:
##########
@@ -328,46 +261,17 @@ 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, so the code stays the builder's
+ // `FeatureUnavailable` (its arm still serves the HTTP caller).
+ FLUSH_UNSAVED_BUFFER_CODE => Err(IggyError::FeatureUnavailable),
+ // A replicated code smuggled inside a `NonReplicated` header keeps the
+ // builder's `FeatureUnavailable`; a table-listed code with no arm
above
+ // and an unknown code are both refused as `InvalidCommand`, never
+ // deferred to the builder's empty-ok catch-all.
Review Comment:
**nit** — with this tail refusing every armless code first,
`build_non_replicated_response`'s `_ => Ok(NonReplicatedResponse::Empty)`
(`responses.rs:598`) is unreachable from **both** callers: the gate refuses it
on the TCP side, and HTTP passes only its eleven named codes. The empty-ok
catch-all this PR set out to kill therefore survives as a fail-open trap for
the next caller that lands.
`Err(IggyError::InvalidCommand)` there is provably safe: all eleven HTTP
codes have named arms above the `_`, and the not-found -> 404 mapping comes
from *inside* `GET_STREAM`/`GET_TOPIC`/`GET_USER`'s own arms via
`response.map_or(Empty, ..)`, not from the catch-all. Sequence it with the
flush arm above so flush's answer does not silently flip from
`FeatureUnavailable` to `InvalidCommand`.
--
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]