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


##########
core/server/src/rewrite.rs:
##########
@@ -0,0 +1,534 @@
+// 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 two pre-consensus request-rewrite chains, side by side.
+//!
+//! [`tcp_chain`] serves the TCP funnel, [`http_chain`] the HTTP submit. Each
+//! step rewrites or validates a request BEFORE consensus, so a rejected
+//! request burns no replicated log entry and no plaintext secret enters
+//! consensus. The chains enter the PAT rewrite through different functions
+//! on purpose: TCP resolves the acting user from the transport
+//! `SessionManager` ([`maybe_rewrite_pat_request`]), HTTP authenticates
+//! against its own session table and passes the resolved `user_id`
+//! ([`rewrite_pat_request_for_user`]).
+//!
+//! Three more links complete the chains but live in their spines, because
+//! they fire partition-read mesh RPCs or are plane-specific. Two are async:
+//! the consumer-group Join/Leave enrichment ([`crate::consumer_group`], the
+//! TCP funnel calls it after [`tcp_chain`]) and the `DeleteSegments` ->
+//! `TruncatePartition` resolution
+//! (`dispatch::partition::resolve_delete_segments_truncate`, called by both
+//! spines). The third, the consumer-offset rewrite on the partition path
+//! (`consumer_group::maybe_rewrite_consumer_offset_request`, called by
+//! `dispatch::partition`), is synchronous.
+
+use crate::pat::{maybe_rewrite_pat_request, rewrite_pat_request_for_user};
+use crate::segment_cleaner::UNENFORCEABLE_TOPIC_SIZE_WARN;
+use crate::session_manager::SessionManager;
+use crate::shell::{ShellBus, ShellShard};
+use crate::users::maybe_rewrite_user_password_request;
+use crate::wire::request_body;
+use consensus::MetadataHandle;
+use iggy_binary_protocol::requests::partitions::{
+    CreatePartitionsRequest, DeletePartitionsRequest,
+};
+use iggy_binary_protocol::requests::streams::{CreateStreamRequest, 
UpdateStreamRequest};
+use iggy_binary_protocol::requests::topics::{CreateTopicRequest, 
UpdateTopicRequest};
+use iggy_binary_protocol::requests::users::{CreateUserRequest, 
UpdateUserRequest};
+use iggy_binary_protocol::{
+    MAX_PARTITIONS_PER_REQUEST, Operation, PrepareHeader, RoutedRequestHeader, 
WireDecode,
+    WireIdentifier, WireOptions,
+};
+use iggy_common::{
+    IggyByteSize, IggyError, MaxTopicSize, TopicCreateOptions, 
UPDATABLE_STREAM_OPTION_KEYS,
+    UPDATABLE_TOPIC_OPTION_KEYS, UPDATABLE_USER_OPTION_KEYS, 
validate_preallocated_topic_bytes,
+    validate_topic_segment_size,
+};
+use journal::superblock::SuperblockStore;
+use journal::{Journal, JournalHandle};
+use metadata::impls::metadata::StreamsFrontend;
+use metadata::stm::stream::Streams;
+use server_common::Message;
+use std::cell::RefCell;
+use std::rc::Rc;
+use tracing::warn;
+
+/// A staged pre-consensus rejection: `stage` labels the chain step for the
+/// deny log line, `error` is the typed code the deny reply carries.
+pub struct RewriteDeny {

Review Comment:
   **nit** — the set of pre-consensus stages is split between this struct and a 
call-site literal. Three live in `stage` (`:110` `"personal-access-token"`, 
`:123` `"user-password"`, `:127` `"static-bounds"`); the fourth is a bare 
`"consumer-group"` at `dispatch/mod.rs:640`, fed to the same `context` 
parameter by the same `send_pre_consensus_deny`. So the type does not own the 
set it appears to define, and a fifth stage has two plausible homes. One enum 
(or one const set) that `mod.rs:640` must also name would close it. (`error` is 
already a typed `IggyError`, so this is about cohesion, not stringly-typed 
errors.)



##########
core/server/src/rewrite.rs:
##########
@@ -0,0 +1,534 @@
+// 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 two pre-consensus request-rewrite chains, side by side.
+//!
+//! [`tcp_chain`] serves the TCP funnel, [`http_chain`] the HTTP submit. Each
+//! step rewrites or validates a request BEFORE consensus, so a rejected
+//! request burns no replicated log entry and no plaintext secret enters
+//! consensus. The chains enter the PAT rewrite through different functions
+//! on purpose: TCP resolves the acting user from the transport
+//! `SessionManager` ([`maybe_rewrite_pat_request`]), HTTP authenticates
+//! against its own session table and passes the resolved `user_id`
+//! ([`rewrite_pat_request_for_user`]).
+//!
+//! Three more links complete the chains but live in their spines, because
+//! they fire partition-read mesh RPCs or are plane-specific. Two are async:
+//! the consumer-group Join/Leave enrichment ([`crate::consumer_group`], the
+//! TCP funnel calls it after [`tcp_chain`]) and the `DeleteSegments` ->
+//! `TruncatePartition` resolution
+//! (`dispatch::partition::resolve_delete_segments_truncate`, called by both
+//! spines). The third, the consumer-offset rewrite on the partition path
+//! (`consumer_group::maybe_rewrite_consumer_offset_request`, called by
+//! `dispatch::partition`), is synchronous.
+
+use crate::pat::{maybe_rewrite_pat_request, rewrite_pat_request_for_user};
+use crate::segment_cleaner::UNENFORCEABLE_TOPIC_SIZE_WARN;
+use crate::session_manager::SessionManager;
+use crate::shell::{ShellBus, ShellShard};
+use crate::users::maybe_rewrite_user_password_request;
+use crate::wire::request_body;
+use consensus::MetadataHandle;
+use iggy_binary_protocol::requests::partitions::{
+    CreatePartitionsRequest, DeletePartitionsRequest,
+};
+use iggy_binary_protocol::requests::streams::{CreateStreamRequest, 
UpdateStreamRequest};
+use iggy_binary_protocol::requests::topics::{CreateTopicRequest, 
UpdateTopicRequest};
+use iggy_binary_protocol::requests::users::{CreateUserRequest, 
UpdateUserRequest};
+use iggy_binary_protocol::{
+    MAX_PARTITIONS_PER_REQUEST, Operation, PrepareHeader, RoutedRequestHeader, 
WireDecode,
+    WireIdentifier, WireOptions,
+};
+use iggy_common::{
+    IggyByteSize, IggyError, MaxTopicSize, TopicCreateOptions, 
UPDATABLE_STREAM_OPTION_KEYS,
+    UPDATABLE_TOPIC_OPTION_KEYS, UPDATABLE_USER_OPTION_KEYS, 
validate_preallocated_topic_bytes,
+    validate_topic_segment_size,
+};
+use journal::superblock::SuperblockStore;
+use journal::{Journal, JournalHandle};
+use metadata::impls::metadata::StreamsFrontend;
+use metadata::stm::stream::Streams;
+use server_common::Message;
+use std::cell::RefCell;
+use std::rc::Rc;
+use tracing::warn;
+
+/// A staged pre-consensus rejection: `stage` labels the chain step for the
+/// deny log line, `error` is the typed code the deny reply carries.
+pub struct RewriteDeny {
+    pub stage: &'static str,
+    pub error: IggyError,
+}
+
+/// The TCP funnel's pre-consensus rewrite chain, in order: the PAT rewrite
+/// (resolving the acting user from the transport `sessions` binding), the
+/// password rewrite, then the static bounds gate. Mirrors [`http_chain`]
+/// plus that bounds step, which the binary wire needs because it has no
+/// `command.validate()` layer. Returns the rewritten request and the raw
+/// PAT token the funnel substitutes into the committed reply; a rejection
+/// names the failing stage for the deny log line.
+pub fn tcp_chain<B, MJ, S, SB>(
+    shard: &Rc<ShellShard<B, MJ, S, SB>>,
+    sessions: &Rc<RefCell<SessionManager>>,
+    transport_client_id: u128,
+    max_tokens_per_user: u32,
+    request: Message<RoutedRequestHeader>,
+) -> Result<(Message<RoutedRequestHeader>, Option<String>), RewriteDeny>
+where
+    B: ShellBus,
+    MJ: JournalHandle + 'static,
+    MJ::Target: Journal<Entry = Message<PrepareHeader>, Header = 
PrepareHeader>,
+    S: 'static,
+    SB: SuperblockStore + 'static,
+{
+    let (request, raw_pat_token) = 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,
+    )
+    // Token cap reached, malformed body, or a lost session binding.
+    .map_err(|error| RewriteDeny {
+        stage: "personal-access-token",
+        error,
+    })?;
+    // 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 =
+        maybe_rewrite_user_password_request(shard, request).map_err(|error| 
RewriteDeny {
+            stage: "user-password",
+            error,
+        })?;
+    static_bounds(shard, &request).map_err(|error| RewriteDeny {
+        stage: "static-bounds",
+        error,
+    })?;
+    Ok((request, raw_pat_token))
+}
+
+/// The HTTP submit's pre-consensus rewrite chain: the PAT rewrite for the
+/// already-authenticated `user_id`, then the password rewrite. Mirrors
+/// [`tcp_chain`] minus the session lookup (the HTTP listener authenticates
+/// against its own session table and resolves the acting user itself) and
+/// minus the static bounds step: HTTP enforces the same bounds in its

