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


##########
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:
   Fixed, with one deviation from the suggested order: the wait is armed 
immediately **before** `broadcast_metadata_bundle`, not after it. The broadcast 
can fail after serving some peers, and by then those peers already hold read 
handles, so arming after it leaves exactly the window you describe one level up.
   
   `IggyMetadata::mux_stm` is now `Rc<M>`, with `new` taking `impl Into<Rc<M>>` 
so the 21 existing call sites are untouched. `shard_main` keeps a clone in a 
binding declared before the wait, and the broadcast moved out of the 
`MetadataHandoff::Owner` arm to just after it. A `recover()` failure inside 
that arm needs no wait at all: no peer holds a handle yet, they are still 
parked in `await_metadata_bundle`.
   
   Drop order is now shard, then `metadata` and its `Rc` clone, then the wait, 
then `shard_main`'s own clone. So the writer outlives every reader on each exit 
from the broadcast onward, the `?` returns inside the callee and a failed 
`.build()` included.



##########
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:
   The certification is gone. The row now names the transient consensus submit 
failure as the one deliberate silence, and calls the `RequestHeader::validate` 
drop what it is: a gap, not a contract. It says why that matters too, that the 
fields decode so a deny could be echoed under the transport id, and that the 
client instead waits out its read timeout.
   
   I did not turn it into a deny in this PR. At that point 
`try_into_typed::<RequestHeader>` has already failed, so only a `GenericHeader` 
exists: building the reply would mean reading `operation` and `request` out of 
a header the server just refused to interpret, and `into_routed()` has not run. 
That wants its own change with its own test rather than a rider here.
   
   A `Scope:` paragraph under the table also states that it covers 
`crate::dispatch` only, and names the two neighbours that answer on their own 
paths, including `stage_transient_deny` shedding the frame when its lifecycle 
queue is full.



##########
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:
   Taken the third way. Instead of adding a row for "undecodable client body 
answered success-shaped", both sites now deny typed, so nothing of that shape 
rides `EmptyFrame` any more. The row says so: a permanent client error never 
rides this channel, because there is nothing to retry.
   
   While there, `handle_delete_segments_request` had the same shape one 
function down and was not in the list: an undecodable body acked status 0 with 
an empty body, which reads as a completed trim. It denies typed now as well, 
and the `Err(_) => None` arm plus the `if let Some(truncate)` wrapper collapsed 
into one `Err(error)` arm that carries the real code.



##########
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:
   Fixed, by a different route than a `reason` field on the send-failure line: 
`FrameChannel::Eviction` now carries the `EvictionReason`, and `FrameChannel` 
grew a `Display` that renders it, so the line reads 
`channel=eviction(MalformedLogin)`.
   
   That keeps one shape for the log line while making the reason part of the 
channel's identity, which is what the taxonomy already claims the channel is. 
It also means `send_host_frame` needs no eviction-specific parameter, and the 
five `"login_rejection"` sites become distinguishable without touching their 
context labels.



##########
core/server/src/dispatch/mod.rs:
##########
@@ -681,353 +512,179 @@ 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,
             )
             .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));
+            // `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;
+        }
+        RequestClass::ReplicatedMetadata => {
+            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!(
+                    // The rewrite only ever fails on an undecodable body, so a
+                    // replay cannot help: deny typed instead of leaving the
+                    // lockstep connection to its read timeout.
+                    send_pre_consensus_deny(
+                        shard,
                         transport_client_id,
-                        error = %error,
-                        "failed to build raw PAT reply"
-                    );
+                        &header,
+                        &error,
+                        "consumer-group",
+                    )
+                    .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"

Review Comment:
   Deny sent. No row added, because with the deny in place this is no longer a 
silence: it rides `TypedDeny`, which the first row already covers.
   
   The comment states the trade instead, that the op committed and only the 
reply could not be rendered, so the client is told the request failed and can 
read the minted token back. Silence on a committed op is the worse of the two.



-- 
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