diff --git a/crates/buzz-acp/src/acp.rs b/crates/buzz-acp/src/acp.rs index ccf0bf3c8..b44a7a51b 100644 --- a/crates/buzz-acp/src/acp.rs +++ b/crates/buzz-acp/src/acp.rs @@ -215,27 +215,6 @@ pub struct AcpClient { /// deltas. Both goose and buzz-agent emit this notification; goose gates /// on client capability advertisement, buzz-agent emits unconditionally. goose_usage: UsageTracker, - /// Accumulated `agent_message_chunk` text for the current turn. Used by the - /// content-delivery fallback: weak local models (e.g. via Buzz shared - /// compute) often answer a conversational prompt in plain assistant - /// `content` instead of calling `buzz messages send`, which would otherwise - /// be silently dropped (buzz-agent's output is its tool calls; streamed - /// text is observability-only). Reset at the start of every turn. - turn_message_text: String, - /// Whether a message publish was CONFIRMED this turn. When true the - /// fallback does NOT fire — the agent delivered its own reply. Confirmation - /// comes from the terminal `tool_call_update` outcome (successful - /// completion, ideally carrying the CLI's `{"accepted":true,...}` envelope - /// in the tool output), never from the tool call's input text alone: - /// intent is not delivery. Reset at the start of every turn. - turn_sent_message: bool, - /// Publish *candidates* for the current turn: `toolCallId`s whose - /// `tool_call` input matched the publish signature, awaiting a terminal - /// `tool_call_update`. A candidate that completes successfully confirms - /// delivery; one that fails (or completes with `isError`, or whose output - /// lacks the publish acknowledgement) is discarded so the fallback stays - /// armed. Reset at the start of every turn. - turn_publish_candidates: std::collections::HashSet, } /// Recursively merge `overlay` into `base`, with `overlay` winning on scalar/shape @@ -586,9 +565,6 @@ impl AcpClient { steering_supported: false, steer_rx: None, goose_usage: UsageTracker::default(), - turn_message_text: String::new(), - turn_sent_message: false, - turn_publish_candidates: std::collections::HashSet::new(), }) } @@ -799,11 +775,6 @@ impl AcpClient { // misattributed to this turn. self.goose_usage.begin_turn(session_id); - // Reset the content-delivery fallback trackers for this turn. - self.turn_message_text.clear(); - self.turn_sent_message = false; - self.turn_publish_candidates.clear(); - self.last_prompt_id = Some(self.next_id); let id = self.next_id; self.next_id += 1; @@ -908,30 +879,6 @@ impl AcpClient { self.goose_usage.take() } - /// Take the accumulated assistant `content` text for the completed turn, - /// if and only if a message publish was NOT confirmed this turn. - /// - /// Returns `Some(trimmed_text)` when the turn produced streamed assistant - /// content but no publish tool call was confirmed delivered (see - /// `publish_outcome_confirms_delivery`) — the caller then delivers it as - /// the channel reply (content-delivery fallback). Returns `None` when the - /// agent's own send was confirmed, when there was no content, or when the - /// content is only a bare acknowledgement (which the base prompt forbids - /// publishing). Clears the buffer either way. - pub fn take_undelivered_turn_message(&mut self) -> Option { - let text = std::mem::take(&mut self.turn_message_text); - let sent = self.turn_sent_message; - self.turn_sent_message = false; - if sent { - return None; - } - let trimmed = text.trim(); - if trimmed.is_empty() || is_bare_acknowledgement(trimmed) { - return None; - } - Some(trimmed.to_string()) - } - /// Install a per-turn steer request channel for goose-native /// non-cancelling mid-turn delivery. /// @@ -1783,10 +1730,6 @@ impl AcpClient { "agent_message_chunk" => { if let Some(text) = update["content"]["text"].as_str() { tracing::info!(target: "acp::stream", "{text}"); - // Accumulate for the content-delivery fallback (see - // `turn_message_text`). Streamed assistant text is otherwise - // observability-only and never posted to the channel. - self.turn_message_text.push_str(text); } false } @@ -1800,28 +1743,6 @@ impl AcpClient { .and_then(|v| v.as_str()) .unwrap_or("unknown"); tracing::info!(target: "acp::tool", "tool_call: {title} ({kind})"); - // Register a message-publish CANDIDATE for the content-delivery - // fallback. Intent is not delivery: the flag that suppresses - // the fallback (`turn_sent_message`) is only set when this - // call's terminal `tool_call_update` confirms success — a send - // that fails must leave the fallback armed, otherwise the - // feature silently drops the exact reply it exists to save. - if tool_call_is_message_publish(update) { - if let Some(id) = update.get("toolCallId").and_then(|v| v.as_str()) { - self.turn_publish_candidates.insert(id.to_string()); - tracing::debug!( - "publish candidate registered ({title}, toolCallId={id}); \ - awaiting terminal outcome" - ); - } - // Some agents emit `tool_call` already carrying a terminal - // status (single-event shape). Handle it like an update. - if let Some(confirmed) = publish_outcome_confirms_delivery(update) { - if confirmed { - self.turn_sent_message = true; - } - } - } true } "tool_call_update" => { @@ -1831,24 +1752,6 @@ impl AcpClient { .unwrap_or("?"); let status = update.get("status").and_then(|v| v.as_str()).unwrap_or("?"); tracing::info!(target: "acp::tool", "tool_call_update: {tool_id} → {status}"); - // Resolve a pending publish candidate on its terminal outcome. - if self.turn_publish_candidates.contains(tool_id) { - if let Some(confirmed) = publish_outcome_confirms_delivery(update) { - self.turn_publish_candidates.remove(tool_id); - if confirmed { - self.turn_sent_message = true; - tracing::debug!( - "publish confirmed (toolCallId={tool_id}); \ - content-delivery fallback stays dormant this turn" - ); - } else { - tracing::debug!( - "publish attempt failed (toolCallId={tool_id}); \ - content-delivery fallback stays armed" - ); - } - } - } false } "plan" => { @@ -2270,181 +2173,6 @@ pub fn model_in_catalog( }) } -/// Return true if a `tool_call` session update *looks like* a Buzz message -/// publish (kind 9 / forum post / comment) from its input. This only registers -/// a publish **candidate** — delivery is confirmed separately from the call's -/// terminal outcome by [`publish_outcome_confirms_delivery`]. -/// -/// The publish path is the dev-mcp `shell` tool running `buzz messages send` -/// (or `buzz social publish`), so the tool name alone is not enough — inspect -/// `rawInput` (the command/args) for the CLI publish signature. The haystack is -/// normalized (quotes/commas/brackets → spaces, whitespace collapsed) so both -/// shell strings (`buzz messages send --channel …`) and argv forms -/// (`['buzz','messages','send',…]`, e.g. Python `subprocess.run`) match. -/// Conservative: only matches an actual send subcommand, not reads like -/// `buzz messages get`. -fn tool_call_is_message_publish(update: &serde_json::Value) -> bool { - // Flatten title + rawInput into one lowercase haystack. rawInput is - // arbitrary JSON (shell command string, or structured args), so serialize - // whatever is there. - let mut haystack = String::new(); - if let Some(title) = update.get("title").and_then(|v| v.as_str()) { - haystack.push_str(title); - haystack.push(' '); - } - if let Some(raw) = update.get("rawInput") { - haystack.push_str(&raw.to_string()); - } - // Normalize away quoting/punctuation so argv-style invocations - // ('buzz','messages','send') match the same signature as shell strings. - let normalized: String = haystack - .to_ascii_lowercase() - .chars() - .map(|c| match c { - '\'' | '"' | '`' | ',' | '[' | ']' | '(' | ')' | '{' | '}' | ':' => ' ', - other => other, - }) - .collect(); - let h = normalized.split_whitespace().collect::>().join(" "); - // Match the publish subcommands that actually post to a channel. Guard - // against read subcommands (get/thread/search/list) sharing the "messages" - // prefix by requiring the send/publish verb. ("messages send-diff" - // contains "messages send", so it is covered.) - h.contains("messages send") || h.contains("social publish") -} - -/// Classify the terminal outcome of a publish tool call. -/// -/// Returns `None` while the call is still pending/in-progress, `Some(true)` -/// when the outcome confirms the message was delivered, and `Some(false)` when -/// the attempt failed (so the content-delivery fallback must stay armed — -/// see the `tool_call`/`tool_call_update` arms in `handle_session_update`). -/// -/// Signals, strongest first: -/// 1. `status: "failed"` (or cancelled) → not delivered. -/// 2. `rawOutput.isError: true` → not delivered (buzz-agent's builtin shape). -/// 3. Visible tool output containing the CLI's response envelope: -/// `"accepted":true` confirms, `"accepted":false` denies. -/// 4. A reported `exit_code` in the output (dev-mcp `shell` completes the -/// *tool* call even when the *command* failed): nonzero → not delivered. -/// 5. Otherwise, a `completed` non-error publish attempt counts as delivered — -/// the status-quo direction (suppressed fallback == today's behavior), -/// chosen over risking a duplicate post when output isn't visible. -fn publish_outcome_confirms_delivery(update: &serde_json::Value) -> Option { - let status = update.get("status").and_then(|v| v.as_str())?; - match status { - "failed" | "cancelled" | "canceled" | "error" => Some(false), - "completed" => { - if update - .get("rawOutput") - .and_then(|r| r.get("isError")) - .and_then(serde_json::Value::as_bool) - == Some(true) - { - return Some(false); - } - // Compact the visible output (content blocks + rawOutput) so the - // envelope matches regardless of pretty-printing — and strip - // backslashes so JSON nested inside a JSON string (rawOutput - // serialization escapes the quotes) matches too. - let compact: String = publish_output_text(update) - .chars() - .filter(|c| !c.is_whitespace() && *c != '\\') - .collect(); - if compact.contains(r#""accepted":true"#) { - return Some(true); - } - if compact.contains(r#""accepted":false"#) { - return Some(false); - } - if let Some(code) = extract_reported_exit_code(&compact) { - return Some(code == 0); - } - Some(true) - } - // "pending" / "in_progress" / anything non-terminal. - _ => None, - } -} - -/// Gather the human-visible output of a tool call update: ACP `content` text -/// blocks plus the serialized `rawOutput`, whichever are present. -fn publish_output_text(update: &serde_json::Value) -> String { - let mut out = String::new(); - if let Some(items) = update.get("content").and_then(|c| c.as_array()) { - for item in items { - // ACP shape: {type:"content", content:{type:"text", text:…}}; - // tolerate a flat {text:…} too. - if let Some(t) = item.pointer("/content/text").and_then(|v| v.as_str()) { - out.push_str(t); - out.push(' '); - } else if let Some(t) = item.get("text").and_then(|v| v.as_str()) { - out.push_str(t); - out.push(' '); - } - } - } - if let Some(raw) = update.get("rawOutput") { - out.push_str(&raw.to_string()); - } - out -} - -/// Extract a `"exit_code": N` value from compacted (whitespace-free) tool -/// output, e.g. the dev-mcp `shell` tool's result JSON. Returns `None` when no -/// exit code is reported. -fn extract_reported_exit_code(compact: &str) -> Option { - const KEY: &str = r#""exit_code":"#; - let idx = compact.find(KEY)?; - let rest = &compact[idx + KEY.len()..]; - let end = rest - .find(|c: char| !(c.is_ascii_digit() || c == '-')) - .unwrap_or(rest.len()); - rest[..end].parse().ok() -} - -/// Return true if `text` is a bare acknowledgement the base prompt forbids -/// publishing ("Got it", "Confirmed", "Standing by", …). Used to keep the -/// content-delivery fallback from posting filler that a capable agent would -/// have suppressed. Deliberately conservative — only short, whole-message -/// acks match, so a substantive reply that merely opens with "Got it, …" -/// still gets delivered. -fn is_bare_acknowledgement(text: &str) -> bool { - // Only consider short messages — a real reply with content is never a bare - // ack even if it starts with one. - if text.chars().count() > 40 { - return false; - } - let normalized: String = text - .to_ascii_lowercase() - .chars() - .filter(|c| c.is_alphanumeric() || c.is_whitespace()) - .collect(); - let normalized = normalized.trim(); - const BARE_ACKS: &[&str] = &[ - "got it", - "confirmed", - "acknowledged", - "ack", - "clear and noted", - "noted", - "aligned", - "standing by", - "parked", - "ok", - "okay", - "will do", - "understood", - "sounds good", - "on it", - "roger", - "roger that", - "i wont reply again", - "i will not reply again", - ]; - BARE_ACKS.contains(&normalized) -} - // ─── Drop: kill child process ───────────────────────────────────────────────── impl Drop for AcpClient { @@ -2505,323 +2233,6 @@ fn configure_no_window(cmd: &mut tokio::process::Command) { mod tests { use super::*; - #[test] - fn tool_call_publish_detection() { - // A `buzz messages send` shell tool call → detected as a publish. - let send = serde_json::json!({ - "title": "shell", - "rawInput": { "command": "buzz messages send --channel abc --content 'hi'" } - }); - assert!(tool_call_is_message_publish(&send)); - - // send-diff variant → detected. - let diff = serde_json::json!({ - "title": "shell", - "rawInput": { "command": "buzz messages send-diff --channel abc" } - }); - assert!(tool_call_is_message_publish(&diff)); - - // social publish → detected. - let social = serde_json::json!({ - "title": "shell", - "rawInput": { "command": "buzz social publish --content x" } - }); - assert!(tool_call_is_message_publish(&social)); - - // Python-argv publish (the pattern agents use for backtick-heavy - // content): serialized rawInput has no "messages send" substring, but - // normalization must still match it. - let argv = serde_json::json!({ - "title": "shell", - "rawInput": { - "command": "python3 - <<'PY'\nimport subprocess\nsubprocess.run(['buzz','messages','send','--channel','abc','--content',content])\nPY" - } - }); - assert!(tool_call_is_message_publish(&argv)); - - // A READ subcommand sharing the "messages" prefix → NOT a publish. - let read = serde_json::json!({ - "title": "shell", - "rawInput": { "command": "buzz messages get --channel abc" } - }); - assert!(!tool_call_is_message_publish(&read)); - - // Unrelated tool → not a publish. - let other = serde_json::json!({ - "title": "read_file", - "rawInput": { "path": "/tmp/foo" } - }); - assert!(!tool_call_is_message_publish(&other)); - } - - #[test] - fn publish_outcome_classification() { - // Non-terminal statuses → None (candidate stays pending). - for status in ["pending", "in_progress"] { - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({ "status": status })), - None, - "{status} is not terminal" - ); - } - // No status at all (e.g. a content-only update) → None. - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({})), - None - ); - - // Failure statuses → Some(false): the fallback must stay armed. - for status in ["failed", "cancelled", "canceled", "error"] { - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({ "status": status })), - Some(false), - "{status} must not confirm delivery" - ); - } - - // completed + rawOutput.isError → not delivered (buzz-agent shape). - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({ - "status": "completed", - "rawOutput": { "isError": true } - })), - Some(false) - ); - - // completed + CLI envelope accepted:true in the content text → delivered. - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({ - "status": "completed", - "content": [{ "type": "content", "content": { "type": "text", - "text": "0 {\"accepted\": true, \"event_id\": \"abc\"}" } }] - })), - Some(true) - ); - - // completed + envelope accepted:false (relay rejected) → not delivered. - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({ - "status": "completed", - "rawOutput": { "stdout": "{\"accepted\": false, \"message\": \"rate limited\"}" } - })), - Some(false) - ); - - // completed, no envelope, dev-mcp shell reports nonzero exit_code → - // the COMMAND failed even though the TOOL completed. Not delivered. - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({ - "status": "completed", - "content": [{ "type": "content", "content": { "type": "text", - "text": "{\"exit_code\": 3, \"stderr\": \"auth failure\"}" } }] - })), - Some(false) - ); - - // completed, exit_code 0, no envelope → delivered. - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({ - "status": "completed", - "content": [{ "type": "content", "content": { "type": "text", - "text": "{\"exit_code\": 0, \"stdout\": \"sent\"}" } }] - })), - Some(true) - ); - - // completed with no inspectable output at all → delivered (status-quo - // direction: suppressing the fallback == today's behavior). - assert_eq!( - publish_outcome_confirms_delivery(&serde_json::json!({ "status": "completed" })), - Some(true) - ); - } - - #[test] - fn bare_acknowledgement_detection() { - // Bare acks the base prompt forbids publishing. - for ack in [ - "Got it", - "confirmed", - "Standing by", - "OK", - " Noted. ", - "will do", - ] { - assert!(is_bare_acknowledgement(ack), "should be bare ack: {ack:?}"); - } - // Substantive replies are NOT bare acks, even if they open with one. - for real in [ - "Got it — I'll start on the migration and report back when the tests pass.", - "I'm doing well, thank you for asking! How are you today?", - "The build failed: missing dependency in Cargo.toml.", - ] { - assert!( - !is_bare_acknowledgement(real), - "should NOT be bare ack: {real:?}" - ); - } - } - - /// Build a `session/update` notification wrapping the given update object. - fn session_update_msg(update: serde_json::Value) -> serde_json::Value { - serde_json::json!({ "params": { "update": update } }) - } - - /// Successful send: candidate registered on `tool_call`, confirmed on the - /// terminal completed update with the CLI envelope → fallback suppressed. - #[tokio::test] - async fn fallback_suppressed_when_send_completes_successfully() { - let mut client = spawn_inert_client().await; - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "agent_message_chunk", - "content": { "type": "text", "text": "Here's my full reply narration." } - }))); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call", - "toolCallId": "tc-1", - "title": "shell", - "status": "pending", - "rawInput": { "command": "buzz messages send --channel abc --content 'hi'" } - }))); - // Not yet confirmed: a crash here must leave the fallback ARMED. - assert!(!client.turn_sent_message); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call_update", - "toolCallId": "tc-1", - "status": "completed", - "rawOutput": { "stdout": "{\"accepted\":true,\"event_id\":\"e1\"}" } - }))); - assert!( - client.turn_sent_message, - "successful send must confirm delivery" - ); - assert_eq!( - client.take_undelivered_turn_message(), - None, - "confirmed delivery suppresses the fallback" - ); - } - - /// Failed send: candidate registered, terminal update is `failed` → the - /// fallback stays armed and the buffered content is released for posting. - /// This is the false-negative path from review: intent is not delivery. - #[tokio::test] - async fn fallback_stays_armed_when_send_fails() { - let mut client = spawn_inert_client().await; - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "agent_message_chunk", - "content": { "type": "text", "text": "The answer is 42." } - }))); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call", - "toolCallId": "tc-2", - "title": "shell", - "status": "pending", - "rawInput": { "command": "buzz messages send --channel abc --content 'The answer is 42.'" } - }))); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call_update", - "toolCallId": "tc-2", - "status": "failed", - "rawOutput": { "error": "relay unreachable" } - }))); - assert!( - !client.turn_sent_message, - "failed send must NOT count as delivery" - ); - assert_eq!( - client.take_undelivered_turn_message().as_deref(), - Some("The answer is 42."), - "failed send leaves the fallback armed with the buffered reply" - ); - } - - /// Cancelled send behaves like failure: fallback stays armed. - #[tokio::test] - async fn fallback_stays_armed_when_send_cancelled() { - let mut client = spawn_inert_client().await; - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "agent_message_chunk", - "content": { "type": "text", "text": "Reply that never made it out." } - }))); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call", - "toolCallId": "tc-3", - "title": "shell", - "status": "pending", - "rawInput": { "command": "buzz messages send --channel abc --content x" } - }))); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call_update", - "toolCallId": "tc-3", - "status": "cancelled" - }))); - assert!(!client.turn_sent_message); - assert!(client.take_undelivered_turn_message().is_some()); - } - - /// A failed attempt followed by a successful retry (new toolCallId) - /// confirms delivery — the fallback must not double-post after a retry. - #[tokio::test] - async fn fallback_suppressed_after_failed_then_successful_retry() { - let mut client = spawn_inert_client().await; - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "agent_message_chunk", - "content": { "type": "text", "text": "narration" } - }))); - for (id, status) in [("tc-4a", "failed"), ("tc-4b", "completed")] { - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call", - "toolCallId": id, - "title": "shell", - "status": "pending", - "rawInput": { "command": "buzz messages send --channel abc --content x" } - }))); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call_update", - "toolCallId": id, - "status": status, - "rawOutput": { "stdout": "{\"accepted\":true,\"event_id\":\"e2\"}" } - }))); - } - assert!( - client.turn_sent_message, - "retry succeeded — delivery confirmed" - ); - assert_eq!(client.take_undelivered_turn_message(), None); - } - - /// Argv-style publish (Python subprocess) is recognized as a candidate and - /// confirmed on success — the normalization regression from review. - #[tokio::test] - async fn fallback_suppressed_for_argv_style_publish() { - let mut client = spawn_inert_client().await; - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "agent_message_chunk", - "content": { "type": "text", "text": "long narration between tool calls" } - }))); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call", - "toolCallId": "tc-5", - "title": "shell", - "status": "pending", - "rawInput": { "command": "python3 - <<'PY'\nimport subprocess\nsubprocess.run(['buzz','messages','send','--channel','abc','--content',content])\nPY" } - }))); - let _ = client.handle_session_update(&session_update_msg(serde_json::json!({ - "sessionUpdate": "tool_call_update", - "toolCallId": "tc-5", - "status": "completed", - "content": [{ "type": "content", "content": { "type": "text", - "text": "0 {\"accepted\": true, \"event_id\": \"abc\"}" } }] - }))); - assert!(client.turn_sent_message); - assert_eq!( - client.take_undelivered_turn_message(), - None, - "argv publish must not double-post the narration" - ); - } - #[test] fn stop_reason_parses_all_known_values() { assert_eq!(StopReason::from_str("end_turn"), Some(StopReason::EndTurn)); diff --git a/crates/buzz-acp/src/pool.rs b/crates/buzz-acp/src/pool.rs index 3befe10f8..bfa20469e 100644 --- a/crates/buzz-acp/src/pool.rs +++ b/crates/buzz-acp/src/pool.rs @@ -1968,28 +1968,6 @@ pub async fn run_prompt_task( None => prompt_sections.iter().map(String::as_str).collect(), }; - // Capture the reply destination for the content-delivery fallback BEFORE - // the prompt runs, so it survives any move of `batch` in the outcome arms. - // Only channel turns with a triggering event can receive a fallback post; - // heartbeats and DMs without a triggering message are skipped (None). - let fallback_reply: Option = batch.as_ref().and_then(|b| { - b.events.last().map(|last| { - let tags = crate::queue::parse_thread_tags(&last.event); - // Thread the reply to the triggering event: if the trigger is - // itself a reply, anchor to its root; otherwise the trigger IS - // the root. Mirrors the CLI `resolve_thread_ref` semantics. - let root_hex = tags - .root_event_id - .clone() - .unwrap_or_else(|| last.event.id.to_hex()); - FallbackReplyTarget { - channel_id: b.channel_id, - root_event_hex: root_hex, - parent_event_hex: last.event.id.to_hex(), - } - }) - }); - // When control_rx is Some (channel tasks), wrap the prompt in select! so // the main loop can cancel, interrupt, or rotate it. Heartbeats // (control_rx=None) take the simple await path — they are not controllable. @@ -2142,15 +2120,6 @@ pub async fn run_prompt_task( Some(buzz_core::agent_turn_metric::StopReason::EndTurn), ) .await; - // Content-delivery fallback (see the main EndTurn arm): - // this rare branch is also a successful turn end, so an - // undelivered plain-text reply still needs posting. - if let (Some(target), Some(content)) = - (&fallback_reply, agent.acp.take_undelivered_turn_message()) - { - post_agent_content_fallback(&ctx.rest_client, target, &content) - .await; - } send_prompt_result( &result_tx, &turn_id, @@ -2214,21 +2183,6 @@ pub async fn run_prompt_task( ) .await; - // Content-delivery fallback: on a normal turn end, if the agent - // produced assistant text but never called a publish tool, post - // that text as the channel reply. Only fires for `EndTurn` (not - // MaxTokens/MaxTurnRequests, which are truncated/aborted turns - // whose partial text shouldn't be treated as a deliberate reply) - // and only when a `fallback_reply` destination was captured - // (channel turns with a triggering event; not heartbeats). - if matches!(stop_reason, StopReason::EndTurn) { - if let (Some(target), Some(content)) = - (&fallback_reply, agent.acp.take_undelivered_turn_message()) - { - post_agent_content_fallback(&ctx.rest_client, target, &content).await; - } - } - send_prompt_result( &result_tx, &turn_id, @@ -3726,81 +3680,6 @@ pub(crate) async fn post_failure_notice( } } -/// Captured reply destination for the content-delivery fallback, taken before -/// the prompt runs so it survives any move of the triggering `batch`. -#[derive(Clone)] -struct FallbackReplyTarget { - channel_id: Uuid, - /// Thread root the reply anchors to (hex). Equals `parent_event_hex` when - /// the trigger was a top-level message. - root_event_hex: String, - /// Immediate parent being replied to (hex) — the triggering event. - parent_event_hex: String, -} - -/// Content-delivery fallback: post an agent's plain-text reply (kind:9) that it -/// generated but never published itself. -/// -/// buzz-agent's output is its tool calls; streamed assistant `content` is -/// observability-only and is normally not posted. Capable models reliably call -/// `buzz messages send`, but weaker local models (e.g. via Buzz shared compute) -/// often answer a conversational prompt in plain content and never call the -/// send tool — silently dropping the reply. When [`AcpClient`] reports such an -/// undelivered turn message, this posts it as a threaded reply, mirroring -/// [`post_failure_notice`]'s build/sign/submit path. Best-effort: any error is -/// logged and swallowed. -async fn post_agent_content_fallback( - rest: &crate::relay::RestClient, - target: &FallbackReplyTarget, - content: &str, -) { - let thread_ref = match ( - nostr::EventId::from_hex(&target.root_event_hex), - nostr::EventId::from_hex(&target.parent_event_hex), - ) { - (Ok(root_id), Ok(parent_id)) => Some(buzz_sdk::ThreadRef { - root_event_id: root_id, - parent_event_id: parent_id, - }), - _ => None, - }; - let builder = match buzz_sdk::build_message( - target.channel_id, - content, - thread_ref.as_ref(), - &[], - false, - &[], - ) { - Ok(b) => b, - Err(e) => { - tracing::warn!(channel = %target.channel_id, "content fallback: build failed: {e}"); - return; - } - }; - let event = match builder.sign_with_keys(&rest.keys) { - Ok(e) => e, - Err(e) => { - tracing::warn!(channel = %target.channel_id, "content fallback: sign failed: {e}"); - return; - } - }; - match tokio::time::timeout(Duration::from_secs(5), rest.submit_event(&event)).await { - Ok(Ok(_)) => { - // WARN (not INFO) and default target (buzz_acp::pool) so it is - // always visible under the harness's `buzz_acp=info` filter — this - // fallback firing is a signal worth surfacing (a model failed to - // call the send tool and we delivered its reply for it). - tracing::warn!( - channel = %target.channel_id, - "content-delivery fallback: posted undelivered agent content as channel reply" - ); - } - Ok(Err(e)) => tracing::warn!(channel = %target.channel_id, "content fallback failed: {e}"), - Err(_) => tracing::warn!(channel = %target.channel_id, "content fallback timed out"), - } -} - /// Best-effort: remove a reaction via a signed kind:5 (NIP-09) deletion event. /// /// Queries kind:7 reactions by our pubkey targeting the event, finds the matching diff --git a/crates/buzz-agent/src/llm.rs b/crates/buzz-agent/src/llm.rs index 23cef24e7..9f59560ec 100644 --- a/crates/buzz-agent/src/llm.rs +++ b/crates/buzz-agent/src/llm.rs @@ -1077,6 +1077,9 @@ fn parse_responses(v: Value) -> Result { Some("completed") => ProviderStop::EndTurn, _ => ProviderStop::Other, }; + if text.is_empty() && !reasoning.is_empty() && tool_calls.is_empty() { + text = reasoning.clone(); + } let input_tokens = sum_usage(&v, &["input_tokens"]); let output_tokens = sum_usage(&v, &["output_tokens"]); Ok(LlmResponse { @@ -1204,7 +1207,7 @@ fn parse_openai(v: Value) -> Result { let msg = choice .get("message") .ok_or_else(|| AgentError::Llm("missing message".into()))?; - let text = str_field(msg, "content"); + let mut text = str_field(msg, "content"); // DeepSeek and vLLM-style OpenAI-compat hosts expose reasoning tokens on the // message object. Prefer `reasoning_content` (DeepSeek's field name); fall // back to `reasoning` (some other providers). Both are absent for standard @@ -1233,6 +1236,9 @@ fn parse_openai(v: Value) -> Result { )?); } } + if text.is_empty() && !reasoning.is_empty() && tool_calls.is_empty() { + text = reasoning.clone(); + } let input_tokens = openai_chat_input_tokens(&v); let output_tokens = sum_usage(&v, &["completion_tokens"]); Ok(LlmResponse { @@ -2407,6 +2413,21 @@ mod tests { assert_eq!(r.stop, ProviderStop::EndTurn); } + #[test] + fn parse_responses_uses_reasoning_when_text_is_empty() { + let v = serde_json::json!({ + "status": "completed", + "output": [{ + "type": "reasoning", + "summary": [{"type": "summary_text", "text": "local model reply"}], + }], + }); + let r = parse_responses(v).unwrap(); + assert_eq!(r.text, "local model reply"); + assert_eq!(r.reasoning, "local model reply"); + assert!(r.tool_calls.is_empty()); + } + #[test] fn parse_responses_completed_with_function_call_is_tool_use() { let v = serde_json::json!({ @@ -3492,6 +3513,23 @@ mod tests { assert_eq!(parse_openai(v).unwrap().input_tokens, Some(123)); } + #[test] + fn parse_openai_uses_reasoning_content_when_content_is_empty() { + let v = serde_json::json!({ + "choices": [{ + "finish_reason": "stop", + "message": { + "content": "", + "reasoning_content": "local model reply" + } + }] + }); + let r = parse_openai(v).unwrap(); + assert_eq!(r.text, "local model reply"); + assert_eq!(r.reasoning, "local model reply"); + assert!(r.tool_calls.is_empty()); + } + #[test] fn parse_openai_databricks_sums_cache_fields() { // Databricks uses the OpenAI chat wire format (prompt_tokens) but also diff --git a/crates/buzz-dev-mcp/src/lib.rs b/crates/buzz-dev-mcp/src/lib.rs index 9b9897480..5d617a877 100644 --- a/crates/buzz-dev-mcp/src/lib.rs +++ b/crates/buzz-dev-mcp/src/lib.rs @@ -12,6 +12,7 @@ use std::sync::Arc; mod paths; mod read_file; +mod reply_delivery; mod rg; mod shell; mod shim; @@ -106,7 +107,17 @@ impl DevMcp { &self, Parameters(_): Parameters, ) -> Result { - todo::text_result(self.todos.stop_objection()) + let objections = [ + self.todos.stop_objection(), + self.state.reply_delivery.stop_objection(), + ] + .into_iter() + .filter(|objection| !objection.is_empty()) + .collect::>(); + if objections.is_empty() { + self.state.reply_delivery.reset_after_allowed_stop(); + } + todo::text_result(objections.join("\n\n")) } /// Hook: called by the agent after context compaction/handoff so the diff --git a/crates/buzz-dev-mcp/src/reply_delivery.rs b/crates/buzz-dev-mcp/src/reply_delivery.rs new file mode 100644 index 000000000..2e892a725 --- /dev/null +++ b/crates/buzz-dev-mcp/src/reply_delivery.rs @@ -0,0 +1,128 @@ +use std::path::Path; +use std::sync::Mutex; + +const STOP_OBJECTION: &str = "You have not posted your reply to Buzz. Do not describe or promise the send. Call the `shell` MCP tool now and run `buzz messages send`. Set `--channel` to the actual channel UUID from `[Context]` and `--content` to your reply; never use literal placeholder text. Do not end the turn until that command succeeds."; + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +enum DeliveryStatus { + #[default] + AwaitingSend, + Sent, +} + +#[derive(Debug, Default)] +pub struct ReplyDeliveryState { + status: Mutex, +} + +impl ReplyDeliveryState { + pub fn record_shell_result(&self, command: &str, exit_code: i32, stdout: &str) { + let is_send = is_buzz_message_send(command); + let accepted = write_was_accepted(stdout); + if exit_code != 0 || !is_send || !accepted { + return; + } + self.with_status(|status| *status = DeliveryStatus::Sent); + } + + pub fn stop_objection(&self) -> String { + self.with_status(|status| match status { + DeliveryStatus::AwaitingSend => STOP_OBJECTION.to_owned(), + DeliveryStatus::Sent => String::new(), + }) + } + + pub fn reset_after_allowed_stop(&self) { + self.with_status(|status| *status = DeliveryStatus::AwaitingSend); + } + + fn with_status(&self, f: impl FnOnce(&mut DeliveryStatus) -> R) -> R { + let mut guard = match self.status.lock() { + Ok(guard) => guard, + Err(poisoned) => poisoned.into_inner(), + }; + f(&mut guard) + } +} + +fn is_buzz_message_send(command: &str) -> bool { + let mut words = command.split_ascii_whitespace(); + let Some(executable) = words.next() else { + return false; + }; + let executable = executable.trim_matches(['\'', '"']); + let is_buzz = Path::new(executable) + .file_name() + .and_then(|name| name.to_str()) + .is_some_and(|name| matches!(name, "buzz" | "buzz.exe")); + is_buzz && words.next() == Some("messages") && words.next() == Some("send") +} + +fn write_was_accepted(stdout: &str) -> bool { + serde_json::from_str::(stdout.trim()) + .ok() + .and_then(|value| value.get("accepted").and_then(serde_json::Value::as_bool)) + == Some(true) +} + +#[cfg(test)] +mod tests { + use super::*; + + const ACCEPTED: &str = r#"{"event_id":"abc","accepted":true,"message":"ok"}"#; + + #[test] + fn successful_buzz_send_satisfies_stop_then_resets() { + let state = ReplyDeliveryState::default(); + assert!(!state.stop_objection().is_empty()); + + state.record_shell_result( + "buzz messages send --channel c --content hello", + 0, + ACCEPTED, + ); + assert!(state.stop_objection().is_empty()); + + state.reset_after_allowed_stop(); + assert!(!state.stop_objection().is_empty()); + } + + #[test] + fn failed_or_rejected_send_does_not_satisfy_stop() { + let state = ReplyDeliveryState::default(); + state.record_shell_result( + "buzz messages send --channel c --content hello", + 1, + ACCEPTED, + ); + state.record_shell_result( + "buzz messages send --channel c --content hello", + 0, + r#"{"accepted":false}"#, + ); + assert!(!state.stop_objection().is_empty()); + } + + #[test] + fn mentioning_send_in_another_command_does_not_satisfy_stop() { + let state = ReplyDeliveryState::default(); + state.record_shell_result("echo buzz messages send", 0, ACCEPTED); + assert!(!state.stop_objection().is_empty()); + } + + #[test] + fn absolute_buzz_path_is_recognized() { + assert!(is_buzz_message_send( + "/Applications/Buzz.app/Contents/MacOS/buzz messages send --channel c --content hi" + )); + } + + #[test] + fn stop_objection_requires_the_shell_tool_and_successful_send() { + assert!(STOP_OBJECTION.contains("`shell` MCP tool")); + assert!(STOP_OBJECTION.contains("`buzz messages send`")); + assert!(STOP_OBJECTION.contains("actual channel UUID")); + assert!(STOP_OBJECTION.contains("never use literal placeholder")); + assert!(STOP_OBJECTION.contains("command succeeds")); + } +} diff --git a/crates/buzz-dev-mcp/src/shell.rs b/crates/buzz-dev-mcp/src/shell.rs index 7aa95b1d8..c9faa251c 100644 --- a/crates/buzz-dev-mcp/src/shell.rs +++ b/crates/buzz-dev-mcp/src/shell.rs @@ -33,6 +33,7 @@ pub struct SharedState { /// bootstrap hint and every `run()` call read the SAME resolution — no drift. pub resolved_shell: Result<(PathBuf, String), String>, pub artifacts: Mutex>, + pub reply_delivery: crate::reply_delivery::ReplyDeliveryState, next_call_id: Mutex, } @@ -58,6 +59,7 @@ impl SharedState { bootstrap_instructions, resolved_shell, artifacts: Mutex::new(VecDeque::with_capacity(ARTIFACT_RING_SIZE)), + reply_delivery: crate::reply_delivery::ReplyDeliveryState::default(), next_call_id: Mutex::new(0), }) } @@ -306,6 +308,10 @@ pub async fn run( let (stderr_text, stderr_truncated, stderr_artifact) = finalize_stream(state, id, "stderr", stderr_cap, &mut notes); + state + .reply_delivery + .record_shell_result(&p.command, exit_code, &stdout_text); + let body = serde_json::json!({ "exit_code": exit_code, "stdout": stdout_text,