From 735a3f4118f252feff733f284b127d527043bbfc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Tue, 22 Sep 2026 01:15:00 -0300 Subject: [PATCH 1/2] redial disconnected cluster peers and remove dead ones --- .../src/cluster/quic_transport.rs | 81 +++++++++- .../src/cluster_agent/event_loop.rs | 13 ++ .../tests/quic_transport_framing.rs | 99 +++++++++++- specs/ClusterRedial.cfg | 15 ++ specs/ClusterRedial.tla | 152 ++++++++++++++++++ specs/ClusterRedial_noredial.cfg | 15 ++ specs/ClusterRedial_reusable.cfg | 13 ++ specs/ClusterRedial_unguarded.cfg | 13 ++ 8 files changed, 392 insertions(+), 9 deletions(-) create mode 100644 specs/ClusterRedial.cfg create mode 100644 specs/ClusterRedial.tla create mode 100644 specs/ClusterRedial_noredial.cfg create mode 100644 specs/ClusterRedial_reusable.cfg create mode 100644 specs/ClusterRedial_unguarded.cfg diff --git a/crates/mqdb-cluster/src/cluster/quic_transport.rs b/crates/mqdb-cluster/src/cluster/quic_transport.rs index 8fe2c96..c9c19df 100644 --- a/crates/mqdb-cluster/src/cluster/quic_transport.rs +++ b/crates/mqdb-cluster/src/cluster/quic_transport.rs @@ -9,7 +9,7 @@ use std::collections::{HashMap, VecDeque}; use std::io::BufReader; use std::net::SocketAddr; use std::path::Path; -use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use tokio::sync::{Notify, RwLock}; use tracing::{debug, error, info, trace, warn}; @@ -22,6 +22,7 @@ struct PeerConnection { _connection: Connection, control_tx: flume::Sender>, bulk_tx: flume::Sender>, + generation: u64, } const MAX_MESSAGE_SIZE: usize = 10 * 1024 * 1024; @@ -38,6 +39,8 @@ pub struct QuicDirectTransport { node_id: NodeId, endpoint: Arc>>, peers: Arc>>, + peer_addrs: Arc>>, + generation: Arc, inbox_tx: flume::Sender, inbox_rx: flume::Receiver, requeue_buffer: Arc>>, @@ -67,6 +70,8 @@ impl Clone for QuicDirectTransport { node_id: self.node_id, endpoint: self.endpoint.clone(), peers: self.peers.clone(), + peer_addrs: self.peer_addrs.clone(), + generation: self.generation.clone(), inbox_tx: self.inbox_tx.clone(), inbox_rx: self.inbox_rx.clone(), requeue_buffer: self.requeue_buffer.clone(), @@ -94,6 +99,8 @@ impl QuicDirectTransport { node_id, endpoint: Arc::new(RwLock::new(None)), peers: Arc::new(RwLock::new(HashMap::new())), + peer_addrs: Arc::new(RwLock::new(HashMap::new())), + generation: Arc::new(AtomicU64::new(0)), inbox_tx, inbox_rx, requeue_buffer: Arc::new(Mutex::new(VecDeque::new())), @@ -175,9 +182,10 @@ impl QuicDirectTransport { let notify = self.message_notify.clone(); let local_node = self.node_id; let peers = self.peers.clone(); + let generation = self.generation.clone(); tokio::spawn(async move { - acceptor_task(endpoint, inbox_tx, notify, local_node, peers).await; + acceptor_task(endpoint, inbox_tx, notify, local_node, peers, generation).await; }); Ok(()) @@ -192,6 +200,8 @@ impl QuicDirectTransport { peer_id: NodeId, peer_addr: SocketAddr, ) -> Result<(), TransportError> { + self.peer_addrs.write().await.insert(peer_id, peer_addr); + let endpoint_guard = self.endpoint.read().await; let endpoint = endpoint_guard .as_ref() @@ -250,14 +260,23 @@ impl QuicDirectTransport { let (control_tx, control_rx) = flume::bounded(PEER_CONTROL_QUEUE_CAPACITY); let (bulk_tx, bulk_rx) = flume::bounded(PEER_BULK_QUEUE_CAPACITY); + let generation = self.generation.fetch_add(1, Ordering::SeqCst) + 1; let peer_conn = PeerConnection { _connection: connection.clone(), control_tx, bulk_tx, + generation, }; self.peers.write().await.insert(peer_id, peer_conn); - tokio::spawn(peer_writer_task(send_stream, control_rx, bulk_rx, peer_id)); + tokio::spawn(peer_writer_task( + send_stream, + control_rx, + bulk_rx, + peer_id, + self.peers.clone(), + generation, + )); let inbox_tx = self.inbox_tx.clone(); let notify = self.message_notify.clone(); @@ -270,6 +289,30 @@ impl QuicDirectTransport { Ok(()) } + /// Re-dial every configured peer that is not currently alive. + /// + /// Driven by the heartbeat liveness view rather than the transport peer + /// map, so a stale entry that has not yet been removed by a send failure + /// does not stop reconnection. `connect_to_peer` replaces any stale entry. + pub async fn redial_unlinked(&self, alive: &[NodeId]) { + let targets: Vec<(NodeId, SocketAddr)> = { + let addrs = self.peer_addrs.read().await; + addrs + .iter() + .filter(|&(node, _)| *node != self.node_id && !alive.contains(node)) + .map(|(node, addr)| (*node, *addr)) + .collect() + }; + + for (node, addr) in targets { + if let Err(e) = self.connect_to_peer(node, addr).await { + debug!(peer = node.get(), addr = %addr, error = %e, "redial to peer failed, will retry"); + } else { + info!(peer = node.get(), addr = %addr, "redialled peer"); + } + } + } + #[must_use] pub fn inbox_rx(&self) -> flume::Receiver { self.inbox_rx.clone() @@ -459,6 +502,7 @@ async fn acceptor_task( notify: Arc, local_node: NodeId, peers: Arc>>, + generation: Arc, ) { info!(node = local_node.get(), "QUIC acceptor task started"); @@ -474,10 +518,13 @@ async fn acceptor_task( let inbox_tx = inbox_tx.clone(); let notify = notify.clone(); let peers = peers.clone(); + let generation = generation.clone(); tokio::spawn(async move { - if let Err(e) = - handle_incoming_connection(connection, inbox_tx, notify, local_node, peers).await + if let Err(e) = handle_incoming_connection( + connection, inbox_tx, notify, local_node, peers, generation, + ) + .await { debug!(error = %e, "incoming connection handler failed"); } @@ -491,6 +538,7 @@ async fn handle_incoming_connection( notify: Arc, local_node: NodeId, peers: Arc>>, + generation: Arc, ) -> Result<(), TransportError> { let (send_stream, mut recv_stream) = connection .accept_bi() @@ -511,11 +559,13 @@ async fn handle_incoming_connection( let (control_tx, control_rx) = flume::bounded(PEER_CONTROL_QUEUE_CAPACITY); let (bulk_tx, bulk_rx) = flume::bounded(PEER_BULK_QUEUE_CAPACITY); + let peer_generation = generation.fetch_add(1, Ordering::SeqCst) + 1; { let peer_conn = PeerConnection { _connection: connection.clone(), control_tx, bulk_tx, + generation: peer_generation, }; peers.write().await.insert(peer_node, peer_conn); } @@ -525,6 +575,8 @@ async fn handle_incoming_connection( control_rx, bulk_rx, peer_node, + peers.clone(), + peer_generation, )); receiver_task(recv_stream, peer_node, inbox_tx, notify, local_node).await; @@ -536,6 +588,8 @@ async fn peer_writer_task( control_rx: flume::Receiver>, bulk_rx: flume::Receiver>, peer_node: NodeId, + peers: Arc>>, + generation: u64, ) { trace!(peer = peer_node.get(), "peer writer task started"); @@ -553,7 +607,8 @@ async fn peer_writer_task( }; if let Err(e) = send_stream.write_all(&frame).await { - warn!(peer = peer_node.get(), error = %e, "peer writer failed, tearing down stream"); + warn!(peer = peer_node.get(), error = %e, "peer writer failed, removing dead peer"); + remove_peer_generation(&peers, peer_node, generation).await; break; } } @@ -561,6 +616,20 @@ async fn peer_writer_task( debug!(peer = peer_node.get(), "peer writer task ended"); } +async fn remove_peer_generation( + peers: &Arc>>, + peer_node: NodeId, + generation: u64, +) { + let mut map = peers.write().await; + if map + .get(&peer_node) + .is_some_and(|peer| peer.generation == generation) + { + map.remove(&peer_node); + } +} + async fn receiver_task( mut recv_stream: RecvStream, peer_node: NodeId, diff --git a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs index 373ce2e..d9bf1cb 100644 --- a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs +++ b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs @@ -148,6 +148,7 @@ impl ClusteredAgent { } _ = mesh_check_interval.tick() => { self.warn_on_unlinked_nodes().await; + self.redial_disconnected_peers().await; } _ = retained_sync_cleanup_interval.tick() => { Self::handle_retained_sync_cleanup(&synced_retained_topics).await; @@ -557,6 +558,18 @@ impl ClusteredAgent { } } + async fn redial_disconnected_peers(&self) { + let (quic, alive) = { + let ctrl = self.controller.read().await; + (ctrl.transport().as_quic().cloned(), ctrl.alive_nodes()) + }; + if let Some(quic) = quic { + tokio::spawn(async move { + quic.redial_unlinked(&alive).await; + }); + } + } + async fn warn_on_unlinked_nodes(&self) { let unlinked = { let ctrl = self.controller.read().await; diff --git a/crates/mqdb-cluster/tests/quic_transport_framing.rs b/crates/mqdb-cluster/tests/quic_transport_framing.rs index c25ef46..0e71b26 100644 --- a/crates/mqdb-cluster/tests/quic_transport_framing.rs +++ b/crates/mqdb-cluster/tests/quic_transport_framing.rs @@ -43,7 +43,7 @@ fn generate_certs() -> Certs { } } -fn stalling_server_endpoint(certs: &Certs) -> Endpoint { +fn server_endpoint(certs: &Certs, stream_window: u32) -> Endpoint { let key = PrivateKeyDer::try_from(certs.leaf_key_der.clone()).unwrap(); let server_crypto = rustls::ServerConfig::builder() .with_no_client_auth() @@ -55,7 +55,7 @@ fn stalling_server_endpoint(certs: &Certs) -> Endpoint { )); let mut transport = TransportConfig::default(); - transport.stream_receive_window(VarInt::from_u32(STALL_WINDOW)); + transport.stream_receive_window(VarInt::from_u32(stream_window)); server_config.transport_config(Arc::new(transport)); Endpoint::server(server_config, "127.0.0.1:0".parse().unwrap()).unwrap() @@ -121,7 +121,7 @@ async fn setup_stalled_peer() -> StalledPeer { std::fs::write(&leaf_path, &certs.leaf_pem).unwrap(); std::fs::write(&leaf_key_path, &certs.leaf_key_pem).unwrap(); - let endpoint = stalling_server_endpoint(&certs); + let endpoint = server_endpoint(&certs, STALL_WINDOW); let far_addr: SocketAddr = endpoint.local_addr().unwrap(); let (release_tx, release_rx) = oneshot::channel(); @@ -309,3 +309,96 @@ async fn control_plane_survives_a_full_bulk_queue() { far_handle.abort(); } + +async fn reconnecting_far_end( + endpoint: Endpoint, + close_first: oneshot::Receiver<()>, + accepted_tx: flume::Sender<()>, +) { + let conn1: Connection = endpoint.accept().await.unwrap().await.unwrap(); + let (_s1, mut r1): (SendStream, RecvStream) = conn1.accept_bi().await.unwrap(); + let mut header = [0u8; 2]; + r1.read_exact(&mut header).await.unwrap(); + accepted_tx.send(()).ok(); + + close_first.await.ok(); + conn1.close(0u32.into(), b"test-close"); + drop(r1); + drop(conn1); + + let conn2: Connection = endpoint.accept().await.unwrap().await.unwrap(); + let (_s2, mut r2): (SendStream, RecvStream) = conn2.accept_bi().await.unwrap(); + let mut header2 = [0u8; 2]; + r2.read_exact(&mut header2).await.unwrap(); + accepted_tx.send(()).ok(); + + std::future::pending::<()>().await; + drop(conn2); +} + +#[tokio::test] +async fn dead_peer_is_removed_and_redialled() { + let _ = rustls::crypto::ring::default_provider().install_default(); + + let certs = generate_certs(); + let dir = tempfile::tempdir().unwrap(); + let ca_path = dir.path().join("ca.pem"); + let leaf_path = dir.path().join("leaf.pem"); + let leaf_key_path = dir.path().join("leaf.key"); + std::fs::write(&ca_path, &certs.ca_pem).unwrap(); + std::fs::write(&leaf_path, &certs.leaf_pem).unwrap(); + std::fs::write(&leaf_key_path, &certs.leaf_key_pem).unwrap(); + + let endpoint = server_endpoint(&certs, 1024 * 1024); + let far_addr: SocketAddr = endpoint.local_addr().unwrap(); + + let (close_tx, close_rx) = oneshot::channel(); + let (accepted_tx, accepted_rx) = flume::unbounded(); + let far_handle = tokio::spawn(reconnecting_far_end(endpoint, close_rx, accepted_tx)); + + let local = NodeId::validated(1).unwrap(); + let peer = NodeId::validated(2).unwrap(); + let transport = QuicDirectTransport::new(local); + transport.set_ca_file(ca_path); + transport + .bind("127.0.0.1:0".parse().unwrap(), &leaf_path, &leaf_key_path) + .await + .unwrap(); + transport.connect_to_peer(peer, far_addr).await.unwrap(); + + tokio::time::timeout(Duration::from_secs(10), accepted_rx.recv_async()) + .await + .expect("far end did not accept the initial connection") + .unwrap(); + assert!( + transport.direct_peers().await.unwrap().contains(&peer), + "peer should be linked after the initial connect" + ); + + close_tx.send(()).ok(); + let mut removed = false; + for _ in 0..200 { + let _ = transport.send(peer, heartbeat(local, 1)).await; + if !transport.direct_peers().await.unwrap().contains(&peer) { + removed = true; + break; + } + tokio::time::sleep(Duration::from_millis(25)).await; + } + assert!( + removed, + "a send-side failure should remove the dead peer from the map" + ); + + transport.redial_unlinked(&[]).await; + tokio::time::timeout(Duration::from_secs(10), accepted_rx.recv_async()) + .await + .expect("far end did not accept the redial") + .unwrap(); + assert!( + transport.direct_peers().await.unwrap().contains(&peer), + "peer should be re-linked after redial_disconnected" + ); + + far_handle.abort(); +} diff --git a/specs/ClusterRedial.cfg b/specs/ClusterRedial.cfg new file mode 100644 index 0000000..c9c51b7 --- /dev/null +++ b/specs/ClusterRedial.cfg @@ -0,0 +1,15 @@ +SPECIFICATION Spec +CONSTANTS + Nodes = {1, 2} + MaxGen = 5 + MaxBreaks = 2 + Guarded = TRUE + MonotonicId = TRUE + Redial = TRUE +INVARIANTS + TypeOK + InvNoLiveDrop +PROPERTIES + Converge + +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterRedial.tla b/specs/ClusterRedial.tla new file mode 100644 index 0000000..ec66bd9 --- /dev/null +++ b/specs/ClusterRedial.tla @@ -0,0 +1,152 @@ +---------------------------- MODULE ClusterRedial ---------------------------- +(***************************************************************************) +(* Design model for the redial half of gh #146 (configured-peers scope), *) +(* baseline main @ 8945b67 (after #145 framing fix, #147 two-lane writer). *) +(* *) +(* WHAT IS MODELLED (see scratchpad/redial-design-facts.md for the code *) +(* facts each abstraction is grounded in): *) +(* * peers map slot per DIRECTED pair a->b: Absent | Live(gen) | Broken *) +(* (gen). A dial opens ONE bidirectional stream, so Connect(a,b) sets *) +(* BOTH a's outbound slot[a][b] and the accept side slot[b][a], each *) +(* UNCONDITIONALLY replacing whatever was there (HashMap::insert *) +(* replaces, dropping the displaced connection). *) +(* * removal is driven by SEND-side death (writer write_all error), which *) +(* emits an in-flight death notice carrying the connection's gen; it is *) +(* NOT triggered by receiver-exit and NOT by clean replacement. *) +(* * removal is a compare-and-remove keyed on the gen id (Guarded); the *) +(* id is monotonic (MonotonicId) or reusable (models quinn stable_id). *) +(* * redial: Connect fires when slot[a][b] is not Live, i.e. absent or *) +(* broken (the impl's mesh tick re-dials configured members the *) +(* heartbeat view reports as not alive, driven by liveness rather than *) +(* the transport peer map; connect_to_peer replaces any stale entry). *) +(* *) +(* The three configs map to gh #140's surviving design rules: *) +(* ClusterRedial.cfg Guarded=TRUE MonotonicId=TRUE -> safety holds + *) +(* broken links re-converge (liveness). *) +(* *_unguarded.cfg Guarded=FALSE -> InvNoLiveDrop *) +(* violated (by-key removal deletes the newer entry). *) +(* *_reusable.cfg Guarded=TRUE MonotonicId=FALSE -> InvNoLiveDrop *) +(* violated (a stale notice matches a reused id). *) +(* *_noredial.cfg Redial=FALSE -> Converge violated (removal without *) +(* re-dial leaves a slot permanently absent). *) +(* *) +(* PREMISE, NOT VERIFIED: that SEND-death (not receiver-exit) is the right *) +(* removal trigger is BAKED IN (SendBreak is the only pd producer; there is *) +(* no receiver-exit action). It is grounded in the code + the #140 audit, *) +(* but this model assumes it rather than proving it. The guard / monotonic *) +(* id / redial rules ARE demonstrated as checker output above. *) +(* IMPLEMENTATION PROOF-OBLIGATION: the safety result holds only if the *) +(* generation counter is node-local and in-memory, co-located with removal *) +(* and the pending-notice queue, so a restart clears all three together. If *) +(* the counter were persisted, derived from a wire value, or outlived by its *) +(* notices, ABA would reappear and this model could not see it. *) +(* *) +(* HONEST SCOPE / NOT MODELLED: node crash/restart incarnations (abstracted *) +(* as SendBreak); address gossip / non-configured discovery (that stays *) +(* #140); the shared-stream framing + two-lane priority (#143/#147); mesh *) +(* single-hop completeness (#140 topology). Connect updates both directed *) +(* slots (dialer insert + accepter insert) in ONE atomic transition, *) +(* collapsing two cross-node async events -- defensible since the slots and *) +(* InvNoLiveDrop are per-directed-pair, but it hides some dialer/accepter *) +(* interleaving. Generations are bounded by MaxGen, so results are *) +(* "verified for these constants", not proven for all sizes; liveness is *) +(* verified up to MaxBreaks disruptions. *) +(***************************************************************************) +EXTENDS Naturals, FiniteSets + +CONSTANTS Nodes, MaxGen, MaxBreaks, Guarded, MonotonicId, Redial + +ASSUME MaxGen \in Nat /\ MaxGen >= 2 +ASSUME MaxBreaks \in Nat +ASSUME Guarded \in BOOLEAN /\ MonotonicId \in BOOLEAN /\ Redial \in BOOLEAN + +VARIABLES slot, gctr, pd, droppedLive, breaksLeft + +vars == <> + +Pairs == {p \in Nodes \X Nodes : p[1] # p[2]} + +Absent == [st |-> "absent", gen |-> 0] + +Notices == [from : Nodes, to : Nodes, gen : 1..MaxGen] + +TypeOK == + /\ slot \in [Nodes -> [Nodes -> [st : {"absent", "live", "broken"}, gen : 0..MaxGen]]] + /\ gctr \in [Nodes -> [Nodes -> 0..MaxGen]] + /\ pd \subseteq Notices + /\ droppedLive \in [Nodes -> [Nodes -> BOOLEAN]] + /\ breaksLeft \in 0..MaxBreaks + +Init == + /\ slot = [a \in Nodes |-> [b \in Nodes |-> Absent]] + /\ gctr = [a \in Nodes |-> [b \in Nodes |-> 0]] + /\ pd = {} + /\ droppedLive = [a \in Nodes |-> [b \in Nodes |-> FALSE]] + /\ breaksLeft = MaxBreaks + +NextGen(cur) == IF MonotonicId THEN cur + 1 ELSE (cur % MaxGen) + 1 + +\* a dials b: only when a's outbound slot to b is absent (redial driver dials +\* the unlinked). One bidirectional stream => both a's outbound and b's accept +\* side are (re)established with fresh generations, replacing any prior entry. +Connect(a, b) == + /\ a # b + /\ slot[a][b].st # "live" \* redial a link that is down (absent or broken) + /\ (Redial \/ gctr[a][b] = 0) \* with redial off, only the first (startup) dial fires + /\ (MonotonicId => (gctr[a][b] < MaxGen /\ gctr[b][a] < MaxGen)) + /\ gctr' = [gctr EXCEPT ![a][b] = NextGen(gctr[a][b]), + ![b][a] = NextGen(gctr[b][a])] + /\ slot' = [slot EXCEPT ![a][b] = [st |-> "live", gen |-> gctr'[a][b]], + ![b][a] = [st |-> "live", gen |-> gctr'[b][a]]] + /\ UNCHANGED <> + +\* a's send to b fails (peer gone / network): the writer emits a death notice +\* carrying this connection's gen. Clean replacement is NOT a SendBreak. +SendBreak(a, b) == + /\ a # b + /\ slot[a][b].st = "live" + /\ breaksLeft > 0 + /\ breaksLeft' = breaksLeft - 1 + /\ slot' = [slot EXCEPT ![a][b].st = "broken"] + /\ pd' = pd \cup {[from |-> a, to |-> b, gen |-> slot[a][b].gen]} + /\ UNCHANGED <> + +\* a death notice is delivered. Guarded: compare-and-remove keyed on the gen id +\* (the real guard checks id equality, NOT a re-checkable "broken" flag). By-key: +\* remove whatever occupies the slot. droppedLive records removing a LIVE slot. +Remove(a, b, g) == + /\ [from |-> a, to |-> b, gen |-> g] \in pd + /\ pd' = pd \ {[from |-> a, to |-> b, gen |-> g]} + /\ LET cur == slot[a][b] + doRemove == IF Guarded THEN (cur.st # "absent" /\ cur.gen = g) + ELSE (cur.st # "absent") + IN IF doRemove + THEN /\ slot' = [slot EXCEPT ![a][b] = Absent] + /\ droppedLive' = IF cur.st = "live" + THEN [droppedLive EXCEPT ![a][b] = TRUE] + ELSE droppedLive + ELSE UNCHANGED <> + /\ UNCHANGED <> + +ConnectStep == \E a, b \in Nodes : Connect(a, b) +BreakStep == \E a, b \in Nodes : SendBreak(a, b) +RemoveStep == \E n \in Notices : Remove(n.from, n.to, n.gen) + +Next == ConnectStep \/ BreakStep \/ RemoveStep + +Fairness == + /\ \A a, b \in Nodes : WF_vars(Connect(a, b)) + /\ \A n \in Notices : WF_vars(Remove(n.from, n.to, n.gen)) + +Spec == Init /\ [][Next]_vars /\ Fairness + +---------------------------------------------------------------------------- +\* SAFETY: a live (healthy) connection is never removed. +InvNoLiveDrop == \A a \in Nodes : \A b \in Nodes : ~droppedLive[a][b] + +\* LIVENESS: every directed pair eventually becomes live and stays live +\* (redial re-establishes broken links; a stale death notice never permanently +\* prevents reconnection). Verified up to MaxBreaks disruptions. +AllLinked == \A p \in Pairs : slot[p[1]][p[2]].st = "live" +Converge == <>[]AllLinked +============================================================================= diff --git a/specs/ClusterRedial_noredial.cfg b/specs/ClusterRedial_noredial.cfg new file mode 100644 index 0000000..4609bdb --- /dev/null +++ b/specs/ClusterRedial_noredial.cfg @@ -0,0 +1,15 @@ +SPECIFICATION Spec +CONSTANTS + Nodes = {1, 2} + MaxGen = 5 + MaxBreaks = 1 + Guarded = TRUE + MonotonicId = TRUE + Redial = FALSE +INVARIANTS + TypeOK + InvNoLiveDrop +PROPERTIES + Converge + +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterRedial_reusable.cfg b/specs/ClusterRedial_reusable.cfg new file mode 100644 index 0000000..9411792 --- /dev/null +++ b/specs/ClusterRedial_reusable.cfg @@ -0,0 +1,13 @@ +SPECIFICATION Spec +CONSTANTS + Nodes = {1, 2} + MaxGen = 2 + MaxBreaks = 4 + Guarded = TRUE + MonotonicId = FALSE + Redial = TRUE +INVARIANTS + TypeOK + InvNoLiveDrop + +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterRedial_unguarded.cfg b/specs/ClusterRedial_unguarded.cfg new file mode 100644 index 0000000..1780032 --- /dev/null +++ b/specs/ClusterRedial_unguarded.cfg @@ -0,0 +1,13 @@ +SPECIFICATION Spec +CONSTANTS + Nodes = {1, 2} + MaxGen = 3 + MaxBreaks = 2 + Guarded = FALSE + MonotonicId = TRUE + Redial = TRUE +INVARIANTS + TypeOK + InvNoLiveDrop + +CHECK_DEADLOCK FALSE From f7e7c877ec1c913e50dce3bf68aa3eda84c8cf88 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Tue, 22 Sep 2026 01:44:41 -0300 Subject: [PATCH 2/2] model: confirm redial over a live connection stays ABA-safe --- specs/ClusterRedial.cfg | 1 + specs/ClusterRedial.tla | 9 +++++++-- specs/ClusterRedial_noredial.cfg | 1 + specs/ClusterRedial_overlive.cfg | 14 ++++++++++++++ specs/ClusterRedial_reusable.cfg | 1 + specs/ClusterRedial_unguarded.cfg | 1 + 6 files changed, 25 insertions(+), 2 deletions(-) create mode 100644 specs/ClusterRedial_overlive.cfg diff --git a/specs/ClusterRedial.cfg b/specs/ClusterRedial.cfg index c9c51b7..e3682f6 100644 --- a/specs/ClusterRedial.cfg +++ b/specs/ClusterRedial.cfg @@ -6,6 +6,7 @@ CONSTANTS Guarded = TRUE MonotonicId = TRUE Redial = TRUE + RedialOverLive = FALSE INVARIANTS TypeOK InvNoLiveDrop diff --git a/specs/ClusterRedial.tla b/specs/ClusterRedial.tla index ec66bd9..5a42a2c 100644 --- a/specs/ClusterRedial.tla +++ b/specs/ClusterRedial.tla @@ -29,6 +29,10 @@ (* violated (a stale notice matches a reused id). *) (* *_noredial.cfg Redial=FALSE -> Converge violated (removal without *) (* re-dial leaves a slot permanently absent). *) +(* *_overlive.cfg RedialOverLive=TRUE -> InvNoLiveDrop still holds: *) +(* the impl re-dials on the HEARTBEAT view, so Connect *) +(* can fire over a still-live slot; the guard keeps it *) +(* safe (a live slot has no pending notice to race). *) (* *) (* PREMISE, NOT VERIFIED: that SEND-death (not receiver-exit) is the right *) (* removal trigger is BAKED IN (SendBreak is the only pd producer; there is *) @@ -54,11 +58,12 @@ (***************************************************************************) EXTENDS Naturals, FiniteSets -CONSTANTS Nodes, MaxGen, MaxBreaks, Guarded, MonotonicId, Redial +CONSTANTS Nodes, MaxGen, MaxBreaks, Guarded, MonotonicId, Redial, RedialOverLive ASSUME MaxGen \in Nat /\ MaxGen >= 2 ASSUME MaxBreaks \in Nat ASSUME Guarded \in BOOLEAN /\ MonotonicId \in BOOLEAN /\ Redial \in BOOLEAN +ASSUME RedialOverLive \in BOOLEAN VARIABLES slot, gctr, pd, droppedLive, breaksLeft @@ -91,7 +96,7 @@ NextGen(cur) == IF MonotonicId THEN cur + 1 ELSE (cur % MaxGen) + 1 \* side are (re)established with fresh generations, replacing any prior entry. Connect(a, b) == /\ a # b - /\ slot[a][b].st # "live" \* redial a link that is down (absent or broken) + /\ (slot[a][b].st # "live" \/ RedialOverLive) \* down link; or (impl's real trigger) over a still-live slot when the heartbeat view lags /\ (Redial \/ gctr[a][b] = 0) \* with redial off, only the first (startup) dial fires /\ (MonotonicId => (gctr[a][b] < MaxGen /\ gctr[b][a] < MaxGen)) /\ gctr' = [gctr EXCEPT ![a][b] = NextGen(gctr[a][b]), diff --git a/specs/ClusterRedial_noredial.cfg b/specs/ClusterRedial_noredial.cfg index 4609bdb..5ff3a52 100644 --- a/specs/ClusterRedial_noredial.cfg +++ b/specs/ClusterRedial_noredial.cfg @@ -6,6 +6,7 @@ CONSTANTS Guarded = TRUE MonotonicId = TRUE Redial = FALSE + RedialOverLive = FALSE INVARIANTS TypeOK InvNoLiveDrop diff --git a/specs/ClusterRedial_overlive.cfg b/specs/ClusterRedial_overlive.cfg new file mode 100644 index 0000000..932fb39 --- /dev/null +++ b/specs/ClusterRedial_overlive.cfg @@ -0,0 +1,14 @@ +SPECIFICATION Spec +CONSTANTS + Nodes = {1, 2} + MaxGen = 5 + MaxBreaks = 2 + Guarded = TRUE + MonotonicId = TRUE + Redial = TRUE + RedialOverLive = TRUE +INVARIANTS + TypeOK + InvNoLiveDrop + +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterRedial_reusable.cfg b/specs/ClusterRedial_reusable.cfg index 9411792..abba361 100644 --- a/specs/ClusterRedial_reusable.cfg +++ b/specs/ClusterRedial_reusable.cfg @@ -6,6 +6,7 @@ CONSTANTS Guarded = TRUE MonotonicId = FALSE Redial = TRUE + RedialOverLive = FALSE INVARIANTS TypeOK InvNoLiveDrop diff --git a/specs/ClusterRedial_unguarded.cfg b/specs/ClusterRedial_unguarded.cfg index 1780032..e80af37 100644 --- a/specs/ClusterRedial_unguarded.cfg +++ b/specs/ClusterRedial_unguarded.cfg @@ -6,6 +6,7 @@ CONSTANTS Guarded = FALSE MonotonicId = TRUE Redial = TRUE + RedialOverLive = FALSE INVARIANTS TypeOK InvNoLiveDrop