1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
//! Write-result resolution of a cross-shard block serve's escrow.
//!
//! When the origin serves a reply to a conn that appears alive, it does not
//! release the target's escrow immediately — a point-in-time "alive" can go
//! stale before the reply is actually written. Instead it records
//! `serve_confirm[conn] = target_shard` and lets the write outcome decide:
//! a clean flush on a live conn releases the escrow, a teardown (the FIN was
//! read, or a write errored) restores the element. Both reactors resolve by
//! the same two entry points here; split from `block_xshard` for the LOC cap.
use crate::Commands;
use crate::message::Inbound;
use crate::shard::Shard;
/// Debug-only signal that a cross-shard serve reply was processed by the
/// origin (i.e. `origin_on_serve_resp` ran). The escrow regression uses it
/// to tell a genuine cross-shard placement from a co-located one: with N
/// shards a random key lands on the conn's own shard ~1/N of the time, and
/// that takes the LOCAL block path, not this cross-shard one — so the test
/// retries until it provably exercised the cross-shard code. Non-I/O, so it
/// does not perturb the timing.
#[cfg(debug_assertions)]
pub mod counters {
use std::sync::atomic::{AtomicU64, Ordering::Relaxed};
/// Counts cross-shard serves since process start. Debug builds only
/// — the test that needs it retries until this proves the cross-shard
/// path actually ran, rather than assuming a pass meant it did.
pub static CROSS_SHARD_SERVES: AtomicU64 = AtomicU64::new(0);
#[inline]
/// Record one cross-shard serve. Relaxed and non-I/O, so arming the
/// counter cannot change the timing it is there to observe.
pub fn note_cross_shard_serve() {
CROSS_SHARD_SERVES.fetch_add(1, Relaxed);
}
/// Cross-shard serves processed since process start.
pub fn cross_shard_serves() -> u64 {
CROSS_SHARD_SERVES.load(Relaxed)
}
}
impl<C: Commands> Shard<C> {
/// The serve reply for `conn` reached the socket (its output flushed
/// without error): release the escrow the target has been holding.
/// Idempotent — a spurious flush after the confirm is a no-op.
pub(crate) fn confirm_serve_delivered(&mut self, conn: u64) {
if let Some(shard) = self.serve_confirm.remove(&conn) {
if shard == self.id {
self.target_release_escrow(self.id, conn);
} else {
self.send_to(shard, Inbound::BlockServeAck { origin: self.id, conn });
}
}
}
/// `conn` is being torn down with a serve reply still unconfirmed — the
/// write never succeeded, so the element never reached a live client.
/// Restore it. Idempotent. Called from the connection-close path.
pub(crate) fn restore_serve_on_teardown(&mut self, conn: u64) {
if let Some(shard) = self.serve_confirm.remove(&conn) {
if shard == self.id {
self.target_apply_escrow(self.id, conn);
} else {
self.send_to(shard, Inbound::BlockServeAbort { origin: self.id, conn });
}
}
}
/// The serve reply came back but the origin's block record is already
/// gone — the conn timed out or disconnected and was torn down first. An
/// empty reply popped nothing; a non-empty one means the target popped an
/// element and holds it in escrow with no record here to tell it the
/// outcome, so the escrow would strand and the element be lost. Route the
/// restore by the key's owning shard (the target that popped it).
/// Idempotent via `escrow_take`, so it cannot double-restore.
pub(crate) fn restore_serve_for_gone_record(&mut self, conn: u64, key: &[u8], reply: &[u8]) {
if reply.is_empty() {
return;
}
let shard = self.shard_of(key);
if shard == self.id {
self.target_apply_escrow(self.id, conn);
} else {
self.send_to(shard, Inbound::BlockServeAbort { origin: self.id, conn });
}
}
/// Resolve a cross-shard block serve's escrow by the write result rather
/// than a point-in-time guess: `closing` (the client's FIN was read, or a
/// write to it errored) means the reply never reached a live client →
/// restore; a clean full drain on a live conn means it did → release.
/// Gated on the map being non-empty so the reply hot path pays only a
/// length check. Idempotent — a later `close_conn` restore is a no-op.
/// The poller calls this from `flush_conn`; io_uring from its own write
/// completion.
pub(crate) fn resolve_serve_by_write(&mut self, conn: u64, closing: bool, drained: bool) {
if self.serve_confirm.is_empty() {
return;
}
if closing {
self.restore_serve_on_teardown(conn);
} else if drained {
self.confirm_serve_delivered(conn);
}
}
/// io_uring twin of the poller's flush-conn escrow resolution: settle a
/// cross-shard block serve by the write outcome. `closing` (the client's
/// FIN was seen, or a write to it errored) → the reply never reached a
/// live client, restore; a clean full drain on a live conn → it did,
/// release. Gated on `serve_confirm` being non-empty, so the write hot
/// path pays a length check. Idempotent — a later teardown restore is a
/// no-op once this has resolved. Linux-only: the io_uring reactor is.
#[cfg(target_os = "linux")]
pub(crate) fn uring_resolve_serve(
&mut self,
cid: u64,
io: &kevy_map::KevyMap<u64, crate::uring_conn::UringConn>,
) {
if self.serve_confirm.is_empty() {
return;
}
let uc = io.get(&cid);
let closing =
uc.is_none_or(|u| u.closing) || self.conns.get(&cid).is_none_or(|c| c.closing);
let drained = uc.is_some_and(|u| u.write_buf.is_empty() && u.write_arcs.is_empty())
&& self.conns.get(&cid).is_some_and(|c| c.output.is_empty());
self.resolve_serve_by_write(cid, closing, drained);
}
}