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


##########
core/sdk/src/quic/quic_client.rs:
##########
@@ -129,7 +147,136 @@ 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 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
+        // 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 =
+                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 = 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
+            // 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 = 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
+                };
+                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;
+                }
+                result = match tokio::time::timeout_at(
+                    overall_deadline,
+                    self.send_raw(code, payload.clone()),
+                )
+                .await
+                {
+                    Ok(result) => result,
+                    Err(_) => Err(IggyError::TransientNotAccepted),

Review Comment:
   This maps an elapsed outer timer to `TransientNotAccepted`, which is the one 
code the PR defines as proof the request never entered a pipeline and is 
therefore safe to re-issue anywhere. At this point the frame has been written 
and the reply is unread, so the outcome is unknown.
   
   `overall_deadline` is pinned at fn entry (`:150`) while `send_raw` mints a 
fresh `RESPONSE_READ_TIMEOUT` per call (`:921`). On the ordinary refusal path 
the inner 2s window returns first, so this is benign. On a **stalled reply** 
the inner budget is a fresh 30s while the outer has `30 - elapsed`, so from hop 
2 onward the outer binds and fires exactly when the request is in flight, 
possibly admitted, replicated and committed.
   
   The fabricated code is then consumed by the SDK's own walk loop, not just 
the app: it satisfies `while matches!(result, Err(TransientNotAccepted))` at 
`:178`, so the walk hops on, `connect_off_leader` mints a new VSR client id, 
and `:271` re-issues the same payload under a session the server's dedup fence 
cannot match. The app-boundary contract breach via error 58 is second-order.
   
   Pre-fix this was a bare `send_raw(...).await` that surfaced `Disconnected` 
on expiry, which the comment at `:915-919` documents as deliberate ("later 
commits would commit it twice"). The outer timeout added for the walk-deadline 
fix silently overrides it.
   
   `TransientNotCommitted` or `Disconnected` is the correct mapping for an 
elapsed outer timer, and either also falsifies the loop condition and stops the 
hop chain, which is the right answer for an unknown outcome. Same shape at 
`:180`. The sibling `Err(_)` arms at `:168`, `:200`, `:232`, `:258` are fine 
and should be left alone — nothing is on the wire there.



##########
core/sdk/src/websocket/websocket_client.rs:
##########
@@ -122,7 +140,136 @@ 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 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
+        // 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 =
+                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 = 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
+            // 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 = 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
+                };
+                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;
+                }
+                result = match tokio::time::timeout_at(
+                    overall_deadline,
+                    self.send_raw(code, payload.clone()),
+                )
+                .await
+                {
+                    Ok(result) => result,
+                    Err(_) => Err(IggyError::TransientNotAccepted),

Review Comment:
   Same defect as `quic_client.rs:275` (and `:173` here): an elapsed outer 
timer is mapped to `TransientNotAccepted`, the code that licenses re-issue 
anywhere, while the request is on the wire with the reply unread. It feeds the 
walk condition at `:179`, so the SDK itself re-issues the payload on the next 
hop under a new session the dedup fence cannot match.
   
   WebSocket is the worse of the two: `send_raw`'s body is a `tokio::spawn`, so 
cancelling the outer timeout drops the `JoinHandle` and **detaches** the task, 
which runs the exchange to completion. The write is guaranteed delivered and 
only the reply is discarded, so the replay double-applies rather than merely 
risking it.
   
   Return `TransientNotCommitted` or `Disconnected` on an elapsed outer timer. 
The `Err(_)` arms wrapping awaits that issue no request are correct as they 
stand.



##########
core/partitions/src/iggy_partition.rs:
##########
@@ -4616,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

Review Comment:
   `fetch_complete` is proved from resident journal headers, but 
`commit_journal` at `:4643` has just evicted them, so a repair that fully 
succeeds reports itself incomplete.
   
   `commit_journal` -> `committed_prefix(commit_max)` bounds the eviction by 
the **live** `commit_max`, not the session's snapshotted `commit_to_op`. 
Delivering the suffix bodies is what lets the group commit past `commit_to_op`, 
so the flush evicts into `(commit_to_op, fetch_to_op]`, and 
`repaired_window_shape` reads only `headers` and never the evicted ring 
(`journal.rs:770-772`). `complete` is false, so `fetch_complete` is false and 
`done` is false at `:4658` even though `commit_min >= commit_to_op`.
   
   The asymmetry is what makes it self-defeating: if the group does *not* 
commit past `commit_to_op` the window sits above `commit_max`, survives 
eviction, and the session closes correctly. It is the successful commit that 
destroys the evidence.
   
   Nothing recovers it on a healthy group. The tick twin at 
`core/shard/src/lib.rs:6534-6541` evaluates the identical expression and is 
blocked identically. The stall retry goes silent because `from_op = commit_min 
+ 1` then exceeds `fetch_to_op` and trips the `from_op <= to_op` guard at 
`shard:6574`, so the wedge is also invisible in the logs. 
`partition.repair.is_some()` (`shard:7479`) then blocks every later re-arm for 
that partition, and every remaining clear path is view-change-gated, so it pins 
until the next election.
   
   The new integration config (`messages_required_to_save: Some(1)`, 
`enforce_fsync: true`) guarantees the flush that triggers this.
   
   Suggest deciding completion from a monotone fact the flush cannot erase: 
record suffix arrivals on the session as they are ingested (a high-water op or 
counter), or use `commit_min >= fetch_to_op`, rather than re-deriving presence 
from a vector `commit_journal` just truncated.



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