Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
33 commits
Select commit Hold shift + click to select a range
afe57e5
Wake agents from verified workflow mentions
loganj Aug 27, 2026
3b4890a
Dispatch workflow mentions only after verification
loganj Aug 27, 2026
4056e28
Enforce current access for workflow wakes
loganj Aug 27, 2026
625ad3b
Allow safe kindless channel search
loganj Aug 27, 2026
e93fe76
Recover workflow wakes across reconnects
loganj Aug 28, 2026
356c3e4
Harden durable workflow wake admission
loganj Aug 28, 2026
99ff9d5
Repair workflow wake lifecycle boundaries
loganj Aug 28, 2026
662af24
Reconcile wake migrations with current foundation
loganj Aug 28, 2026
d6f4c78
Configure lifecycle fixtures before state construction
loganj Aug 28, 2026
0372b25
Keep lifecycle regressions in backend integration gate
loganj Aug 28, 2026
0ef5041
Create workflow fixture owner in tenant user table
loganj Aug 28, 2026
819ed22
Exercise authority admission with real Redis in lifecycle tests
loganj Aug 28, 2026
470285d
Use bridge filter array in wake removal regression
loganj Aug 28, 2026
23b9faa
Exercise wake revocation across WebSocket read and fanout paths
loganj Aug 28, 2026
6473304
Distinguish unavailable wake authority from revocation
loganj Aug 28, 2026
0a6889b
Pin PostgreSQL timeout recovery and clear bridge lint
loganj Aug 28, 2026
0938253
Exercise exhausted authority recovery through transport replay
loganj Aug 28, 2026
ae903c5
Handle background heartbeat in wake replay fixture
loganj Aug 28, 2026
5ec6929
Recover workflow wakes after interrupted authority bodies
loganj Aug 28, 2026
71056bb
Keep replay-guard outages distinct from authentication denials
loganj Aug 28, 2026
4db89eb
Verify captured revision revocation through deletion ingress
loganj Aug 28, 2026
929c697
Supply authenticated deletion scope and check fixture reads
loganj Aug 28, 2026
8f56bc4
Integrate workflow delivery with domain datastore tracing
loganj Aug 28, 2026
dc0a020
Place wake migrations after updated workflow foundation
loganj Aug 28, 2026
0d19d13
Point FTS migration fixture at renumbered wake migration
loganj Aug 28, 2026
ac456ab
test(search): assert raw NULL and negated-query privacy for workflow …
loganj Sep 1, 2026
2a2f263
fix(search): skip proven-safe wake FTS policies and preserve generate…
loganj Sep 1, 2026
e707882
fix(workflow): bind wake admission and resumes to captured authority
loganj Sep 1, 2026
db1e8c2
fix(workflow): retain deletion revocation on approval resume
loganj Sep 1, 2026
520989f
fix(acp): distinguish pending and absent wake signing identities
loganj Sep 1, 2026
8f7809c
test(acp): ignore heartbeat frames in wake replay assertion
loganj Sep 1, 2026
4a0dda2
ci: test wake migrations on the supported PostgreSQL 17
loganj Sep 1, 2026
15826a9
fix(workflows): rely on read authorization without rewriting event se…
loganj Sep 1, 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
3 changes: 2 additions & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -448,7 +448,8 @@ jobs:
contents: read
services:
postgres:
image: postgres:16
# Match the PostgreSQL 17 contract in VISION.md and docker-compose.yml.
image: postgres:17
env:
POSTGRES_USER: buzz
POSTGRES_PASSWORD: ${{ env.BUZZ_TEST_POSTGRES_PASSWORD }}
Expand Down
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.

1 change: 1 addition & 0 deletions crates/buzz-acp/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ reqwest = { workspace = true }
# Serialization
serde = { workspace = true }
serde_json = { workspace = true }
serde_yaml = { workspace = true }

