Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
34 commits
Select commit Hold shift + click to select a range
69beeac
Reconstruct mission truth from signed observer evidence
3stepwin Aug 21, 2026
dbeeba3
Close Activity Ledger durability and authority gaps
3stepwin Aug 21, 2026
350ade4
Bind Activity Ledger proof to durable owner evidence
3stepwin Aug 21, 2026
fe46de9
Reconcile Activity Ledger with archive retention mainline
3stepwin Aug 21, 2026
4cfe7c3
Satisfy strict archive command gates
3stepwin Aug 21, 2026
117da0b
Reconcile Activity Ledger with latest mobile mainline
3stepwin Aug 21, 2026
e98684c
Close final Activity Ledger review gaps
3stepwin Aug 21, 2026
4669e7f
Close fresh Activity Ledger review findings
3stepwin Aug 21, 2026
170aa87
Keep owner Today activity current and publishable
3stepwin Aug 21, 2026
4a58983
Keep Activity Ledger failure paths fail-closed
3stepwin Aug 21, 2026
3358d47
Make archive byte-cap overflow restart-durable
3stepwin Aug 21, 2026
818c9c9
Keep Activity Ledger current under late evidence
3stepwin Aug 21, 2026
f391efc
Keep Activity Ledger fail-closed across lifecycle boundaries
3stepwin Aug 22, 2026
fbb892d
Keep Activity Ledger bounded without losing proof
3stepwin Aug 22, 2026
025c74b
Merge remote-tracking branch 'origin/main' into codex/buzz-activity-l…
3stepwin Aug 22, 2026
0ddaf69
Keep Today available across transient boundaries
3stepwin Aug 22, 2026
affc7eb
Merge remote-tracking branch 'origin/main' into codex/buzz-activity-l…
3stepwin Aug 22, 2026
5bd9af4
Keep Activity Ledger authority scoped and continuous
3stepwin Aug 22, 2026
4502822
Keep journal verification current across every frame
3stepwin Aug 22, 2026
10b651e
Bind Activity Ledger authority to the exact agent
3stepwin Aug 22, 2026
f8d18a6
Keep Today continuous across local midnight
3stepwin Aug 22, 2026
63d3b0c
Keep Activity Ledger proof and Today fail-closed
3stepwin Aug 22, 2026
0b8fbba
Keep Activity Ledger evidence durable across boundaries
3stepwin Aug 22, 2026
d77597f
Close Activity Ledger's final fail-closed gaps
3stepwin Aug 22, 2026
260df0c
Close Activity Ledger's final review gaps
3stepwin Aug 22, 2026
515ede4
Disclose observer frames without attributable time
3stepwin Aug 22, 2026
95a4d66
Fence Today publication against concurrent exclusions
3stepwin Aug 22, 2026
3e1dc85
Keep Today stable across concurrent archive changes
3stepwin Aug 22, 2026
16a4c98
Reject impossible observer calendar timestamps
3stepwin Aug 22, 2026
e75848a
Keep Today truthful under sustained archive activity
3stepwin Aug 22, 2026
bf94d16
Keep missing telemetry from becoming completion proof
3stepwin Aug 22, 2026
d3167a2
Keep final observer evidence through relay shutdown
3stepwin Aug 22, 2026
0ae72be
Keep archive concurrency failures observable
3stepwin Aug 22, 2026
8201c27
Keep archive opens read-only after migration
3stepwin Aug 22, 2026
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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

299 changes: 96 additions & 203 deletions crates/buzz-acp/src/lib.rs

Large diffs are not rendered by default.

46 changes: 46 additions & 0 deletions crates/buzz-acp/src/observer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ struct ObserverInner {
tx: broadcast::Sender<ObserverEvent>,
buffer: Mutex<VecDeque<ObserverEvent>>,
seq: AtomicU64,
replay_dropped: AtomicU64,
}

fn new_observer_handle() -> ObserverHandle {
Expand All @@ -49,6 +50,7 @@ fn new_observer_handle() -> ObserverHandle {
tx,
buffer: Mutex::new(VecDeque::with_capacity(OBSERVER_BUFFER_CAP)),
seq: AtomicU64::new(1),
replay_dropped: AtomicU64::new(0),
}),
}
}
Expand Down Expand Up @@ -85,11 +87,13 @@ impl ObserverHandle {
}

