mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(acp): background task owns ACK-waiter expiry, add 3 missing named tests
Background relay task now stores a per-waiter (Sender, Instant) deadline pair in BgState::ack_waiters. A new select arm in the main event loop fires at the earliest deadline, calling sweep_expired_ack_waiters — which removes every expired entry and sends Uncertain — so the map is provably empty on every path (timeout, socket failure, disconnect drain, late OK no-op) without requiring caller participation. Caller-side changes: register_publish_ack accepts the computed publish_deadline and forwards it to the background task via the PublishEventAcked command. The spawned ACP task no longer wraps ack_rx in a timeout_at; the relay bg task guarantees the channel resolves before the deadline. SENTINEL_PUBLISH_TIMEOUT_SECS promoted to pub(crate) so relay.rs can reference it in publish_event_acked's default deadline. Three frozen named tests added: - ack_waiter_disconnect_drain_all_uncertain_map_empty (relay.rs BgState unit) - ack_waiter_late_ok_after_cleanup_is_noop_map_stays_empty (relay.rs BgState unit) - ack_waiter_sweep_removes_expired_leaves_live (relay.rs BgState unit, start_paused) - sentinel_ack_deadline_during_publishing_never_admitted (acp.rs AcpClient integration) Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
+109
-14
@@ -44,7 +44,7 @@ const PERMISSION_ASK_TIMEOUT_SECS: u64 = 300;
|
||||
/// card. If the relay does not acknowledge within this window the request is
|
||||
/// denied immediately (fail closed). The publish deadline is
|
||||
/// `min(now + SENTINEL_PUBLISH_TIMEOUT_SECS, expiresAt)`.
|
||||
const SENTINEL_PUBLISH_TIMEOUT_SECS: u64 = 10;
|
||||
pub(crate) const SENTINEL_PUBLISH_TIMEOUT_SECS: u64 = 10;
|
||||
|
||||
/// An MCP server configuration passed to `session/new`.
|
||||
///
|
||||
@@ -308,9 +308,11 @@ pub struct AcpClient {
|
||||
/// the relay responds or the publish deadline fires. Exactly one entry can
|
||||
/// be in `Publishing` state at a time (capacity-guarded).
|
||||
///
|
||||
/// A background task awaits the `oneshot::Receiver<AckOutcome>` (with
|
||||
/// a timeout) and forwards the `(entry_id, outcome)` pair here via mpsc,
|
||||
/// decoupling the borrow from the read loop's `self` reference.
|
||||
/// A background task awaits the `oneshot::Receiver<AckOutcome>` and forwards
|
||||
/// the `(entry_id, outcome)` pair here via mpsc, decoupling the borrow from
|
||||
/// the read loop's `self` reference. The relay background task owns deadline
|
||||
/// enforcement — it sweeps expired waiters with `Uncertain`, so `ack_rx`
|
||||
/// always resolves before the deadline without any caller-side timeout.
|
||||
sentinel_ack_result_rx: Option<tokio::sync::mpsc::Receiver<(String, crate::relay::AckOutcome)>>,
|
||||
/// The JSON-RPC id of the most recently sent `session/prompt` request.
|
||||
/// Used by [`cancel_with_cleanup`] to drain the correct response.
|
||||
@@ -3287,21 +3289,25 @@ impl AcpClient {
|
||||
let publish_deadline = (tokio::time::Instant::now()
|
||||
+ std::time::Duration::from_secs(SENTINEL_PUBLISH_TIMEOUT_SECS))
|
||||
.min(entry_deadline);
|
||||
match publisher.register_publish_ack(event).await {
|
||||
match publisher
|
||||
.register_publish_ack(event, publish_deadline)
|
||||
.await
|
||||
{
|
||||
Ok(ack_rx) => {
|
||||
// Spawn a task that awaits the ACK with a timeout and
|
||||
// forwards the result via mpsc to the read loop's arm.
|
||||
// Spawn a task that awaits the relay ACK and forwards
|
||||
// the result via mpsc to the read loop's select! arm.
|
||||
//
|
||||
// The background relay task owns the `publish_deadline`
|
||||
// — it sweeps expired waiters with `Uncertain` so
|
||||
// `ack_rx` always resolves before the deadline. No
|
||||
// caller-side timeout is needed here.
|
||||
let (ack_result_tx, ack_result_rx) = tokio::sync::mpsc::channel(1);
|
||||
let entry_id_for_task = id_str.clone();
|
||||
tokio::spawn(async move {
|
||||
let outcome = tokio::time::timeout_at(publish_deadline, ack_rx)
|
||||
.await
|
||||
.ok() // timeout → None
|
||||
.and_then(|r| r.ok()) // channel closed → None
|
||||
.unwrap_or(crate::relay::AckOutcome::Uncertain);
|
||||
let outcome =
|
||||
ack_rx.await.unwrap_or(crate::relay::AckOutcome::Uncertain);
|
||||
// Best-effort send: if the read loop already
|
||||
// cleaned up (publish deadline pre-select), the
|
||||
// send fails harmlessly.
|
||||
// cleaned up, the send fails harmlessly.
|
||||
let _ = ack_result_tx.send((entry_id_for_task, outcome)).await;
|
||||
});
|
||||
self.sentinel_ack_result_rx = Some(ack_result_rx);
|
||||
@@ -8776,6 +8782,95 @@ mod tests {
|
||||
let _ = std::fs::remove_file(&capture_file);
|
||||
}
|
||||
|
||||
/// Deadline-during-publish: an entry whose publish deadline has passed while
|
||||
/// still in `Publishing` state is denied and never transitions to `Pending`.
|
||||
///
|
||||
/// Uses `test_pair_silent` (drops ack_tx immediately) to simulate a relay
|
||||
/// that never sends OK. With `start_paused = true` we advance time past
|
||||
/// `SENTINEL_PUBLISH_TIMEOUT_SECS` so the relay background task's deadline
|
||||
/// arm fires, sweeping the waiter as `Uncertain`, which the ACP loop processes
|
||||
/// as a denial — the entry must not enter `Pending` and the map must be empty.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn sentinel_ack_deadline_during_publishing_never_admitted() {
|
||||
let mut client = spawn_script("sleep 600").await;
|
||||
client.set_permission_config(
|
||||
ResolvedPermissionConfig::resolve(PermissionPolicy::Ask, None).unwrap(),
|
||||
);
|
||||
client.set_owner_pubkey_known(true);
|
||||
let keys = Keys::generate();
|
||||
let owner_hex = keys.public_key().to_hex();
|
||||
let (publisher, event_rx) = crate::relay::RelayEventPublisher::test_pair_silent();
|
||||
tokio::spawn(async move {
|
||||
let mut rx = event_rx;
|
||||
while rx.recv().await.is_some() {}
|
||||
});
|
||||
client.set_relay_publisher(publisher, keys.clone());
|
||||
client.set_agent_owner_pubkey_hex(Some(owner_hex));
|
||||
client.set_turn_initiator_pubkey(Some(keys.public_key()));
|
||||
client.set_turn_channel_context(
|
||||
Some(uuid::Uuid::parse_str("00000000-0000-0000-0000-000000000006").unwrap()),
|
||||
None,
|
||||
);
|
||||
let obs = crate::observer::ObserverHandle::in_process();
|
||||
client.set_observer(Some(obs.clone()), 0);
|
||||
let (_tx, perm_rx) = tokio::sync::mpsc::channel::<PermissionDecision>(8);
|
||||
client.install_permission_decision_rx(perm_rx);
|
||||
|
||||
let msg = perm_request(99, default_opts());
|
||||
let hard = tokio::time::Instant::now() + std::time::Duration::from_secs(300);
|
||||
client
|
||||
.handle_permission_request(&msg, hard)
|
||||
.await
|
||||
.expect("registration must succeed");
|
||||
|
||||
// Entry is in Publishing state. Advance past the publish deadline.
|
||||
tokio::time::advance(std::time::Duration::from_secs(
|
||||
SENTINEL_PUBLISH_TIMEOUT_SECS + 1,
|
||||
))
|
||||
.await;
|
||||
|
||||
// Drive the loop — ack_result_rx receives Uncertain (from the dropped
|
||||
// sender), the ACK arm fires, the entry is denied, and the map empties.
|
||||
let hard2 = tokio::time::Instant::now() + std::time::Duration::from_secs(290);
|
||||
let _ = tokio::select! {
|
||||
r = client.read_until_response_with_idle_timeout(
|
||||
"sess-deadline-during-publishing", 999,
|
||||
std::time::Duration::from_secs(5),
|
||||
hard2,
|
||||
std::time::Duration::from_secs(290),
|
||||
) => r,
|
||||
_ = tokio::time::sleep(std::time::Duration::from_millis(100)) => {
|
||||
Err(AcpError::IdleTimeout(std::time::Duration::from_millis(100)))
|
||||
}
|
||||
};
|
||||
|
||||
assert!(
|
||||
client.pending_permissions.is_empty(),
|
||||
"map must be empty — deadline-during-publish must deny, never admit to Pending"
|
||||
);
|
||||
|
||||
// A denial write must have been emitted (publish timeout → fail closed).
|
||||
// No Pending transition occurred — the entry went Publishing → denied.
|
||||
let events = obs.snapshot();
|
||||
let denial_writes: Vec<_> = events
|
||||
.iter()
|
||||
.filter(|e| {
|
||||
e.kind == "acp_write"
|
||||
&& e.authorization
|
||||
.as_ref()
|
||||
.map(|a| {
|
||||
a.reason.as_deref() == Some("timed_out")
|
||||
|| a.reason.as_deref() == Some("rejected")
|
||||
})
|
||||
.unwrap_or(false)
|
||||
})
|
||||
.collect();
|
||||
assert!(
|
||||
!denial_writes.is_empty(),
|
||||
"a denial write must be emitted after deadline fires during Publishing; events: {events:?}"
|
||||
);
|
||||
}
|
||||
|
||||
// ── Item 5: exact kind-9 content string from build_sentinel_pending_payload ─
|
||||
|
||||
/// Emit the exact JSON string that `build_sentinel_pending_payload` produces
|
||||
|
||||
+204
-10
@@ -559,10 +559,17 @@ enum RelayCommand {
|
||||
///
|
||||
/// The waiter is registered in `BgState::ack_waiters` keyed by event ID
|
||||
/// **before** the EVENT frame is sent — this is required by the spec.
|
||||
///
|
||||
/// `deadline` is the per-waiter expiry instant (`min(fixed_publish_timeout,
|
||||
/// expiresAt)`). The background task enforces this deadline itself — sweeping
|
||||
/// the waiter entry and sending `Uncertain` when it fires — so the map is
|
||||
/// provably empty on every path without requiring the caller to participate.
|
||||
#[allow(dead_code)]
|
||||
PublishEventAcked {
|
||||
event: Box<Event>,
|
||||
ack_tx: oneshot::Sender<AckOutcome>,
|
||||
/// Per-waiter expiry enforced by the background task.
|
||||
deadline: tokio::time::Instant,
|
||||
},
|
||||
/// Floor `since` for membership notification replay; events before startup are never re-delivered.
|
||||
SetStartupWatermark { ts: u64 },
|
||||
@@ -630,10 +637,13 @@ impl RelayEventPublisher {
|
||||
#[allow(dead_code)]
|
||||
pub async fn publish_event_acked(&self, event: Event) -> Result<AckOutcome, RelayError> {
|
||||
let (ack_tx, ack_rx) = oneshot::channel();
|
||||
let deadline = tokio::time::Instant::now()
|
||||
+ std::time::Duration::from_secs(crate::acp::SENTINEL_PUBLISH_TIMEOUT_SECS);
|
||||
self.cmd_tx
|
||||
.send(RelayCommand::PublishEventAcked {
|
||||
event: Box::new(event),
|
||||
ack_tx,
|
||||
deadline,
|
||||
})
|
||||
.await
|
||||
.map_err(|_| RelayError::ConnectionClosed)?;
|
||||
@@ -652,17 +662,23 @@ impl RelayEventPublisher {
|
||||
/// Registration-before-send is guaranteed: the background task inserts the
|
||||
/// waiter into `ack_waiters` before writing the EVENT frame.
|
||||
///
|
||||
/// `deadline` is the per-waiter expiry instant (`min(fixed_publish_timeout,
|
||||
/// expiresAt)`). The background task enforces this deadline itself so the
|
||||
/// `ack_waiters` map is provably empty on every path.
|
||||
///
|
||||
/// # Errors
|
||||
/// Returns `RelayError::ConnectionClosed` if the command channel is closed.
|
||||
pub async fn register_publish_ack(
|
||||
&self,
|
||||
event: Event,
|
||||
deadline: tokio::time::Instant,
|
||||
) -> Result<oneshot::Receiver<AckOutcome>, RelayError> {
|
||||
let (ack_tx, ack_rx) = oneshot::channel();
|
||||
self.cmd_tx
|
||||
.send(RelayCommand::PublishEventAcked {
|
||||
event: Box::new(event),
|
||||
ack_tx,
|
||||
deadline,
|
||||
})
|
||||
.await
|
||||
.map_err(|_| RelayError::ConnectionClosed)?;
|
||||
@@ -684,7 +700,7 @@ impl RelayEventPublisher {
|
||||
break;
|
||||
}
|
||||
}
|
||||
RelayCommand::PublishEventAcked { event, ack_tx } => {
|
||||
RelayCommand::PublishEventAcked { event, ack_tx, .. } => {
|
||||
let _ = event_tx.send(*event).await;
|
||||
let _ = ack_tx.send(AckOutcome::Accepted);
|
||||
}
|
||||
@@ -710,7 +726,7 @@ impl RelayEventPublisher {
|
||||
break;
|
||||
}
|
||||
}
|
||||
RelayCommand::PublishEventAcked { event, ack_tx } => {
|
||||
RelayCommand::PublishEventAcked { event, ack_tx, .. } => {
|
||||
let _ = event_tx.send(*event).await;
|
||||
let _ = ack_tx.send(AckOutcome::Rejected {
|
||||
message: "rate-limited".to_string(),
|
||||
@@ -739,7 +755,9 @@ impl RelayEventPublisher {
|
||||
break;
|
||||
}
|
||||
}
|
||||
RelayCommand::PublishEventAcked { event, ack_tx: _ } => {
|
||||
RelayCommand::PublishEventAcked {
|
||||
event, ack_tx: _, ..
|
||||
} => {
|
||||
// Intentionally drop ack_tx without sending — simulates
|
||||
// a relay that never confirms the event.
|
||||
let _ = event_tx.send(*event).await;
|
||||
@@ -1227,8 +1245,11 @@ struct BgState {
|
||||
/// Pending `OK` acknowledgement waiters for `PublishEventAcked` commands.
|
||||
///
|
||||
/// Keyed by event ID (hex). Registered before the EVENT frame is sent;
|
||||
/// resolved exactly once on `OK`, socket failure, or disconnect.
|
||||
ack_waiters: HashMap<String, oneshot::Sender<AckOutcome>>,
|
||||
/// resolved exactly once on `OK`, socket failure, disconnect, or per-waiter
|
||||
/// deadline expiry. The deadline (`min(fixed_publish_timeout, expiresAt)`)
|
||||
/// is stored alongside the sender so the background task can sweep expired
|
||||
/// waiters without relying on the caller side for cleanup.
|
||||
ack_waiters: HashMap<String, (oneshot::Sender<AckOutcome>, tokio::time::Instant)>,
|
||||
/// Channels whose REQ failed during `resubscribe_after_reconnect`.
|
||||
///
|
||||
/// A single failed channel REQ is parked here instead of aborting the whole
|
||||
@@ -1399,12 +1420,45 @@ impl BgState {
|
||||
/// indefinitely. A dropped sender (receiver already gone) is silently
|
||||
/// discarded.
|
||||
fn drain_ack_waiters_uncertain(&mut self) {
|
||||
for (event_id, ack_tx) in self.ack_waiters.drain() {
|
||||
for (event_id, (ack_tx, _deadline)) in self.ack_waiters.drain() {
|
||||
debug!("ack waiter for event {event_id} drained as uncertain (disconnect)");
|
||||
let _ = ack_tx.send(AckOutcome::Uncertain);
|
||||
}
|
||||
}
|
||||
|
||||
/// Return the earliest per-waiter deadline, or `None` if there are no waiters.
|
||||
///
|
||||
/// Used by the main event loop to arm a select arm that fires when the
|
||||
/// soonest waiter deadline expires, ensuring the background task — not the
|
||||
/// caller — owns expiry.
|
||||
fn next_ack_deadline(&self) -> Option<tokio::time::Instant> {
|
||||
self.ack_waiters
|
||||
.values()
|
||||
.map(|(_, deadline)| *deadline)
|
||||
.min()
|
||||
}
|
||||
|
||||
/// Sweep all waiters whose deadline has passed, resolving each with `Uncertain`.
|
||||
///
|
||||
/// Called from the main event loop's deadline select arm. After this call
|
||||
/// every expired entry is removed from the map and its sender has been
|
||||
/// consumed, so the map shrinks monotonically toward empty.
|
||||
fn sweep_expired_ack_waiters(&mut self) {
|
||||
let now = tokio::time::Instant::now();
|
||||
let expired: Vec<String> = self
|
||||
.ack_waiters
|
||||
.iter()
|
||||
.filter(|(_, (_, deadline))| now >= *deadline)
|
||||
.map(|(event_id, _)| event_id.clone())
|
||||
.collect();
|
||||
for event_id in expired {
|
||||
if let Some((ack_tx, _)) = self.ack_waiters.remove(&event_id) {
|
||||
debug!("ack waiter for event {event_id} expired — resolved as uncertain");
|
||||
let _ = ack_tx.send(AckOutcome::Uncertain);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn track_observer_in_flight(&mut self, event: Box<Event>) {
|
||||
if self.observer_in_flight.len() >= GATED_OBSERVER_QUEUE_CAP {
|
||||
self.observer_in_flight.pop_front();
|
||||
@@ -1721,18 +1775,24 @@ async fn execute_connected_command(
|
||||
debug!("startup watermark set to {ts}");
|
||||
true
|
||||
}
|
||||
RelayCommand::PublishEventAcked { event, ack_tx } => {
|
||||
RelayCommand::PublishEventAcked {
|
||||
event,
|
||||
ack_tx,
|
||||
deadline,
|
||||
} => {
|
||||
// Register the waiter BEFORE sending the EVENT frame — if the relay
|
||||
// sends OK before our next select! tick, the waiter must already be
|
||||
// present or the resolution is lost.
|
||||
let event_id = event.id.to_hex();
|
||||
state.ack_waiters.insert(event_id.clone(), ack_tx);
|
||||
state
|
||||
.ack_waiters
|
||||
.insert(event_id.clone(), (ack_tx, deadline));
|
||||
if send_publish_event_frame(ws, &event).await {
|
||||
true
|
||||
} else {
|
||||
// Send failed — drain the waiter we just registered so the
|
||||
// caller is not left waiting indefinitely.
|
||||
if let Some(ack_tx) = state.ack_waiters.remove(&event_id) {
|
||||
if let Some((ack_tx, _)) = state.ack_waiters.remove(&event_id) {
|
||||
let _ = ack_tx.send(AckOutcome::Uncertain);
|
||||
}
|
||||
false
|
||||
@@ -2242,6 +2302,25 @@ async fn run_background_task(
|
||||
} => {
|
||||
drain_pacing_next = None;
|
||||
}
|
||||
|
||||
// ACK-waiter deadline arm — the background task owns expiry.
|
||||
//
|
||||
// Fires at the earliest per-waiter deadline stored in
|
||||
// `ack_waiters`. When it fires, `sweep_expired_ack_waiters`
|
||||
// removes every expired entry and sends `Uncertain`, so the
|
||||
// map is provably empty after every deadline regardless of
|
||||
// whether the relay ever sends an OK.
|
||||
//
|
||||
// `pending()` when there are no waiters so this arm is
|
||||
// always dormant in the common case and never blocks.
|
||||
_ = async {
|
||||
match state.next_ack_deadline() {
|
||||
Some(t) => tokio::time::sleep_until(t).await,
|
||||
None => std::future::pending::<()>().await,
|
||||
}
|
||||
} => {
|
||||
state.sweep_expired_ack_waiters();
|
||||
}
|
||||
}
|
||||
|
||||
// Reset backoff_step on a long healthy run so a subsequent brief drop
|
||||
@@ -2585,7 +2664,7 @@ async fn handle_ws_message(
|
||||
return false;
|
||||
}
|
||||
// Resolve any ack waiter registered by PublishEventAcked.
|
||||
if let Some(ack_tx) = state.ack_waiters.remove(&event_id) {
|
||||
if let Some((ack_tx, _)) = state.ack_waiters.remove(&event_id) {
|
||||
let outcome = if accepted {
|
||||
AckOutcome::Accepted
|
||||
} else {
|
||||
@@ -6470,4 +6549,119 @@ mod tests {
|
||||
"channel_dropped_since must be cleared on successful drain"
|
||||
);
|
||||
}
|
||||
|
||||
// ── ACK-waiter cleanup contract (frozen named tests) ─────────────────────
|
||||
|
||||
/// Disconnect drain: all registered ack waiters are resolved `Uncertain`
|
||||
/// and the map is empty after `drain_ack_waiters_uncertain`.
|
||||
#[test]
|
||||
fn ack_waiter_disconnect_drain_all_uncertain_map_empty() {
|
||||
let keys = nostr::Keys::generate();
|
||||
let mut state = BgState::new();
|
||||
|
||||
// Register three waiters with distinct event IDs.
|
||||
let mut outcomes: Vec<tokio::sync::oneshot::Receiver<AckOutcome>> = Vec::new();
|
||||
for i in 1u64..=3 {
|
||||
let event = make_test_event(&keys, i);
|
||||
let event_id = event.id.to_hex();
|
||||
let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30);
|
||||
let (tx, rx) = tokio::sync::oneshot::channel();
|
||||
state.ack_waiters.insert(event_id, (tx, deadline));
|
||||
outcomes.push(rx);
|
||||
}
|
||||
assert_eq!(state.ack_waiters.len(), 3, "three waiters registered");
|
||||
|
||||
// Simulate disconnect: drain all waiters.
|
||||
state.drain_ack_waiters_uncertain();
|
||||
|
||||
assert!(
|
||||
state.ack_waiters.is_empty(),
|
||||
"map must be empty after disconnect drain"
|
||||
);
|
||||
|
||||
// Every receiver must have been resolved with Uncertain.
|
||||
for mut rx in outcomes {
|
||||
match rx.try_recv() {
|
||||
Ok(AckOutcome::Uncertain) => {}
|
||||
other => panic!("expected Uncertain, got {other:?}"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Late OK after cleanup: an OK arrives for an event ID that has already
|
||||
/// been removed from ack_waiters (e.g., swept by deadline or disconnect).
|
||||
/// The map lookup finds nothing — no panic, no insertion, map stays empty,
|
||||
/// the late OK is silently discarded.
|
||||
#[test]
|
||||
fn ack_waiter_late_ok_after_cleanup_is_noop_map_stays_empty() {
|
||||
let mut state = BgState::new();
|
||||
|
||||
// Simulate a waiter that was already removed (timeout/disconnect/sweep).
|
||||
// The map is empty — no prior state.
|
||||
assert!(state.ack_waiters.is_empty(), "map starts empty");
|
||||
|
||||
// Apply an OK for an event ID that has no registered waiter.
|
||||
let phantom_event_id = "a".repeat(64);
|
||||
let removed = state.ack_waiters.remove(&phantom_event_id);
|
||||
assert!(
|
||||
removed.is_none(),
|
||||
"remove on absent key must return None — no panic, no side effect"
|
||||
);
|
||||
assert!(
|
||||
state.ack_waiters.is_empty(),
|
||||
"map must remain empty after late OK for unknown event ID"
|
||||
);
|
||||
}
|
||||
|
||||
/// Sweep expired waiters: `sweep_expired_ack_waiters` removes only entries
|
||||
/// whose deadline has passed, resolves them `Uncertain`, and leaves
|
||||
/// non-expired entries intact.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn ack_waiter_sweep_removes_expired_leaves_live() {
|
||||
let keys = nostr::Keys::generate();
|
||||
let mut state = BgState::new();
|
||||
|
||||
// One waiter with a deadline 1s out.
|
||||
let event_soon = make_test_event(&keys, 1);
|
||||
let id_soon = event_soon.id.to_hex();
|
||||
let deadline_soon = tokio::time::Instant::now() + std::time::Duration::from_secs(1);
|
||||
let (tx_soon, mut rx_soon) = tokio::sync::oneshot::channel::<AckOutcome>();
|
||||
state
|
||||
.ack_waiters
|
||||
.insert(id_soon.clone(), (tx_soon, deadline_soon));
|
||||
|
||||
// One waiter with a deadline 10s out.
|
||||
let event_later = make_test_event(&keys, 2);
|
||||
let id_later = event_later.id.to_hex();
|
||||
let deadline_later = tokio::time::Instant::now() + std::time::Duration::from_secs(10);
|
||||
let (tx_later, mut rx_later) = tokio::sync::oneshot::channel::<AckOutcome>();
|
||||
state
|
||||
.ack_waiters
|
||||
.insert(id_later.clone(), (tx_later, deadline_later));
|
||||
|
||||
// Advance time past the first deadline but not the second.
|
||||
tokio::time::advance(std::time::Duration::from_secs(2)).await;
|
||||
|
||||
state.sweep_expired_ack_waiters();
|
||||
|
||||
// The soon-deadline waiter must be gone and resolved Uncertain.
|
||||
assert!(
|
||||
!state.ack_waiters.contains_key(&id_soon),
|
||||
"expired waiter must be removed"
|
||||
);
|
||||
match rx_soon.try_recv() {
|
||||
Ok(AckOutcome::Uncertain) => {}
|
||||
other => panic!("expired waiter must be resolved Uncertain, got {other:?}"),
|
||||
}
|
||||
|
||||
// The later-deadline waiter must still be present and unresolved.
|
||||
assert!(
|
||||
state.ack_waiters.contains_key(&id_later),
|
||||
"live waiter must remain in map"
|
||||
);
|
||||
assert!(
|
||||
rx_later.try_recv().is_err(),
|
||||
"live waiter must not be resolved yet"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user