use std::time::{Duration, Instant};
use mecha_core::agent::{AgentEvent, RunOutcome};
use mecha_slack::chat::{self, Chunk, TaskStatus};
use mecha_slack::{blocks, Slack};
use tokio::sync::mpsc::UnboundedReceiver;
pub struct PumpConfig {
pub flush_chars: usize,
pub flush_ms: u64,
}
pub async fn pump(
slack: &Slack,
channel: &str,
thread_ts: &str,
mut events: UnboundedReceiver<AgentEvent>,
cfg: &PumpConfig,
) -> Option<String> {
let stream_ts = match chat::start_stream(slack, channel, thread_ts).await {
Ok(ts) => ts,
Err(e) => {
tracing::warn!("could not open a Slack stream: {e}");
while events.recv().await.is_some() {}
return None;
}
};
let mut buffer = String::new();
let mut last_flush = Instant::now();
let mut titles: std::collections::HashMap<String, String> = std::collections::HashMap::new();
let mut notice = 0u32;
let mut outcome: Option<Box<RunOutcome>> = None;
while let Some(event) = events.recv().await {
match event {
AgentEvent::TextDelta(text) => {
buffer.push_str(&text);
if buffer.chars().count() >= cfg.flush_chars
|| last_flush.elapsed() >= Duration::from_millis(cfg.flush_ms)
{
flush(slack, channel, &stream_ts, &mut buffer, &mut last_flush).await;
}
}
AgentEvent::ThinkingDelta(_) | AgentEvent::AssistantText(_) => {}
AgentEvent::ToolCall { id, name, input } => {
flush(slack, channel, &stream_ts, &mut buffer, &mut last_flush).await;
let title = describe(&name, &input);
titles.insert(id.clone(), title.clone());
task(
slack,
channel,
&stream_ts,
&id,
&title,
TaskStatus::InProgress,
)
.await;
}
AgentEvent::ToolResult {
id, name, is_error, ..
} => {
let status = if is_error {
TaskStatus::Error
} else {
TaskStatus::Complete
};
let title = titles.remove(&id).unwrap_or(name);
task(slack, channel, &stream_ts, &id, &title, status).await;
}
AgentEvent::ToolDenied { name, reason } => {
flush(slack, channel, &stream_ts, &mut buffer, &mut last_flush).await;
notice += 1;
task(
slack,
channel,
&stream_ts,
&format!("notice-{notice}"),
&format!("{name} refused — {reason}"),
TaskStatus::Error,
)
.await;
}
AgentEvent::Compacted { .. } => {
notice += 1;
task(
slack,
channel,
&stream_ts,
&format!("notice-{notice}"),
"the transcript was summarised to fit the context window",
TaskStatus::Complete,
)
.await;
}
AgentEvent::QueuedInput(text) => {
notice += 1;
task(
slack,
channel,
&stream_ts,
&format!("notice-{notice}"),
&format!("steering: {text}"),
TaskStatus::Complete,
)
.await;
}
AgentEvent::Done(done) => outcome = Some(done),
_ => {}
}
}
flush(slack, channel, &stream_ts, &mut buffer, &mut last_flush).await;
let footer = outcome.as_deref().map(footer_blocks);
if let Err(e) = chat::stop_stream(slack, channel, &stream_ts, footer).await {
tracing::warn!("could not close the Slack stream: {e}");
}
Some(stream_ts)
}
async fn flush(
slack: &Slack,
channel: &str,
ts: &str,
buffer: &mut String,
last_flush: &mut Instant,
) {
if buffer.is_empty() {
return;
}
let chunk = Chunk::Markdown(std::mem::take(buffer));
if let Err(e) = chat::append_stream(slack, channel, ts, &chunk).await {
tracing::warn!("could not append to the Slack stream: {e}");
}
*last_flush = Instant::now();
}
async fn task(slack: &Slack, channel: &str, ts: &str, id: &str, title: &str, status: TaskStatus) {
let chunk = Chunk::Task {
id: id.to_string(),
title: title.to_string(),
status,
};
if let Err(e) = chat::append_stream(slack, channel, ts, &chunk).await {
tracing::warn!("could not append a task update: {e}");
}
}
fn describe(name: &str, input: &serde_json::Value) -> String {
let detail = input
.get("command")
.or_else(|| input.get("path"))
.or_else(|| input.get("url"))
.and_then(|v| v.as_str())
.map(|s| s.chars().take(120).collect::<String>());
match detail {
Some(d) => format!("{name}: {d}"),
None => name.to_string(),
}
}
fn footer_blocks(outcome: &RunOutcome) -> Vec<serde_json::Value> {
vec![blocks::context(&footer_parts(outcome).join(" · "))]
}
fn footer_parts(outcome: &RunOutcome) -> Vec<String> {
let mut parts = vec![format!("{} turns", outcome.turns)];
parts.push(format!(
"{} in / {} out{}",
outcome.usage.input_tokens,
outcome.usage.output_tokens,
if outcome.usage_complete {
""
} else {
" (partial)"
}
));
if let Some(cost) = outcome.cost_usd {
parts.push(format!("${cost:.4}"));
}
parts.push(format!("{:?}", outcome.stop_cause).to_lowercase());
if outcome.compactions > 0 {
parts.push(format!("{} compaction(s)", outcome.compactions));
}
if outcome.blocked_sends > 0 {
parts.push(format!(
"{} send(s) refused by the interlock",
outcome.blocked_sends
));
}
if outcome.exhausted {
parts.push("*hit the turn limit — the answer is probably incomplete*".into());
}
parts
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn a_call_is_described_by_what_it_will_do() {
assert_eq!(
describe("shell", &json!({"command": "cargo test"})),
"shell: cargo test"
);
assert_eq!(describe("todo", &json!({})), "todo");
}
#[test]
fn a_long_command_is_cut_before_it_reaches_slacks_task_limit() {
let d = describe("shell", &json!({"command": "x".repeat(1000)}));
assert!(d.chars().count() <= 120 + "shell: ".len());
}
fn outcome() -> RunOutcome {
RunOutcome {
text: String::new(),
stop_reason: mecha_core::message::StopReason::EndTurn,
usage: mecha_core::Usage::default(),
turns: 3,
refusal: None,
exhausted: false,
ended_on_failed_call: false,
tool_calls: Vec::new(),
malformed_tool_args: 0,
blocked_sends: 0,
taint: Default::default(),
stop_cause: mecha_core::agent::StopCause::Completed,
cost_usd: None,
compactions: 0,
usage_complete: true,
}
}
#[test]
fn the_footer_says_how_it_ended_and_what_it_cost() {
let mut o = outcome();
o.usage.input_tokens = 100;
o.usage.output_tokens = 20;
o.blocked_sends = 2;
o.exhausted = true;
o.compactions = 1;
let parts = footer_parts(&o).join(" · ");
assert!(parts.contains("3 turns"), "{parts}");
assert!(parts.contains("100 in / 20 out"), "{parts}");
assert!(parts.contains("1 compaction"), "{parts}");
assert!(parts.contains("refused by the interlock"), "{parts}");
assert!(parts.contains("turn limit"), "{parts}");
}
#[test]
fn partial_usage_says_so_rather_than_reading_as_a_measurement() {
let mut o = outcome();
o.usage_complete = false;
assert!(footer_parts(&o).join(" ").contains("(partial)"));
assert!(!footer_parts(&outcome()).join(" ").contains("(partial)"));
}
}