/// Subscribe to live observer events.
#[cfg(test)]
pub fn subscribe(&self) -> broadcast::Receiver<ObserverEvent> {
self.inner.tx.subscribe()
}

/// Return the current replay buffer.
#[cfg(test)]
pub fn snapshot(&self) -> Vec<ObserverEvent> {
match self.inner.buffer.lock() {
Ok(buffer) => buffer.iter().cloned().collect(),
Expand All @@ -100,6 +104,23 @@ impl ObserverHandle {
}
}

/// Atomically subscribe, snapshot, and take replay-buffer overflow count.
pub(crate) fn subscribe_with_snapshot(
&self,
) -> (Vec<ObserverEvent>, broadcast::Receiver<ObserverEvent>, u64) {
let buffer = match self.inner.buffer.lock() {
Ok(buffer) => buffer,
Err(error) => {
tracing::warn!(target: "observer", "observer replay buffer lock poisoned: {error}");
error.into_inner()
}
};
let receiver = self.inner.tx.subscribe();
let snapshot = buffer.iter().cloned().collect();
let dropped = self.inner.replay_dropped.swap(0, Ordering::AcqRel);
(snapshot, receiver, dropped)
}

/// Emit a local observer event.
pub fn emit(
&self,
Expand All @@ -124,6 +145,7 @@ impl ObserverHandle {
Ok(mut buffer) => {
if buffer.len() >= OBSERVER_BUFFER_CAP {
buffer.pop_front();
self.inner.replay_dropped.fetch_add(1, Ordering::Relaxed);
}
buffer.push_back(event.clone());
}
Expand Down Expand Up @@ -164,3 +186,27 @@ pub fn context_for_turn(
started_at: Some(started_at),
}
}

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

#[test]
fn subscribe_snapshot_reports_replay_overflow_in_source_units() {
let observer = ObserverHandle::in_process();
for seq in 0..OBSERVER_BUFFER_CAP + 3 {
observer.emit(
"test",
None,
&ObserverContext::default(),
serde_json::json!({ "seq": seq }),
);
}

let (snapshot, _receiver, dropped) = observer.subscribe_with_snapshot();
assert_eq!(snapshot.len(), OBSERVER_BUFFER_CAP);
assert_eq!(dropped, 3);
let (_, _, dropped_again) = observer.subscribe_with_snapshot();
assert_eq!(dropped_again, 0, "reported gaps are not double counted");
}
}
295 changes: 295 additions & 0 deletions crates/buzz-acp/src/observer_gap.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,295 @@
//! Signed, in-band accounting for observer events lost before relay archival.

use crate::observer;
use buzz_core::observer::{encrypt_observer_payload, OBSERVER_FRAME_TELEMETRY};
use nostr::{Event, Keys, PublicKey};

pub(crate) const OBSERVER_TELEMETRY_GAP_KIND: &str = "observer_telemetry_gap";

#[derive(Clone, Copy, Debug)]
pub(crate) enum ObserverGapReason {
ReplayBufferOverflow,
PublishQueueEviction,
BroadcastLag,
PublishFailure,
RelayQueueEviction,
}

#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub(crate) struct ObserverGapCounts {
replay_buffer_overflow: u64,
publish_queue_eviction: u64,
broadcast_lag: u64,
publish_failure: u64,
relay_queue_eviction: u64,
first_observed_at: Option<String>,
last_observed_at: Option<String>,
}

pub(crate) struct PendingObserverFrame {
pub(crate) event: observer::ObserverEvent,
pub(crate) source_events: u64,
pub(crate) reported_gaps: ObserverGapCounts,
}

