Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
288 changes: 235 additions & 53 deletions crates/buzz-acp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ mod usage;

pub use usage::TurnUsage;

use std::collections::{HashMap, HashSet};
use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::Arc;
use std::time::Duration;

Expand Down Expand Up @@ -287,6 +287,53 @@ async fn check_sibling_via_profile(
false
}

const OBSERVER_PUBLISH_INTERVAL: Duration = Duration::from_millis(167);
const OBSERVER_PUBLISH_LIMIT_PER_MINUTE: usize = 90;

struct ObserverPublishPacer {
next_publish: tokio::time::Instant,
published: VecDeque<tokio::time::Instant>,
}

impl ObserverPublishPacer {
fn new() -> Self {
Self {
// No initial burst: even the first snapshot frame waits for its slot.
next_publish: tokio::time::Instant::now() + OBSERVER_PUBLISH_INTERVAL,
published: VecDeque::with_capacity(OBSERVER_PUBLISH_LIMIT_PER_MINUTE),
}
}

async fn wait(&mut self) {
loop {
let now = tokio::time::Instant::now();
while self
.published
.front()
.is_some_and(|sent| now.duration_since(*sent) >= Duration::from_secs(60))
{
self.published.pop_front();
}

let minute_slot = self.published.front().and_then(|sent| {
(self.published.len() >= OBSERVER_PUBLISH_LIMIT_PER_MINUTE)
.then_some(*sent + Duration::from_secs(60))
});
let publish_at =
minute_slot.map_or(self.next_publish, |slot| slot.max(self.next_publish));
if publish_at > now {
tokio::time::sleep_until(publish_at).await;
continue;
}

let published_at = tokio::time::Instant::now();
self.published.push_back(published_at);
self.next_publish = published_at + OBSERVER_PUBLISH_INTERVAL;
return;
}
}
}

fn spawn_relay_observer_publisher(
observer: observer::ObserverHandle,
publisher: RelayEventPublisher,
Expand All @@ -296,68 +343,102 @@ fn spawn_relay_observer_publisher(
owner_pubkey: PublicKey,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
let mut coalescer = ObserverChunkCoalescer::default();
for event in observer.snapshot() {
for event in coalescer.ingest(event) {
publish_relay_observer_event(
&publisher,
&keys,
&agent_pubkey_hex,
&owner_pubkey_hex,
&owner_pubkey,
event,
)
.await;
}
// Subscribe BEFORE snapshotting so an event emitted between the two
// calls is never lost: it lands in the snapshot, the live receiver, or
// both. The overlap is deduped in the run loop via the snapshot's
// high-water `seq` (monotonic, assigned at emit).
let rx = observer.subscribe();
let snapshot = observer.snapshot();
run_relay_observer_publisher(
snapshot,
rx,
publisher,
keys,
agent_pubkey_hex,
owner_pubkey_hex,
owner_pubkey,
)
.await;
})
}

async fn run_relay_observer_publisher(
snapshot: Vec<observer::ObserverEvent>,
mut rx: tokio::sync::broadcast::Receiver<observer::ObserverEvent>,
publisher: RelayEventPublisher,
keys: nostr::Keys,
agent_pubkey_hex: String,
owner_pubkey_hex: String,
owner_pubkey: PublicKey,
) {
let mut coalescer = ObserverChunkCoalescer::default();
let mut pacer = ObserverPublishPacer::new();
let max_snapshot_seq = snapshot.iter().map(|event| event.seq).max().unwrap_or(0);
for event in snapshot {
for event in coalescer.ingest(event) {
publish_relay_observer_event(
&publisher,
&keys,
&agent_pubkey_hex,
&owner_pubkey_hex,
&owner_pubkey,
&mut pacer,
event,
)
.await;
}
}

let mut rx = observer.subscribe();
let mut flush_interval = tokio::time::interval(std::time::Duration::from_millis(500));
flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
result = rx.recv() => {
match result {
Ok(event) => {
for event in coalescer.ingest(event) {
publish_relay_observer_event(
&publisher, &keys, &agent_pubkey_hex,
&owner_pubkey_hex, &owner_pubkey, event,
).await;
}
let mut flush_interval = tokio::time::interval(std::time::Duration::from_millis(500));
flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
result = rx.recv() => {
match result {
Ok(event) => {
// Skip live events already delivered via the snapshot
// (the subscribe-before-snapshot overlap).
if event.seq <= max_snapshot_seq {
continue;
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(count)) => {
for event in coalescer.flush() {
publish_relay_observer_event(
&publisher, &keys, &agent_pubkey_hex,
&owner_pubkey_hex, &owner_pubkey, event,
).await;
}
tracing::warn!(dropped = count, "relay observer publisher lagged");
for event in coalescer.ingest(event) {
publish_relay_observer_event(
&publisher, &keys, &agent_pubkey_hex,
&owner_pubkey_hex, &owner_pubkey, &mut pacer, event,
).await;
}
Err(tokio::sync::broadcast::error::RecvError::Closed) => {
for event in coalescer.flush() {
publish_relay_observer_event(
&publisher, &keys, &agent_pubkey_hex,
&owner_pubkey_hex, &owner_pubkey, event,
).await;
}
break;
}
Err(tokio::sync::broadcast::error::RecvError::Lagged(count)) => {
for event in coalescer.flush() {
publish_relay_observer_event(
&publisher, &keys, &agent_pubkey_hex,
&owner_pubkey_hex, &owner_pubkey, &mut pacer, event,
).await;
}
tracing::warn!(dropped = count, "relay observer publisher lagged");
}
}
_ = flush_interval.tick() => {
// Periodic flush ensures live streaming even during continuous chunk delivery.
for event in coalescer.flush() {
publish_relay_observer_event(
&publisher, &keys, &agent_pubkey_hex,
&owner_pubkey_hex, &owner_pubkey, event,
).await;
Err(tokio::sync::broadcast::error::RecvError::Closed) => {
for event in coalescer.flush() {
publish_relay_observer_event(
&publisher, &keys, &agent_pubkey_hex,
&owner_pubkey_hex, &owner_pubkey, &mut pacer, event,
).await;
}
break;
}
}
}
_ = flush_interval.tick() => {
// Periodic flush ensures live streaming even during continuous chunk delivery.
for event in coalescer.flush() {
publish_relay_observer_event(
&publisher, &keys, &agent_pubkey_hex,
&owner_pubkey_hex, &owner_pubkey, &mut pacer, event,
).await;
}
}
}
})
}
}

#[derive(Default)]
Expand Down Expand Up @@ -638,8 +719,10 @@ async fn publish_relay_observer_event(
agent_pubkey_hex: &str,
owner_pubkey_hex: &str,
owner_pubkey: &PublicKey,
pacer: &mut ObserverPublishPacer,
mut event: observer::ObserverEvent,
) {
pacer.wait().await;
// Trim oversized frames to fit the plaintext cap rather than letting
// encrypt_observer_payload reject and drop them whole (silent telemetry loss).
fit_observer_event_to_budget(&mut event);
Expand Down Expand Up @@ -4070,6 +4153,105 @@ mod author_gate_tests {
}
}

#[cfg(test)]
mod observer_snapshot_race_tests {
use super::*;
use nostr::Keys;

fn emit_marker(observer: &observer::ObserverHandle, marker: &str) {
observer.emit(
"test_event",
None,
&observer::context_for(None, None, None),
serde_json::json!({ "marker": marker }),
);
}

/// An event emitted between `subscribe()` and `snapshot()` lands in BOTH
/// the snapshot and the live receiver; the seq high-water dedupe must
/// deliver it exactly once — and never lose events on either side of it.
#[tokio::test(start_paused = true)]
async fn overlap_between_subscribe_and_snapshot_publishes_exactly_once() {
let observer = observer::ObserverHandle::in_process();
let agent_keys = Keys::generate();
let owner_keys = Keys::generate();
let (publisher, mut published_rx) = RelayEventPublisher::test_pair();

// Before the publisher starts: replay-buffer only.
emit_marker(&observer, "before");
// The race window: emitted after subscribe() but before snapshot(),
// so it is present in the snapshot AND queued on the receiver.
let rx = observer.subscribe();
emit_marker(&observer, "overlap");
let snapshot = observer.snapshot();
assert_eq!(snapshot.len(), 2, "overlap event must be in the snapshot");
// After the snapshot: live receiver only.
emit_marker(&observer, "after");
// Close the broadcast channel so the run loop drains and exits.
drop(observer);

run_relay_observer_publisher(
snapshot,
rx,
publisher,
agent_keys.clone(),
agent_keys.public_key().to_hex(),
owner_keys.public_key().to_hex(),
owner_keys.public_key(),
)
.await;

// The run loop has exited, dropping the publisher; drain the forwarded
// events until the channel closes (deterministic — no try_recv race
// with the test_pair forwarding task).
let mut markers = Vec::new();
while let Some(event) = published_rx.recv().await {
let payload: serde_json::Value =
decrypt_observer_payload(&owner_keys, &event).expect("decrypt published frame");
markers.push(payload["payload"]["marker"].as_str().unwrap().to_string());
}
assert_eq!(
markers,
["before", "overlap", "after"],
"each event must be published exactly once, in order"
);
}
}

#[cfg(test)]
mod observer_publish_pacer_tests {
use super::*;

#[tokio::test(start_paused = true)]
async fn starts_without_a_burst_and_spaces_frames() {
let started = tokio::time::Instant::now();
let mut pacer = ObserverPublishPacer::new();

pacer.wait().await;
let first = tokio::time::Instant::now();
pacer.wait().await;
let second = tokio::time::Instant::now();

assert_eq!(first.duration_since(started), OBSERVER_PUBLISH_INTERVAL);
assert_eq!(second.duration_since(first), OBSERVER_PUBLISH_INTERVAL);
}

#[tokio::test(start_paused = true)]
async fn limits_frames_in_each_rolling_minute() {
let mut pacer = ObserverPublishPacer::new();
pacer.wait().await;
let first = tokio::time::Instant::now();
for _ in 1..OBSERVER_PUBLISH_LIMIT_PER_MINUTE {
pacer.wait().await;
}

pacer.wait().await;
let ninety_first = tokio::time::Instant::now();

assert_eq!(ninety_first.duration_since(first), Duration::from_secs(60));
}
}

#[cfg(test)]
mod observer_chunk_coalescer_tests {
use super::*;
Expand Down
Loading
Loading