# IDs
uuid = { workspace = true }
Expand Down
36 changes: 32 additions & 4 deletions crates/buzz-acp/src/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1324,6 +1324,7 @@ pub fn resolve_channel_filters(
) -> HashMap<Uuid, ChannelFilter> {
use buzz_core::kind::{
KIND_STREAM_MESSAGE, KIND_STREAM_REMINDER, KIND_WORKFLOW_APPROVAL_REQUESTED,
KIND_WORKFLOW_MENTION_WAKE,
};

let target_channels: Vec<Uuid> = if let Some(ref overrides) = config.channels_override {
Expand All @@ -1343,6 +1344,7 @@ pub fn resolve_channel_filters(
let kinds = config.kinds_override.clone().unwrap_or_else(|| {
vec![
KIND_STREAM_MESSAGE,
KIND_WORKFLOW_MENTION_WAKE,
KIND_WORKFLOW_APPROVAL_REQUESTED,
KIND_STREAM_REMINDER,
]
Expand Down Expand Up @@ -1426,6 +1428,7 @@ pub fn resolve_dynamic_channel_filter(
) -> Option<ChannelFilter> {
use buzz_core::kind::{
KIND_STREAM_MESSAGE, KIND_STREAM_REMINDER, KIND_WORKFLOW_APPROVAL_REQUESTED,
KIND_WORKFLOW_MENTION_WAKE,
};

// In Mentions/All mode, if the operator explicitly constrained channels
Expand All @@ -1448,6 +1451,7 @@ pub fn resolve_dynamic_channel_filter(
kinds: Some(config.kinds_override.clone().unwrap_or_else(|| {
vec![
KIND_STREAM_MESSAGE,
KIND_WORKFLOW_MENTION_WAKE,
KIND_WORKFLOW_APPROVAL_REQUESTED,
KIND_STREAM_REMINDER,
]
Expand Down Expand Up @@ -1596,13 +1600,37 @@ mod tests {
for ch in &channels {
let f = result.get(ch).expect("channel should be present");
assert!(f.require_mention, "mentions mode requires mention");
let kinds = f.kinds.as_ref().expect("should have kinds");
assert!(kinds.contains(&buzz_core::kind::KIND_STREAM_MESSAGE));
assert!(kinds.contains(&buzz_core::kind::KIND_WORKFLOW_APPROVAL_REQUESTED));
assert!(kinds.contains(&buzz_core::kind::KIND_STREAM_REMINDER));
assert_eq!(
f.kinds,
Some(vec![
buzz_core::kind::KIND_STREAM_MESSAGE,
buzz_core::kind::KIND_WORKFLOW_MENTION_WAKE,
buzz_core::kind::KIND_WORKFLOW_APPROVAL_REQUESTED,
buzz_core::kind::KIND_STREAM_REMINDER,
])
);
}
}

#[test]
fn test_mentions_mode_dynamic_default_kinds_include_workflow_wake() {
let config = test_config(SubscribeMode::Mentions);
let channel = Uuid::new_v4();
let filter = resolve_dynamic_channel_filter(&config, channel, &[])
.expect("dynamic channel should be subscribed");

assert!(filter.require_mention);
assert_eq!(
filter.kinds,
Some(vec![
buzz_core::kind::KIND_STREAM_MESSAGE,
buzz_core::kind::KIND_WORKFLOW_MENTION_WAKE,
buzz_core::kind::KIND_WORKFLOW_APPROVAL_REQUESTED,
buzz_core::kind::KIND_STREAM_REMINDER,
])
);
}

#[test]
fn test_mentions_mode_custom_kinds() {
let mut config = test_config(SubscribeMode::Mentions);
Expand Down
237 changes: 229 additions & 8 deletions crates/buzz-acp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ mod relay;
mod scope;
mod setup_mode;
mod usage;
mod workflow_wake;

pub use usage::TurnUsage;

Expand Down Expand Up @@ -414,6 +415,15 @@ mod inbound_author_gate {
.await
}

/// Durable wakes require a current-generation identity, unlike ordinary
/// admission's availability-oriented retained-key policy.
#[derive(Debug, PartialEq, Eq)]
pub(crate) enum WakeIdentity {
Ready(nostr::PublicKey),
Unavailable,
Retry,
}

pub(crate) struct InboundAuthorGate {
agent_pubkey_hex: String,
relay_self: Option<String>,
Expand Down Expand Up @@ -456,6 +466,47 @@ mod inbound_author_gate {
self.relay_self.as_deref()
}

/// Resolve the signing key through the same generation-fenced lifecycle
/// used by ordinary author admission. Durable wakes must not pin a
/// separate startup identity across relay reconnects.
pub(crate) async fn relay_identity_for_generation(
&mut self,
rest_client: &relay::RestClient,
event_generation: u64,
) -> Option<nostr::PublicKey> {
if refresh_needed(self.refreshed_generation, event_generation) {
let (relay_self, completed) =
refresh_relay_self(rest_client, self.relay_self.take(), "listener").await;
self.relay_self = relay_self;
if completed {
self.refreshed_generation = Some(event_generation);
}
}
self.relay_self
.as_deref()
.and_then(|key| nostr::PublicKey::from_hex(key).ok())
}

/// Resolve durable-wake identity without consuming a wake against a
/// stale key. Missing identity in a complete document is terminal for
/// this generation. Transient failures remain replayable, paced even
/// when discovery fails immediately (for example HTTP 500).
pub(crate) async fn wake_identity_for_generation(
&mut self,
rest_client: &relay::RestClient,
event_generation: u64,
) -> WakeIdentity {
let key = self
.relay_identity_for_generation(rest_client, event_generation)
.await;
if refresh_needed(self.refreshed_generation, event_generation) {
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
WakeIdentity::Retry
} else {
key.map_or(WakeIdentity::Unavailable, WakeIdentity::Ready)
}
}

/// Refresh relay identity, resolve channel trust, and apply trusted
/// workflow attribution and author policy for one listener event.
///
Expand All @@ -474,14 +525,8 @@ mod inbound_author_gate {
// Retry failed startup discovery on generation 0 as well as failed
// reconnect refreshes. Only an authoritative result completes the
// generation; transient failure retains the last verified key.
if refresh_needed(self.refreshed_generation, buzz_event.connection_generation) {
let (relay_self, completed) =
refresh_relay_self(rest_client, self.relay_self.take(), "listener").await;
self.relay_self = relay_self;
if completed {
self.refreshed_generation = Some(buzz_event.connection_generation);
}
}
self.relay_identity_for_generation(rest_client, buzz_event.connection_generation)
.await;
let is_dm = is_dm_channel(buzz_event.channel_id, channel_info).await;
self.evaluate_with_channel_trust(
&buzz_event.event,
Expand Down Expand Up @@ -2605,6 +2650,7 @@ async fn tokio_main() -> Result<()> {
kinds: config.kinds_override.clone().unwrap_or_else(|| {
vec![
KIND_STREAM_MESSAGE,
buzz_core::kind::KIND_WORKFLOW_MENTION_WAKE,
KIND_WORKFLOW_APPROVAL_REQUESTED,
KIND_STREAM_REMINDER,
]
Expand Down Expand Up @@ -3135,6 +3181,74 @@ async fn tokio_main() -> Result<()> {
match buzz_event {
Some(buzz_event) => {
let kind_u32 = buzz_event.event.kind.as_u16() as u32;
// Revision-labelled messages dispatch only through their
// durable wake, even while identity discovery is unavailable.
if workflow_wake::requires_verified_wake(&buzz_event.event) {
continue;
}

let buzz_event = if kind_u32
== buzz_core::kind::KIND_WORKFLOW_MENTION_WAKE
{
let Some((wake, workflow_relay_pubkey)) = workflow_wake::authenticate_for_listener(
&mut author_gate_ctx, &relay, &ctx.rest_client, &buzz_event,
).await else {
continue;
};
let authority = match ctx
.rest_client
.workflow_wake_authority(wake.run_id(), &wake.message_event_id())
.await
{
Ok(authority) => authority,
Err(error) if error.is_transient() => {
// HTTP-status failures exhaust bounded retries; body
// interruptions also return transient after pacing.
// Transport dedup recorded this relay-signed wake, but
// dispatch has not occurred. Re-admit it for filtered
// replay rather than losing it or bypassing verification.
if let Err(replay_error) = relay
.replay_event(
buzz_event.channel_id,
buzz_event.event.id.to_hex(),
buzz_event.event.created_at.as_secs(),
)
.await
{
tracing::warn!(
%replay_error,
"failed to arrange workflow wake authority replay"
);
}
tracing::warn!(%error, "workflow wake authority unavailable; replay queued");
continue;
}
Err(error) => {
// 403/404 and malformed authority bundles are terminal:
// replays cannot make a rejected or invalid authority safe.
tracing::warn!(%error, "workflow wake authority rejected");
continue;
}
};
let Some((message, _signed_author)) = workflow_wake::verify(
&buzz_event.event,
authority,
workflow_relay_pubkey,
config.keys.public_key(),
buzz_event.channel_id,
) else {
tracing::warn!("workflow wake authority verification failed");
continue;
};
relay::BuzzEvent {
channel_id: buzz_event.channel_id,
connection_generation: buzz_event.connection_generation,
event: message,
}
} else {
buzz_event
};
let kind_u32 = buzz_event.event.kind.as_u16() as u32;

if kind_u32 == KIND_MEMBER_ADDED_NOTIFICATION
|| kind_u32 == KIND_MEMBER_REMOVED_NOTIFICATION
Expand Down Expand Up @@ -7008,6 +7122,113 @@ mod author_gate_tests {
/// The first authorized event after reconnect must restore attribution
/// through the same decision boundary both listeners use, without a
/// separate identity-refresh call.
#[tokio::test]
async fn durable_wake_identity_tracks_listener_generation_rotation() {
let old = nostr::Keys::generate();
let new = nostr::Keys::generate();
let agent = nostr::Keys::generate();
let (rest, server) = nip11_scripted_server(std::collections::VecDeque::from([
Ok(serde_json::json!({"self": old.public_key().to_hex()})),
Ok(serde_json::json!({"self": new.public_key().to_hex()})),
]))
.await;
let mut gate =
InboundAuthorGate::connect(&rest, &agent.public_key().to_hex(), "test").await;
assert_eq!(
gate.relay_identity_for_generation(&rest, 0).await,
Some(old.public_key())
);
let channel = Uuid::new_v4();
let wake = buzz_core::workflow_wake::WorkflowMentionWake::new(
agent.public_key(),
channel,
Uuid::new_v4(),
nostr::EventId::from_byte_array([1; 32]),
nostr::EventId::from_byte_array([2; 32]),
)
.sign(&new)
.unwrap();
let key = gate.relay_identity_for_generation(&rest, 1).await.unwrap();
assert_eq!(key, new.public_key());
assert!(workflow_wake::authenticate(&wake, key).is_some());
assert!(workflow_wake::authenticate(&wake, old.public_key()).is_none());
// The ordinary gate consumes that exact identity, not a second startup cache.
assert_eq!(
gate.relay_identity_for_test(),
Some(new.public_key().to_hex().as_str())
);
server.abort();
}

#[tokio::test]
async fn wake_identity_missing_is_terminal_until_a_new_generation() {
use inbound_author_gate::WakeIdentity;
let keys = nostr::Keys::generate();
let (rest, server) = nip11_scripted_server(std::collections::VecDeque::from([
Ok(serde_json::json!({})),
Ok(serde_json::json!({})), // /info fallback also has no identity
Ok(serde_json::json!({"self": keys.public_key().to_hex()})),
]))
.await;
let mut gate = InboundAuthorGate::connect(&rest, "agent", "test").await;
// The scripted valid key must remain unread: repeated deliveries on a
// completed generation cannot make progress and must not request replay.
for _ in 0..3 {
assert_eq!(
gate.wake_identity_for_generation(&rest, 0).await,
WakeIdentity::Unavailable
);
}
assert_eq!(
gate.wake_identity_for_generation(&rest, 1).await,
WakeIdentity::Ready(keys.public_key())
);
server.abort();
}

#[tokio::test]
async fn wake_identity_transient_rotation_is_paced_and_replayable() {
use inbound_author_gate::WakeIdentity;
let old = nostr::Keys::generate();
let new = nostr::Keys::generate();
let agent = nostr::Keys::generate();
let (rest, server) = nip11_scripted_server(std::collections::VecDeque::from([
Ok(serde_json::json!({"self": old.public_key().to_hex()})),
Err(()),
Err(()), // /info fallback also fails
Ok(serde_json::json!({"self": new.public_key().to_hex()})),
]))
.await;
let mut gate =
InboundAuthorGate::connect(&rest, &agent.public_key().to_hex(), "test").await;
let wake = buzz_core::workflow_wake::WorkflowMentionWake::new(
agent.public_key(),
Uuid::new_v4(),
Uuid::new_v4(),
nostr::EventId::from_byte_array([1; 32]),
nostr::EventId::from_byte_array([2; 32]),
)
.sign(&new)
.unwrap();
let started = tokio::time::Instant::now();
assert_eq!(
gate.wake_identity_for_generation(&rest, 1).await,
WakeIdentity::Retry
);
assert!(started.elapsed() >= std::time::Duration::from_secs(1));
// Ordinary admission may retain A, but the durable-wake boundary must
// not authenticate against A and permanently consume B's wake.
assert_eq!(
gate.relay_identity_for_test(),
Some(old.public_key().to_hex().as_str())
);
let WakeIdentity::Ready(key) = gate.wake_identity_for_generation(&rest, 1).await else {
panic!("replayed wake must recover the new signing identity");
};
assert!(workflow_wake::authenticate(&wake, key).is_some());
server.abort();
}

#[tokio::test]
async fn test_gate_refresh_arms_attribution_after_reconnect() {
let relay_keys = nostr::Keys::generate();
Expand Down
Loading
Loading