Review Comment:
   **nit** — this sentence newly asserts cross-transport bounds parity, and the 
two caps it rests on are defined twice with no cross-check: 
`MAX_PARTITIONS_COUNT = 1000` (`common/src/lib.rs:162`, used only by the HTTP 
DTO validators) and `MAX_PARTITIONS_PER_REQUEST = 1000` 
(`binary_protocol/src/lib.rs:107`, used only on the binary path at `:182`). 
Equivalent today, so no live divergence — but the parity is now documented and 
nothing enforces it. This repo's own idiom for exactly this drift class is a 
`const _: () = assert!(..)` (see `boot/recovery.rs:301-350`).



##########
core/integration/tests/server/raw_tcp.rs:
##########
@@ -0,0 +1,192 @@
+// 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.
+
+//! Raw TCP framing for the server suites that hand-craft client frames the
+//! SDK cannot emit: connect to the harness server, write one request frame,
+//! read the header the server answers with, and register root so a frame
+//! can ride a bound session.
+
+use std::mem::offset_of;
+use std::time::Duration;
+
+use iggy::prelude::*;
+use iggy_binary_protocol::codec::{WireDecode, WireEncode};
+use iggy_binary_protocol::consensus::{
+    Command, Operation, ReplyHeader, RequestHeader, read_size_field, 
result_code,
+    result_section_len,
+};
+use iggy_binary_protocol::requests::users::LoginRegisterRequest;
+use iggy_binary_protocol::responses::users::LoginRegisterResponse;
+use iggy_binary_protocol::{
+    ClientVersionInfo, EvictionHeader, HEADER_SIZE, IGGY_PROTOCOL_VERSION, 
WireName,
+};
+use integration::harness::TestHarness;
+use secrecy::SecretString;
+use tokio::io::{AsyncReadExt, AsyncWriteExt};
+use tokio::net::TcpStream;
+use tokio::time::{Instant, sleep, timeout};
+
+/// Per-frame reply wait. A server that drops the frame answers nothing at
+/// all, so an unanswered read is a verdict, not a reason to wait longer.
+const REPLY_WAIT: Duration = Duration::from_secs(5);
+
+/// Budget for the register to commit: right after boot the single node may
+/// still be electing itself and answers transient rejections meanwhile.
+const COMMIT_BUDGET: Duration = Duration::from_secs(15);
+
+const RETRY_PAUSE: Duration = Duration::from_millis(100);
+
+pub(crate) async fn connect(harness: &TestHarness) -> TcpStream {
+    let addr = harness
+        .server()
+        .tcp_addr()
+        .expect("server must expose a TCP address");
+    TcpStream::connect(addr).await.unwrap()
+}
+
+pub(crate) fn request_header(
+    operation: Operation,
+    client: u128,
+    session: u64,
+    request: u64,
+    body_len: usize,
+) -> RequestHeader {
+    RequestHeader {
+        command: Command::Request,
+        operation,
+        size: u32::try_from(HEADER_SIZE + body_len).unwrap(),
+        client,
+        session,
+        request,
+        ..Default::default()
+    }
+}
+
+/// A header-only `NonReplicated` frame; the command code travels in the
+/// first 4 reserved bytes.
+pub(crate) fn non_replicated_header(
+    client: u128,
+    session: u64,
+    request: u64,
+    code: u32,
+) -> RequestHeader {
+    let mut header = request_header(Operation::NonReplicated, client, session, 
request, 0);
+    header.reserved[..4].copy_from_slice(&code.to_le_bytes());
+    header
+}
+
+pub(crate) async fn write_frame(stream: &mut TcpStream, header: 
&RequestHeader, body: &[u8]) {
+    stream.write_all(bytemuck::bytes_of(header)).await.unwrap();
+    if !body.is_empty() {
+        stream.write_all(body).await.unwrap();
+    }
+}
+
+/// Read the header of the next server frame, within [`REPLY_WAIT`].
+pub(crate) async fn read_frame_header(stream: &mut TcpStream) -> [u8; 
HEADER_SIZE] {
+    let mut header = [0u8; HEADER_SIZE];
+    timeout(REPLY_WAIT, stream.read_exact(&mut header))
+        .await
+        .expect("server must answer within the reply wait, not stall")
+        .expect("reply header read failed");
+    header
+}
+
+/// Write one frame and read one Reply off the lockstep connection, the body
+/// sized by the reply's size field.
+pub(crate) async fn exchange(
+    stream: &mut TcpStream,
+    header: &RequestHeader,
+    body: &[u8],
+) -> ([u8; HEADER_SIZE], Vec<u8>) {
+    write_frame(stream, header, body).await;
+    let reply_header = read_frame_header(stream).await;
+    let command = frame_command(&reply_header);
+    assert_eq!(
+        command,
+        Command::Reply as u8,
+        "expected a Reply frame, got command byte {command} (an Eviction 
carries reason {})",
+        reply_header[offset_of!(EvictionHeader, reason)]
+    );
+
+    let total_size = read_size_field(&reply_header).expect("reply size field") 
as usize;
+    let mut reply_body = vec![0u8; total_size - HEADER_SIZE];

