fix(acp): publish the setup nudge over REST so a dropped one is not silent

The setup nudge is a durable kind:9 reply, but it was published through
the WS path, which is built for ephemera. While the rate-limit gate is
armed, or while disconnected, that path drops any non-observer publish
and returns Ok: the command is queued and the caller never learns what
happened to it. The nudge was therefore lost while being logged as
"nudge published" — the user asked how to set the agent up and got
nothing back.

Dedup made the loss permanent. should_nudge_for_event recorded the event
id by insert-on-check, before the publish was attempted, so the retry on
the next delivery of the same event was suppressed by an entry written
for a message that was never sent. One dropped nudge meant no nudge ever.

Both halves are fixed by subtraction rather than new machinery:

The nudge now goes through RestClient::submit_event, joining the four
other durable publishes on the HTTP bridge. The rest_client was already
in scope at the call site. submit_event surfaces the relay's actual
response, so a refusal is an Err.

The gate stops mutating: it takes &HashSet and only checks. The caller
inserts the id in the Ok arm of the publish, so a nudge that failed to
send stays retryable.

An earlier draft of this change instead widened the WS drop guard from
kind == KIND_AGENT_OBSERVER_FRAME to is_ephemeral(kind). That is wrong,
and the comment now says why: observer frames are themselves
range-ephemeral, so testing is_ephemeral would route them into the drop
path and undo the parking directly above it. The boundary is encoded as
a debug_assert instead. Enumerating the callers shows it holds: presence
(20001), typing (20002), observer frames (24200, parked before the
assert), and HarnessRelay::publish_event, which has no in-repo callers.

That assertion immediately caught publish_during_replay_pacing_is_sent_
on_live_socket publishing a kind:1 through the publish path. It was a
fixture artifact, not a defect — the test asserts on the event id and
never on the kind, so it reached for the generic make_test_event. The
fixture now builds a typing indicator, and both helpers' docs say which
path they belong to.

Mutation results. Reverting the dedup insert back into the gate is
killed by test_failed_nudge_remains_retryable; note the pre-existing
test_same_event_id_twice_nudges_exactly_once survives it, so the new
test is the discriminator rather than decoration. Reverting the
transport to publisher.publish_event initially survived the entire
suite: the dedup tests are pure and structurally cannot see which
transport is used. That is what the three transport tests are for --
they drive publish_setup_nudge against a real socket, and they kill that
mutant 3/3. Removing the debug_assert is caught by restoring the kind:1
fixture.

