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(¤t, &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(¤t, &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]