Review Comment:
   **warning** — `total_size - HEADER_SIZE` has no floor, so a reply whose size 
field is under `HEADER_SIZE` fails in a way that hides the regression under 
test. `read_size_field` (`binary_protocol/src/consensus/header.rs:73-78`) is a 
bare 4-byte LE read with no `>= HEADER_SIZE` guard — the SDK adds its own at 
`sdk/src/vsr.rs:181-183`, which is why this helper cannot inherit it. Debug 
builds panic "attempt to subtract with overflow"; release test runs have no 
`overflow-checks` (`[profile.release]` sets only `lto` and `codegen-units`), so 
it wraps to a ~1.8e19 `vec![0u8; n]` and aborts on allocation failure. Assert 
`total_size >= HEADER_SIZE` with a message naming the frame.



##########
core/integration/tests/server/legacy_login_vsr.rs:
##########
@@ -60,39 +59,20 @@ async fn 
given_legacy_pat_login_code_when_sent_raw_should_evict_malformed_login(
 /// eviction. The reject runs before the session gate, so this unbound socket
 /// exercises the same path a bound connection would.
 async fn assert_legacy_login_code_evicted(harness: &TestHarness, code: u32) {
-    let mut header = RequestHeader {
-        command: Command::Request,
-        operation: Operation::NonReplicated,
-        size: u32::try_from(HEADER_SIZE).unwrap(),
-        // NonReplicated leaves session / request unchecked, but the header
-        // validator still requires a nonzero client id.
-        client: 0xC0FFEE,
-        session: 0,
-        request: 0,
-        ..Default::default()
-    };
-    // A non-replicated command code travels in the first 4 reserved bytes.
-    header.reserved[..4].copy_from_slice(&code.to_le_bytes());
+    // NonReplicated leaves session / request unchecked, but the header
+    // validator still requires a nonzero client id.
+    let header = non_replicated_header(0xC0FFEE, 0, 0, code);
 
-    let addr = harness
-        .server()
-        .tcp_addr()
-        .expect("server must expose a TCP address");
-    let mut stream = TcpStream::connect(addr).await.unwrap();
-    stream.write_all(bytemuck::bytes_of(&header)).await.unwrap();
+    let mut stream = connect(harness).await;
+    write_frame(&mut stream, &header, &[]).await;
 
-    // Eviction is header-only: exactly 256 bytes. The timeout makes the
+    // Eviction is header-only: exactly 256 bytes. The bounded read makes the
     // fail-fast contract explicit -- a regression that silently drops the
     // frame trips this instead of hanging until the test wall clock.
-    let mut reply = [0u8; HEADER_SIZE];
-    timeout(Duration::from_secs(5), stream.read_exact(&mut reply))
-        .await
-        .expect("server must answer a legacy login code within 5s, not stall")
-        .expect("reading the eviction frame must succeed");
+    let reply = read_frame_header(&mut stream).await;
 
-    let command_offset = offset_of!(RequestHeader, command);
     assert_eq!(
-        reply[command_offset],
+        frame_command(&reply),

Review Comment:
   **nit** — this file still reads the reason byte as `reply[HEADER_SIZE - 1]` 
against a hardcoded `EVICTION_REASON_MALFORMED_LOGIN: u8 = 15`, while the new 
harness next door reads the same field as `offset_of!(EvictionHeader, reason)` 
(`raw_tcp.rs:123`). Both are correct today — `header.rs:755` statically asserts 
`reason` is the last header byte — but renumber `EvictionReason` and the 
hardcoded `15` silently checks a different reason while still passing. Export 
an `eviction_reason()` accessor from `raw_tcp` and compare against 
`EvictionReason::MalformedLogin as u8`.



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