impl ObserverGapCounts {
pub(crate) fn record(&mut self, reason: ObserverGapReason, count: u64) {
if count == 0 {
return;
}
let timestamp = chrono::Utc::now().to_rfc3339();
self.first_observed_at
.get_or_insert_with(|| timestamp.clone());
self.last_observed_at = Some(timestamp);
let bucket = match reason {
ObserverGapReason::ReplayBufferOverflow => &mut self.replay_buffer_overflow,
ObserverGapReason::PublishQueueEviction => &mut self.publish_queue_eviction,
ObserverGapReason::BroadcastLag => &mut self.broadcast_lag,
ObserverGapReason::PublishFailure => &mut self.publish_failure,
ObserverGapReason::RelayQueueEviction => &mut self.relay_queue_eviction,
};
*bucket = bucket.saturating_add(count);
}

pub(crate) fn merge(&mut self, other: Self) {
if other.is_empty() {
return;
}
if self.first_observed_at.is_none() {
self.first_observed_at = other.first_observed_at.clone();
}
self.last_observed_at = other.last_observed_at.clone();
self.replay_buffer_overflow = self
.replay_buffer_overflow
.saturating_add(other.replay_buffer_overflow);
self.publish_queue_eviction = self
.publish_queue_eviction
.saturating_add(other.publish_queue_eviction);
self.broadcast_lag = self.broadcast_lag.saturating_add(other.broadcast_lag);
self.publish_failure = self.publish_failure.saturating_add(other.publish_failure);
self.relay_queue_eviction = self
.relay_queue_eviction
.saturating_add(other.relay_queue_eviction);
}

pub(crate) fn is_empty(&self) -> bool {
self.total() == 0
}

pub(crate) fn total(&self) -> u64 {
self.replay_buffer_overflow
.saturating_add(self.publish_queue_eviction)
.saturating_add(self.broadcast_lag)
.saturating_add(self.publish_failure)
.saturating_add(self.relay_queue_eviction)
}

pub(crate) fn into_event(
self,
template: Option<&observer::ObserverEvent>,
) -> observer::ObserverEvent {
let total = self.total();
observer::ObserverEvent {
seq: template.map_or(0, |event| event.seq),
timestamp: chrono::Utc::now().to_rfc3339(),
kind: OBSERVER_TELEMETRY_GAP_KIND.to_string(),
agent_index: template.and_then(|event| event.agent_index),
channel_id: template.and_then(|event| event.channel_id.clone()),
session_id: None,
turn_id: None,
started_at: None,
payload: serde_json::json!({
"droppedEvents": total,
"reasonCounts": {
"replayBufferOverflow": self.replay_buffer_overflow,
"publishQueueEviction": self.publish_queue_eviction,
"broadcastLag": self.broadcast_lag,
"publishFailure": self.publish_failure,
"relayQueueEviction": self.relay_queue_eviction,
},
"firstObservedAt": self.first_observed_at,
"lastObservedAt": self.last_observed_at,
"scope": "publisher_global",
}),
}
}
}

fn observer_owner(event: &Event) -> Result<PublicKey, String> {
let owner = event
.tags
.iter()
.find_map(|tag| {
let parts = tag.as_slice();
(parts.first().map(String::as_str) == Some("p"))
.then(|| parts.get(1))
.flatten()
})
.ok_or("observer frame is missing owner p tag")?;
PublicKey::from_hex(owner).map_err(|error| format!("invalid observer owner: {error}"))
}

/// Count decrypted observer events represented by one signed relay frame.
pub(crate) fn represented_event_count(keys: Option<&Keys>, event: &Event) -> u64 {
let Some(keys) = keys else { return 1 };
let Ok(owner) = observer_owner(event) else {
return 1;
};
let Ok(plaintext) =
nostr::nips::nip44::decrypt(keys.secret_key(), &owner, event.content.as_str())
else {
return 1;
};
let Ok(value) = serde_json::from_str::<serde_json::Value>(&plaintext) else {
return 1;
};
represented_value_count(&value)
}

fn represented_value_count(value: &serde_json::Value) -> u64 {
let Some(object) = value.as_object() else {
return 1;
};
if object.get("kind").and_then(serde_json::Value::as_str) == Some(OBSERVER_TELEMETRY_GAP_KIND) {
return object
.get("payload")
.and_then(|payload| payload.get("droppedEvents"))
.and_then(serde_json::Value::as_u64)
.filter(|count| *count > 0)
.unwrap_or(1);
}
if let Some(events) = object
.get("payload")
.and_then(|payload| payload.get("events"))
.and_then(serde_json::Value::as_array)
{
return events
.iter()
.map(represented_value_count)
.sum::<u64>()
.max(1);
}
1
}

