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


##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -248,11 +263,24 @@ impl BinaryTransport for TcpClient {
         // Login and register are the exception: the server stays deliberately
         // silent on a transient register failure and relies on the client
         // replaying, so that replay is the protocol rather than a retry.
-        let replay_after_reconnect = replay_is_safe(code, &error);
+        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 _routing_guard = if nested_connect {
+            None
+        } else {
+            Some(self.routing_lock.lock().await)

Review Comment:
   `_routing_guard` is a named binding, so it lives to the end of the function, 
and both replay sites below (`:280` and `:316`) call `self.send_raw(...)` while 
it is still held. `send_raw` re-acquires the same lock at `:1265` when the 
server answers `TransientNotAccepted`. Same task, non-reentrant `tokio::Mutex`, 
no timeout on `lock()`.
   
   Reachable on this PR's own scenario: `SEND_MESSAGES` fails with 
`NotConnected`, `replay_after_session_reset_is_safe` returns true, the 
reconnect lands on the restarted node which is now a partition backup, the 
backup answers `TransientNotAccepted`, and the second acquire never returns. No 
error, no log, and the request timeout does not apply. The guard is never 
released, so every later refused request on the same client blocks behind it.
   
   QUIC and WebSocket are immune: their guard is scoped to the walk block, and 
their `send_raw` never takes it.
   
   Dropping the guard before the replay, or lifting the walk out of `send_raw` 
the way QUIC/WS do, would close this.



##########
core/sdk/src/tcp/tcp_client.rs:
##########
@@ -248,11 +263,24 @@ impl BinaryTransport for TcpClient {
         // Login and register are the exception: the server stays deliberately
         // silent on a transient register failure and relies on the client
         // replaying, so that replay is the protocol rather than a retry.
-        let replay_after_reconnect = replay_is_safe(code, &error);
+        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();

Review Comment:
   `nested_connect` is decided from a global flag with no task identity, so it 
cannot distinguish the owner's own sign-in from an unrelated caller.
   
   Any task calling the public `login_user()` while another task owns the 
connect takes this branch: it skips `routing_lock`, skips the coordinator, runs 
`disconnect_transport()` at `:282` unconditionally (tearing down the owner's 
in-flight stream), then calls `connect_inner()` at `:298` concurrently with the 
owner's. Two connect sequences race `stream`, `state` and 
`current_server_address`, which is what `ConnectCoordinator` was added in this 
PR to prevent.
   
   Same shape at `quic_client.rs:251` and `websocket_client.rs:241`.
   
   An explicit ownership token threaded down the connect chain would also fix 
the `settle_off_leader_once` leak: waiters in `ConnectCoordinator::run` return 
the owner's result without running the closure, so `connect_inner`'s 
`swap(false)` never runs and the flag leaks into a later unrelated reconnect.



##########
core/shard/src/lib.rs:
##########
@@ -7433,16 +7456,33 @@ where
         B: MessageBus,
     {
         let consensus = partition.consensus();
-        if !consensus.is_normal()
-            || consensus.is_transferring()
-            || consensus.commit_min() >= consensus.commit_max()
-            || partition.repair.is_some()
-        {
+        if !consensus.is_normal() || consensus.is_transferring() || 
partition.repair.is_some() {
+            return;
+        }
+        // The window ends at the group head when suffix bodies are missing,
+        // not at the commit point. A backup that adopted a StartView holds
+        // suffix HEADERS above `commit_max` whose bodies it may never have
+        // received: its ack for them is withheld until the body is journaled,
+        // and the primary's retransmit is dropped by the backup gap check
+        // because adoption already advanced the sequencer to the head. With a
+        // 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 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 && !missing_suffix {
             return;
         }
         let nonce = iggy_common::random_id::get_uuid();
         let from_op = consensus.commit_min() + 1;
-        let to_op = consensus.commit_max();
+        let to_op = if missing_suffix {

Review Comment:
   Widening `to_op` to the sequencer head means repair now carries ops above 
`commit_max`, but the partition ingest arm is fence-free precisely because that 
could not happen before.
   
   `apply_repaired_prepare` (`core/partitions/src/iggy_partition.rs:4464`) 
gates only on `commit_min < op <= session.to_op` plus 
`verify_prepare_integrity`, which is frame self-consistency. No view check, no 
merged-log comparison, no peer check, and `RepairPrepareHeader` carries no 
nonce so the ingest could not filter by session even if it wanted to. The 
metadata arm has all of these (`repair_op_in_scope` plus the checksum 
`disagrees` test); the comment at `:4589-4591` states the gap outright.
   
   The session also survives a view change (`maybe_request_partition_repair` 
no-ops while `repair.is_some()`), so a late or duplicated `RepairPrepare` from 
the previous window can refill a truncated op with a non-canonical entry. The 
primary's retransmit is then rejected by `journaled_prepare_matches_retransmit` 
and never acked, and once `commit_max` passes that op the divergent bytes are 
flushed to segment storage. `reconcile_partition_view_divergence` runs only at 
StartView adoption, never at repair ingest.
   
   Two further consequences of the same line. `repaired_window_is_complete` can 
never turn true against a peer with an empty memory-only journal, so the 
`FloorRefused` to state-transfer escape stops firing exactly when a fast rejoin 
needs it. And `to_op = head` is view-scoped rather than monotone, so after a 
view change discards the suffix `commit_min >= session.to_op` is unreachable 
and the new clear at `:6521` never fires.
   
   Bounding the completeness and clear verdict at `commit_max` while widening 
only the fetch would address all three.



##########
core/shard/src/lib.rs:
##########
@@ -7433,16 +7456,33 @@ where
         B: MessageBus,
     {
         let consensus = partition.consensus();
-        if !consensus.is_normal()
-            || consensus.is_transferring()
-            || consensus.commit_min() >= consensus.commit_max()
-            || partition.repair.is_some()
-        {
+        if !consensus.is_normal() || consensus.is_transferring() || 
partition.repair.is_some() {
+            return;
+        }
+        // The window ends at the group head when suffix bodies are missing,
+        // not at the commit point. A backup that adopted a StartView holds
+        // suffix HEADERS above `commit_max` whose bodies it may never have
+        // received: its ack for them is withheld until the body is journaled,
+        // and the primary's retransmit is dropped by the backup gap check
+        // because adoption already advanced the sequencer to the head. With a
+        // 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 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| {

Review Comment:
   `header_by_op` is an unindexed `headers.iter().find` 
(`core/partitions/src/journal.rs:760`), so this is O(window x resident 
headers), synchronous, on the shard pump. `journal.rs:768-775` documents this 
exact op-by-op shape as what can "turn one rejoin into an election storm", 
which is why `repaired_window_shape` exists.
   
   The cheap `commit_min() >= commit_max()` early-out that used to guard this 
function now sits below the scan at `:7476`, so a caught-up backup that is 
merely pipelining pays the full walk where it previously returned on three 
boolean loads. Reached from `on_commit` / `CommitOutcome::Advanced`, so roughly 
once per `commit_broadcast_interval` per partition group per backup.
   
   `!partition.log.journal().inner.repaired_window_shape(commit_max, 
head).complete` (`journal.rs:779`) is a single pass and is semantically 
identical, empty window included.



##########
core/sdk/src/quic/quic_client.rs:
##########
@@ -129,7 +152,60 @@ impl BinaryTransport for QuicClient {
     }
 
     async fn send_raw_with_response(&self, code: u32, payload: Bytes) -> 
Result<Bytes, IggyError> {
-        let result = self.send_raw(code, payload.clone()).await;
+        let mut result = self.send_raw(code, payload.clone()).await;
+
+        // A persistent not-admitted refusal is a verdict about who leads the
+        // TARGET group, which the metadata leader check alone cannot repair:
+        // metadata and partition consensus groups elect independently. Recheck
+        // the leader once, then walk the roster, one visit per endpoint. Only
+        // recoverable with a session to re-establish, hence the auto-login
+        // gate, and a login/register replay stays on its own connection.
+        if matches!(result, Err(IggyError::TransientNotAccepted))
+            && !is_login_register_code(code)
+            && code != GET_CLUSTER_METADATA_CODE
+            && self.config.reconnection.enabled
+            && !matches!(self.config.auto_login, AutoLogin::Disabled)
+        {
+            let _routing_guard = self.routing_lock.lock().await;
+            // A concurrent refused request may have completed the movement
+            // while this request waited for the gate.
+            result = self.send_raw(code, payload.clone()).await;
+            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
+            // partition replica refused it.
+            let mut checked_metadata_leader = false;
+            while matches!(result, Err(IggyError::TransientNotAccepted)) {
+                let current = self.current_server_address.lock().await.clone();
+                let redirected = if checked_metadata_leader {
+                    false
+                } else {
+                    checked_metadata_leader = true;
+                    let redirected = 
matches!(self.handle_leader_redirection().await, Ok(true));
+                    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 target = 
self.current_server_address.lock().await.clone();
+                    if let Some(walk) = roster_walk.as_mut() {
+                        walk.record_attempt(&target);
+                    }
+                }
+                self.connect().await?;

Review Comment:
   `self.connect().await?` propagates `CannotEstablishConnection` out of 
`send_raw_with_response`, so the first undialable roster hop aborts the whole 
walk. `settle_on_endpoint(next).await?` a few lines above has the same escape.
   
   TCP handles it instead: `tcp_client.rs:1310` swallows 
`Err(IggyError::CannotEstablishConnection)`, pulls the next `RosterWalk::next`, 
and returns `TransientNotAccepted` only once the walk is exhausted.
   
   This is guaranteed to bite in the scenario the PR targets, since a stopped 
or restarting primary is exactly a roster entry that will not dial. QUIC and 
WebSocket then never reach the partition primary, and the caller sees a connect 
error instead of the typed refusal.
   
   Related: this loop has no overall deadline, where TCP bounds its walk by 
`overall_deadline`. Each hop mints a fresh `RESPONSE_READ_TIMEOUT` plus a full 
connect retry ladder.



##########
core/sdk/src/websocket/websocket_client.rs:
##########
@@ -122,7 +145,60 @@ impl BinaryTransport for WebSocketClient {
     }
 
     async fn send_raw_with_response(&self, code: u32, payload: Bytes) -> 
Result<Bytes, IggyError> {
-        let result = self.send_raw(code, payload.clone()).await;
+        let mut result = self.send_raw(code, payload.clone()).await;
+
+        // A persistent not-admitted refusal is a verdict about who leads the
+        // TARGET group, which the metadata leader check alone cannot repair:
+        // metadata and partition consensus groups elect independently. Recheck
+        // the leader once, then walk the roster, one visit per endpoint. Only
+        // recoverable with a session to re-establish, hence the auto-login
+        // gate, and a login/register replay stays on its own connection.
+        if matches!(result, Err(IggyError::TransientNotAccepted))
+            && !is_login_register_code(code)
+            && code != GET_CLUSTER_METADATA_CODE
+            && self.config.reconnection.enabled
+            && !matches!(self.config.auto_login, AutoLogin::Disabled)
+        {
+            let _routing_guard = self.routing_lock.lock().await;
+            // A concurrent refused request may have completed the movement
+            // while this request waited for the gate.
+            result = self.send_raw(code, payload.clone()).await;
+            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
+            // partition replica refused it.
+            let mut checked_metadata_leader = false;
+            while matches!(result, Err(IggyError::TransientNotAccepted)) {
+                let current = self.current_server_address.lock().await.clone();
+                let redirected = if checked_metadata_leader {
+                    false
+                } else {
+                    checked_metadata_leader = true;
+                    let redirected = 
matches!(self.handle_leader_redirection().await, Ok(true));
+                    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 target = 
self.current_server_address.lock().await.clone();
+                    if let Some(walk) = roster_walk.as_mut() {
+                        walk.record_attempt(&target);
+                    }
+                }
+                self.connect().await?;

Review Comment:
   Same as `quic_client.rs:200`. `self.connect().await?`, and 
`settle_on_endpoint(next).await?` above it, abort the roster walk on the first 
undialable hop, where TCP skips and continues at `tcp_client.rs:1310`. A 
stopped primary in the roster is this PR's own scenario, so the walk never 
reaches the partition primary.
   
   Also unbounded in wall time: no caller deadline is threaded through the loop.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to