numinnex commented on code in PR #4036:
URL: https://github.com/apache/iggy/pull/4036#discussion_r3926680721
##########
core/server/src/dispatch/partition.rs:
##########
@@ -509,11 +511,13 @@ pub(in crate::dispatch) async fn handle_poll_messages<B,
MJ, S, SB>(
{
let Ok(wire) = PollMessagesRequest::decode_from(request_body(request))
else {
// Undecodable poll: keep the fail-fast empty-poll shape.
+ let (body, channel) = empty_poll_fallback(0);
Review Comment:
**warning** — this is the last fail-open of exactly the class the PR set out
to close, and it is on the hottest read path. An undecodable `POLL_MESSAGES`
body answers status 0 plus a *valid* 16-byte `PolledMessages` body, which every
SDK decodes as a successful 0-message poll, so a version-skewed consumer loops
forever with no error — strictly worse than a wedge, which at least surfaces a
timeout. Every sibling undecodable-body path in this same diff denies typed
(`rewrite.rs` `static_bounds` -> `InvalidCommand`; `mod.rs:640`, whose commit
is titled "deny malformed consumer-group requests instead of dropping").
No SDK breaks on the typed deny, and it is not a new wire shape: the server
already sends nonzero-status `POLL_MESSAGES` replies at `:539` (authz) and
`:566-577` (`PartitionNotFound`/`StreamIdNotFound`/`TopicIdNotFound`), pinned
end-to-end in
`integration/tests/server/poll_semantics_vsr.rs:70-73,113,145,188-192`. All
five native SDKs peek `status` before any body decode — Rust
`sdk/src/vsr.rs:204,234`, Go `reply.go:81`, Node `reply.ts:79`, Java
`VsrResponseHandler.java:204`, .NET `VsrReplyDecoder.cs:60` — and
Python/C++/PHP are FFI over Rust.
The rustdoc at `:496-497` ("empty body so the SDK fails fast on decode")
argues for the current shape but contradicts its own function: `:578-580`
explains the body is a valid 16-byte poll precisely so the decoder does *not*
fail. And `failure.rs:562-593` never enters `handle_poll_messages`, so this
change leaves that snapshot byte-identical.
Fix: `send_non_replicated_deny(shard, request, transport_client_id,
IggyError::InvalidCommand.as_code())`. Scope it to this undecodable-body branch
only — leave the `ConsumerGroupPartitionNotOwned` resync sentinel at `:585-593`
alone.
##########
core/server/src/dispatch/partition.rs:
##########
@@ -642,6 +695,7 @@ pub(in crate::dispatch) async fn
handle_get_consumer_offset<B, MJ, S, SB>(
request,
transport_client_id,
Bytes::new(),
+ FrameChannel::Reply,
Review Comment:
**nit** — same shape as the undecodable poll: a malformed request body gets
a success-shaped status-0 empty frame, which the SDK reads as "no stored
offset". It is also labelled `FrameChannel::Reply` with the identical
`"get_consumer_offset"` context as the legitimate no-offset success at `:757`,
so the two are indistinguishable in the send-failure log by every field. Deny
typed here too, or at minimum give this branch its own `context`.
##########
core/server/src/dispatch/reads.rs:
##########
@@ -154,17 +158,14 @@ pub(in crate::dispatch) async fn
handle_non_replicated_request<B, MJ, S, SB>(
request.header().session,
commit,
);
- if let Err(error) = shard
- .bus
- .send_to_client(transport_client_id,
reply.into_generic().into_frozen())
- .await
- {
- warn!(
- transport_client_id,
- error = %error,
- "failed to send non-replicated ping reply"
- );
- }
+ send_host_frame(
Review Comment:
**nit** — `record_heartbeat` was already called unconditionally in the
funnel (`dispatch/mod.rs:457`) before `classify` routed here, so the PING arm's
own call at `:153` is redundant: two `Instant::now()` and two `get_mut` per
ping. The funnel call strictly covers it. Pre-existing, not introduced here.
##########
core/server/src/dispatch/reads.rs:
##########
@@ -324,18 +333,14 @@ async fn handle_default_non_replicated<B, MJ, S, SB>(
request.header().session,
commit,
);
- if let Err(error) = shard
- .bus
- .send_to_client(transport_client_id,
reply.into_generic().into_frozen())
- .await
- {
- warn!(
- transport_client_id,
- code,
- error = %error,
- "failed to send non-replicated VSR reply"
- );
- }
+ send_host_frame(
Review Comment:
**warning** — the gate-deny branch at `:308-311` emits no log at all, while
this builder-`Err` branch 25 lines below `warn!`s with `code` and `error`.
`send_non_replicated_deny` only logs on *send* failure, so a successful refusal
leaves zero server-side breadcrumb — and the refusal this PR introduces (every
armless or unknown non-replicated code, via `authorize_default_read`'s new
tail) is exactly the one an operator would need to see. Add a `warn!` with
`code` + `error` on the gate-deny branch.
##########
core/server/src/dispatch/mod.rs:
##########
@@ -615,45 +457,34 @@ async fn handle_client_request<B, MJ, S, SB>(
sessions.borrow_mut().record_heartbeat(transport_client_id);
let header = *request.header();
- if header.operation == Operation::NonReplicated {
- // Auth bypass guard: `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.
- let nr_code =
u32::from_le_bytes(request.header().reserved[..4].try_into().unwrap());
- // Legacy (pre-register) login codes. The server authenticates only via
- // the Register handshake (LOGIN_REGISTER / LOGIN_REGISTER_WITH_PAT,
- // Operation::Register); the vsr SDK funnels both logins there and
never
- // emits these. Reject them uniformly with a typed MalformedLogin (the
- // SDK maps it to InvalidFormat) before the session gate, so a legacy
or
- // foreign client fails fast instead of getting the generic
- // Unauthenticated deny the pre-auth guard would send unbound, or the
- // silent empty-ok Reply the bound non-replicated path would send.
- if matches!(
- nr_code,
- LOGIN_USER_CODE | LOGIN_WITH_PERSONAL_ACCESS_TOKEN_CODE
- ) {
+ let bound = sessions.borrow().get_session(transport_client_id);
+ match classify(&header, bound.is_some()) {
+ RequestClass::LegacyLogin => {
+ // Legacy (pre-register) login codes. The server authenticates
only via
+ // the Register handshake (LOGIN_REGISTER /
LOGIN_REGISTER_WITH_PAT,
+ // Operation::Register); the vsr SDK funnels both logins there and
never
+ // emits these. Reject them uniformly with a typed MalformedLogin
(the
+ // SDK maps it to InvalidFormat) before the session gate, so a
legacy or
+ // foreign client fails fast instead of getting the generic
+ // Unauthenticated deny the pre-auth guard would send unbound, or
the
+ // silent empty-ok Reply the bound non-replicated path would send.
+ let nr_code = non_replicated_code(&header);
warn!(
transport_client_id,
code = nr_code,
"rejecting legacy login code; server requires the register
handshake"
);
- send_login_eviction(
+ send_eviction(
shard,
transport_client_id,
header.client,
EvictionReason::MalformedLogin,
+ "legacy login rejection",
)
.await;
- return;
}
- let allowed_pre_auth = nr_code == PING_CODE;
- if !allowed_pre_auth &&
sessions.borrow().get_session(transport_client_id).is_none() {
+ RequestClass::UnauthenticatedRead => {
+ let nr_code = non_replicated_code(&header);
Review Comment:
**warning** — the rationale just below (`:488-490`, "Foreign SDKs still
probe `GET_CLUSTER_METADATA` before login until they are fixed") is stale for
every SDK in tree, so a genuine unauthenticated read of the private replica
topology now hides at `debug` (`:491-495`).
All seven read the roster only *after* login: Node `client.socket.ts:773`
(comment says "deliberately absent: the server auth-gates it"), Go
`internal/util/leader_aware.go:52-88`, .NET `TcpMessageStream.Vsr.cs:489`, Java
`AsyncIggyTcpClient.java:854-856` (`loginAndSettleOnLeader` composes the roster
read onto the *completed* login future), and Python/C++/PHP inherit Rust.
Java's poll is the most aggressive of the seven — `LEADERLESS_WAIT_BUDGET = 5s`
at `LEADERLESS_POLL_INTERVAL = 250ms` (`LeaderAwareness.java:60-62`), up to 20
reads per login — but all 20 are authenticated, and this arm fires only for an
unbound transport, so dropping the special case cannot cause a warn flood.
Fix: drop the special case and let it `warn!` like every other pre-auth
denial.
##########
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:
**nit** — this `Err` returns with no frame at all, *after* the metadata op
already committed. The lockstep connection then wedges to its read timeout on
an op that succeeded server-side, and `failure.rs:31` declares exactly two
deliberate silences, of which this is not one.
Both `Err` arms look practically unreachable (`responses.rs:1566` already
gates `command == Command::Reply` on a server-built frame, and
`WireName::new(raw)` caps at 255 bytes against the fixed short token from
`pat.rs:131`), so this is cheap either way: it is reachable only for a PAT
mint, an admin-frequency path. Send the typed deny and add the row.
##########
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
Review Comment:
**nit** — worth stating what the reader needs here: the reason a replay
cannot help is that both of `maybe_rewrite_consumer_group_request`'s own errors
are `InvalidCommand` decode failures (`consumer_group.rs:80`, `:94`). Its third
error path, `rewrite_request_body`'s `InvalidConfiguration` (`wire.rs:76-77`),
needs a body past `u32::MAX` against a 64 MiB `MAX_MESSAGE_SIZE`, so it is
unreachable for any client frame — the deny is correct for all three
regardless. Fine as-is if you would rather not expand it; flagging only because
the sentence is load-bearing for why this is terminal rather than transient.
##########
core/server/src/dispatch/mod.rs:
##########
@@ -615,45 +457,34 @@ async fn handle_client_request<B, MJ, S, SB>(
sessions.borrow_mut().record_heartbeat(transport_client_id);
let header = *request.header();
- if header.operation == Operation::NonReplicated {
- // Auth bypass guard: `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.
- let nr_code =
u32::from_le_bytes(request.header().reserved[..4].try_into().unwrap());
- // Legacy (pre-register) login codes. The server authenticates only via
- // the Register handshake (LOGIN_REGISTER / LOGIN_REGISTER_WITH_PAT,
- // Operation::Register); the vsr SDK funnels both logins there and
never
- // emits these. Reject them uniformly with a typed MalformedLogin (the
- // SDK maps it to InvalidFormat) before the session gate, so a legacy
or
- // foreign client fails fast instead of getting the generic
- // Unauthenticated deny the pre-auth guard would send unbound, or the
- // silent empty-ok Reply the bound non-replicated path would send.
- if matches!(
- nr_code,
- LOGIN_USER_CODE | LOGIN_WITH_PERSONAL_ACCESS_TOKEN_CODE
- ) {
+ let bound = sessions.borrow().get_session(transport_client_id);
Review Comment:
**nit (perf)** — `get_session` is now unconditional here, where pre-PR it
ran only inside the `NonReplicated` non-pre-auth branch and again after the
Register/Logout returns. PING, Register and Logout each pay one extra `RefCell`
borrow plus a map lookup. This is the only per-frame cost this diff adds, and
it is small — PING is per-heartbeat-interval, Register/Logout are once per
connection — while the diff removes more than it adds (five `Rc`/`Arc` clones
now sit behind the `active.insert` guard instead of ahead of it, one fewer
checked header cast, one bus clone moved to handler-build time). Net per-frame
cost is negative. Not blocking; noting so it is a deliberate trade.
Two pre-existing items on the same line region, for a follow-up rather than
this PR: `let header = *request.header()` at `:459` copies 256 bytes on every
frame when only the `ReplicatedMetadata` arm needs the pre-`transmute_header`
snapshot — bind the *reference* once here and keep the by-value copy inside
that arm (do not replace it with per-arm `request.header()` calls:
`RequestBacking::header` runs `bytemuck::checked::try_from_bytes` + `.expect`
per call at `server_common/src/consensus_message.rs:159-163`, so N accessor
calls can cost more than the one memcpy — same reason `failure.rs:168-170`,
`:325-327` and `reads.rs:157-160` should bind once). And five same-key lookups
happen per data-plane frame (`bus.client_meta` plus an `Rc` clone at `:451`,
`ensure_connection`, `record_heartbeat` at `:457`, this `get_session`, then
`get_user_id`/`read_context`) where `ensure_transport_connection` is a no-op
after frame 1, so the `client_meta` lookup is dead work — one `connections`
walk mirroring the existing `read_context` would collapse them.
##########
core/server/src/dispatch/mod.rs:
##########
@@ -299,8 +230,19 @@ fn enqueue_client_request<B, MJ, S, SB>(
return;
}
- let bus = shard.bus.clone();
+ let shard_handle = Rc::clone(shard_handle);
+ let sessions = Rc::clone(sessions);
+ let system_config = Arc::clone(system_config);
+ let queues = Rc::clone(queues);
+ let active = Rc::clone(active);
bus.spawn(async move {
+ // The handle is set once the shard is built. A frame that beats
+ // it stays queued and the client's next frame drains both, so the
+ // active slot must be released here or that next frame never spawns.
+ let Some(shard) = upgrade_shard_handle(&shard_handle) else {
+ active.borrow_mut().remove(&client_id);
Review Comment:
**nit** — releasing the active slot here is right (otherwise the client's
next frame finds the slot taken and nothing ever drains), but the queued frame
stays in `queues[client_id]` forever if the client sends nothing more: the
connection-lost hook clears neither `queues` nor `active`. Unreachable in
production, since listeners bind after the backfill (`boot/mod.rs:912-920`
asserts it).
Three pre-existing items share this state machine and want one fix in the
connection-lost hook: a panic inside `handle_client_request` leaves `client_id`
in `active` forever (compio catches the detached task's panic), so every later
frame for that client queues and never drains; `pop_next_client_request`
removes the map entry whenever a queue drains to empty (`:303`), so a lockstep
SDK pays a `HashMap` insert plus a `VecDeque` alloc/free per request — keeping
the entry is safe, `pop_next_client_request` already handles empty-but-present
correctly; and the per-client `VecDeque` is unbounded with no end-to-end
backpressure, because the bounded socket->dispatch channel drains at spin speed
into it (`enqueue_client_request` is synchronous). Also consider `AHashMap` for
`queues`/`active` and `SessionManager`'s maps — ids are server-minted so there
is no HashDoS surface, `ahash` is already a direct dep, and
`message_bus/src/lib.rs:702` sets the precedent (check `iter_clients` orderi
ng first).
##########
core/server/src/dispatch/mod.rs:
##########
@@ -366,193 +308,93 @@ fn pop_next_client_request(
message
}
-/// Per-request partitions-count cap, shared by create-topic, create-partitions
-/// and delete-partitions admission. Runs pre-consensus like
-/// [`validate_topic_bounds`]: an oversized count must not burn a replicated
-/// log entry (create-partitions admission would also allocate that many
-/// consensus-group ids before replicating).
-///
-/// Zero passes here because a zero-partition TOPIC is legal (legacy
-/// `create_topic` admits `0..=MAX`); the add/remove requests reject it in
-/// [`validate_partitions_change_count`].
-const fn validate_partitions_count(partitions_count: u32) -> Result<(),
IggyError> {
- if partitions_count > MAX_PARTITIONS_PER_REQUEST {
- return Err(IggyError::TooManyPartitions);
- }
- Ok(())
+/// Where the funnel routes a client request. Derived by [`classify`], which
+/// IS the routing: [`handle_client_request`] matches on its result. The
+/// variant ORDER mirrors the order of the checks inside [`classify`], and
+/// that order is semantics (documented there).
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum RequestClass {
Review Comment:
**nit** — `RequestClass` and `classify` (`:362`) are `pub` but referenced
only inside this file, its own tests included; sibling dispatch items use
`pub(in crate::dispatch)` and the submodules are private `mod`. Tighten both.
##########
core/server/src/dispatch/mod.rs:
##########
@@ -172,6 +114,12 @@ where
})
}
+/// Build the shard's one client-request handler: per-client FIFO queues
+/// drained one task per client, and the bus connection-lost hook that
+/// logs a dropped connection out. Every transport on the shard must
+/// share the instance (shard 0 hands it to its local QUIC, TCP-TLS and
+/// WSS listeners as well), or a client's ordering guarantee and the
+/// disconnect hook split by transport.
Review Comment:
**nit** — the connection-lost hook below puts `remove_connection(client_id)`
as the *first* operand of the `&&` chain, so it mutates before the weak upgrade
is attempted: a failed upgrade strips the `SessionManager` entry without
submitting the replicated `Logout`, leaking the `ClientTable` entry and its
consumer-group memberships. The `&&` chain is unchanged from pre-PR, but shard
0 only now reaches it — the deleted `make_client_request_handler` held a strong
`Rc<shard>` whose upgrade could not fail. Window is pre-build /
post-runtime-drop only, so effectively unreachable, and it self-heals via
`recover()`'s `remove_consumer_group_member` on next boot. Cheap fix: upgrade
first, then `remove_connection`.
##########
core/server/src/dispatch/session_ops.rs:
##########
@@ -552,35 +475,18 @@ async fn evict_stale_client<B, MJ, S, SB>(
if let Some((vsr_client_id, session)) = bound {
submit_disconnect_logout(Rc::clone(shard), vsr_client_id, session);
}
- let ctx = shard.plane.metadata().consensus.as_ref().map_or(
- consensus::EvictionContext {
- cluster: 0,
- view: 0,
- replica: 0,
- },
- consensus::EvictionContext::from_consensus,
- );
- let eviction = consensus::build_eviction_message(
- ctx,
+ warn!(
Review Comment:
**nit** — this `warn!` now runs unconditionally *before* `send_eviction`,
where pre-PR it sat in the `else` of the send result and fired only on success.
The log now asserts a delivery it did not observe, and the contradicting line
says only "failed to send host frame" with `context = "stale-client eviction"`.
##########
core/server/src/dispatch/session_ops.rs:
##########
@@ -331,58 +331,28 @@ async fn surface_login_failure<B, MJ, S, SB>(
SB: SuperblockStore + 'static,
{
if error.is_terminal() {
- send_login_eviction(
+ send_eviction(
shard,
transport_client_id,
request_header.client,
eviction_reason_for(error),
+ "login rejection",
Review Comment:
**nit** — this label is passed at four more sites (`:1186`, `:1204`,
`:1302`, `:1317`) covering four distinct `EvictionReason`s, so `context` cannot
distinguish them; see the `send_host_frame` comment about the dropped `reason`
field. More broadly, the one `context` field now carries three naming
conventions — snake_case opcodes (`reads.rs:124` `"get_me"`, `partition.rs:551`
`"poll_messages"`), kebab-case stages (`rewrite.rs:110`
`"personal-access-token"`), prose (`failure.rs:117` `"request denial"`, this
line) and hybrids (`partition.rs:1088` `"delete_segments reply"`) — which makes
filtering by it impossible. Repo rule 7 wants labels matching literal API
names; pick one shape.
--
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]