This is an automated email from the ASF dual-hosted git repository.

hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to refs/heads/master by this push:
     new c5fec1ff6 test(integration): spec clients-table durability across node 
restart (#3736)
c5fec1ff6 is described below

commit c5fec1ff646e7513619c771cea909188526a7141
Author: Hubert Gruszecki <[email protected]>
AuthorDate: Mon Jul 27 09:48:51 2026 +0200

    test(integration): spec clients-table durability across node restart (#3736)
---
 .../tests/cluster/client_table_restart.rs          | 382 +++++++++++++++++++++
 core/integration/tests/cluster/mod.rs              |  22 +-
 core/integration/tests/mod.rs                      |   5 +-
 3 files changed, 385 insertions(+), 24 deletions(-)

diff --git a/core/integration/tests/cluster/client_table_restart.rs 
b/core/integration/tests/cluster/client_table_restart.rs
new file mode 100644
index 000000000..3eb6195ce
--- /dev/null
+++ b/core/integration/tests/cluster/client_table_restart.rs
@@ -0,0 +1,382 @@
+// 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.
+
+//! Spec tests for clients-table durability across a node restart (IGGY-137).
+//!
+//! A client that keeps its `(client, session, request)` identity across a
+//! node crash must be able to continue: a retry of an already-committed
+//! request id must be answered from the dedup cache (never re-applied,
+//! never silently dropped), and the next request id must be admitted.
+//! Today the table lives only in memory, so a rebooted node has no record
+//! of the session or its request watermark and both scenarios fail.
+//!
+//! The Rust SDK cannot drive this: it resets its `ConsensusSession` on every
+//! disconnect and re-registers under a fresh identity. The frames are
+//! therefore hand-crafted on a raw TCP socket, same technique as the
+//! protocol-version gate tests.
+//!
+//! What has to land for these tests to go green, in order:
+//!
+//! 1. Persist the clients table (IGGY-137, standalone): include the
+//!    (client id, last request id, cached reply) entries in the checkpoint
+//!    and recover them on boot from WAL replay, so a rebooted node
+//!    remembers where each client left off.
+//! 2. The client stops forgetting itself on disconnect: keep client id,
+//!    session id and request counter across reconnects and present the old
+//!    identity instead of a fresh Register.
+//! 3. The server accepts a resumed identity: look the session up in the
+//!    replicated table, rebind the new transport to it, and answer with
+//!    the last committed request id so the client knows whether its
+//!    in-doubt request went through. Define the conflict rule for the same
+//!    session arriving on two connections (evict the older).
+//! 4. SDK retry rule change: a replicated write may only be retried under
+//!    the same (client id, request id); the path that re-issues an
+//!    in-doubt write under a fresh session after failover goes away.
+//!
+//! Steps 2+3 must ship together; 1 is standalone. An alternative to 2-4 is
+//! a per-request idempotency key that is independent of the session, the
+//! way TigerBeetle does it: no session resume at all, retries from a fresh
+//! client session stay safe because dedup keys off the request, not the
+//! (client, session) pair.
+//!
+//! These tests pin the implicit-rebind contract: a resumed client simply
+//! keeps sending under its old `(client, session)` on a fresh connection and
+//! the server rebinds the transport from the persisted table. If the
+//! session-resume work settles on an explicit resume handshake instead,
+//! adjust `resume_request` to speak it.
+
+#![cfg(feature = "vsr")]
+
+use bytes::Bytes;
+use iggy::prelude::*;
+use iggy_binary_protocol::codec::{WireDecode, WireEncode};
+use iggy_binary_protocol::consensus::{
+    Command2, Operation, ReplyHeader, RequestHeader, read_size_field, 
result_code,
+    result_section_len,
+};
+use iggy_binary_protocol::namespace::METADATA_CONSENSUS_NAMESPACE;
+use iggy_binary_protocol::requests::streams::CreateStreamRequest;
+use iggy_binary_protocol::requests::users::LoginRegisterRequest;
+use iggy_binary_protocol::responses::users::LoginRegisterResponse;
+use iggy_binary_protocol::{ClientVersionInfo, HEADER_SIZE, 
IGGY_PROTOCOL_VERSION, WireName};
+use integration::harness::TestHarness;
+use integration::iggy_harness;
+use secrecy::SecretString;
+use std::mem::offset_of;
+use std::net::SocketAddr;
+use std::time::Duration;
+use tokio::io::{AsyncReadExt, AsyncWriteExt};
+use tokio::net::TcpStream;
+use tokio::time::{Instant, sleep, timeout};
+
+/// Fixed wire identity so the post-restart frames are byte-identical to the
+/// pre-restart ones; the SDK would randomize this on reconnect.
+const CLIENT_ID: u128 = 0x1337_C0FFEE;
+
+/// Budget for one committed round-trip (covers transient replays while the
+/// single node elects itself after boot).
+const COMMIT_BUDGET: Duration = Duration::from_secs(15);
+
+/// Budget for the post-restart continuation attempts. Longer than
+/// `COMMIT_BUDGET`: it also absorbs the listener coming back up.
+const RESUME_BUDGET: Duration = Duration::from_secs(20);
+
+/// Per-attempt reply wait. A server that silently drops the frame (the
+/// `RequestGap` failure mode) 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);
+
+const RETRY_PAUSE: Duration = Duration::from_millis(100);
+
+#[iggy_harness]
+#[ignore = "red until clients-table persistence + session resume land"]
+async fn given_committed_request_when_node_restarts_should_dedup_same_id_retry(
+    harness: &mut TestHarness,
+) {
+    let addr = tcp_addr(harness);
+    let (mut stream, session) = register(addr).await;
+    let create_stream = create_stream_payload("iggy137-dedup");
+    commit_request(&mut stream, session, 1, &create_stream).await;
+    drop(stream);
+
+    harness.restart_server().await.unwrap();
+
+    // The reply for request 1 was already delivered, but the client cannot
+    // know that in the crash window; retrying the same id must converge on
+    // the cached reply, never on a second apply or a silent drop.
+    let addr = tcp_addr(harness);
+    resume_request(addr, session, 1, &create_stream).await;
+}
+
+#[iggy_harness]
+#[ignore = "red until clients-table persistence + session resume land"]
+async fn given_bound_session_when_node_restarts_should_accept_next_request_id(
+    harness: &mut TestHarness,
+) {
+    let addr = tcp_addr(harness);
+    let (mut stream, session) = register(addr).await;
+    commit_request(
+        &mut stream,
+        session,
+        1,
+        &create_stream_payload("iggy137-first"),
+    )
+    .await;
+    drop(stream);
+
+    harness.restart_server().await.unwrap();
+
+    // Continuation, not retry: the session advances to the next id. A node
+    // that forgot the watermark sees request 2 on an unknown session and
+    // either drops it as a gap or bounces the session entirely.
+    let addr = tcp_addr(harness);
+    resume_request(addr, session, 2, 
&create_stream_payload("iggy137-second")).await;
+}
+
+fn tcp_addr(harness: &TestHarness) -> SocketAddr {
+    harness
+        .server()
+        .tcp_addr()
+        .expect("server must expose a TCP address")
+}
+
+fn create_stream_payload(name: &str) -> Bytes {
+    CreateStreamRequest {
+        name: WireName::new(name).unwrap(),
+    }
+    .to_bytes()
+}
+
+fn request_header(
+    operation: Operation,
+    session: u64,
+    request: u64,
+    body_len: usize,
+) -> RequestHeader {
+    RequestHeader {
+        command: Command2::Request,
+        operation,
+        size: u32::try_from(HEADER_SIZE + body_len).unwrap(),
+        client: CLIENT_ID,
+        session,
+        request,
+        namespace: match operation {
+            Operation::Register => METADATA_CONSENSUS_NAMESPACE,
+            _ => 0,
+        },
+        ..Default::default()
+    }
+}
+
+/// Register `CLIENT_ID` as root and return the connection with its bound
+/// session id. The session binds to THIS transport connection server-side,
+/// so the pre-restart request must reuse the returned stream. Replays on
+/// transient rejections: right after boot the single node may not have
+/// elected itself yet.
+async fn register(addr: SocketAddr) -> (TcpStream, u64) {
+    let body = LoginRegisterRequest {
+        version_info: ClientVersionInfo {
+            protocol_version: IGGY_PROTOCOL_VERSION,
+            sdk_name: WireName::new("iggy137-raw").unwrap(),
+            sdk_version: WireName::new("0.0.1").unwrap(),
+        },
+        username: WireName::new(DEFAULT_ROOT_USERNAME).unwrap(),
+        password: SecretString::from(DEFAULT_ROOT_PASSWORD),
+        client_context: None,
+    }
+    .to_bytes();
+    let header = request_header(Operation::Register, 0, 0, body.len());
+
+    let mut stream = TcpStream::connect(addr).await.unwrap();
+    let deadline = Instant::now() + COMMIT_BUDGET;
+    loop {
+        match exchange(&mut stream, &header, &body).await.verdict() {
+            Verdict::Success(payload) => {
+                let response = LoginRegisterResponse::decode_from(&payload)
+                    .expect("register payload must decode");
+                assert_ne!(response.session, 0, "server must bind a nonzero 
session");
+                return (stream, response.session);
+            }
+            Verdict::Rejected(code) if is_transient(code) && Instant::now() < 
deadline => {
+                sleep(RETRY_PAUSE).await;
+            }
+            other => panic!("register did not commit: {other:?}"),
+        }
+    }
+}
+
+/// Send one replicated metadata request on the registered connection and
+/// require a committed success within `COMMIT_BUDGET`.
+async fn commit_request(stream: &mut TcpStream, session: u64, request: u64, 
body: &Bytes) {
+    let header = request_header(Operation::CreateStream, session, request, 
body.len());
+    let deadline = Instant::now() + COMMIT_BUDGET;
+    loop {
+        match exchange(stream, &header, body).await.verdict() {
+            Verdict::Success(_) => return,
+            Verdict::Rejected(code) if is_transient(code) && Instant::now() < 
deadline => {
+                sleep(RETRY_PAUSE).await;
+            }
+            other => panic!("request {request} did not commit: {other:?}"),
+        }
+    }
+}
+
+/// Post-restart continuation: keep presenting the old identity until the
+/// server commits (or serves the cached reply for) the request. Every
+/// attempt uses a fresh connection, both because the old one died with the
+/// node and so an unanswered frame cannot desync the next attempt. Panics
+/// with the last observed failure mode when the budget runs out.
+async fn resume_request(addr: SocketAddr, session: u64, request: u64, body: 
&Bytes) {
+    let header = request_header(Operation::CreateStream, session, request, 
body.len());
+    let deadline = Instant::now() + RESUME_BUDGET;
+    let mut last_failure = "the listener never came back".to_string();
+    while Instant::now() < deadline {
+        let Ok(mut stream) = TcpStream::connect(addr).await else {
+            sleep(RETRY_PAUSE).await;
+            continue;
+        };
+        match exchange(&mut stream, &header, body).await.verdict() {
+            Verdict::Success(_) => return,
+            Verdict::Rejected(code) if is_transient(code) => {
+                last_failure = format!("still transient (code {code})");
+            }
+            Verdict::Rejected(code) => {
+                last_failure = format!(
+                    "request {request} answered with committed code {code} (a \
+                     duplicate-apply rejection means the dedup cache was lost)"
+                );
+            }
+            Verdict::NoResultSection => {
+                last_failure = format!(
+                    "request {request} got the unbound-transport empty Reply: 
the \
+                     restarted node does not recognize session {session}"
+                );
+            }
+            Verdict::Ignored => {
+                last_failure = format!(
+                    "request {request} on session {session} was silently 
ignored for \
+                     {REPLY_WAIT:?} (RequestGap-style drop: the restarted node 
lost \
+                     the request watermark)"
+                );
+            }
+            Verdict::Evicted(reason) => {
+                last_failure = format!(
+                    "session {session} was evicted with reason {reason} 
instead of \
+                     being rebound from the persisted table"
+                );
+            }
+        }
+        sleep(RETRY_PAUSE).await;
+    }
+    panic!(
+        "session {session} did not survive the restart within 
{RESUME_BUDGET:?}: {last_failure}"
+    );
+}
+
+/// Everything one request/reply exchange can end in, spelled out so the
+/// red-test panic names the exact failure mode instead of a decode error.
+#[derive(Debug)]
+enum Verdict {
+    /// Committed success; carries the payload after the result section.
+    Success(Bytes),
+    /// Committed (or pre-consensus transient) rejection code.
+    Rejected(u32),
+    /// A Reply with no result section, i.e. the empty Reply the server emits
+    /// for a replicated request on a transport it has no session for.
+    NoResultSection,
+    /// No frame within `REPLY_WAIT`.
+    Ignored,
+    /// Session-terminal Eviction frame; carries the wire reason byte.
+    Evicted(u8),
+}
+
+enum Exchange {
+    Reply { status: u32, body: Bytes },
+    Eviction { reason: u8 },
+    Ignored,
+}
+
+impl Exchange {
+    fn verdict(self) -> Verdict {
+        match self {
+            Self::Ignored => Verdict::Ignored,
+            Self::Eviction { reason } => Verdict::Evicted(reason),
+            // A nonzero status is the pre-commit deny channel (authz etc.);
+            // fold it into the rejection space, the codes are shared.
+            Self::Reply { status, .. } if status != 0 => 
Verdict::Rejected(status),
+            Self::Reply { body, .. } => match result_code(&body) {
+                None => Verdict::NoResultSection,
+                Some(0) => {
+                    let payload_start = result_section_len(&body).unwrap();
+                    Verdict::Success(body.slice(payload_start..))
+                }
+                Some(code) => Verdict::Rejected(code),
+            },
+        }
+    }
+}
+
+/// Write one frame and read one frame off the lockstep connection.
+async fn exchange(stream: &mut TcpStream, header: &RequestHeader, body: 
&Bytes) -> Exchange {
+    stream.write_all(bytemuck::bytes_of(header)).await.unwrap();
+    if !body.is_empty() {
+        stream.write_all(body).await.unwrap();
+    }
+
+    let mut reply_header = [0u8; HEADER_SIZE];
+    match timeout(REPLY_WAIT, stream.read_exact(&mut reply_header)).await {
+        Err(_elapsed) => return Exchange::Ignored,
+        Ok(read) => {
+            read.expect("reply header read failed");
+        }
+    }
+
+    let command_offset = offset_of!(RequestHeader, command);
+    if reply_header[command_offset] == Command2::Eviction as u8 {
+        return Exchange::Eviction {
+            reason: reply_header[HEADER_SIZE - 1],
+        };
+    }
+    assert_eq!(
+        reply_header[command_offset],
+        Command2::Reply as u8,
+        "expected a Reply frame"
+    );
+
+    let status_offset = offset_of!(ReplyHeader, status);
+    let status = u32::from_le_bytes(
+        reply_header[status_offset..status_offset + 4]
+            .try_into()
+            .unwrap(),
+    );
+
+    let total_size = read_size_field(&reply_header).expect("reply size field") 
as usize;
+    let mut body = vec![0u8; total_size - HEADER_SIZE];
+    timeout(REPLY_WAIT, stream.read_exact(&mut body))
+        .await
+        .expect("reply body timed out")
+        .expect("reply body read failed");
+    Exchange::Reply {
+        status,
+        body: body.into(),
+    }
+}
+
+fn is_transient(code: u32) -> bool {
+    code == IggyError::TransientNotCommitted.as_code()
+        || code == IggyError::TransientNotAccepted.as_code()
+}
diff --git a/core/integration/tests/cluster/mod.rs 
b/core/integration/tests/cluster/mod.rs
index 0351adaa1..93c631bd0 100644
--- a/core/integration/tests/cluster/mod.rs
+++ b/core/integration/tests/cluster/mod.rs
@@ -15,24 +15,4 @@
 // specific language governing permissions and limitations
 // under the License.
 
-use iggy::prelude::*;
-use integration::iggy_harness;
-
-#[iggy_harness(cluster_nodes = [3, 5, 7],
-    test_client_transport = [ Tcp, Http, WebSocket, Quic, TcpTlsSelfSigned, 
TcpTlsGenerated, WebSocketTlsSelfSigned, WebSocketTlsGenerated],
-    server(segment.size = ["1MiB", "2MiB"],
-           segment.cache_indexes = ["open_segment", "all"]))]
-#[ignore]
-async fn should_ping_all_cluster_nodes(harness: TestHarness) {
-    for i in 0..harness.cluster_size() {
-        let client = harness
-            .node(i)
-            .test_client()
-            .unwrap()
-            .with_root_login()
-            .connect()
-            .await
-            .unwrap();
-        client.ping().await.unwrap();
-    }
-}
+mod client_table_restart;
diff --git a/core/integration/tests/mod.rs b/core/integration/tests/mod.rs
index dd3cf48f6..efb0cc8ca 100644
--- a/core/integration/tests/mod.rs
+++ b/core/integration/tests/mod.rs
@@ -36,9 +36,8 @@ use tracing_subscriber::{EnvFilter, fmt};
 // design: flush returns FeatureUnavailable, the session-timeout message
 // differs, and purge is eventually consistent so server state is polled.
 mod cli;
-// A single `#[ignore]`d multi-node ping matrix stub; none of its cells run in
-// either mode today.
-#[cfg(not(feature = "vsr"))]
+// Raw-wire spec tests for VSR session continuity across a node restart
+// (IGGY-137); the module is vsr-only by construction (file-level cfg).
 mod cluster;
 mod config_provider;
 mod connectors;

Reply via email to