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

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


The following commit(s) were added to refs/heads/recover_writes by this push:
     new 266602968 fix
266602968 is described below

commit 2666029680cfef9b74b40e022383fb31771217b1
Author: spetz <[email protected]>
AuthorDate: Fri Aug 28 16:03:35 2026 +0200

    fix
---
 .../src/harness/handle/client_builder.rs           |  61 +++++-
 .../tests/cluster/fast_primary_rejoin.rs           | 146 +++++++++++--
 core/partitions/src/iggy_partition.rs              | 148 ++++++++++---
 core/partitions/src/journal.rs                     |  14 ++
 core/partitions/src/types.rs                       |  16 +-
 core/sdk/src/leader_aware.rs                       | 228 ++++++++++++++++++++-
 core/sdk/src/quic/quic_client.rs                   | 198 ++++++++++++++----
 core/sdk/src/tcp/tcp_client.rs                     |  96 +++++++--
 core/sdk/src/websocket/websocket_client.rs         | 221 ++++++++++++++++----
 core/server/src/partition_reconciler.rs            |   4 +-
 core/shard/src/lib.rs                              | 117 +++++++++--
 .../AsyncIggyTcpClientTransientFailoverTest.java   |  43 ++--
 12 files changed, 1077 insertions(+), 215 deletions(-)

diff --git a/core/integration/src/harness/handle/client_builder.rs 
b/core/integration/src/harness/handle/client_builder.rs
index 040056c7e..4dd449a04 100644
--- a/core/integration/src/harness/handle/client_builder.rs
+++ b/core/integration/src/harness/handle/client_builder.rs
@@ -38,12 +38,12 @@ use crate::harness::config::{AutoLoginConfig, TlsConfig};
 use crate::harness::error::TestBinaryError;
 use iggy::http::http_client::HttpClient;
 use iggy::prelude::{
-    Client, HttpClientConfig, IggyClient, QuicClientConfig, TcpClient, 
TcpClientConfig, UserClient,
-    WebSocketClientConfig,
+    Client, HttpClientConfig, IggyClient, IggyDuration, QuicClientConfig, 
TcpClient,
+    TcpClientConfig, UserClient, WebSocketClientConfig,
 };
 use iggy::quic::quic_client::QuicClient;
 use iggy::websocket::websocket_client::WebSocketClient;
-use iggy_common::TransportProtocol;
+use iggy_common::{AutoLogin, Credentials, TransportProtocol};
 use std::net::SocketAddr;
 use std::path::PathBuf;
 use std::sync::Arc;
@@ -69,7 +69,9 @@ pub struct ClientBuilder {
     transport: TransportProtocol,
     connection: ServerConnection,
     auto_login: Option<AutoLoginConfig>,
+    reconnecting_login: bool,
     tcp_nodelay: bool,
+    reestablish_after: Option<IggyDuration>,
     encryptor: Option<Arc<iggy_common::EncryptorKind>>,
 }
 
@@ -79,7 +81,9 @@ impl ClientBuilder {
             transport,
             connection,
             auto_login: None,
+            reconnecting_login: false,
             tcp_nodelay: false,
+            reestablish_after: None,
             encryptor: None,
         }
     }
@@ -90,6 +94,14 @@ impl ClientBuilder {
         self
     }
 