/// Build a replacement gap frame using the dropped frame's owner scope.
pub(crate) fn signed_relay_gap(
keys: &Keys,
template: &Event,
dropped_events: u64,
) -> Result<Event, String> {
let owner = observer_owner(template)?;
let mut gaps = ObserverGapCounts::default();
gaps.record(ObserverGapReason::RelayQueueEviction, dropped_events);
let payload = gaps.into_event(None);
let encrypted = encrypt_observer_payload(keys, &owner, &payload)
.map_err(|error| format!("encrypt relay gap: {error}"))?;
buzz_sdk::build_agent_observer_frame(
&owner.to_hex(),
&keys.public_key().to_hex(),
OBSERVER_FRAME_TELEMETRY,
&encrypted,
)
.map_err(|error| format!("build relay gap: {error}"))?
.sign_with_keys(keys)
.map_err(|error| format!("sign relay gap: {error}"))
}

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

#[test]
fn gap_counts_merge_and_build_a_frame_payload() {
let mut gaps = ObserverGapCounts::default();
gaps.record(ObserverGapReason::BroadcastLag, 2);
let mut later = ObserverGapCounts::default();
later.record(ObserverGapReason::PublishFailure, 3);
gaps.merge(later);

let event = gaps.into_event(None);
assert_eq!(event.kind, OBSERVER_TELEMETRY_GAP_KIND);
assert_eq!(event.payload["droppedEvents"], 5);
assert_eq!(event.payload["reasonCounts"]["broadcastLag"], 2);
assert_eq!(event.payload["reasonCounts"]["publishFailure"], 3);
assert_eq!(event.payload["scope"], "publisher_global");
}

#[test]
fn relay_gap_counts_batch_members_and_remains_owner_decryptable() {
let agent = Keys::generate();
let owner = Keys::generate();
let batch = serde_json::json!({
"seq": 2,
"timestamp": "2026-08-22T12:00:00Z",
"kind": "batch",
"payload": {"events": [{"seq": 1}, {"seq": 2}]},
});
let encrypted = encrypt_observer_payload(&agent, &owner.public_key(), &batch).unwrap();
let template = buzz_sdk::build_agent_observer_frame(
&owner.public_key().to_hex(),
&agent.public_key().to_hex(),
OBSERVER_FRAME_TELEMETRY,
&encrypted,
)
.unwrap()
.sign_with_keys(&agent)
.unwrap();

assert_eq!(represented_event_count(Some(&agent), &template), 2);
let gap = signed_relay_gap(&agent, &template, 2).unwrap();
let payload: serde_json::Value =
buzz_core::observer::decrypt_observer_payload(&owner, &gap).unwrap();
assert_eq!(payload["kind"], OBSERVER_TELEMETRY_GAP_KIND);
assert_eq!(payload["payload"]["droppedEvents"], 2);
assert_eq!(payload["payload"]["reasonCounts"]["relayQueueEviction"], 2);
}

#[test]
fn represented_event_count_recursively_preserves_nested_gap_totals() {
let agent = Keys::generate();
let owner = Keys::generate();
let nested = serde_json::json!({
"seq": 9,
"timestamp": "2026-08-22T12:00:00Z",
"kind": "batch",
"payload": {
"events": [
{
"seq": 7,
"timestamp": "2026-08-22T11:59:58Z",
"kind": OBSERVER_TELEMETRY_GAP_KIND,
"payload": {
"droppedEvents": 11,
"reasonCounts": {
"relayQueueEviction": 11
}
}
},
{
"seq": 8,
"timestamp": "2026-08-22T11:59:59Z",
"kind": "tool_call",
"payload": { "ok": true }
}
]
}
});
let encrypted = encrypt_observer_payload(&agent, &owner.public_key(), &nested).unwrap();
let template = buzz_sdk::build_agent_observer_frame(
&owner.public_key().to_hex(),
&agent.public_key().to_hex(),
OBSERVER_FRAME_TELEMETRY,
&encrypted,
)
.unwrap()
.sign_with_keys(&agent)
.unwrap();

assert_eq!(
represented_event_count(Some(&agent), &template),
12,
"batch accounting must preserve signed gap droppedEvents recursively"
);
}
}
Loading