Co-authored-by: Sami <f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68@buzz.block.builderlab.xyz>
Signed-off-by: Sami <f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68@buzz.block.builderlab.xyz>
This commit is contained in:
Sami
2026-08-05 12:17:10 -04:00
parent 186f7e71ad
commit 9281d7b313
2 changed files with 311 additions and 35 deletions
+38 -7
View File
@@ -1515,11 +1515,28 @@ async fn execute_connected_command(
// indicators are worthless and sending them would consume admission
// budget the relay already rejected us on.
//
// INVARIANT: apart from observer frames (parked above), the WS publish
// path carries only ephemeral kinds (typing indicators). The silent
// drop-while-gated relies on that invariant. If a future caller
// publishes durable events through this path, it must extend the
// kind guard above to avoid silently discarding user data.
// INVARIANT: the WS publish path carries only ephemeral kinds
// (presence 20001, typing 20002, observer frames 24200 — the last
// parked above rather than dropped). The silent drop-while-gated
// below is only correct under that invariant: dropping a durable
// event here would discard user data while reporting success to
// the caller, since `publish_event` returns `Ok` once the command
// is queued and never learns what happened to it.
//
// The assertion encodes the boundary rather than widening the
// guard. Note that `is_ephemeral` is not usable as the drop
// predicate here: observer frames are themselves range-ephemeral,
// so testing it would send them down the drop path and undo the
// parking above. A durable event reaching this point is a bug at
// the *caller* — it should use `RestClient::submit_event`, which
// surfaces the relay's real response (see `publish_setup_nudge`).
debug_assert!(
buzz_core::kind::is_ephemeral(u32::from(event.kind.as_u16())),
"durable kind {} reached the WS publish path, which silently \
drops while rate-gated or disconnected; publish durable events \
via RestClient::submit_event instead",
event.kind.as_u16()
);
if state.check_rate_gate().is_some() {
debug!("rate-gated: dropping ephemeral PublishEvent (typing indicator)");
return true;
@@ -4365,11 +4382,15 @@ mod tests {
assert_eq!(sub_id, "ch-12345678-1234-5678-1234-567812345678");
}
/// Build a real signed Nostr event for testing BgState.
/// Build a real signed durable kind:1 note for testing BgState.
///
/// Uses `custom_created_at` so tests can control the timestamp.
/// The event ID is determined by the nostr signing process — we don't
/// control it, but we return it so callers can use it for dedup tests.
///
/// Not valid as a `RelayCommand::PublishEvent` payload for
/// `execute_connected_command`: that path carries ephemeral kinds only and
/// asserts as much. Use [`make_test_typing_event`] for publish-path tests.
fn make_test_event(keys: &nostr::Keys, created_at_secs: u64) -> Event {
let ts = nostr::Timestamp::from(created_at_secs);
EventBuilder::new(nostr::Kind::TextNote, "test")
@@ -4379,6 +4400,15 @@ mod tests {
.expect("signing should succeed")
}
/// A kind:20002 typing indicator — a realistic payload for the WS publish
/// path, which carries ephemeral kinds only.
fn make_test_typing_event(keys: &nostr::Keys) -> Event {
EventBuilder::new(Kind::Custom(KIND_TYPING_INDICATOR as u16), "")
.tags([Tag::parse(["h", &Uuid::new_v4().to_string()]).expect("h tag")])
.sign_with_keys(keys)
.expect("signing should succeed")
}
async fn test_ws_pair() -> (WsStream, WebSocketStream<tokio::net::TcpStream>) {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
@@ -4542,7 +4572,8 @@ mod tests {
let mut state = BgState::new();
let channel_id = Uuid::new_v4();
seed_test_subscription(&mut state, channel_id);
let event = make_test_event(&nostr::Keys::generate(), 2_000);
// A typing indicator: the WS publish path carries ephemeral kinds only.
let event = make_test_typing_event(&nostr::Keys::generate());
let event_id = event.id.to_hex();
let task = tokio::spawn(async move {
+273 -28
View File
@@ -74,7 +74,7 @@ use crate::{
author_allowed,
config::Config,
event_mentions_agent, filter,
relay::{HarnessRelay, RelayEventPublisher},
relay::{HarnessRelay, RestClient},
};
// ── Payload ───────────────────────────────────────────────────────────────────
@@ -380,7 +380,6 @@ pub(crate) async fn run_setup_listener(config: Config, payload: SetupPayload) ->
}
}
let publisher = relay.event_publisher();
let rest_client = relay.rest_client();
let channel_info = crate::pool::ChannelInfoResolver::new(channel_info_map, rest_client.clone());
@@ -450,19 +449,27 @@ pub(crate) async fn run_setup_listener(config: Config, payload: SetupPayload) ->
.await
.is_some();
// Pure gate: author gate verdict + event-id dedup.
// Pure gate: author gate verdict + event-id dedup check.
if !should_nudge_for_event(
buzz_event.event.id,
allowed,
filter_matched,
&mut nudged_event_ids,
&nudged_event_ids,
) {
continue;
}
// Build and publish the setup nudge.
if let Err(e) = publish_setup_nudge(
&publisher,
// Build and publish the setup nudge over the HTTP bridge.
//
// The nudge is a durable kind:9 message. It must not ride the WS
// publish path: that path silently discards non-observer publishes
// while the rate-limit gate is armed and while disconnected, reporting
// success to the caller either way. For an ephemeral typing indicator
// that is correct; for a user-visible reply it is permanent data loss
// logged as "nudge published". `submit_event` surfaces the relay's
// actual response instead.
match publish_setup_nudge(
&rest_client,
&config.keys,
buzz_event.channel_id,
&buzz_event.event,
@@ -470,13 +477,21 @@ pub(crate) async fn run_setup_listener(config: Config, payload: SetupPayload) ->
)
.await
{
tracing::warn!("setup-mode: failed to publish nudge: {e}");
} else {
tracing::info!(
channel_id = %buzz_event.channel_id,
event_id = %buzz_event.event.id,
"setup-mode: nudge published"
);
Err(e) => {
// Dedup is NOT recorded: the nudge was not delivered, so a
// later retry for this event must still be allowed. Recording
// before the send is what made a dropped nudge permanent.
tracing::warn!("setup-mode: failed to publish nudge: {e}");
}
Ok(()) => {
// Record dedup only after the relay accepted the nudge.
nudged_event_ids.insert(buzz_event.event.id);
tracing::info!(
channel_id = %buzz_event.channel_id,
event_id = %buzz_event.event.id,
"setup-mode: nudge published"
);
}
}
}
@@ -487,8 +502,14 @@ pub(crate) async fn run_setup_listener(config: Config, payload: SetupPayload) ->
///
/// Callers compute the async gates (`author_allowed`, `filter::match_event`)
/// up-front, then pass the boolean results here. This helper handles
/// everything that is synchronous and stateful: the author gate verdict
/// and event-id dedup.
/// everything that is synchronous: the author gate verdict and the dedup
/// *check*.
///
/// Recording the dedup entry is deliberately NOT done here. The caller inserts
/// the event id only after the relay has accepted the nudge, so a nudge that
/// fails to send can still be retried. Inserting at check time — as this
/// function used to — made any dropped nudge permanently un-retryable, because
/// the event was marked as nudged before anything was delivered.
///
/// Returns `true` when the event should produce a nudge.
#[must_use]
@@ -496,7 +517,7 @@ pub(crate) fn should_nudge_for_event(
event_id: EventId,
author_allowed: bool,
filter_matched: bool,
nudged_event_ids: &mut HashSet<EventId>,
nudged_event_ids: &HashSet<EventId>,
) -> bool {
if !author_allowed {
tracing::debug!("setup-mode: event filtered by author gate");
@@ -505,7 +526,7 @@ pub(crate) fn should_nudge_for_event(
if !filter_matched {
return false;
}
if !nudged_event_ids.insert(event_id) {
if nudged_event_ids.contains(&event_id) {
tracing::debug!(%event_id, "setup-mode: skipping already-nudged event");
return false;
}
@@ -592,8 +613,14 @@ async fn handle_setup_membership(
///
/// Threading: flat reply to the thread root if one exists; otherwise reply
/// to the triggering event itself. P-tags the asker.
///
/// Published over the HTTP bridge rather than the WS publish path. The nudge
/// is a durable kind:9 message, and the WS path is documented to carry only
/// ephemeral kinds: it drops non-observer publishes while rate-gated or
/// disconnected and still returns `Ok`. Going through `submit_event` means a
/// failure is actually reported to the caller.
async fn publish_setup_nudge(
publisher: &RelayEventPublisher,
rest_client: &RestClient,
keys: &nostr::Keys,
channel_id: Uuid,
triggering_event: &nostr::Event,
@@ -637,8 +664,8 @@ async fn publish_setup_nudge(
.sign_with_keys(keys)
.map_err(|e| anyhow::anyhow!("failed to sign setup nudge: {e}"))?;
publisher
.publish_event(signed)
rest_client
.submit_event(&signed)
.await
.map_err(|e| anyhow::anyhow!("failed to publish setup nudge: {e}"))?;
@@ -1000,13 +1027,13 @@ mod tests {
#[test]
fn test_non_allowlisted_author_returns_no_nudge() {
// author_allowed = false → should return false regardless of other args.
let mut dedup: HashSet<EventId> = HashSet::new();
let dedup: HashSet<EventId> = HashSet::new();
let event_id = fake_event_id(0xAA);
let result = should_nudge_for_event(
event_id, false, // author NOT allowed
true, // filter matched — would otherwise nudge
&mut dedup,
&dedup,
);
assert!(!result, "non-allowlisted author must not produce a nudge");
@@ -1017,29 +1044,247 @@ mod tests {
);
}
/// Dedup suppresses a replay only once the send has been recorded.
///
/// The gate now *checks* the dedup set without mutating it; the caller
/// records the id after the relay accepts the nudge. So a replay is
/// accepted until that recording happens, and rejected afterwards.
#[test]
fn test_same_event_id_twice_nudges_exactly_once() {
// The first call with a given event-id should return true; the second
// call with the identical id must return false (replay dedup).
let mut dedup: HashSet<EventId> = HashSet::new();
let event_id = fake_event_id(0xBB);
let first = should_nudge_for_event(
event_id, true, // allowed
true, // matched
&mut dedup,
&dedup,
);
assert!(first, "first occurrence must be accepted");
// The caller records the id only after a successful publish.
dedup.insert(event_id);
// Simulate reconnect replay: same event arrives again.
let second = should_nudge_for_event(
event_id, true, // allowed
true, // matched
&mut dedup,
&dedup,
);
assert!(
!second,
"replay of the same event-id must be rejected (dedup)"
"replay of the same event-id must be rejected once recorded"
);
}
/// A nudge that failed to send stays retryable.
///
/// This is the F9/F11 regression: dedup used to be recorded by the gate
/// itself, before the publish was attempted. Combined with a WS publish
/// path that silently drops durable events while rate-gated or
/// disconnected, that made a dropped nudge permanent — the user never got
/// a reply, and the retry was suppressed by an entry recorded for a
/// message that was never delivered. The caller must leave the set
/// untouched on failure.
#[test]
fn test_failed_nudge_remains_retryable() {
let mut dedup: HashSet<EventId> = HashSet::new();
let event_id = fake_event_id(0xCC);
assert!(
should_nudge_for_event(event_id, true, true, &dedup),
"first attempt must be accepted"
);
// Publish fails: the caller records nothing.
assert!(
dedup.is_empty(),
"the gate must not record dedup on its own — recording is the \
caller's job, after the relay accepts the nudge"
);
assert!(
should_nudge_for_event(event_id, true, true, &dedup),
"a nudge that was never delivered must still be retryable"
);
// Now it succeeds and the caller records it.
dedup.insert(event_id);
assert!(
!should_nudge_for_event(event_id, true, true, &dedup),
"once delivered, the replay must be suppressed"
);
}
// ── nudge transport tests ─────────────────────────────────────────────────
//
// The dedup tests above are pure: they prove the gate stopped recording,
// but they cannot see which transport the nudge rides. That distinction is
// the other half of the fix, and reverting `publish_setup_nudge` to the WS
// publisher leaves every pure test green.
//
// These drive `publish_setup_nudge` against a real socket so the transport
// is observable: the request must arrive over HTTP, and a relay refusal
// must come back as `Err` rather than being swallowed.
/// A one-shot HTTP server that records the request line and body, and
/// replies with `status`. Returns the bound base URL and a handle to the
/// captured request.
async fn capturing_bridge(
status: &'static str,
body: &'static str,
) -> (
String,
std::sync::Arc<tokio::sync::Mutex<Vec<String>>>,
tokio::task::JoinHandle<()>,
) {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind test HTTP server");
let base_url = format!("http://{}", listener.local_addr().unwrap());
let captured = std::sync::Arc::new(tokio::sync::Mutex::new(Vec::new()));
let server_captured = captured.clone();
let server = tokio::spawn(async move {
while let Ok((mut socket, _)) = listener.accept().await {
let mut buf = vec![0u8; 16384];
let n = socket.read(&mut buf).await.unwrap_or(0);
server_captured
.lock()
.await
.push(String::from_utf8_lossy(&buf[..n]).to_string());
let response = format!(
"HTTP/1.1 {status}\r\nContent-Type: application/json\r\n\
Content-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
);
let _ = socket.write_all(response.as_bytes()).await;
}
});
(base_url, captured, server)
}
fn test_rest_client(base_url: String) -> crate::relay::RestClient {
crate::relay::RestClient {
http: reqwest::Client::new(),
base_url,
keys: nostr::Keys::generate(),
auth_tag_json: None,
}
}
fn trigger_event(keys: &nostr::Keys) -> nostr::Event {
nostr::EventBuilder::new(
nostr::Kind::Custom(KIND_STREAM_MESSAGE as u16),
"how do I set you up?",
)
.tags([])
.sign_with_keys(keys)
.expect("sign trigger event")
}
fn test_payload() -> SetupPayload {
SetupPayload {
agent_name: "Fizz".into(),
agent_pubkey: "aabbccddeeff0011".into(),
requirements: vec![],
}
}
/// The nudge goes out over the HTTP bridge, as a durable kind:9.
///
/// This is the transport half of F9/F11. `publish_event` (the WS path)
/// returns `Ok` as soon as the command is queued and never touches a
/// socket in-process, so if this fix were reverted no HTTP request would
/// arrive and this test fails on the request count.
#[tokio::test]
async fn nudge_is_published_over_the_http_bridge() {
let (base_url, captured, server) = capturing_bridge("200 OK", "{}").await;
let keys = nostr::Keys::generate();
let trigger = trigger_event(&keys);
publish_setup_nudge(
&test_rest_client(base_url),
&keys,
Uuid::new_v4(),
&trigger,
&test_payload(),
)
.await
.expect("a 200 from the bridge must be reported as success");
let requests = captured.lock().await;
assert_eq!(requests.len(), 1, "the nudge must reach the HTTP bridge");
let request = &requests[0];
assert!(
request.starts_with("POST /events "),
"nudge must POST to the events endpoint, got: {}",
request.lines().next().unwrap_or_default()
);
assert!(
request.contains(&format!("\"kind\":{KIND_STREAM_MESSAGE}")),
"the nudge is a durable kind:{KIND_STREAM_MESSAGE} message"
);
server.abort();
}
/// A relay refusal is surfaced to the caller, not swallowed.
///
/// This is what makes the dedup-after-`Ok` ordering meaningful: the caller
/// can only decline to record the dedup entry if the failure is actually
/// reported. The WS path returned `Ok` for a dropped publish, which is why
/// a rate-gated nudge was lost while being logged as "nudge published".
#[tokio::test]
async fn a_refused_nudge_is_reported_as_an_error() {
// 403 is non-retriable, so the retry ladder returns immediately.
let (base_url, captured, server) = capturing_bridge("403 Forbidden", "denied").await;
let keys = nostr::Keys::generate();
let trigger = trigger_event(&keys);
let result = publish_setup_nudge(
&test_rest_client(base_url),
&keys,
Uuid::new_v4(),
&trigger,
&test_payload(),
)
.await;
assert!(
result.is_err(),
"a relay refusal must be an Err — reporting Ok is the data-loss bug"
);
assert_eq!(captured.lock().await.len(), 1, "the attempt was made");
server.abort();
}
/// With no relay listening at all, the nudge fails rather than reporting
/// success. This is the disconnected case: the WS path accepted the
/// publish into a queue nobody was draining and told the caller `Ok`.
#[tokio::test]
async fn a_nudge_with_no_relay_listening_fails() {
// Bind and immediately drop the listener so the port is closed.
let base_url = {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind to claim a port");
format!("http://{}", listener.local_addr().unwrap())
};
let keys = nostr::Keys::generate();
let trigger = trigger_event(&keys);
let result = publish_setup_nudge(
&test_rest_client(base_url),
&keys,
Uuid::new_v4(),
&trigger,
&test_payload(),
)
.await;
assert!(
result.is_err(),
"a nudge that never reached a relay must not be reported as published"
);
}