+    /// Configure the binary transport itself to restore the root session
+    /// after reconnecting instead of running a one-time harness login.
+    pub fn with_reconnecting_root_login(mut self) -> Self {
+        self.auto_login = Some(AutoLoginConfig::root());
+        self.reconnecting_login = true;
+        self
+    }
+
     /// Enable automatic login with custom credentials after connection.
     pub fn with_login(mut self, username: impl Into<String>, password: impl 
Into<String>) -> Self {
         self.auto_login = Some(AutoLoginConfig::new(username, password));
@@ -102,6 +114,13 @@ impl ClientBuilder {
         self
     }
 
+    /// Override how long a binary transport prefers its previous endpoint
+    /// before rotating through the cluster roster after a disconnect.
+    pub fn with_reestablish_after(mut self, reestablish_after: IggyDuration) 
-> Self {
+        self.reestablish_after = Some(reestablish_after);
+        self
+    }
+
     /// Set the client-side encryptor for encrypting/decrypting message 
payloads and headers.
     pub fn with_encryptor(mut self, encryptor: 
Arc<iggy_common::EncryptorKind>) -> Self {
         self.encryptor = Some(encryptor);
@@ -117,7 +136,9 @@ impl ClientBuilder {
             TransportProtocol::WebSocket => 
self.create_websocket_client().await?,
         };
 
-        if let Some(ref login) = self.auto_login {
+        if let Some(ref login) = self.auto_login
+            && (!self.reconnecting_login || self.transport == 
TransportProtocol::Http)
+        {
             client
                 .login_user(&login.username, &login.password)
                 .await
@@ -142,7 +163,7 @@ impl ClientBuilder {
         let tls_enabled = self.connection.tls.is_some();
         let tls_validate = self.connection.tls.as_ref().is_some_and(|t| 
!t.self_signed);
 
-        let config = TcpClientConfig {
+        let mut config = TcpClientConfig {
             server_address: addr.to_string(),
             nodelay: self.tcp_nodelay,
             tls_enabled,
@@ -153,8 +174,12 @@ impl ClientBuilder {
                 .as_ref()
                 .map(|p| p.to_string_lossy().to_string()),
             tls_validate_certificate: tls_validate,
+            auto_login: self.binary_auto_login(),
             ..TcpClientConfig::default()
         };
+        if let Some(reestablish_after) = self.reestablish_after {
+            config.reconnection.reestablish_after = reestablish_after;
+        }
 
         let client =
             TcpClient::create(Arc::new(config)).map_err(|e| 
TestBinaryError::ClientCreation {
@@ -213,11 +238,15 @@ impl ClientBuilder {
                 message: "QUIC transport not available".to_string(),
             })?;
 
-        let config = QuicClientConfig {
+        let mut config = QuicClientConfig {
             server_address: addr.to_string(),
             max_idle_timeout: 2_000_000,
+            auto_login: self.binary_auto_login(),
             ..QuicClientConfig::default()
         };
+        if let Some(reestablish_after) = self.reestablish_after {
+            config.reconnection.reestablish_after = reestablish_after;
+        }
 
         let client =
             QuicClient::create(Arc::new(config)).map_err(|e| 
TestBinaryError::ClientCreation {
@@ -256,7 +285,7 @@ impl ClientBuilder {
             .as_ref()
             .is_some_and(|t| !t.self_signed);
 
-        let config = WebSocketClientConfig {
+        let mut config = WebSocketClientConfig {
             server_address: addr.to_string(),
             tls_enabled,
             tls_domain: "localhost".to_string(),
@@ -266,8 +295,12 @@ impl ClientBuilder {
                 .as_ref()
                 .map(|p| p.to_string_lossy().to_string()),
             tls_validate_certificate: tls_validate,
+            auto_login: self.binary_auto_login(),
             ..WebSocketClientConfig::default()
         };
+        if let Some(reestablish_after) = self.reestablish_after {
+            config.reconnection.reestablish_after = reestablish_after;
+        }
 
         let client = WebSocketClient::create(Arc::new(config)).map_err(|e| {
             TestBinaryError::ClientCreation {
@@ -292,6 +325,20 @@ impl ClientBuilder {
         ))
     }
 
+    fn binary_auto_login(&self) -> AutoLogin {
+        if !self.reconnecting_login {
+            return AutoLogin::Disabled;
+        }
+        self.auto_login
+            .as_ref()
+            .map_or(AutoLogin::Disabled, |login| {
+                AutoLogin::Enabled(Credentials::UsernamePassword(
+                    login.username.clone(),
+                    login.password.clone().into(),
+                ))
+            })
+    }
+
     fn get_address_string(&self) -> String {
         match self.transport {
             TransportProtocol::Tcp => self
diff --git a/core/integration/tests/cluster/fast_primary_rejoin.rs 
b/core/integration/tests/cluster/fast_primary_rejoin.rs
index 343476740..435c4ce5e 100644
--- a/core/integration/tests/cluster/fast_primary_rejoin.rs
+++ b/core/integration/tests/cluster/fast_primary_rejoin.rs
@@ -94,17 +94,34 @@ async fn create_stream_and_topic(harness: &TestHarness) {
 
 /// Connect one client pinned to the current primary's own endpoint, the way a
 /// leader-aware SDK ends up connected to whichever node answers as leader.
-async fn pinned_producer(harness: &TestHarness, leader: usize) -> (IggyClient, 
String) {
-    let primary_endpoint = harness
-        .node(leader)
-        .tcp_addr()
-        .expect("leader exposes a TCP endpoint")
-        .to_string();
-    let producer = harness
-        .node(leader)
-        .tcp_client()
-        .expect("leader exposes a TCP endpoint")
-        .with_root_login()
+async fn pinned_producer(
+    harness: &TestHarness,
+    leader: usize,
+    transport: TransportProtocol,
+    reestablish_after: Option<IggyDuration>,
+) -> (IggyClient, String) {
+    let node = harness.node(leader);
+    let primary_endpoint = match transport {
+        TransportProtocol::Tcp => node.tcp_addr(),
+        TransportProtocol::Quic => node.quic_addr(),
+        TransportProtocol::WebSocket => node.websocket_addr(),
+        TransportProtocol::Http => panic!("HTTP does not expose a persistent 
Iggy client"),
+    }
+    .unwrap_or_else(|| panic!("leader exposes a {transport} endpoint"))
+    .to_string();
+    let builder = match transport {
+        TransportProtocol::Tcp => node.tcp_client(),
+        TransportProtocol::Quic => node.quic_client(),
+        TransportProtocol::WebSocket => node.websocket_client(),
+        TransportProtocol::Http => panic!("HTTP does not expose a persistent 
Iggy client"),
+    }
+    .unwrap_or_else(|error| panic!("leader exposes a {transport} client: 
{error}"));
+    let builder = match reestablish_after {
+        Some(duration) => builder.with_reestablish_after(duration),
+        None => builder,
+    };
+    let producer = builder
+        .with_reconnecting_root_login()
         .connect()
         .await
         .expect("connect the producer to the primary");
@@ -223,10 +240,13 @@ async fn require_readback(producer: &IggyClient, 
expected: &[(u64, String)]) {
 /// primary it is pinned to.
 async fn warmed_producer_past_follower_restart(
     harness: &mut TestHarness,
+    transport: TransportProtocol,
+    reestablish_after: Option<IggyDuration>,
 ) -> (IggyClient, usize, String, Vec<(u64, String)>) {
     create_stream_and_topic(harness).await;
     let leader = disk::leader_node_index(harness).await;
-    let (producer, primary_endpoint) = pinned_producer(harness, leader).await;
+    let (producer, primary_endpoint) =
+        pinned_producer(harness, leader, transport, reestablish_after).await;
     let mut acked = require_acked_sends(&producer, "warmup", WARM_ACKS, 
WARMUP_TIMEOUT).await;
 
     let follower = (0..harness.cluster_size())
@@ -240,14 +260,12 @@ async fn warmed_producer_past_follower_restart(
     (producer, leader, primary_endpoint, acked)
 }
 
-/// Baseline sibling: the primary stops gracefully and stays down. The same
-/// client must settle on a survivor, write, and read back.
-#[iggy_harness(cluster_nodes = 3)]
-async fn 
given_a_pinned_producer_when_its_primary_stops_and_stays_down_should_resume(
+async fn require_resume_after_primary_stop(
     harness: &mut TestHarness,
+    transport: TransportProtocol,
 ) {
     let (producer, leader, primary_endpoint, mut acked) =
-        warmed_producer_past_follower_restart(harness).await;
+        warmed_producer_past_follower_restart(harness, transport, None).await;
 
     harness.stop_node(leader).expect("stop the primary");
 
@@ -260,6 +278,48 @@ async fn 
given_a_pinned_producer_when_its_primary_stops_and_stays_down_should_re
     require_readback(&producer, &acked).await;
 }
 
+async fn require_resume_after_fast_primary_rejoin(
+    harness: &mut TestHarness,
+    transport: TransportProtocol,
+) {
+    let (producer, leader, _primary_endpoint, mut acked) =
+        warmed_producer_past_follower_restart(harness, transport, None).await;
+
+    harness
+        .restart_node(leader)
+        .expect("restart the primary with its data intact");
+
+    acked.extend(require_acked_sends(&producer, "post-primary-restart", 1, 
RESUME_BUDGET).await);
+    require_readback(&producer, &acked).await;
+}
+
+async fn require_resume_after_fast_primary_rejoin_with_dead_roster_hop(
+    harness: &mut TestHarness,
+    transport: TransportProtocol,
+) {
+    let (producer, leader, _primary_endpoint, mut acked) =
+        warmed_producer_past_follower_restart(harness, transport, None).await;
+    let dead_hop = (leader + 1) % harness.cluster_size();
+    harness
+        .stop_node(dead_hop)
+        .expect("stop the first roster hop after the original primary");
+    harness
+        .restart_node(leader)
+        .expect("restart the primary with its data intact");
+
+    acked.extend(require_acked_sends(&producer, "post-dead-roster-hop", 1, 
RESUME_BUDGET).await);
+    require_readback(&producer, &acked).await;
+}
+
+/// Baseline sibling: the primary stops gracefully and stays down. The same
+/// client must settle on a survivor, write, and read back.
+#[iggy_harness(cluster_nodes = 3)]
+async fn 
given_a_pinned_producer_when_its_primary_stops_and_stays_down_should_resume(
+    harness: &mut TestHarness,
+) {
+    require_resume_after_primary_stop(harness, TransportProtocol::Tcp).await;
+}
+
 /// Fast-rejoin sibling: the primary stops gracefully and is started again
 /// immediately, so its endpoint answers TCP and metadata as a rejoining
 /// follower while the group primaryship settles elsewhere. The same client
@@ -268,17 +328,61 @@ async fn 
given_a_pinned_producer_when_its_primary_stops_and_stays_down_should_re
 async fn 
given_a_pinned_producer_when_its_primary_restarts_quickly_should_resume(
     harness: &mut TestHarness,
 ) {
-    let (producer, leader, _primary_endpoint, mut acked) =
-        warmed_producer_past_follower_restart(harness).await;
+    require_resume_after_fast_primary_rejoin(harness, 
TransportProtocol::Tcp).await;
+}
+
+#[iggy_harness(cluster_nodes = 3)]
+async fn 
given_a_zero_cooldown_tcp_producer_when_replay_lands_on_a_partition_backup_should_resume(
+    harness: &mut TestHarness,
+) {
+    let (producer, leader, _primary_endpoint, mut acked) = 
warmed_producer_past_follower_restart(
+        harness,
+        TransportProtocol::Tcp,
+        Some(IggyDuration::from(0u64)),
+    )
+    .await;
 
     harness
         .restart_node(leader)
         .expect("restart the primary with its data intact");
 
-    acked.extend(require_acked_sends(&producer, "post-primary-restart", 1, 
RESUME_BUDGET).await);
+    acked.extend(require_acked_sends(&producer, "zero-cooldown-replay", 1, 
RESUME_BUDGET).await);
     require_readback(&producer, &acked).await;
 }
 
+#[iggy_harness(cluster_nodes = 3)]
+async fn given_a_quic_producer_when_its_primary_restarts_quickly_should_resume(
+    harness: &mut TestHarness,
+) {
+    require_resume_after_fast_primary_rejoin(harness, 
TransportProtocol::Quic).await;
+}
+
+#[iggy_harness(cluster_nodes = 3)]
+async fn 
given_a_quic_producer_when_its_first_roster_hop_is_down_should_reach_the_partition_primary(
+    harness: &mut TestHarness,
+) {
+    require_resume_after_fast_primary_rejoin_with_dead_roster_hop(harness, 
TransportProtocol::Quic)
+        .await;
+}
+
+#[iggy_harness(cluster_nodes = 3)]
+async fn 
given_a_websocket_producer_when_its_primary_restarts_quickly_should_resume(
+    harness: &mut TestHarness,
+) {
+    require_resume_after_fast_primary_rejoin(harness, 
TransportProtocol::WebSocket).await;
+}
+
+#[iggy_harness(cluster_nodes = 3)]
+async fn 
given_a_websocket_producer_when_its_first_roster_hop_is_down_should_reach_the_partition_primary(
+    harness: &mut TestHarness,
+) {
+    require_resume_after_fast_primary_rejoin_with_dead_roster_hop(
+        harness,
+        TransportProtocol::WebSocket,
+    )
+    .await;
+}
+
 /// The stateless HTTP transport has no client-side partition routing. After
 /// the old primary rejoins as a backup, its listener must walk the bounded
 /// server roster for acknowledged partition writes.
@@ -293,7 +397,7 @@ async fn 
given_http_writes_on_a_rejoined_backup_when_the_primary_moved_should_fo
     harness: &mut TestHarness,
 ) {
     let (producer, leader, primary_endpoint, mut acked) =
-        warmed_producer_past_follower_restart(harness).await;
+        warmed_producer_past_follower_restart(harness, TransportProtocol::Tcp, 
None).await;
     harness
         .restart_node(leader)
         .expect("restart the primary with its data intact");
diff --git a/core/partitions/src/iggy_partition.rs 
b/core/partitions/src/iggy_partition.rs
index 85c4a0906..202d0fb79 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -4463,10 +4463,29 @@ where
     /// commit walk runs at `RepairDone`, after the floor is known.
     pub async fn apply_repaired_prepare(&mut self, message: 
Message<PrepareHeader>) {
         let header = *message.header();
-        let Some(session) = &self.repair else {
+        let Some(session) = self.repair else {
             return;
         };
-        if header.op <= self.consensus().commit_min() || header.op > 
session.to_op {
+        let consensus = self.consensus();
+        if !consensus.is_normal() || consensus.view() != session.view {
+            self.repair = None;
+            return;
+        }
+        if header.op <= consensus.commit_min() || header.op > 
session.fetch_to_op {
+            return;
+        }
+        let canonical_checksum = consensus
+            .with_pending_view_log(|pending| {
+                pending
+                    .headers
+                    .iter()
+                    .find(|expected| expected.op == header.op)
+                    .map(|expected| expected.checksum)
+            })
+            .flatten();
+        if canonical_checksum.is_some_and(|expected| expected != 
header.checksum)
+            || (header.op > session.commit_to_op && 
canonical_checksum.is_none())
+        {
             return;
         }
         // Any in-window frame proves the stream is alive; only silence
@@ -4480,7 +4499,9 @@ where
         let applied = if header.operation == Operation::SendMessages {
             match self.append_repaired_send_messages(message).await {
                 Ok(base_offset) => {
-                    if let (Some(base_offset), Some(session)) = (base_offset, 
self.repair.as_mut())
+                    if header.op <= session.commit_to_op
+                        && let (Some(base_offset), Some(session)) =
+                            (base_offset, self.repair.as_mut())
                     {
                         session.first_batch_offset = Some(
                             session
@@ -4535,6 +4556,10 @@ where
         let Some(session) = self.repair else {
             return RepairConclusion::Done;
         };
+        if !self.consensus().is_normal() || self.consensus().view() != 
session.view {
+            self.repair = None;
+            return RepairConclusion::Done;
+        }
         if let Some(floor) = session.floor {
             // A peer may have evicted past this replica's commit frontier;
             // an unclamped floor would drive commit_min above commit_max and
@@ -4555,6 +4580,11 @@ where
             let stand_in = durable_end
                 .map(|durable| durable.saturating_add(1))
                 .max(self.installed_frontier);
+            let committed_shape = self
+                .log
+                .journal()
+                .inner
+                .repaired_window_shape(floor, session.commit_to_op);
             let connected = match (session.first_batch_offset, stand_in) {
                 (Some(first), Some(bound)) => first <= bound,
                 (Some(first), None) => first == 0,
@@ -4566,7 +4596,11 @@ where
                 // a fully evicted window -- is indistinguishable from a
                 // message range below the floor that this replica does not
                 // durably own, and accepting it would serve a holed log.
-                (None, _) => self.repaired_window_is_offsets_only(floor, 
session.to_op),
+                (None, _) => {
+                    floor < session.commit_to_op
+                        && committed_shape.complete
+                        && !committed_shape.holds_messages
+                }
             };
             if !connected {
                 tracing::error!(
@@ -4591,11 +4625,11 @@ where
                 // Both are the state-transfer trigger; the session is
                 // dropped here so the caller's arming funnel starts clean,
                 // and a transfer-unavailable fallback re-arms repair fresh.
-                if self.repaired_window_is_complete(floor, session.to_op) {
+                if committed_shape.complete {
                     self.repair = None;
                     return RepairConclusion::FloorRefused {
                         floor,
-                        to_op: session.to_op,
+                        to_op: session.commit_to_op,
                     };
                 }
                 return RepairConclusion::InProgress;
@@ -4614,7 +4648,14 @@ where
         // that reached the requested frontier closes the session; anything
         // less keeps it armed and the stall retry re-requests the remains
         // (`commit_min + 1..`), converging over rounds.
-        let done = commit_min >= session.to_op;
+        let fetch_complete = session.fetch_to_op <= session.commit_to_op
+            || self
+                .log
+                .journal()
+                .inner
+                .repaired_window_shape(session.commit_to_op, 
session.fetch_to_op)
+                .complete;
+        let done = commit_min >= session.commit_to_op && fetch_complete;
         if done {
             self.repair = None;
         }
@@ -4625,7 +4666,9 @@ where
             commit_min_before = before,
             commit_min_after = commit_min,
             commit_max = self.consensus().commit_max(),
-            to_op = session.to_op,
+            commit_to_op = session.commit_to_op,
+            fetch_to_op = session.fetch_to_op,
+            fetch_complete,
             done,
             "repair window commit walk finished"
         );
@@ -4636,31 +4679,6 @@ where
         }
     }
 
-    /// Whether every op in `(floor, to_op]` is journaled. An empty window
-    /// (`floor >= to_op`) counts as complete: there is nothing left that
-    /// could arrive and change the floor verdict.
-    fn repaired_window_is_complete(&self, floor: u64, to_op: u64) -> bool {
-        self.log
-            .journal()
-            .inner
-            .repaired_window_shape(floor, to_op)
-            .complete
-    }
-
-    /// Whether the served repair window `(floor, to_op]` arrived complete and
-    /// holds no `SendMessages` op. Only then may a commit floor be accepted
-    /// without a batch anchor: the window demonstrably moved no messages, so
-    /// the consumer-offset table on disk stands in below the floor. An empty
-    /// window (`floor >= to_op`) carries no evidence at all and never
-    /// qualifies.
-    fn repaired_window_is_offsets_only(&self, floor: u64, to_op: u64) -> bool {
-        if floor >= to_op {
-            return false;
-        }
-        let shape = self.log.journal().inner.repaired_window_shape(floor, 
to_op);
-        shape.complete && !shape.holds_messages
-    }
-
     /// Journal a repaired `SendMessages` prepare, preserving its embedded
     /// batch stamps. A stored prepare was stamped by `append_messages` on
     /// the serving replica BEFORE it was journaled, so its `base_offset` /
@@ -6630,9 +6648,20 @@ mod tests {
     }
 
     fn armed_session(to_op: u64, floor: u64, first_batch_offset: Option<u64>) 
-> RepairSession {
+        armed_fetch_session(to_op, to_op, floor, first_batch_offset)
+    }
+
+    fn armed_fetch_session(
+        to_op: u64,
+        fetch_to_op: u64,
+        floor: u64,
+        first_batch_offset: Option<u64>,
+    ) -> RepairSession {
         RepairSession {
             nonce: 1,
-            to_op,
+            view: 0,
+            commit_to_op: to_op,
+            fetch_to_op,
             floor: Some(floor),
             peer: 0,
             first_batch_offset,
@@ -6759,6 +6788,26 @@ mod tests {
         );
     }
 
+    #[compio::test]
+    async fn 
given_prior_view_repair_when_a_new_view_started_should_discard_it() {
+        let mut partition = test_partition();
+        partition.repair = Some(armed_fetch_session(0, 1, 0, None));
+        partition.consensus.set_view(1);
+
+        partition
+            .apply_repaired_prepare(repaired_send_prepare(1, 0, 0x11))
+            .await;
+
+        assert!(
+            partition.repair.is_none(),
+            "the prior-view session is obsolete"
+        );
+        assert!(
+            partition.log.journal().inner.header_by_op(1).is_none(),
+            "a delayed prior-view body must not enter the new view's journal"
+        );
+    }
+
     #[compio::test]
     async fn 
given_session_remint_when_attempts_burned_should_survive_on_partition() {
         let mut partition = test_partition();
@@ -6880,6 +6929,37 @@ mod tests {
         assert!(partition.repair.is_none());
     }
 
+    #[compio::test]
+    async fn 
given_empty_committed_window_with_a_suffix_fetch_should_escape_to_state_transfer()
 {
+        let mut partition = test_partition();
+        partition.consensus().advance_commit_max(5);
+        partition.repair = Some(armed_fetch_session(5, 9, 5, None));
+
+        let conclusion = partition.complete_repair(&repair_config()).await;
+
+        assert_eq!(
+            conclusion,
+            RepairConclusion::FloorRefused { floor: 5, to_op: 5 },
+            "the uncommitted fetch ceiling must not postpone a definitive 
committed-floor refusal"
+        );
+        assert!(partition.repair.is_none());
+    }
+
+    #[compio::test]
+    async fn 
given_suffix_fetch_when_its_view_is_discarded_should_clear_the_session() {
+        let mut partition = test_partition();
+        partition.repair = Some(armed_fetch_session(0, 3, 0, None));
+        partition.consensus.set_view(1);
+
+        let conclusion = partition.complete_repair(&repair_config()).await;
+
+        assert_eq!(conclusion, RepairConclusion::Done);
+        assert!(
+            partition.repair.is_none(),
+            "a discarded view must not leave its suffix fetch blocking future 
repair"
+        );
+    }
+
     #[compio::test]
     async fn 
given_repaired_batch_above_durable_end_when_floor_arrives_should_refuse_commit_floor()
     {
diff --git a/core/partitions/src/journal.rs b/core/partitions/src/journal.rs
index ec45aefb9..fcb3a9af7 100644
--- a/core/partitions/src/journal.rs
+++ b/core/partitions/src/journal.rs
@@ -1329,6 +1329,20 @@ mod tests {
         assert_eq!(journal.last_op(), Some(4));
     }
 
+    #[compio::test]
+    async fn 
repaired_window_shape_rejects_unbounded_sparse_window_before_allocation() {
+        let journal = 
PartitionJournal::<PartitionJournalMemStorage>::default();
+        journal
+            .append(build_prepare(7, HEADER_SIZE + 16).into_frozen())
+            .await
+            .expect("append");
+
+        let shape = journal.repaired_window_shape(0, u64::MAX);
+
+        assert!(!shape.complete);
+        assert!(!shape.holds_messages);
+    }
+
     #[compio::test]
     async fn repair_headers_in_serves_the_commit_point_from_the_evicted_ring() 
{
         // Blank AT the commit point is the one slot a merge can neither adopt 
nor
diff --git a/core/partitions/src/types.rs b/core/partitions/src/types.rs
index df15b4393..bf9d936f6 100644
--- a/core/partitions/src/types.rs
+++ b/core/partitions/src/types.rs
@@ -213,10 +213,20 @@ pub const REPAIR_RETRY_TICKS: u32 = 100;
 /// One in-flight journal-repair stream for a partition group.
 #[derive(Debug, Clone, Copy)]
 pub struct RepairSession {
-    /// Fences stale repair frames from an earlier attempt.
+    /// Fences range replies from an earlier attempt. Repair bodies carry the
+    /// stored prepare header instead, so [`Self::view`] and canonical suffix
+    /// checks fence their ingest.
     pub nonce: u128,
-    /// Last op the stream is expected to serve (the frontier at request time).
-    pub to_op: u64,
+    /// Consensus view in which this session was armed. A later view discards
+    /// the session before any delayed repair body can enter its journal.
+    pub view: u32,
+    /// Committed frontier this repair must make locally walkable. Floor
+    /// completeness and session completion are bounded here.
+    pub commit_to_op: u64,
+    /// Highest op requested from the peer. This may extend above
+    /// [`Self::commit_to_op`] only for the canonical suffix carried by the
+    /// adopted `StartView`.
+    pub fetch_to_op: u64,
     /// Commit floor learned from `RangeEvicted { retained_from }`:
     /// `retained_from - 1`. `None` until (unless) the serving peer reports a
     /// truncated prefix.
diff --git a/core/sdk/src/leader_aware.rs b/core/sdk/src/leader_aware.rs
index 6aa67f202..2d1c61be8 100644
--- a/core/sdk/src/leader_aware.rs
+++ b/core/sdk/src/leader_aware.rs
@@ -413,8 +413,11 @@ impl RosterWalk {
 /// that exact result instead of treating `Connecting` as success.
 #[derive(Debug)]
 pub(crate) struct ConnectCoordinator {
+    id: u64,
     active: AtomicBool,
     abandoned: AtomicBool,
+    active_token: AtomicU64,
+    next_token: AtomicU64,
     generation: AtomicU64,
     result: StdMutex<Option<(u64, Result<(), IggyError>)>>,
     changed: Notify,
@@ -423,8 +426,11 @@ pub(crate) struct ConnectCoordinator {
 impl ConnectCoordinator {
     pub(crate) fn new() -> Self {
         Self {
+            id: NEXT_CONNECT_COORDINATOR_ID.fetch_add(1, Ordering::SeqCst),
             active: AtomicBool::new(false),
             abandoned: AtomicBool::new(false),
+            active_token: AtomicU64::new(0),
+            next_token: AtomicU64::new(1),
             generation: AtomicU64::new(0),
             result: StdMutex::new(None),
             changed: Notify::new(),
@@ -435,9 +441,37 @@ impl ConnectCoordinator {
         self.active.load(Ordering::SeqCst)
     }
 
+    pub(crate) fn current_owner_context(&self) -> Option<ConnectOwnerContext> {
+        let context = CONNECT_OWNER_CONTEXT.try_with(|context| *context).ok()?;
+        (context.coordinator_id == self.id
+            && context.token == self.active_token.load(Ordering::SeqCst))
+        .then_some(context)
+    }
+
+    pub(crate) fn owner_context(
+        &self,
+        token: ConnectOwnerToken,
+        settle_off_leader: bool,
+        single_attempt: bool,
+    ) -> ConnectOwnerContext {
+        ConnectOwnerContext {
+            coordinator_id: self.id,
+            token: token.0,
+            settle_off_leader,
+            single_attempt,
+        }
+    }
+
+    pub(crate) async fn scope_owner<Fut, T>(&self, context: 
ConnectOwnerContext, future: Fut) -> T
+    where
+        Fut: Future<Output = T>,
+    {
+        CONNECT_OWNER_CONTEXT.scope(context, future).await
+    }
+
     pub(crate) async fn run<F, Fut>(&self, operation: F) -> Result<(), 
IggyError>
     where
-        F: FnOnce(bool) -> Fut,
+        F: FnOnce(bool, ConnectOwnerToken) -> Fut,
         Fut: Future<Output = Result<(), IggyError>>,
     {
         let observed_generation = self.generation.load(Ordering::SeqCst);
@@ -447,11 +481,14 @@ impl ConnectCoordinator {
             .is_ok()
         {
             let abandoned = self.abandoned.swap(false, Ordering::SeqCst);
+            let token = self.next_token.fetch_add(1, Ordering::SeqCst).max(1);
+            self.active_token.store(token, Ordering::SeqCst);
             let mut owner = ConnectOwner {
                 coordinator: self,
+                token,
                 completed: false,
             };
-            let result = operation(abandoned).await;
+            let result = operation(abandoned, ConnectOwnerToken(token)).await;
             owner.complete(result.clone());
             return result;
         }
@@ -476,7 +513,7 @@ impl ConnectCoordinator {
             .unwrap_or(Err(IggyError::Disconnected))
     }
 
-    fn finish(&self, result: Result<(), IggyError>, abandoned: bool) {
+    fn finish(&self, token: u64, result: Result<(), IggyError>, abandoned: 
bool) {
         if abandoned {
             self.abandoned.store(true, Ordering::SeqCst);
         }
@@ -486,6 +523,9 @@ impl ConnectCoordinator {
             .expect("connect result mutex poisoned")
             .replace((generation, result));
         self.generation.store(generation, Ordering::SeqCst);
+        if self.active_token.load(Ordering::SeqCst) == token {
+            self.active_token.store(0, Ordering::SeqCst);
+        }
         self.active.store(false, Ordering::SeqCst);
         self.changed.notify_waiters();
     }
@@ -499,12 +539,13 @@ impl Default for ConnectCoordinator {
 
 struct ConnectOwner<'a> {
     coordinator: &'a ConnectCoordinator,
+    token: u64,
     completed: bool,
 }
 
 impl ConnectOwner<'_> {
     fn complete(&mut self, result: Result<(), IggyError>) {
-        self.coordinator.finish(result, false);
+        self.coordinator.finish(self.token, result, false);
         self.completed = true;
     }
 }
@@ -512,11 +553,39 @@ impl ConnectOwner<'_> {
 impl Drop for ConnectOwner<'_> {
     fn drop(&mut self) {
         if !self.completed {
-            self.coordinator.finish(Err(IggyError::Disconnected), true);
+            self.coordinator
+                .finish(self.token, Err(IggyError::Disconnected), true);
         }
     }
 }
 
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub(crate) struct ConnectOwnerToken(u64);
+
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub(crate) struct ConnectOwnerContext {
+    coordinator_id: u64,
+    token: u64,
+    settle_off_leader: bool,
+    single_attempt: bool,
+}
+
+impl ConnectOwnerContext {
+    pub(crate) fn settle_off_leader(self) -> bool {
+        self.settle_off_leader
+    }
+
+    pub(crate) fn single_attempt(self) -> bool {
+        self.single_attempt
+    }
+}
+
+static NEXT_CONNECT_COORDINATOR_ID: AtomicU64 = AtomicU64::new(1);
+
+tokio::task_local! {
+    static CONNECT_OWNER_CONTEXT: ConnectOwnerContext;
+}
+
 /// Struct to track leader redirection state
 #[derive(Debug, Clone)]
 pub struct LeaderRedirectionState {
@@ -753,7 +822,7 @@ mod tests {
             let release = Arc::clone(&release);
             tokio::spawn(async move {
                 coordinator
-                    .run(|abandoned| async move {
+                    .run(|abandoned, _token| async move {
                         assert!(!abandoned);
                         operations.fetch_add(1, Ordering::SeqCst);
                         started.notify_one();
@@ -769,7 +838,7 @@ mod tests {
             let operations = Arc::clone(&operations);
             tokio::spawn(async move {
                 coordinator
-                    .run(|_| async move {
+                    .run(|_, _token| async move {
                         operations.fetch_add(1, Ordering::SeqCst);
                         Ok(())
                     })
@@ -799,7 +868,7 @@ mod tests {
             let started = Arc::clone(&started);
             tokio::spawn(async move {
                 coordinator
-                    .run(|_| async move {
+                    .run(|_, _token| async move {
                         started.notify_one();
                         std::future::pending::<Result<(), IggyError>>().await
                     })
@@ -809,7 +878,7 @@ mod tests {
         started.notified().await;
         let waiter = {
             let coordinator = Arc::clone(&coordinator);
-            tokio::spawn(async move { coordinator.run(|_| async { Ok(()) 
}).await })
+            tokio::spawn(async move { coordinator.run(|_, _token| async { 
Ok(()) }).await })
         };
         tokio::task::yield_now().await;
         owner.abort();
@@ -821,7 +890,7 @@ mod tests {
         let observed_abandoned = Arc::new(AtomicBool::new(false));
         let marker = Arc::clone(&observed_abandoned);
         coordinator
-            .run(|abandoned| async move {
+            .run(|abandoned, _token| async move {
                 marker.store(abandoned, Ordering::SeqCst);
                 Ok(())
             })
@@ -830,6 +899,145 @@ mod tests {
         assert!(observed_abandoned.load(Ordering::SeqCst));
     }
 
+    #[tokio::test]
+    async fn connect_owner_context_is_visible_only_to_the_owner_task() {
+        let coordinator = Arc::new(ConnectCoordinator::new());
+        let started = Arc::new(Notify::new());
+        let release = Arc::new(Notify::new());
+        let owner = {
+            let coordinator = Arc::clone(&coordinator);
+            let owner_coordinator = Arc::clone(&coordinator);
+            let started = Arc::clone(&started);
+            let release = Arc::clone(&release);
+            tokio::spawn(async move {
+                coordinator
+                    .run(move |_, token| async move {
+                        let context = owner_coordinator.owner_context(token, 
true, true);
+                        owner_coordinator
+                            .scope_owner(context, async {
+                                assert_eq!(
+                                    owner_coordinator.current_owner_context(),
+                                    Some(context)
+                                );
+                                started.notify_one();
+                                release.notified().await;
+                                Ok(())
+                            })
+                            .await
+                    })
+                    .await
+            })
+        };
+        started.notified().await;
+
+        assert!(coordinator.is_active());
+        assert_eq!(coordinator.current_owner_context(), None);
+        release.notify_one();
+        owner.await.unwrap().unwrap();
+        assert_eq!(coordinator.current_owner_context(), None);
+    }
+
+    #[tokio::test]
+    async fn owner_tokens_cannot_collide_across_connect_coordinators() {
+        let first = Arc::new(ConnectCoordinator::new());
+        let second = Arc::new(ConnectCoordinator::new());
+        let second_started = Arc::new(Notify::new());
+        let second_release = Arc::new(Notify::new());
+        let second_owner = {
+            let second = Arc::clone(&second);
+            let owner_coordinator = Arc::clone(&second);
+            let second_started = Arc::clone(&second_started);
+            let second_release = Arc::clone(&second_release);
+            tokio::spawn(async move {
+                second
+                    .run(move |_, token| async move {
+                        let context = owner_coordinator.owner_context(token, 
false, false);
+                        owner_coordinator
+                            .scope_owner(context, async {
+                                second_started.notify_one();
+                                second_release.notified().await;
+                                Ok(())
+                            })
+                            .await
+                    })
+                    .await
+            })
+        };
+        second_started.notified().await;
+
+        let first_owner = Arc::clone(&first);
+        let second_from_first = Arc::clone(&second);
+        first
+            .run(move |_, token| async move {
+                let context = first_owner.owner_context(token, true, true);
+                first_owner
+                    .scope_owner(context, async {
+                        assert_eq!(first_owner.current_owner_context(), 
Some(context));
+                        assert_eq!(second_from_first.current_owner_context(), 
None);
+                        Ok(())
+                    })
+                    .await
+            })
+            .await
+            .unwrap();
+
+        second_release.notify_one();
+        second_owner.await.unwrap().unwrap();
+    }
+
+    #[tokio::test]
+    async fn a_waiters_settlement_mode_cannot_leak_into_the_owner() {
+        let coordinator = Arc::new(ConnectCoordinator::new());
+        let started = Arc::new(Notify::new());
+        let release = Arc::new(Notify::new());
+        let waiter_ran = Arc::new(AtomicBool::new(false));
+        let owner = {
+            let coordinator = Arc::clone(&coordinator);
+            let owner_coordinator = Arc::clone(&coordinator);
+            let started = Arc::clone(&started);
+            let release = Arc::clone(&release);
+            tokio::spawn(async move {
+                coordinator
+                    .run(move |_, token| async move {
+                        let context = owner_coordinator.owner_context(token, 
false, false);
+                        owner_coordinator
+                            .scope_owner(context, async {
+                                assert!(!context.settle_off_leader());
+                                assert!(!context.single_attempt());
+                                started.notify_one();
+                                release.notified().await;
+                                Ok(())
+                            })
+                            .await
+                    })
+                    .await
+            })
+        };
+        started.notified().await;
+        let waiter = {
+            let coordinator = Arc::clone(&coordinator);
+            let waiter_coordinator = Arc::clone(&coordinator);
+            let waiter_ran = Arc::clone(&waiter_ran);
+            tokio::spawn(async move {
+                coordinator
+                    .run(move |_, token| async move {
+                        waiter_ran.store(true, Ordering::SeqCst);
+                        let context = waiter_coordinator.owner_context(token, 
true, true);
+                        assert!(context.settle_off_leader());
+                        assert!(context.single_attempt());
+                        Ok(())
+                    })
+                    .await
+            })
+        };
+        tokio::task::yield_now().await;
+        release.notify_one();
+
+        owner.await.unwrap().unwrap();
+        waiter.await.unwrap().unwrap();
+        assert!(!waiter_ran.load(Ordering::SeqCst));
+    }
+
     #[test]
     fn test_normalize_address() {
         assert_eq!(normalize_address("localhost:8090"), "127.0.0.1:8090");
diff --git a/core/sdk/src/quic/quic_client.rs b/core/sdk/src/quic/quic_client.rs
index 88b4bb18f..1cb1da19f 100644
--- a/core/sdk/src/quic/quic_client.rs
+++ b/core/sdk/src/quic/quic_client.rs
@@ -16,8 +16,8 @@
 // under the License.
 
 use crate::leader_aware::{
-    ConnectCoordinator, LeaderRedirectionState, RosterWalk, 
check_and_redirect_to_leader,
-    is_unauthenticated_metadata_probe,
+    ConnectCoordinator, ConnectOwnerContext, LeaderRedirectionState, 
RosterWalk,
+    check_and_redirect_to_leader, is_unauthenticated_metadata_probe,
 };
 use crate::prelude::AutoLogin;
 use crate::session::ConsensusSession;
@@ -47,7 +47,6 @@ use std::net::{SocketAddr, ToSocketAddrs};
 use std::str::FromStr;
 use std::sync::Arc;
 use std::sync::Mutex as StdMutex;
-use std::sync::atomic::{AtomicBool, Ordering};
 use std::time::Duration;
 use tokio::sync::Mutex;
 use tokio::time::sleep;
@@ -96,10 +95,6 @@ pub struct QuicClient {
     /// as walk candidates for a request the current node keeps refusing to
     /// admit (its replica of the target partition group is not the primary).
     roster_endpoints: Mutex<Vec<String>>,
-    /// Set when the next connect must stay on the endpoint it dialed instead
-    /// of settling on the metadata leader: the roster walk has to survive its
-    /// own sign-in's leader check. See the TCP client's twin field.
-    settle_off_leader_once: AtomicBool,
     /// Serializes leader checks and roster walks after refused requests, so
     /// concurrent QUIC streams cannot tear down each other's new connection.
     routing_lock: Mutex<()>,
@@ -152,6 +147,7 @@ impl BinaryTransport for QuicClient {
     }
 
     async fn send_raw_with_response(&self, code: u32, payload: Bytes) -> 
Result<Bytes, IggyError> {
+        let roster_deadline = tokio::time::Instant::now() + 
RESPONSE_READ_TIMEOUT;
         let mut result = self.send_raw(code, payload.clone()).await;
 
         // A persistent not-admitted refusal is a verdict about who leads the
@@ -166,10 +162,23 @@ impl BinaryTransport for QuicClient {
             && self.config.reconnection.enabled
             && !matches!(self.config.auto_login, AutoLogin::Disabled)
         {
-            let _routing_guard = self.routing_lock.lock().await;
+            let _routing_guard =
+                match tokio::time::timeout_at(roster_deadline, 
self.routing_lock.lock()).await {
+                    Ok(guard) => guard,
+                    Err(_) => return Err(IggyError::TransientNotAccepted),
+                };
+            let overall_deadline = roster_deadline;
             // A concurrent refused request may have completed the movement
             // while this request waited for the gate.
-            result = self.send_raw(code, payload.clone()).await;
+            result = match tokio::time::timeout_at(
+                overall_deadline,
+                self.send_raw(code, payload.clone()),
+            )
+            .await
+            {
+                Ok(result) => result,
+                Err(_) => Err(IggyError::TransientNotAccepted),
+            };
             let mut roster_walk: Option<RosterWalk> = None;
             // Once the walk starts it keeps walking: a leader recheck between
             // hops would put the request straight back on the node whose
@@ -181,28 +190,90 @@ impl BinaryTransport for QuicClient {
                     false
                 } else {
                     checked_metadata_leader = true;
-                    let redirected = 
matches!(self.handle_leader_redirection().await, Ok(true));
+                    let redirected = match tokio::time::timeout_at(
+                        overall_deadline,
+                        self.handle_leader_redirection(),
+                    )
+                    .await
+                    {
+                        Ok(result) => matches!(result, Ok(true)),
+                        Err(_) => return Err(IggyError::TransientNotAccepted),
+                    };
                     let roster = self.roster_endpoints.lock().await.clone();
                     roster_walk = Some(RosterWalk::new(&current, &roster));
                     redirected
                 };
-                if !redirected {
-                    let Some(next) = 
roster_walk.as_mut().and_then(RosterWalk::next) else {
-                        break;
-                    };
-                    self.settle_on_endpoint(next).await?;
-                } else {
+                let (mut target, mut needs_settle) = if redirected {
                     let target = 
self.current_server_address.lock().await.clone();
                     if let Some(walk) = roster_walk.as_mut() {
                         walk.record_attempt(&target);
                     }
+                    (target, false)
+                } else if let Some(next) = 
roster_walk.as_mut().and_then(RosterWalk::next) {
+                    (next, true)
+                } else {
+                    break;
+                };
+
+                loop {
+                    if tokio::time::Instant::now() >= overall_deadline {
+                        return Err(IggyError::TransientNotAccepted);
+                    }
+                    let settled = if needs_settle {
+                        match tokio::time::timeout_at(
+                            overall_deadline,
+                            self.settle_on_endpoint(target.clone()),
+                        )
+                        .await
+                        {
+                            Ok(Ok(())) => true,
+                            Ok(Err(IggyError::CannotEstablishConnection)) => 
false,
+                            Ok(Err(error)) => return Err(error),
+                            Err(_) => return 
Err(IggyError::TransientNotAccepted),
+                        }
+                    } else {
+                        true
+                    };
+                    if settled {
+                        let connect_result = if needs_settle {
+                            tokio::time::timeout_at(overall_deadline, 
self.connect_off_leader())
+                                .await
+                        } else {
+                            tokio::time::timeout_at(overall_deadline, 
self.connect()).await
+                        };
+                        match connect_result {
+                            Ok(Ok(())) => {
+                                let connected = 
self.current_server_address.lock().await.clone();
+                                let first_visit = roster_walk
+                                    .as_mut()
+                                    .is_some_and(|walk| 
walk.record_attempt(&connected));
+                                if 
crate::leader_aware::is_same_spelling(&connected, &target)
+                                    || first_visit
+                                {
+                                    break;
+                                }
+                            }
+                            Ok(Err(IggyError::CannotEstablishConnection)) => {}
+                            Ok(Err(error)) => return Err(error),
+                            Err(_) => return 
Err(IggyError::TransientNotAccepted),
+                        }
+                    }
+
+                    let Some(next) = 
roster_walk.as_mut().and_then(RosterWalk::next) else {
+                        return Err(IggyError::TransientNotAccepted);
+                    };
+                    target = next;
+                    needs_settle = true;
                 }
-                self.connect().await?;
-                let connected = 
self.current_server_address.lock().await.clone();
-                if let Some(walk) = roster_walk.as_mut() {
-                    walk.record_attempt(&connected);
-                }
-                result = self.send_raw(code, payload.clone()).await;
+                result = match tokio::time::timeout_at(
+                    overall_deadline,
+                    self.send_raw(code, payload.clone()),
+                )
+                .await
+                {
+                    Ok(result) => result,
+                    Err(_) => Err(IggyError::TransientNotAccepted),
+                };
             }
         }
 
@@ -248,7 +319,10 @@ impl BinaryTransport for QuicClient {
 
         let replay_after_reconnect = replay_after_session_reset_is_safe(code, 
&error);
         let skip_auto_login = is_login_register_code(code);
-        let nested_connect = skip_auto_login && 
self.connect_coordinator.is_active();
+        let owner_context = skip_auto_login
+            .then(|| self.connect_coordinator.current_owner_context())
+            .flatten();
+        let nested_connect = owner_context.is_some();
         let _routing_guard = if nested_connect {
             None
         } else {
@@ -271,7 +345,8 @@ impl BinaryTransport for QuicClient {
             server_address, self.config.client_address
         );
         let reconnect = if nested_connect {
-            self.connect_inner().await
+            self.connect_inner(owner_context.expect("owner context checked 
above"))
+                .await
         } else {
             self.connect().await
         };
@@ -405,7 +480,6 @@ impl QuicClient {
             consensus_session: 
Arc::new(StdMutex::new(ConsensusSession::new())),
             skip_auto_login_once: Mutex::new(false),
             roster_endpoints: Mutex::new(Vec::new()),
-            settle_off_leader_once: AtomicBool::new(false),
             routing_lock: Mutex::new(()),
             connect_coordinator: ConnectCoordinator::new(),
             consumer_group_state: 
Arc::new(iggy_common::ConsumerGroupClientState::new()),
@@ -446,20 +520,36 @@ impl QuicClient {
     }
 
     async fn connect(&self) -> Result<(), IggyError> {
+        self.connect_with_settlement(false).await
+    }
+
+    pub(crate) async fn connect_off_leader(&self) -> Result<(), IggyError> {
+        self.connect_with_settlement(true).await
+    }
+
+    async fn connect_with_settlement(&self, settle_off_leader: bool) -> 
Result<(), IggyError> {
         self.connect_coordinator
-            .run(|abandoned| async move {
-                if abandoned {
-                    self.clear_abandoned_connect().await?;
-                }
-                self.connect_inner().await
+            .run(|abandoned, token| async move {
+                let context = self.connect_coordinator.owner_context(
+                    token,
+                    settle_off_leader,
+                    settle_off_leader,
+                );
+                self.connect_coordinator
+                    .scope_owner(context, async move {
+                        if abandoned {
+                            self.clear_abandoned_connect().await?;
+                        }
+                        self.connect_inner(context).await
+                    })
+                    .await
             })
             .await
     }
 
-    async fn connect_inner(&self) -> Result<(), IggyError> {
-        // Consume before fallible connection work. A walk target that cannot
-        // connect must not make an unrelated later reconnect skip settlement.
-        let settle_off_leader = self.settle_off_leader_once.swap(false, 
Ordering::SeqCst);
+    async fn connect_inner(&self, context: ConnectOwnerContext) -> Result<(), 
IggyError> {
+        let settle_off_leader = context.settle_off_leader();
+        let single_attempt = context.single_attempt();
         loop {
             match self.get_state().await {
                 ClientState::Shutdown => {
@@ -480,7 +570,7 @@ impl QuicClient {
             }
 
             self.set_state(ClientState::Connecting).await;
-            if let Some(connected_at) = 
self.connected_at.lock().await.as_ref() {
+            if !single_attempt && let Some(connected_at) = 
self.connected_at.lock().await.as_ref() {
                 let now = IggyTimestamp::now();
                 let elapsed = now.as_micros() - connected_at.as_micros();
                 let interval = 
self.config.reconnection.reestablish_after.as_micros();
@@ -518,14 +608,26 @@ impl QuicClient {
                     "{NAME} client is connecting to server: {}...",
                     server_address
                 );
-                let connection_result = self
+                let connection_result = match self
                     .endpoint
                     .connect(server_address, &self.config.server_name)
-                    .unwrap()
-                    .await;
+                {
+                    Ok(connecting) => connecting.await,
+                    Err(error) => {
+                        error!("Failed to start QUIC connection: {error}");
+                        self.set_state(ClientState::Disconnected).await;
+                        
self.publish_event(DiagnosticEvent::Disconnected).await;
+                        return Err(IggyError::CannotEstablishConnection);
+                    }
+                };
 
                 if connection_result.is_err() {
                     error!("Failed to connect to server: {}", server_address);
+                    if single_attempt {
+                        self.set_state(ClientState::Disconnected).await;
+                        
self.publish_event(DiagnosticEvent::Disconnected).await;
+                        return Err(IggyError::CannotEstablishConnection);
+                    }
                     if !self.config.reconnection.enabled {
                         warn!("Automatic reconnection is disabled.");
                         return Err(IggyError::CannotEstablishConnection);
@@ -717,7 +819,6 @@ impl QuicClient {
         self.connected_at.lock().await.take();
         self.disconnect().await?;
         *self.current_server_address.lock().await = next;
-        self.settle_off_leader_once.store(true, Ordering::SeqCst);
         Ok(())
     }
 
@@ -979,6 +1080,25 @@ fn configure(config: &QuicClientConfig) -> 
Result<ClientConfig, IggyError> {
 mod tests {
     use super::*;
 
+    #[tokio::test]
+    async fn a_roster_hop_does_not_enter_the_reconnect_ladder() {
+        let socket = std::net::UdpSocket::bind("127.0.0.1:0").expect("reserve 
UDP address");
+        let server_address = socket.local_addr().unwrap().to_string();
+        drop(socket);
+        let client = QuicClient::create(Arc::new(QuicClientConfig {
+            server_address,
+            max_idle_timeout: 100,
+            ..QuicClientConfig::default()
+        }))
+        .expect("create QUIC client");
+
+        let result = tokio::time::timeout(Duration::from_secs(10), 
client.connect_off_leader())
+            .await
+            .expect("one QUIC dial must not enter unlimited reconnect");
+        assert!(matches!(result, Err(IggyError::CannotEstablishConnection)));
+        assert_eq!(client.get_state().await, ClientState::Disconnected);
+    }
+
     #[tokio::test]
     async fn should_fail_with_a_zero_heartbeat_interval() {
         let value = 
"iggy+quic://user:[email protected]:1234?heartbeat_interval=none";
diff --git a/core/sdk/src/tcp/tcp_client.rs b/core/sdk/src/tcp/tcp_client.rs
index 634c4ddda..23f70f791 100644
--- a/core/sdk/src/tcp/tcp_client.rs
+++ b/core/sdk/src/tcp/tcp_client.rs
@@ -16,8 +16,9 @@
 // under the License.
 
 use crate::leader_aware::{
-    ConnectCoordinator, LeaderRedirectionState, RosterWalk, 
check_and_redirect_to_leader,
-    is_same_spelling, is_unauthenticated_metadata_probe, 
read_transport_endpoints,
+    ConnectCoordinator, ConnectOwnerContext, LeaderRedirectionState, 
RosterWalk,
+    check_and_redirect_to_leader, is_same_spelling, 
is_unauthenticated_metadata_probe,
+    read_transport_endpoints,
 };
 use crate::prelude::Client;
 use crate::prelude::TcpClientConfig;
@@ -125,12 +126,6 @@ pub struct TcpClient {
     // contention with zero correctness benefit.
     consensus_session: Arc<StdMutex<ConsensusSession>>,
     skip_auto_login_once: Mutex<bool>,
-    /// Set when the next connect must stay on the endpoint it dialed instead
-    /// of settling on the metadata leader. Metadata and partition consensus
-    /// groups elect independently, so the metadata leader can hold a follower
-    /// replica of the partition a request targets; the failover that walks the
-    /// roster past it has to survive its own sign-in's leader check.
-    settle_off_leader_once: AtomicBool,
     /// Serializes connection movement after a refused request. The stream is
     /// lockstep, but the refusal releases it before the leader check and
     /// reconnect, where another request could otherwise run a competing walk.
@@ -266,7 +261,10 @@ impl BinaryTransport for TcpClient {
         let replay_after_reconnect = replay_after_session_reset_is_safe(code, 
&error);
 
         let skip_auto_login = is_login_register_code(code);
-        let nested_connect = skip_auto_login && 
self.connect_coordinator.is_active();
+        let owner_context = skip_auto_login
+            .then(|| self.connect_coordinator.current_owner_context())
+            .flatten();
+        let nested_connect = owner_context.is_some();
         let _routing_guard = if nested_connect {
             None
         } else {
@@ -277,6 +275,7 @@ impl BinaryTransport for TcpClient {
             if !replay_after_reconnect {
                 return Err(error);
             }
+            drop(_routing_guard);
             return self.send_raw(code, payload).await;
         }
         self.disconnect_transport().await?;
@@ -295,7 +294,8 @@ impl BinaryTransport for TcpClient {
         }
 
         let reconnect = if nested_connect {
-            self.connect_inner().await
+            self.connect_inner(owner_context.expect("owner context checked 
above"))
+                .await
         } else {
             self.connect().await
         };
@@ -313,6 +313,7 @@ impl BinaryTransport for TcpClient {
             return Err(error);
         }
 
+        drop(_routing_guard);
         self.send_raw(code, payload).await
     }
 
@@ -498,7 +499,6 @@ impl TcpClient {
             session_credentials: Mutex::new(None),
             consensus_session: 
Arc::new(StdMutex::new(ConsensusSession::new())),
             skip_auto_login_once: Mutex::new(false),
-            settle_off_leader_once: AtomicBool::new(false),
             routing_lock: Mutex::new(()),
             connect_coordinator: ConnectCoordinator::new(),
             consumer_group_state: 
Arc::new(iggy_common::ConsumerGroupClientState::new()),
@@ -506,20 +506,33 @@ impl TcpClient {
     }
 
     async fn connect(&self) -> Result<(), IggyError> {
+        self.connect_with_settlement(false).await
+    }
+
+    async fn connect_off_leader(&self) -> Result<(), IggyError> {
+        self.connect_with_settlement(true).await
+    }
+
+    async fn connect_with_settlement(&self, settle_off_leader: bool) -> 
Result<(), IggyError> {
         self.connect_coordinator
-            .run(|abandoned| async move {
-                if abandoned {
-                    self.clear_abandoned_connect().await?;
-                }
-                self.connect_inner().await
+            .run(|abandoned, token| async move {
+                let context =
+                    self.connect_coordinator
+                        .owner_context(token, settle_off_leader, false);
+                self.connect_coordinator
+                    .scope_owner(context, async move {
+                        if abandoned {
+                            self.clear_abandoned_connect().await?;
+                        }
+                        self.connect_inner(context).await
+                    })
+                    .await
             })
             .await
     }
 
-    async fn connect_inner(&self) -> Result<(), IggyError> {
-        // Consume before fallible connection work. A walk target that cannot
-        // connect must not make an unrelated later reconnect skip settlement.
-        let settle_off_leader = self.settle_off_leader_once.swap(false, 
Ordering::SeqCst);
+    async fn connect_inner(&self, context: ConnectOwnerContext) -> Result<(), 
IggyError> {
+        let settle_off_leader = context.settle_off_leader();
         loop {
             // Read and claimed under one lock acquisition. Apart, two callers
             // both find `Disconnected` and both sweep: the loser's
@@ -887,7 +900,6 @@ impl TcpClient {
         self.connected_at.lock().await.take();
         self.disconnect_transport().await?;
         *self.current_server_address.lock().await = next;
-        self.settle_off_leader_once.store(true, Ordering::SeqCst);
         Ok(())
     }
 
@@ -1312,7 +1324,12 @@ impl TcpClient {
                         if needs_settle {
                             self.settle_on_endpoint(target.clone()).await?;
                         }
-                        match self.connect().await {
+                        let connect = if needs_settle {
+                            self.connect_off_leader().await
+                        } else {
+                            self.connect().await
+                        };
+                        match connect {
                             Ok(()) => {
                                 let connected = 
self.current_server_address.lock().await.clone();
                                 let first_visit = roster_walk
@@ -1650,6 +1667,41 @@ mod tests {
         ));
     }
 
+    #[tokio::test]
+    async fn an_unrelated_public_login_cannot_take_over_a_connect_owner() {
+        let (_listener, silent) = live_endpoint().await;
+        let client = Arc::new(
+            TcpClient::create(Arc::new(TcpClientConfig {
+                server_address: silent,
+                tls_enabled: true,
+                tls_validate_certificate: false,
+                ..TcpClientConfig::default()
+            }))
+            .expect("create the client"),
+        );
+        let owner = {
+            let client = Arc::clone(&client);
+            tokio::spawn(async move { TcpClient::connect(&client).await })
+        };
+        while client.get_state().await != ClientState::Connecting {
+            tokio::task::yield_now().await;
+        }
+        let mut login = {
+            let client = Arc::clone(&client);
+            tokio::spawn(async move { client.login_user("iggy", "iggy").await 
})
+        };
+
+        assert!(
+            tokio::time::timeout(std::time::Duration::from_millis(100), &mut 
login)
+                .await
+                .is_err(),
+            "an unrelated login bypassed the connect owner instead of waiting"
+        );
+        assert_eq!(client.get_state().await, ClientState::Connecting);
+        owner.abort();
+        assert!(login.await.unwrap().is_err());
+    }
+
     // With reconnection off there are no retries, but the endpoints the roster
     // named are still there to be tried and each gets its one turn.
     #[tokio::test]
diff --git a/core/sdk/src/websocket/websocket_client.rs 
b/core/sdk/src/websocket/websocket_client.rs
index c808d5fca..f9573e79a 100644
--- a/core/sdk/src/websocket/websocket_client.rs
+++ b/core/sdk/src/websocket/websocket_client.rs
@@ -16,8 +16,8 @@
 // under the License.
 
 use crate::leader_aware::{
-    ConnectCoordinator, LeaderRedirectionState, RosterWalk, 
check_and_redirect_to_leader,
-    is_unauthenticated_metadata_probe,
+    ConnectCoordinator, ConnectOwnerContext, LeaderRedirectionState, 
RosterWalk,
+    check_and_redirect_to_leader, is_unauthenticated_metadata_probe,
 };
 use crate::session::ConsensusSession;
 use crate::vsr::replay_after_session_reset_is_safe;
@@ -44,7 +44,6 @@ use secrecy::ExposeSecret;
 use std::net::SocketAddr;
 use std::sync::Arc;
 use std::sync::Mutex as StdMutex;
-use std::sync::atomic::{AtomicBool, Ordering};
 use tokio::net::TcpStream;
 use tokio::sync::Mutex;
 use tokio::time::sleep;
@@ -92,10 +91,6 @@ pub struct WebSocketClient {
     /// as walk candidates for a request the current node keeps refusing to
     /// admit (its replica of the target partition group is not the primary).
     roster_endpoints: Mutex<Vec<String>>,
-    /// Set when the next connect must stay on the endpoint it dialed instead
-    /// of settling on the metadata leader: the roster walk has to survive its
-    /// own sign-in's leader check. See the TCP client's twin field.
-    settle_off_leader_once: AtomicBool,
     /// Serializes leader checks and roster walks after refused requests, so
     /// concurrent callers cannot tear down each other's new connection.
     routing_lock: Mutex<()>,
@@ -145,6 +140,7 @@ impl BinaryTransport for WebSocketClient {
     }
 
     async fn send_raw_with_response(&self, code: u32, payload: Bytes) -> 
Result<Bytes, IggyError> {
+        let roster_deadline = tokio::time::Instant::now() + 
RESPONSE_READ_TIMEOUT;
         let mut result = self.send_raw(code, payload.clone()).await;
 
         // A persistent not-admitted refusal is a verdict about who leads the
@@ -159,10 +155,23 @@ impl BinaryTransport for WebSocketClient {
             && self.config.reconnection.enabled
             && !matches!(self.config.auto_login, AutoLogin::Disabled)
         {
-            let _routing_guard = self.routing_lock.lock().await;
+            let _routing_guard =
+                match tokio::time::timeout_at(roster_deadline, 
self.routing_lock.lock()).await {
+                    Ok(guard) => guard,
+                    Err(_) => return Err(IggyError::TransientNotAccepted),
+                };
+            let overall_deadline = roster_deadline;
             // A concurrent refused request may have completed the movement
             // while this request waited for the gate.
-            result = self.send_raw(code, payload.clone()).await;
+            result = match tokio::time::timeout_at(
+                overall_deadline,
+                self.send_raw(code, payload.clone()),
+            )
+            .await
+            {
+                Ok(result) => result,
+                Err(_) => Err(IggyError::TransientNotAccepted),
+            };
             let mut roster_walk: Option<RosterWalk> = None;
             // Once the walk starts it keeps walking: a leader recheck between
             // hops would put the request straight back on the node whose
@@ -174,28 +183,90 @@ impl BinaryTransport for WebSocketClient {
                     false
                 } else {
                     checked_metadata_leader = true;
-                    let redirected = 
matches!(self.handle_leader_redirection().await, Ok(true));
+                    let redirected = match tokio::time::timeout_at(
+                        overall_deadline,
+                        self.handle_leader_redirection(),
+                    )
+                    .await
+                    {
+                        Ok(result) => matches!(result, Ok(true)),
+                        Err(_) => return Err(IggyError::TransientNotAccepted),
+                    };
                     let roster = self.roster_endpoints.lock().await.clone();
                     roster_walk = Some(RosterWalk::new(&current, &roster));
                     redirected
                 };
-                if !redirected {
-                    let Some(next) = 
roster_walk.as_mut().and_then(RosterWalk::next) else {
-                        break;
-                    };
-                    self.settle_on_endpoint(next).await?;
-                } else {
+                let (mut target, mut needs_settle) = if redirected {
                     let target = 
self.current_server_address.lock().await.clone();
                     if let Some(walk) = roster_walk.as_mut() {
                         walk.record_attempt(&target);
                     }
+                    (target, false)
+                } else if let Some(next) = 
roster_walk.as_mut().and_then(RosterWalk::next) {
+                    (next, true)
+                } else {
+                    break;
+                };
+
+                loop {
+                    if tokio::time::Instant::now() >= overall_deadline {
+                        return Err(IggyError::TransientNotAccepted);
+                    }
+                    let settled = if needs_settle {
+                        match tokio::time::timeout_at(
+                            overall_deadline,
+                            self.settle_on_endpoint(target.clone()),
+                        )
+                        .await
+                        {
+                            Ok(Ok(())) => true,
+                            Ok(Err(IggyError::CannotEstablishConnection)) => 
false,
+                            Ok(Err(error)) => return Err(error),
+                            Err(_) => return 
Err(IggyError::TransientNotAccepted),
+                        }
+                    } else {
+                        true
+                    };
+                    if settled {
+                        let connect_result = if needs_settle {
+                            tokio::time::timeout_at(overall_deadline, 
self.connect_off_leader())
+                                .await
+                        } else {
+                            tokio::time::timeout_at(overall_deadline, 
self.connect()).await
+                        };
+                        match connect_result {
+                            Ok(Ok(())) => {
+                                let connected = 
self.current_server_address.lock().await.clone();
+                                let first_visit = roster_walk
+                                    .as_mut()
+                                    .is_some_and(|walk| 
walk.record_attempt(&connected));
+                                if 
crate::leader_aware::is_same_spelling(&connected, &target)
+                                    || first_visit
+                                {
+                                    break;
+                                }
+                            }
+                            Ok(Err(IggyError::CannotEstablishConnection)) => {}
+                            Ok(Err(error)) => return Err(error),
+                            Err(_) => return 
Err(IggyError::TransientNotAccepted),
+                        }
+                    }
+
+                    let Some(next) = 
roster_walk.as_mut().and_then(RosterWalk::next) else {
+                        return Err(IggyError::TransientNotAccepted);
+                    };
+                    target = next;
+                    needs_settle = true;
                 }
-                self.connect().await?;
-                let connected = 
self.current_server_address.lock().await.clone();
-                if let Some(walk) = roster_walk.as_mut() {
-                    walk.record_attempt(&connected);
-                }
-                result = self.send_raw(code, payload.clone()).await;
+                result = match tokio::time::timeout_at(
+                    overall_deadline,
+                    self.send_raw(code, payload.clone()),
+                )
+                .await
+                {
+                    Ok(result) => result,
+                    Err(_) => Err(IggyError::TransientNotAccepted),
+                };
             }
         }
 
@@ -238,7 +309,10 @@ impl BinaryTransport for WebSocketClient {
 
         let replay_after_reconnect = replay_after_session_reset_is_safe(code, 
&error);
         let skip_auto_login = is_login_register_code(code);
-        let nested_connect = skip_auto_login && 
self.connect_coordinator.is_active();
+        let owner_context = skip_auto_login
+            .then(|| self.connect_coordinator.current_owner_context())
+            .flatten();
+        let nested_connect = owner_context.is_some();
         let _routing_guard = if nested_connect {
             None
         } else {
@@ -266,7 +340,8 @@ impl BinaryTransport for WebSocketClient {
         }
 
         let reconnect = if nested_connect {
-            self.connect_inner().await
+            self.connect_inner(owner_context.expect("owner context checked 
above"))
+                .await
         } else {
             self.connect().await
         };
@@ -352,7 +427,6 @@ impl WebSocketClient {
             consensus_session: 
Arc::new(StdMutex::new(ConsensusSession::new())),
             skip_auto_login_once: Mutex::new(false),
             roster_endpoints: Mutex::new(Vec::new()),
-            settle_off_leader_once: AtomicBool::new(false),
             routing_lock: Mutex::new(()),
             connect_coordinator: ConnectCoordinator::new(),
             consumer_group_state: 
Arc::new(iggy_common::ConsumerGroupClientState::new()),
@@ -376,20 +450,36 @@ impl WebSocketClient {
     }
 
     async fn connect(&self) -> Result<(), IggyError> {
+        self.connect_with_settlement(false).await
+    }
+
+    pub(crate) async fn connect_off_leader(&self) -> Result<(), IggyError> {
+        self.connect_with_settlement(true).await
+    }
+
+    async fn connect_with_settlement(&self, settle_off_leader: bool) -> 
Result<(), IggyError> {
         self.connect_coordinator
-            .run(|abandoned| async move {
-                if abandoned {
-                    self.clear_abandoned_connect().await?;
-                }
-                self.connect_inner().await
+            .run(|abandoned, token| async move {
+                let context = self.connect_coordinator.owner_context(
+                    token,
+                    settle_off_leader,
+                    settle_off_leader,
+                );
+                self.connect_coordinator
+                    .scope_owner(context, async move {
+                        if abandoned {
+                            self.clear_abandoned_connect().await?;
+                        }
+                        self.connect_inner(context).await
+                    })
+                    .await
             })
             .await
     }
 
-    async fn connect_inner(&self) -> Result<(), IggyError> {
-        // Consume before fallible connection work. A walk target that cannot
-        // connect must not make an unrelated later reconnect skip settlement.
-        let settle_off_leader = self.settle_off_leader_once.swap(false, 
Ordering::SeqCst);
+    async fn connect_inner(&self, context: ConnectOwnerContext) -> Result<(), 
IggyError> {
+        let settle_off_leader = context.settle_off_leader();
+        let single_attempt = context.single_attempt();
         loop {
             if self.get_state().await == ClientState::Connected {
                 return Ok(());
@@ -441,7 +531,10 @@ impl WebSocketClient {
                     })?;
 
                 let connection_stream = if self.config.tls_enabled {
-                    match self.connect_tls(server_addr, &mut 
retry_count).await {
+                    match self
+                        .connect_tls(server_addr, &mut retry_count, 
single_attempt)
+                        .await
+                    {
                         Ok(stream) => stream,
                         Err(IggyError::CannotEstablishConnection) => {
                             return Err(IggyError::CannotEstablishConnection);
@@ -449,7 +542,10 @@ impl WebSocketClient {
                         Err(_) => continue, // retry
                     }
                 } else {
-                    match self.connect_plain(server_addr, &mut 
retry_count).await {
+                    match self
+                        .connect_plain(server_addr, &mut retry_count, 
single_attempt)
+                        .await
+                    {
                         Ok(stream) => stream,
                         Err(IggyError::CannotEstablishConnection) => {
                             return Err(IggyError::CannotEstablishConnection);
@@ -491,6 +587,7 @@ impl WebSocketClient {
         &self,
         server_addr: SocketAddr,
         retry_count: &mut u32,
+        single_attempt: bool,
     ) -> Result<WebSocketStreamKind, IggyError> {
         let tcp_stream = match TcpStream::connect(&server_addr).await {
             Ok(stream) => stream,
@@ -499,7 +596,9 @@ impl WebSocketClient {
                     "Failed to connect to server: {}. Error: {}",
                     self.config.server_address, error
                 );
-                return self.handle_connection_error(retry_count).await;
+                return self
+                    .handle_connection_error(retry_count, single_attempt)
+                    .await;
             }
         };
 
@@ -516,7 +615,9 @@ impl WebSocketClient {
                 Ok(result) => result,
                 Err(error) => {
                     error!("WebSocket handshake failed: {}", error);
-                    return self.handle_connection_error(retry_count).await;
+                    return self
+                        .handle_connection_error(retry_count, single_attempt)
+                        .await;
                 }
             };
 
@@ -533,12 +634,15 @@ impl WebSocketClient {
         &self,
         server_addr: SocketAddr,
         retry_count: &mut u32,
+        single_attempt: bool,
     ) -> Result<WebSocketStreamKind, IggyError> {
         let tcp_stream = match TcpStream::connect(server_addr).await {
             Ok(stream) => stream,
             Err(error) => {
                 error!("Failed to connect to server: {server_addr}. Error: 
{error}");
-                return self.handle_connection_error(retry_count).await;
+                return self
+                    .handle_connection_error(retry_count, single_attempt)
+                    .await;
             }
         };
         let tls_config = self.build_tls_config()?;
@@ -570,7 +674,9 @@ impl WebSocketClient {
             Ok(result) => result,
             Err(error) => {
                 error!("WebSocket TLS handshake failed: {}", error);
-                return self.handle_connection_error(retry_count).await;
+                return self
+                    .handle_connection_error(retry_count, single_attempt)
+                    .await;
             }
         };
 
@@ -629,7 +735,16 @@ impl WebSocketClient {
         Ok(config)
     }
 
-    async fn handle_connection_error<T>(&self, retry_count: &mut u32) -> 
Result<T, IggyError> {
+    async fn handle_connection_error<T>(
+        &self,
+        retry_count: &mut u32,
+        single_attempt: bool,
+    ) -> Result<T, IggyError> {
+        if single_attempt {
+            self.set_state(ClientState::Disconnected).await;
+            self.publish_event(DiagnosticEvent::Disconnected).await;
+            return Err(IggyError::CannotEstablishConnection);
+        }
         if !self.config.reconnection.enabled {
             warn!("Automatic reconnection is disabled.");
             return Err(IggyError::CannotEstablishConnection);
@@ -745,7 +860,6 @@ impl WebSocketClient {
         self.connected_at.lock().await.take();
         self.disconnect().await?;
         *self.current_server_address.lock().await = next;
-        self.settle_off_leader_once.store(true, Ordering::SeqCst);
         Ok(())
     }
 
@@ -970,6 +1084,29 @@ mod tests {
     use super::*;
     use std::str::FromStr;
 
+    #[tokio::test]
+    async fn a_roster_hop_does_not_enter_the_reconnect_ladder() {
+        let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
+            .await
+            .expect("reserve TCP address");
+        let server_address = listener.local_addr().unwrap().to_string();
+        drop(listener);
+        let client = WebSocketClient::create(Arc::new(WebSocketClientConfig {
+            server_address,
+            ..WebSocketClientConfig::default()
+        }))
+        .expect("create WebSocket client");
+
+        let result = tokio::time::timeout(
+            std::time::Duration::from_secs(1),
+            client.connect_off_leader(),
+        )
+        .await
+        .expect("one WebSocket dial must not enter unlimited reconnect");
+        assert!(matches!(result, Err(IggyError::CannotEstablishConnection)));
+        assert_eq!(client.get_state().await, ClientState::Disconnected);
+    }
+
     #[test]
     fn should_be_created_with_default_config() {
         let client = WebSocketClient::default();
diff --git a/core/server/src/partition_reconciler.rs 
b/core/server/src/partition_reconciler.rs
index 699e9c94d..39df9c2ee 100644
--- a/core/server/src/partition_reconciler.rs
+++ b/core/server/src/partition_reconciler.rs
@@ -2545,7 +2545,9 @@ mod tests {
             .expect("partition is materialised")
             .repair = Some(RepairSession {
             nonce: NONCE,
-            to_op: 5,
+            view: 0,
+            commit_to_op: 5,
+            fetch_to_op: 5,
             floor: None,
             peer: 1,
             first_batch_offset: None,
diff --git a/core/shard/src/lib.rs b/core/shard/src/lib.rs
index 22fa69015..62ecd3d30 100644
--- a/core/shard/src/lib.rs
+++ b/core/shard/src/lib.rs
@@ -3976,6 +3976,10 @@ where
                 .handle_start_view(PlaneKind::Partitions, &header, 
suffix_body);
         let adopted = !actions.is_empty();
         if adopted {
+            // Any stream armed before this adoption belongs to the superseded
+            // view. Repair bodies carry no nonce, so drop the receiving 
session
+            // before reconciling or arming the new view's canonical range.
+            partition.repair = None;
             // Ahead of the local dispatch, which rebuilds the pipeline out of 
the
             // journal this rewrites. Same position as the metadata arm's 
twin, and
             // like it, pending-less adoptions (empty StartView suffix) still 
sweep
@@ -4725,6 +4729,10 @@ where
         if header.nonce != session.nonce {
             return;
         }
+        if !partition.consensus().is_normal() || partition.consensus().view() 
!= session.view {
+            partition.repair = None;
+            return;
+        }
         // Receiver half of the serve-side purge gate: while a committed purge
         // is not yet locally applied, this replica's 
`recovered_durable_offset`
         // still describes the PRE-purge segments, so a floor from a peer that
@@ -4834,7 +4842,7 @@ where
                 } else {
                     let commit_min = partition.consensus().commit_min();
                     let next = partition.repair.as_ref().and_then(|live| {
-                        (commit_min > before).then_some((live.peer, 
live.nonce, live.to_op))
+                        (commit_min > before).then_some((live.peer, 
live.nonce, live.fetch_to_op))
                     });
                     let cluster = partition.consensus().cluster();
                     let self_id = partition.consensus().replica();
@@ -6515,19 +6523,31 @@ where
                     continue;
                 };
                 let consensus_normal = partition.consensus().is_normal();
+                let consensus_view = partition.consensus().view();
                 let commit_min = partition.consensus().commit_min();
                 let cluster = partition.consensus().cluster();
                 let self_id = partition.consensus().replica();
-                if partition
-                    .repair
-                    .is_some_and(|session| commit_min >= session.to_op)
-                {
+                let repair_finished = partition.repair.is_some_and(|session| {
+                    if !consensus_normal || consensus_view != session.view {
+                        return true;
+                    }
+                    let fetch_complete = session.fetch_to_op <= 
session.commit_to_op
+                        || partition
+                            .log
+                            .journal()
+                            .inner
+                            .repaired_window_shape(session.commit_to_op, 
session.fetch_to_op)
+                            .complete;
+                    commit_min >= session.commit_to_op && fetch_complete
+                });
+                if repair_finished {
                     partition.repair = None;
                     tracing::info!(
                         shard = self.id,
                         namespace_raw = namespace.inner(),
                         commit_min,
-                        "partition journal repair completed after the commit 
frontier advanced"
+                        consensus_view,
+                        "partition journal repair completed or was superseded"
                     );
                     continue;
                 }
@@ -6543,8 +6563,8 @@ where
                     Some((
                         session.peer,
                         session.nonce,
-                        commit_min + 1,
-                        session.to_op,
+                        commit_min.saturating_add(1),
+                        session.fetch_to_op,
                         cluster,
                         self_id,
                     ))
@@ -7468,27 +7488,36 @@ where
         // commit-bounded window nothing ever delivers those bodies, the
         // primary cannot gather quorum for the suffix, and the group wedges
         // one op below its head with the client write never confirmed.
+        let commit_to_op = consensus.commit_max();
+        let commit_lag = consensus.commit_min() < commit_to_op;
         let head = consensus.sequencer().current_sequence();
-        let commit_lag = consensus.commit_min() < consensus.commit_max();
-        let missing_suffix = 
consensus.commit_max().checked_add(1).is_some_and(|first| {
-            (first..=head).any(|op| 
partition.log.journal().inner.header_by_op(op).is_none())
-        });
+        if !commit_lag && head <= commit_to_op {
+            return;
+        }
+        let canonical_suffix = consensus
+            .with_pending_view_log(|pending| pending_covers_suffix(pending, 
commit_to_op, head))
+            .unwrap_or(false);
+        let missing_suffix = canonical_suffix
+            && !partition
+                .log
+                .journal()
+                .inner
+                .repaired_window_shape(commit_to_op, head)
+                .complete;
         if !commit_lag && !missing_suffix {
             return;
         }
         let nonce = iggy_common::random_id::get_uuid();
         let from_op = consensus.commit_min() + 1;
-        let to_op = if missing_suffix {
-            head
-        } else {
-            consensus.commit_max()
-        };
+        let fetch_to_op = if missing_suffix { head } else { commit_to_op };
         let cluster = consensus.cluster();
         let self_id = consensus.replica();
         let namespace = consensus.group();
         partition.repair = Some(partitions::RepairSession {
             nonce,
-            to_op,
+            view: consensus.view(),
+            commit_to_op,
+            fetch_to_op,
             floor: None,
             peer,
             first_batch_offset: None,
@@ -7498,11 +7527,20 @@ where
             shard = self.id,
             namespace_raw = namespace,
             from_op,
-            to_op,
+            commit_to_op,
+            fetch_to_op,
             "partition behind the group frontier; requesting repair"
         );
-        self.send_request_prepares(cluster, self_id, peer, nonce, from_op, 
to_op, namespace)
-            .await;
+        self.send_request_prepares(
+            cluster,
+            self_id,
+            peer,
+            nonce,
+            from_op,
+            fetch_to_op,
+            namespace,
+        )
+        .await;
     }
 
     /// Receiver side of a partition descriptor: accept the manifest, adopt
@@ -8754,6 +8792,27 @@ fn repair_serve_ceiling(requested_to_op: u64, 
commit_max: u64, head: u64) -> u64
     requested_to_op.min(commit_max.max(head))
 }
 
+/// Whether the parked `StartView` log names every op in the uncommitted
+/// suffix `(commit_max, head]`, in descending order. Only this canonical list
+/// makes fetching bodies above the commit point safe.
+fn pending_covers_suffix(pending: &MergedLog, commit_max: u64, head: u64) -> 
bool {
+    if head <= commit_max || pending.commit_max != commit_max || 
pending.op_head != head {
+        return false;
+    }
+    let mut expected = head;
+    for header in pending
+        .headers
+        .iter()
+        .filter(|header| header.op > commit_max)
+    {
+        if header.op != expected {
+            return false;
+        }
+        expected -= 1;
+    }
+    expected == commit_max
+}
+
 /// Read this replica's uncommitted suffix out of the metadata journal, for the
 /// window `commit..=op`.
 ///
@@ -9525,7 +9584,7 @@ mod repair_scope_tests {
 
     use iggy_binary_protocol::{Command, PrepareHeader};
 
-    use super::{MergedLog, repair_op_in_scope, repair_serve_ceiling};
+    use super::{MergedLog, pending_covers_suffix, repair_op_in_scope, 
repair_serve_ceiling};
 
     fn header(op: u64) -> PrepareHeader {
         PrepareHeader {
@@ -9596,6 +9655,20 @@ mod repair_scope_tests {
         // `commit_max` above the local head still counts: heartbeats outrun 
prepares.
         assert_eq!(repair_serve_ceiling(u64::MAX, 120, 90), 120);
     }
+
+    #[test]
+    fn 
given_a_parked_view_when_fetching_above_commit_should_require_dense_canonical_suffix()
 {
+        let pending = parked();
+        assert!(pending_covers_suffix(&pending, 98, 100));
+
+        let mut missing = pending.clone();
+        missing.headers.retain(|header| header.op != 99);
+        assert!(!pending_covers_suffix(&missing, 98, 100));
+
+        let mut wrong_frontier = pending;
+        wrong_frontier.commit_max = 97;
+        assert!(!pending_covers_suffix(&wrong_frontier, 98, 100));
+    }
 }
 
 #[cfg(test)]
diff --git 
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java
 
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java
index b4ea3b57b..abc7a3e6b 100644
--- 
a/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java
+++ 
b/foreign/java/java-sdk/src/test/java/org/apache/iggy/client/async/tcp/AsyncIggyTcpClientTransientFailoverTest.java
@@ -457,23 +457,38 @@ class AsyncIggyTcpClientTransientFailoverTest {
         return serve(server, 1, handler);
     }
 
+    /**
+     * Runs blocking socket I/O on one dedicated daemon thread per mock node.
+     * The common fork-join pool has only cores minus one workers on small CI
+     * runners, so three blocking nodes can starve the client continuations the
+     * test is waiting for when the full suite runs concurrently.
+     */
     private static CompletableFuture<Void> serve(ServerSocket server, int 
connectionCount, RequestHandler handler) {
-        return CompletableFuture.runAsync(() -> {
-            try {
-                for (int connection = 0; connection < connectionCount; 
connection++) {
-                    try (Socket socket = server.accept()) {
-                        InputStream input = socket.getInputStream();
-                        OutputStream output = socket.getOutputStream();
-                        Request request;
-                        while ((request = readRequest(input)) != null) {
-                            writeResponse(output, request, 
handler.handle(request));
+        CompletableFuture<Void> serving = new CompletableFuture<>();
+        Thread serverThread = new Thread(
+                () -> {
+                    try {
+                        for (int connection = 0; connection < connectionCount; 
connection++) {
+                            try (Socket socket = server.accept()) {
+                                InputStream input = socket.getInputStream();
+                                OutputStream output = socket.getOutputStream();
+                                Request request;
+                                while ((request = readRequest(input)) != null) 
{
+                                    writeResponse(output, request, 
handler.handle(request));
+                                }
+                            }
                         }
+                        serving.complete(null);
+                    } catch (IOException error) {
+                        serving.completeExceptionally(new 
IllegalStateException("Mock VSR server failed", error));
+                    } catch (RuntimeException error) {
+                        serving.completeExceptionally(error);
                     }
-                }
-            } catch (IOException error) {
-                throw new IllegalStateException("Mock VSR server failed", 
error);
-            }
-        });
+                },
+                "transient-failover-server-" + server.getLocalPort());
+        serverThread.setDaemon(true);
+        serverThread.start();
+        return serving;
     }
 
     private static Request readRequest(InputStream input) throws IOException {

Reply via email to