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