From 3daeecc5c21cb80ebca9fbf2875bc50d658a0246 Mon Sep 17 00:00:00 2001 From: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz> Date: Mon, 17 Aug 2026 10:43:57 -0400 Subject: [PATCH] fix(acp): address thread context review findings Signed-off-by: Larry <8cf5a83f590ec0955b11647d1c88f796a98e088c30a492c58e0e46c3026ae7a4@buzz.block.builderlab.xyz> --- crates/buzz-acp/src/pool.rs | 5 +++ crates/buzz-acp/src/queue.rs | 51 ++++++++++++++---------- crates/buzz-acp/src/setup_mode.rs | 19 ++++++++- crates/buzz-cli/src/commands/messages.rs | 45 ++++++++++++--------- 4 files changed, 79 insertions(+), 41 deletions(-) diff --git a/crates/buzz-acp/src/pool.rs b/crates/buzz-acp/src/pool.rs index d15f5e2e3..27d8a0ef2 100644 --- a/crates/buzz-acp/src/pool.rs +++ b/crates/buzz-acp/src/pool.rs @@ -2975,12 +2975,14 @@ fn conversation_context_delta( messages, total, truncated, + root_kind, } => { let messages = filter(messages); (!messages.is_empty()).then_some(ConversationContext::Thread { messages, total, truncated, + root_kind, }) } ConversationContext::Dm { @@ -6451,6 +6453,7 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" ], total: 3, truncated: false, + root_kind: Some(buzz_core::kind::KIND_FORUM_POST), }; let delta = conversation_context_delta(Some(context), &delivered, &triggering) @@ -6460,11 +6463,13 @@ printf '%s\n' '{{"jsonrpc":"2.0","id":0,"result":{{"stopReason":"end_turn"}}}}'" messages, total, truncated, + root_kind, } => { assert_eq!(messages.len(), 1); assert_eq!(messages[0].event_id, "new"); assert_eq!(total, 3); assert!(!truncated); + assert_eq!(root_kind, Some(buzz_core::kind::KIND_FORUM_POST)); } _ => panic!("expected thread context"), } diff --git a/crates/buzz-acp/src/queue.rs b/crates/buzz-acp/src/queue.rs index c69996748..e1dbd0a92 100644 --- a/crates/buzz-acp/src/queue.rs +++ b/crates/buzz-acp/src/queue.rs @@ -1317,15 +1317,20 @@ fn append_channel_description(s: &mut String, channel_info: Option<&PromptChanne /// replies; in the channel branch a `Some` anchor means a human-facing /// top-level mention whose reply should open a new thread rooted at the /// triggering event. +#[derive(Clone, Copy)] +struct ContextHintState<'a> { + is_dm: bool, + has_conversation_context: bool, + had_delivered_events: bool, + reply_anchor: Option<&'a str>, + thread_root_kind: Option, +} + fn format_context_hints( channel_id: Uuid, channel_info: Option<&PromptChannelInfo>, thread_tags: &ThreadTags, - is_dm: bool, - has_conversation_context: bool, - conversation_context_had_delivered_events: bool, - reply_anchor: Option<&str>, - thread_root_kind: Option, + state: ContextHintState<'_>, ) -> String { let channel_display = match channel_info { Some(ci) => format!("{} (#{channel_id})", ci.name), @@ -1334,17 +1339,17 @@ fn format_context_hints( // DM check comes first — a DM reply has both thread tags AND is_dm=true, // and the scope should be "dm" (not "thread") because the agent is in a DM. - if is_dm { + if state.is_dm { let is_reply = thread_tags.root_event_id.is_some(); // DM replies use thread command because /messages excludes thread replies. // DM non-replies use get for recent conversation. - let ctx_hint = if has_conversation_context && is_reply { + let ctx_hint = if state.has_conversation_context && is_reply { "Thread context included below. Use `buzz messages thread --channel --event ` for full history if truncated." - } else if has_conversation_context { + } else if state.has_conversation_context { "Conversation context included below. Use `buzz messages get --channel ` for full history if truncated." - } else if conversation_context_had_delivered_events && is_reply { + } else if state.had_delivered_events && is_reply { "Earlier thread context was already delivered in this session. Use `buzz messages thread --channel --event ` to re-read the reply chain." - } else if conversation_context_had_delivered_events { + } else if state.had_delivered_events { "Earlier conversation context was already delivered in this session. Use `buzz messages get --channel ` to re-read it." } else if is_reply { "Use `buzz messages thread --channel --event ` to fetch the reply chain." @@ -1360,7 +1365,7 @@ fn format_context_hints( // If this is a DM reply, include thread structural info as supplementary. if let Some(ref root) = thread_tags.root_event_id { s.push_str(&format!("\nThread root: {root}")); - if let Some(kind) = thread_root_kind { + if let Some(kind) = state.thread_root_kind { s.push_str(&format!("\nThread root kind: {kind}")); } if let Some(ref parent) = thread_tags.parent_event_id { @@ -1368,15 +1373,15 @@ fn format_context_hints( s.push_str(&format!("\nParent: {parent}")); } } - if let Some(event_id) = reply_anchor { + if let Some(event_id) = state.reply_anchor { append_reply_instruction(&mut s, event_id); } } s } else if let Some(ref root) = thread_tags.root_event_id { - let ctx_hint = if has_conversation_context { + let ctx_hint = if state.has_conversation_context { "Thread context included below. Use `buzz messages thread --channel --event ` for full history if truncated." - } else if conversation_context_had_delivered_events { + } else if state.had_delivered_events { "Earlier thread context was already delivered in this session. Use `buzz messages thread --channel --event ` to re-read it." } else { "Use `buzz messages thread --channel --event ` to fetch thread context." @@ -1388,7 +1393,7 @@ fn format_context_hints( ); append_channel_description(&mut s, channel_info); s.push_str(&format!("\nThread root: {root}")); - if let Some(kind) = thread_root_kind { + if let Some(kind) = state.thread_root_kind { s.push_str(&format!("\nThread root kind: {kind}")); } if let Some(ref parent) = thread_tags.parent_event_id { @@ -1397,7 +1402,7 @@ fn format_context_hints( } } s.push_str(&format!("\n{ctx_hint}")); - if let Some(event_id) = reply_anchor { + if let Some(event_id) = state.reply_anchor { append_reply_instruction(&mut s, event_id); } s @@ -1411,7 +1416,7 @@ fn format_context_hints( s.push_str( "\nHint: Use `buzz messages get --channel ` for recent messages if needed.", ); - if let Some(event_id) = reply_anchor { + if let Some(event_id) = state.reply_anchor { append_new_thread_reply_instruction(&mut s, event_id); } s @@ -1643,11 +1648,13 @@ pub fn format_prompt(batch: &FlushBatch, args: &FormatPromptArgs<'_>) -> Vec continue; } - // Build and publish the setup nudge. + // Build and publish the setup nudge. Dedup records successful delivery, + // not an attempt: a transient lookup or relay failure must remain retryable. if let Err(e) = publish_setup_nudge( &publisher, &rest_client, @@ -472,6 +473,7 @@ pub(crate) async fn run_setup_listener(config: Config, payload: SetupPayload) -> ) .await { + record_nudge_failure(&mut nudged_event_ids, buzz_event.event.id); tracing::warn!("setup-mode: failed to publish nudge: {e}"); } else { tracing::info!( @@ -514,6 +516,10 @@ pub(crate) fn should_nudge_for_event( true } +fn record_nudge_failure(nudged_event_ids: &mut HashSet, event_id: EventId) { + nudged_event_ids.remove(&event_id); +} + fn is_setup_message_kind(kind: u32) -> bool { matches!( kind, @@ -1197,6 +1203,17 @@ mod tests { ); } + #[test] + fn failed_nudge_attempt_remains_retryable() { + let mut dedup = HashSet::new(); + let event_id = fake_event_id(0xCC); + assert!(should_nudge_for_event(event_id, true, true, &mut dedup)); + + record_nudge_failure(&mut dedup, event_id); + + assert!(should_nudge_for_event(event_id, true, true, &mut dedup)); + } + // ── availability round-trip tests ───────────────────────────────────────── // // These tests prove the desktop→buzz-acp→sentinel path preserves the diff --git a/crates/buzz-cli/src/commands/messages.rs b/crates/buzz-cli/src/commands/messages.rs index 40a9ae80b..86cb4a865 100644 --- a/crates/buzz-cli/src/commands/messages.rs +++ b/crates/buzz-cli/src/commands/messages.rs @@ -391,6 +391,20 @@ pub async fn cmd_get_messages( Ok(()) } +fn thread_query_filters( + channel_id: &str, + event_id: &str, + limit: u32, + depth_limit: Option, +) -> [serde_json::Value; 2] { + let mut replies = serde_json::json!({"kinds": [9, 40002, 40003, 40008, 45003], "#h": [channel_id], "#e": [event_id], "limit": limit}); + if let Some(depth) = depth_limit { + replies["depth_limit"] = serde_json::json!(depth); + } + let root = serde_json::json!({"ids": [event_id], "#h": [channel_id], "limit": 1}); + [replies, root] +} + pub async fn cmd_get_thread( client: &BuzzClient, channel_id: &str, @@ -403,23 +417,8 @@ pub async fn cmd_get_thread( validate_hex64(event_id)?; let limit = limit.unwrap_or(100).min(500); - // Two filters ORed in a single HTTP call: - // 1. Replies referencing this event via e-tag (no kind restriction) - // 2. The root event itself by ID - let mut reply_filter = serde_json::json!({ - "kinds": [9, 40002, 40003, 40008, 45003], - "#h": [channel_id], - "#e": [event_id], - "limit": limit - }); - if let Some(d) = depth_limit { - reply_filter["depth_limit"] = serde_json::json!(d); - } - let root_filter = serde_json::json!({ - "ids": [event_id], - "limit": 1 - }); - let resp = client.query_multi(&[reply_filter, root_filter]).await?; + let filters = thread_query_filters(channel_id, event_id, limit, depth_limit); + let resp = client.query_multi(&filters).await?; let mut events: Vec = serde_json::from_str(&resp).unwrap_or_default(); events.sort_by_key(|e| e.get("created_at").and_then(|v| v.as_u64()).unwrap_or(0)); let normalized = normalize_events(&events); @@ -995,13 +994,23 @@ mod tests { use super::{ event_mention_pubkeys, find_root_from_tags, match_profiles_by_name, merge_message_mentions, missing_members, normalize_explicit_mentions, parse_member_pubkeys, - resolve_names_to_pubkeys, + resolve_names_to_pubkeys, thread_query_filters, }; use buzz_sdk::mentions::{ extract_at_mentions_with_known, extract_at_names, match_names_to_profiles, MentionProfile, }; use serde_json::json; + #[test] + fn thread_root_query_is_scoped_to_requested_channel() { + let channel = "aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee"; + let event = "11".repeat(32); + let [replies, root] = thread_query_filters(channel, &event, 50, Some(3)); + assert_eq!(replies["#h"], json!([channel])); + assert_eq!(root["#h"], json!([channel])); + assert_eq!(root["ids"], json!([event])); + } + const ID_A: &str = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; const ID_B: &str = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; const PUBKEY: &str = "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc";