use std::collections::{HashMap, VecDeque};
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use anyhow::{bail, Context, Result};
use mecha_core::agent::{Agent, AgentEvent, Budget, Conversation, RunOutcome};
use mecha_core::message::Message;
use mecha_core::outbox::OutboxRoute;
use mecha_core::session::{Record, Session, SessionMeta};
use mecha_core::tool::ToolCtx;
use mecha_slack::binding::{self, Binding, Credentials, Gate, SlackStore};
use mecha_slack::envelope::{FileRef, Inbound, Interaction, SlackEvent};
use mecha_slack::{blocks, chat, Slack, SocketMode, SocketOptions};
use tokio::sync::{mpsc, oneshot};
use tokio_util::sync::CancellationToken;
use super::actions::{self, Action, ActionLedger, Executor};
use super::approve::{self, Answer, Mode, SlackApprover};
use super::pump::{pump, PumpConfig};
use super::review::{self, ReviewMode};
use super::threads::{Event, RunMarker, ThreadRecord, ThreadStore};
use crate::{setup, GlobalOpts};
const SEEN_EVENTS: usize = 512;
struct Live {
cancel: CancellationToken,
queue: Arc<Mutex<VecDeque<String>>>,
mode: Arc<Mutex<Mode>>,
unanswered: Arc<AtomicBool>,
}
struct Completion {
key: String,
conversation: Conversation,
outcome: Result<Box<RunOutcome>, String>,
}
struct PendingApproval {
reply: oneshot::Sender<Answer>,
channel: String,
message_ts: String,
tool: String,
thread_key: String,
expires_at: std::time::Instant,
}
struct ConnectorLock {
_file: std::fs::File,
}
impl ConnectorLock {
fn take(path: &std::path::Path) -> Result<Self> {
let file = std::fs::OpenOptions::new()
.create(true)
.truncate(false)
.write(true)
.open(path)
.with_context(|| format!("opening {}", path.display()))?;
let rc = unsafe {
libc::flock(
std::os::unix::io::AsRawFd::as_raw_fd(&file),
libc::LOCK_EX | libc::LOCK_NB,
)
};
if rc != 0 {
bail!(
"another `mecha slack connect` is already running (lock held on {}). \
Two connectors would both answer the same message.",
path.display()
);
}
Ok(Self { _file: file })
}
}
pub async fn run(global: &GlobalOpts) -> Result<()> {
let store = SlackStore::open(mecha_core::work::mecha_home()?.join("slack"))?;
let _lock = ConnectorLock::take(&store.root().join("connector.lock"))?;
let creds: Credentials = store
.credentials()?
.context("no Slack tokens stored — run `mecha slack auth` first")?;
let binding: Binding = store
.binding()?
.context("nothing is bound — run `mecha slack link` first")?;
let global_cfg = mecha_core::config::Config::load_global()?;
let outbox_root = match global_cfg.outbox.dir.clone() {
Some(dir) => dir,
None => mecha_core::outbox::OutboxStore::default_root()?,
};
let cfg = global_cfg.slack;
let threads = ThreadStore::open(
mecha_core::work::mecha_home()?
.join("slack")
.join("threads"),
)?;
let slack = Slack::new(&creds.bot_token);
let me: serde_json::Value = slack.call("auth.test", serde_json::json!({})).await?;
let my_user_id = me["user_id"].as_str().unwrap_or_default().to_string();
for orphan in threads.sweep()? {
if let Some(ts) = &orphan.controls_ts {
let _ = chat::update(
&slack,
&orphan.channel_id,
ts,
"✗ Run lost to a restart",
Some(vec![blocks::context("✗ Run lost to a restart")]),
)
.await;
}
let _ = chat::post_message(
&slack,
&orphan.channel_id,
Some(&orphan.thread_ts),
"This run did not survive a restart of the connector. Nothing was lost that \
was written down — send another message to pick it up.",
None,
)
.await;
let _ = threads.apply(&orphan.key, Event::OrphanAnnounced);
}
let prepared = build_agent(global, &cfg).await?;
let provider = prepared.provider_name.clone();
let model = prepared.model.clone();
let prepared_config = prepared.config.clone();
let agent = Arc::new(prepared.agent);
let (inbound_tx, mut inbound_rx) = mpsc::channel(64);
let (approval_tx, mut approval_rx) = mpsc::channel::<approve::Request>(32);
let (completion_tx, mut completion_rx) = mpsc::channel::<Completion>(32);
let socket = SocketMode::new(
slack.clone(),
SocketOptions {
app_token: creds.app_token.clone(),
debug_reconnects: false,
},
);
let socket_task = tokio::spawn(async move { socket.run(inbound_tx, || false).await });
let mut state = State {
slack,
binding,
threads,
cfg,
agent,
my_user_id,
live: HashMap::new(),
conversations: HashMap::new(),
pending: HashMap::new(),
seen: VecDeque::new(),
approval_seq: 0,
staged_before: HashMap::new(),
files_before: HashMap::new(),
outbox_root,
ledger: Arc::new(ActionLedger::open_default()),
review: HashMap::new(),
provider,
model,
config: prepared_config,
approval_tx,
completion_tx,
};
println!(
"Connected to {} as {}. {} owner(s), {} thread(s) known.",
me["team"].as_str().unwrap_or("slack"),
me["user"].as_str().unwrap_or("mecha"),
state.binding.owners.len(),
state.threads.all().map(|t| t.len()).unwrap_or(0),
);
println!("Waiting for a direct message. Ctrl-C or SIGTERM to stop.");
loop {
tokio::select! {
inbound = inbound_rx.recv() => match inbound {
Some(inbound) => state.on_inbound(inbound).await,
None => break,
},
request = approval_rx.recv() => if let Some(r) = request {
state.on_approval_request(r).await;
},
done = completion_rx.recv() => if let Some(c) = done {
state.on_completion(c).await;
},
_ = tokio::time::sleep(Duration::from_secs(15)) => {
state.retire_expired_approvals().await;
}
_ = shutdown_signal() => {
println!("Stopping; in-flight runs cancel at their next safe point.");
for live in state.live.values() {
live.cancel.cancel();
}
break;
}
}
}
if socket_task.is_finished() {
match socket_task.await {
Ok(Err(e)) => return Err(anyhow::anyhow!("slack socket stopped: {e}")),
Ok(Ok(())) => {}
Err(e) if !e.is_cancelled() => {
return Err(anyhow::anyhow!("slack socket task failed: {e}"))
}
Err(_) => {}
}
} else {
socket_task.abort();
}
Ok(())
}
async fn shutdown_signal() {
#[cfg(unix)]
{
use tokio::signal::unix::{signal, SignalKind};
let mut term = match signal(SignalKind::terminate()) {
Ok(s) => s,
Err(_) => {
let _ = tokio::signal::ctrl_c().await;
return;
}
};
tokio::select! {
_ = tokio::signal::ctrl_c() => {}
_ = term.recv() => {}
}
}
#[cfg(not(unix))]
{
let _ = tokio::signal::ctrl_c().await;
}
}
async fn build_agent(
global: &GlobalOpts,
cfg: &mecha_core::config::SlackConfig,
) -> Result<setup::Prepared> {
let opts = GlobalOpts {
global_config_only: true,
provider: global.provider.clone(),
model: global.model.clone(),
tools: cfg.tools.clone(),
workspace: Some(producer_root()?),
..GlobalOpts::default()
};
setup::prepare(&opts, false).await
}
struct State {
slack: Slack,
binding: Binding,
threads: ThreadStore,
cfg: mecha_core::config::SlackConfig,
agent: Arc<Agent>,
my_user_id: String,
live: HashMap<String, Live>,
conversations: HashMap<String, Conversation>,
pending: HashMap<String, PendingApproval>,
seen: VecDeque<String>,
approval_seq: u64,
staged_before: HashMap<String, std::collections::HashSet<String>>,
files_before: HashMap<String, HashMap<PathBuf, (u64, std::time::SystemTime)>>,
outbox_root: std::path::PathBuf,
ledger: Arc<ActionLedger>,
review: HashMap<String, review::Setting>,
provider: String,
model: String,
config: mecha_core::config::Config,
approval_tx: mpsc::Sender<approve::Request>,
completion_tx: mpsc::Sender<Completion>,
}
impl State {
async fn on_inbound(&mut self, inbound: Inbound) {
match inbound {
Inbound::Event { event, .. } => self.on_event(*event).await,
Inbound::Interactive { interaction, .. } => self.on_interaction(*interaction).await,
_ => {}
}
}
fn allowed(&self, user: Option<&str>, team: Option<&str>) -> Gate {
binding::check(Some(&self.binding), user, team)
}
async fn on_event(&mut self, event: SlackEvent) {
if event.kind != "message" || !event.is_from_a_human() {
return;
}
if event.user.as_deref() == Some(self.my_user_id.as_str()) {
return;
}
if !self.first_time(&event.event_id) {
return;
}
let gate = self.allowed(event.user.as_deref(), event.team_id.as_deref());
if !gate.is_allowed() {
println!("ignored a message: {}", gate.reason());
return;
}
let (Some(channel), Some(thread_ts), Some(text)) = (
event.channel.clone(),
event.thread_key(),
event.text.clone(),
) else {
return;
};
if super::doctor::is_doctor_command(&text) {
let slack = self.slack.clone();
let (channel, thread_ts) = (channel.clone(), thread_ts.clone());
tokio::spawn(async move {
let (text, blocks) = super::doctor::report().await;
let _ = chat::post_message(&slack, &channel, Some(&thread_ts), &text, blocks).await;
});
return;
}
if super::triggers::is_triggers_command(&text) {
let slack = self.slack.clone();
let (channel, thread_ts) = (channel.clone(), thread_ts.clone());
tokio::spawn(async move {
let (text, blocks) = super::triggers::listing().await;
let _ = chat::post_message(&slack, &channel, Some(&thread_ts), &text, blocks).await;
});
return;
}
if let Some(asked) = review::command(&text) {
let reply = match asked {
Some(mode) => {
let (key, scope) = review::scope_for(&channel, event.thread_ts.as_deref());
let who = event.user.clone().unwrap_or_default();
self.review
.insert(key, review::Setting { mode, set_by: who });
format!(
"Review mode for {scope} is now `{}` — {}. Tainted drafts \
always stop for review, and the mode lasts only while the \
connector runs.",
mode.name(),
mode.describe()
)
}
None => {
let key = super::threads::key_for(&channel, &thread_ts);
let mode = review::effective(&self.review, &key, &channel)
.map(|s| s.mode)
.unwrap_or(ReviewMode::Now);
format!(
"Review mode is `{}` — {}. Set it with `review now|later|auto`.",
mode.name(),
mode.describe()
)
}
};
let _ = chat::post_message(&self.slack, &channel, Some(&thread_ts), &reply, None).await;
return;
}
let record = match self
.threads
.ensure(&channel, &thread_ts, &self.cfg.default_mode)
{
Ok(r) => r,
Err(e) => {
tracing::warn!("could not open the thread record: {e}");
return;
}
};
if let Some(live) = self.live.get(&record.key) {
if let Ok(mut queue) = live.queue.lock() {
queue.push_back(text);
}
live.unanswered.store(false, Ordering::Relaxed);
return;
}
if self.live.len() >= self.cfg.max_concurrent {
let _ = chat::post_message(
&self.slack,
&channel,
Some(&thread_ts),
&format!(
"Already running {} threads, which is the configured limit. \
Send this again when one finishes.",
self.cfg.max_concurrent
),
None,
)
.await;
return;
}
self.start_run(record, text, event.files).await;
}
fn first_time(&mut self, event_id: &str) -> bool {
if event_id.is_empty() {
return true;
}
if self.seen.iter().any(|s| s == event_id) {
return false;
}
self.seen.push_back(event_id.to_string());
if self.seen.len() > SEEN_EVENTS {
self.seen.pop_front();
}
true
}
async fn start_run(&mut self, record: ThreadRecord, prompt: String, files: Vec<FileRef>) {
let key = record.key.clone();
let channel = record.channel_id.clone();
let thread_ts = record.thread_ts.clone();
let workspace = match thread_workspace(&key) {
Ok(w) => w,
Err(e) => {
tracing::warn!("no workspace for {key}: {e}");
return;
}
};
let mode = Arc::new(Mutex::new(Mode::parse(&record.mode).unwrap_or(Mode::Ask)));
let cancel = CancellationToken::new();
let queue = Arc::new(Mutex::new(VecDeque::new()));
let session = match self.session_for(&record).await {
Some(s) => s,
None => return,
};
let mut cx = (**self.agent.context()).clone();
cx.tools = Arc::new(ToolCtx {
workspace,
..(*self.agent.ctx()).clone()
});
let approver = SlackApprover::new(
key.clone(),
Arc::clone(&mode),
self.approval_tx.clone(),
Duration::from_secs(self.cfg.approval_timeout_secs),
);
let unanswered = approver.unanswered_latch();
cx.approver = Arc::new(approver);
cx.budget = Budget {
max_turns: Some(self.cfg.max_turns),
max_cost_usd: self.cfg.max_cost_usd,
..Budget::default()
};
cx.cancel = Some(cancel.clone());
cx.queued_input = Some(Arc::clone(&queue));
if let Some(shared) = &self.agent.context().outbox {
let Ok(store) = mecha_core::outbox::OutboxStore::open(&self.outbox_root) else {
return;
};
let mine = OutboxRoute::new(
store,
shared.routed().map(String::from).collect::<Vec<_>>(),
shared.publishes().map(String::from).collect::<Vec<_>>(),
);
mine.set_session_id(&session.meta.id);
cx.outbox = Some(Arc::new(mine));
}
let mut prompt = format!(
"[This thread's workspace is {}. Relative paths resolve there.]\n\n{prompt}",
cx.tools.workspace.display()
);
if !files.is_empty() {
let landed = self.fetch_attachments(&files, &cx.tools.workspace).await;
if !landed.is_empty() {
prompt.push_str("\n\nThe user attached:\n");
for path in &landed {
prompt.push_str(&format!("- {path}\n"));
}
}
}
let staged_before = pending_outbox_ids(&self.outbox_root);
let files_before = workspace_snapshot(&cx.tools.workspace);
let mut conversation = self.conversations.remove(&key).unwrap_or_default();
conversation.messages.push(Message::user(&prompt));
let controls_channel = channel.clone();
let controls_thread_ts = thread_ts.clone();
let (events_tx, events_rx) = mpsc::unbounded_channel::<AgentEvent>();
let agent = Arc::clone(&self.agent);
let slack = self.slack.clone();
let completion_tx = self.completion_tx.clone();
let pump_cfg = PumpConfig {
flush_chars: self.cfg.stream_flush_chars,
flush_ms: self.cfg.stream_flush_ms,
};
let key_for_task = key.clone();
let session_for_task = session;
self.staged_before.insert(key.clone(), staged_before);
self.files_before.insert(key.clone(), files_before);
tokio::spawn(async move {
let renderer = {
let slack = slack.clone();
let channel = channel.clone();
let thread_ts = thread_ts.clone();
tokio::spawn(async move {
pump(&slack, &channel, &thread_ts, events_rx, &pump_cfg).await
})
};
let before = conversation.messages.clone();
let outcome = agent
.run_in(&cx, &mut conversation, Some(events_tx))
.await
.map(Box::new)
.map_err(|e| e.to_string());
let _ = renderer.await;
let _ = session_for_task.record_run(&before, &conversation);
if let Ok(o) = &outcome {
let _ = session_for_task.record_outcome(o);
}
let _ = completion_tx
.send(Completion {
key: key_for_task,
conversation,
outcome,
})
.await;
});
let mode_label = mode.lock().map(|m| m.as_str()).unwrap_or("ask");
let controls = chat::post_message(
&self.slack,
&controls_channel,
Some(&controls_thread_ts),
"Working…",
Some(controls_blocks(&key, mode_label)),
)
.await;
println!("[{key}] run started");
self.live.insert(
key.clone(),
Live {
cancel,
queue,
mode,
unanswered,
},
);
let _ = self.threads.apply(&key, Event::OwnerSpoke);
if let Ok(Some(mut r)) = self.threads.get(&key) {
r.run = Some(RunMarker::here());
r.controls_ts = controls.ok();
let _ = self.threads.put(&r);
}
}
async fn fetch_attachments(
&self,
files: &[FileRef],
workspace: &std::path::Path,
) -> Vec<String> {
let inbox = workspace.join("inbox");
if let Err(e) = std::fs::create_dir_all(&inbox) {
tracing::warn!("could not make an inbox directory: {e}");
return Vec::new();
}
let mut landed = Vec::new();
for file in files {
let name = safe_filename(file.name.as_deref().unwrap_or(&file.id));
match mecha_slack::files::download(
&self.slack,
file,
self.cfg.max_upload_mb.saturating_mul(1024 * 1024),
)
.await
{
Ok(bytes) => {
let path = inbox.join(&name);
match std::fs::write(&path, &bytes) {
Ok(()) => landed.push(format!("./inbox/{name}")),
Err(e) => tracing::warn!("could not save {name}: {e}"),
}
}
Err(e) => tracing::warn!("could not fetch {}: {e}", file.id),
}
}
landed
}
async fn on_completion(&mut self, done: Completion) {
match &done.outcome {
Ok(o) => println!("[{}] run finished — {} turns", done.key, o.turns),
Err(e) => println!("[{}] run failed — {e}", done.key),
}
self.live.remove(&done.key);
self.conversations
.insert(done.key.clone(), done.conversation);
if let Ok(Some(r)) = self.threads.get(&done.key) {
if let Some(ts) = &r.controls_ts {
let ended = if done.outcome.is_ok() {
"✓ Run complete"
} else {
"✗ Run failed"
};
let _ = chat::update(
&self.slack,
&r.channel_id,
ts,
ended,
Some(vec![blocks::context(&format!(
"{ended} · mode `{}`",
r.mode
))]),
)
.await;
}
}
match done.outcome {
Ok(outcome) => {
let staged = outcome.tool_calls.iter().any(|c| c.staged);
let finished_clean = !outcome.stop_cause.is_early();
let _ = self.threads.apply(&done.key, Event::Finished { staged });
self.post_artifacts(&done.key).await;
if staged {
self.offer_drafts(&done.key, finished_clean).await;
}
}
Err(error) => {
let _ = self.threads.apply(&done.key, Event::Errored);
if let Ok(Some(r)) = self.threads.get(&done.key) {
let _ = chat::post_message(
&self.slack,
&r.channel_id,
Some(&r.thread_ts),
&format!("The run failed: {error} — send `doctor` for a health report"),
None,
)
.await;
}
}
}
}
async fn cycle_mode(&mut self, key: &str) {
let Ok(Some(mut record)) = self.threads.get(key) else {
return;
};
let next = match Mode::parse(&record.mode).unwrap_or(Mode::Ask) {
Mode::Ask => Mode::Allow,
Mode::Allow => Mode::ReadOnly,
Mode::ReadOnly => Mode::Ask,
};
record.mode = next.as_str().to_string();
let _ = self.threads.put(&record);
if let Some(live) = self.live.get(key) {
if let Ok(mut m) = live.mode.lock() {
*m = next;
}
}
if let Some(ts) = &record.controls_ts {
let _ = chat::update(
&self.slack,
&record.channel_id,
ts,
"Working…",
Some(controls_blocks(key, next.as_str())),
)
.await;
}
}
async fn session_for(&mut self, record: &ThreadRecord) -> Option<Session> {
let dir = Session::default_dir().ok()?;
let session = Session::create(
&dir,
SessionMeta {
id: Session::new_id(),
created_at: chrono::Utc::now(),
provider: self.provider.clone(),
model: self.model.clone(),
workspace: record.workspace.clone().unwrap_or_default(),
title: Some(format!("slack: {}", record.key)),
},
)
.ok()?;
let _ = session.append(&Record::Config(mecha_core::session::RunConfig::of(
&self.agent,
&self.config,
&self.provider,
)));
if let Ok(Some(mut r)) = self.threads.get(&record.key) {
r.session_id = Some(session.meta.id.clone());
let _ = self.threads.put(&r);
}
Some(session)
}
async fn post_artifacts(&mut self, key: &str) {
const MAX_FILES: usize = 5;
let Ok(Some(record)) = self.threads.get(key) else {
return;
};
let before = self.files_before.remove(key).unwrap_or_default();
let Some(workspace) = record
.workspace
.clone()
.or_else(|| thread_workspace(key).ok())
else {
return;
};
let after = workspace_snapshot(&workspace);
let mut changed: Vec<PathBuf> = after
.iter()
.filter(|(path, stamp)| before.get(*path).is_none_or(|old| old != *stamp))
.map(|(path, _)| path.clone())
.collect();
if changed.is_empty() {
return;
}
changed.sort();
let limit = self.cfg.max_upload_mb.saturating_mul(1024 * 1024);
let (send, skipped): (Vec<_>, Vec<_>) = changed.iter().partition(|p| {
std::fs::metadata(p)
.map(|m| m.len() <= limit)
.unwrap_or(false)
});
let mut named = Vec::new();
for path in send.iter().take(MAX_FILES) {
let Ok(bytes) = std::fs::read(path) else {
continue;
};
let name = path
.file_name()
.map(|n| n.to_string_lossy().to_string())
.unwrap_or_else(|| "artifact".into());
let share = mecha_slack::files::Share {
channel_id: Some(&record.channel_id),
thread_ts: Some(&record.thread_ts),
title: Some(&name),
..Default::default()
};
if let Err(e) = mecha_slack::files::upload(&self.slack, &name, &bytes, &share).await {
tracing::warn!("could not upload {name}: {e}");
} else {
named.push(name);
}
}
let over = send.len().saturating_sub(MAX_FILES);
if over > 0 || !skipped.is_empty() {
let mut note = String::new();
if over > 0 {
note.push_str(&format!("{over} more file(s) changed and were not sent. "));
}
if !skipped.is_empty() {
note.push_str(&format!(
"{} file(s) are over the {} MB limit: {}.",
skipped.len(),
self.cfg.max_upload_mb,
skipped
.iter()
.filter_map(|p| p.file_name().map(|n| n.to_string_lossy().to_string()))
.collect::<Vec<_>>()
.join(", ")
));
}
let _ = chat::post_message(
&self.slack,
&record.channel_id,
Some(&record.thread_ts),
note.trim(),
None,
)
.await;
}
if !named.is_empty() {
println!("[{key}] posted {} artifact(s)", named.len());
}
}
async fn offer_drafts(&mut self, key: &str, finished_clean: bool) {
let Ok(Some(record)) = self.threads.get(key) else {
return;
};
let Some(session_id) = record.session_id.clone() else {
return;
};
let before = self.staged_before.remove(key).unwrap_or_default();
let Ok(store) = mecha_core::outbox::OutboxStore::open(&self.outbox_root) else {
return;
};
let Ok(items) = store.items() else { return };
let fresh: Vec<_> = items
.into_iter()
.filter(|i| {
i.status == "pending"
&& i.session_id.as_deref() == Some(session_id.as_str())
&& !before.contains(&i.id)
})
.collect();
if fresh.is_empty() {
return;
}
let setting = review::effective(&self.review, key, &record.channel_id).cloned();
let mode = setting.as_ref().map(|s| s.mode).unwrap_or(ReviewMode::Now);
if mode == ReviewMode::Later {
let _ = chat::post_message(
&self.slack,
&record.channel_id,
Some(&record.thread_ts),
&format!(
"{} draft(s) staged — waiting in the outbox (`mecha outbox review`, \
or send `review now` and re-run).",
fresh.len()
),
None,
)
.await;
return;
}
if mode == ReviewMode::Auto && !finished_clean {
let _ = chat::post_message(
&self.slack,
&record.channel_id,
Some(&record.thread_ts),
"The run stopped early — nothing auto-releases. Its drafts are \
carded below for review.",
None,
)
.await;
}
for item in fresh {
let armed = item.taint.private && item.taint.untrusted;
if releases_without_card(setting.as_ref(), armed, finished_clean) {
let setter = setting
.as_ref()
.map(|s| s.set_by.clone())
.unwrap_or_default();
self.dispatch_action(
Action::OutboxSend {
id: item.id.clone(),
},
&setter,
"review-auto",
ActionCard::Reply {
channel: record.channel_id.clone(),
thread_ts: record.thread_ts.clone(),
},
);
continue;
}
let _ = chat::post_message(
&self.slack,
&record.channel_id,
Some(&record.thread_ts),
"A draft is waiting for review.",
Some(draft_offer_blocks(&item)),
)
.await;
}
}
async fn outbox_item(&self, id: &str) -> Option<mecha_core::outbox::OutboxItem> {
let root = self.outbox_root.clone();
let id = id.to_string();
tokio::task::spawn_blocking(move || {
mecha_core::outbox::OutboxStore::open(&root)
.ok()?
.item_exact(&id)
.ok()
.flatten()
})
.await
.ok()
.flatten()
}
async fn ask_to_confirm_tainted(
&mut self,
id: &str,
item: Option<&mecha_core::outbox::OutboxItem>,
channel: &str,
ts: &str,
) {
let (text, card) = tainted_confirm_card(id, item);
let _ = chat::update(&self.slack, channel, ts, &text, Some(card)).await;
}
fn dispatch_action(
&mut self,
action: Action,
who: &str,
surface: &'static str,
target: ActionCard,
) {
let tap_id = actions::new_tap_id();
let slack = self.slack.clone();
let ledger = Arc::clone(&self.ledger);
let executor = Executor {
outbox_root: self.outbox_root.clone(),
};
let who = who.to_string();
tokio::spawn(async move {
ledger.dispatched(&tap_id, &who, &action, surface);
let pending = format!("⏳ {} — <@{}>", action.describe(), who);
let (channel, card_ts, reply_ts) = match &target {
ActionCard::Rewrite {
channel,
ts,
thread_ts,
} => {
let _ = chat::update(
&slack,
channel,
ts,
&pending,
Some(vec![blocks::context(&pending)]),
)
.await;
(
channel.clone(),
Some(ts.clone()),
thread_ts.clone().unwrap_or_else(|| ts.clone()),
)
}
ActionCard::Reply { channel, thread_ts } => {
let ts = chat::post_message(&slack, channel, Some(thread_ts), &pending, None)
.await
.ok();
(channel.clone(), ts, thread_ts.clone())
}
};
let started = std::time::Instant::now();
let outcome = executor.run(&action).await;
ledger.resolved(&tap_id, &outcome.status, &outcome.line);
let line = format!("{} · <@{}>", outcome.line, who);
let (update_card, fresh_reply) = outcome_delivery(
card_ts.is_some(),
started.elapsed() > Duration::from_secs(60),
);
if update_card {
if let Some(ts) = &card_ts {
let _ = chat::update(
&slack,
&channel,
ts,
&line,
Some(vec![blocks::context(&line)]),
)
.await;
}
}
if fresh_reply {
let _ = chat::post_message(&slack, &channel, Some(&reply_ts), &line, None).await;
}
if let Action::OutboxSend { id } = &action {
let item = executor.item(id).await;
if failed_auto_release_needs_card(surface, item.as_ref().map(|i| i.status.as_str()))
{
let item = item.expect("checked by the predicate");
let _ = chat::post_message(
&slack,
&channel,
Some(&reply_ts),
"The auto-release failed and the draft is still pending — review it \
here.",
Some(draft_offer_blocks(&item)),
)
.await;
}
}
});
}
async fn retire_expired_approvals(&mut self) {
let now = std::time::Instant::now();
let expired: Vec<String> = self
.pending
.iter()
.filter(|(_, p)| p.expires_at <= now)
.map(|(id, _)| id.clone())
.collect();
for id in expired {
if let Some(p) = self.pending.remove(&id) {
let line = format!(
"`{}` was not approved — nobody answered in time, and the call was \
refused.",
p.tool
);
let _ = chat::update(
&self.slack,
&p.channel,
&p.message_ts,
&line,
Some(vec![blocks::context(&line)]),
)
.await;
let _ = self.threads.apply(&p.thread_key, Event::InputSettled);
}
}
}
async fn refuse_pending_for(&mut self, thread_key: &str) {
let theirs: Vec<String> = self
.pending
.iter()
.filter(|(_, p)| p.thread_key == thread_key)
.map(|(id, _)| id.clone())
.collect();
for id in theirs {
if let Some(p) = self.pending.remove(&id) {
drop(p.reply);
let line = format!("`{}` was not approved — the run was stopped.", p.tool);
let _ = chat::update(
&self.slack,
&p.channel,
&p.message_ts,
&line,
Some(vec![blocks::context(&line)]),
)
.await;
}
}
}
async fn on_approval_request(&mut self, request: approve::Request) {
let Ok(Some(record)) = self.threads.get(&request.thread_key) else {
return;
};
self.approval_seq += 1;
let id = format!("{}-{}", request.thread_key, self.approval_seq);
let card = vec![
blocks::section(&format!("*Approve this call?*\n`{}`", request.summary)),
blocks::context(&format!("thread {} · {}", record.thread_ts, request.tool)),
blocks::actions(vec![
blocks::button("slack_approve", "Approve", &id, Some("primary")),
blocks::button("slack_approve_run", "Allow for this run", &id, None),
blocks::button("slack_reject", "Reject", &id, Some("danger")),
]),
];
match chat::post_message(
&self.slack,
&record.channel_id,
Some(&record.thread_ts),
"An approval is waiting.",
Some(card),
)
.await
{
Ok(ts) => {
let _ = self
.threads
.apply(&request.thread_key, Event::AskedForInput);
self.pending.insert(
id,
PendingApproval {
reply: request.reply,
channel: record.channel_id,
message_ts: ts,
tool: request.tool,
thread_key: request.thread_key.clone(),
expires_at: std::time::Instant::now()
+ Duration::from_secs(self.cfg.approval_timeout_secs),
},
);
}
Err(e) => {
tracing::warn!("could not post an approval card: {e}");
}
}
}
async fn on_interaction(&mut self, interaction: Interaction) {
let gate = self.allowed(
interaction.user_id.as_deref(),
interaction.team_id.as_deref(),
);
if !gate.is_allowed() {
tracing::warn!(
"dropped an interaction ({}): user {}, channel {}",
gate.reason(),
interaction.user_id.as_deref().unwrap_or("<none>"),
interaction.channel_id.as_deref().unwrap_or("<none>"),
);
return;
}
if interaction.kind == "view_submission" {
self.on_view_submission(&interaction);
return;
}
for action in &interaction.actions {
let value = action.value.clone().unwrap_or_default();
match action.action_id.as_str() {
"slack_stop" => {
if let Some(live) = self.live.get(&value) {
live.cancel.cancel();
self.refuse_pending_for(&value).await;
let _ = self.threads.apply(&value, Event::StopPressed);
}
}
"slack_mode" => self.cycle_mode(&value).await,
"slack_approve" | "slack_approve_run" | "slack_reject" => {
if let Some((key, _)) = value.rsplit_once('-') {
if let Some(live) = self.live.get(key) {
live.unanswered.store(false, Ordering::Relaxed);
}
}
if let Some(pending) = self.pending.remove(&value) {
let answer = match action.action_id.as_str() {
"slack_approve" => Answer::Approve,
"slack_approve_run" => Answer::ApproveForRun,
_ => Answer::Reject("rejected from Slack".into()),
};
let verb = match action.action_id.as_str() {
"slack_approve" => "approved",
"slack_approve_run" => "allowed for this run",
_ => "rejected",
};
let landed = pending.reply.send(answer).is_ok();
if !landed {
let line = format!(
"`{}` was already refused before this was pressed — nothing ran.",
pending.tool
);
let _ = chat::update(
&self.slack,
&pending.channel,
&pending.message_ts,
&line,
Some(vec![blocks::context(&line)]),
)
.await;
continue;
}
let who = interaction.user_id.as_deref().unwrap_or("someone");
let _ = chat::update(
&self.slack,
&pending.channel,
&pending.message_ts,
&format!("`{}` {verb} by <@{who}>", pending.tool),
Some(vec![blocks::context(&format!(
"`{}` {verb} by <@{who}>",
pending.tool
))]),
)
.await;
}
if let Some(key) = value.rsplit_once('-').map(|(k, _)| k.to_string()) {
let _ = self.threads.apply(&key, Event::InputSettled);
}
}
super::doctor::REVIEW_HERE_OUTBOX | super::doctor::REVIEW_HERE_FRONTDOOR => {
let Some(channel) = interaction.channel_id.clone() else {
continue;
};
let Some(thread_ts) = interaction
.thread_ts
.clone()
.or_else(|| interaction.message_ts.clone())
else {
continue;
};
let slack = self.slack.clone();
if action.action_id == super::doctor::REVIEW_HERE_OUTBOX {
let root = self.outbox_root.clone();
tokio::spawn(async move {
review_here_drafts(&slack, &root, &channel, &thread_ts).await;
});
} else {
tokio::spawn(async move {
review_here_requests(&slack, &channel, &thread_ts).await;
});
}
}
super::frontdoor::CLOSE_OPEN | super::frontdoor::NEEDS_INFO_OPEN => {
let Some(trigger) = interaction.trigger_id.clone() else {
continue;
};
let Some(seq) = value.parse::<i64>().ok().filter(|s| *s > 0) else {
continue;
};
let meta = super::frontdoor::metadata(
seq,
interaction.channel_id.as_deref(),
interaction.message_ts.as_deref(),
interaction.thread_ts.as_deref(),
);
let view = if action.action_id == super::frontdoor::CLOSE_OPEN {
super::frontdoor::close_modal(seq, &meta)
} else {
super::frontdoor::needs_info_modal(seq, &meta)
};
if let Err(e) = mecha_slack::views::open(&self.slack, &trigger, view).await {
tracing::warn!("could not open the modal for request {seq}: {e}");
if let (Some(channel), Some(ts)) = (
interaction.channel_id.as_deref(),
interaction
.thread_ts
.as_deref()
.or(interaction.message_ts.as_deref()),
) {
let _ = chat::post_message(
&self.slack,
channel,
Some(ts),
&format!(
"The form could not open. At a terminal: \
`mecha frontdoor close {seq} --reason ...` or \
`mecha frontdoor needs-info {seq} --note ...`"
),
None,
)
.await;
}
}
}
actions::ids::OUTBOX_SEND_CONFIRM => {
let (Some(channel), Some(ts)) = (
interaction.channel_id.clone(),
interaction.message_ts.clone(),
) else {
continue;
};
let Some((id, _)) = actions::parse_confirm_value(&value) else {
continue;
};
let item = self.outbox_item(id).await;
match confirm_press(&value, item.as_ref()) {
ConfirmPress::Send { id } => {
let who = interaction
.user_id
.as_deref()
.unwrap_or("someone")
.to_string();
self.dispatch_action(
Action::OutboxSend { id },
&who,
"draft-card",
ActionCard::Rewrite {
channel,
ts,
thread_ts: interaction.thread_ts.clone(),
},
);
}
ConfirmPress::Recard { id } => {
let (_, mut card) = tainted_confirm_card(&id, item.as_ref());
card.insert(
0,
blocks::context(
"⚠️ The draft changed after this card was composed — \
nothing was sent. Re-read it.",
),
);
let _ = chat::update(
&self.slack,
&channel,
&ts,
"This draft changed after its card was composed — nothing \
was sent.",
Some(card),
)
.await;
}
ConfirmPress::Unreadable { id } => {
let (text, card) = tainted_confirm_card(&id, None);
let _ =
chat::update(&self.slack, &channel, &ts, &text, Some(card)).await;
}
}
}
other => {
let Some(act) = Action::from_payload(other, &value) else {
continue;
};
let who = interaction
.user_id
.as_deref()
.unwrap_or("someone")
.to_string();
let (Some(channel), Some(ts)) = (
interaction.channel_id.clone(),
interaction.message_ts.clone(),
) else {
continue;
};
if other == actions::ids::OUTBOX_SEND {
let item = self.outbox_item(&value).await;
if item
.as_ref()
.is_some_and(|i| i.taint.private && i.taint.untrusted)
{
self.ask_to_confirm_tainted(&value, item.as_ref(), &channel, &ts)
.await;
continue;
}
}
let (surface, target) = match &act {
Action::OutboxSend { .. } | Action::OutboxReject { .. } => (
"draft-card",
ActionCard::Rewrite {
channel,
ts,
thread_ts: interaction.thread_ts.clone(),
},
),
_ => (
"doctor",
ActionCard::Reply {
thread_ts: interaction
.thread_ts
.clone()
.unwrap_or_else(|| ts.clone()),
channel,
},
),
};
self.dispatch_action(act, &who, surface, target);
}
}
}
}
fn on_view_submission(&mut self, interaction: &Interaction) {
let Some(view) = &interaction.view else {
return;
};
let Some(meta) = super::frontdoor::parse_metadata(&view.private_metadata) else {
tracing::warn!("a view submission carried unusable metadata; nothing ran");
return;
};
let Some(text) = view.values.get(super::frontdoor::TEXT_INPUT) else {
tracing::warn!("a view submission carried no text; nothing ran");
return;
};
let Some(act) = Action::from_submission(&view.callback_id, &meta.seq.to_string(), text)
else {
tracing::warn!(
"a view submission ({}) parsed to no action; nothing ran",
view.callback_id
);
return;
};
let who = interaction
.user_id
.as_deref()
.unwrap_or("someone")
.to_string();
let target = match (meta.channel, meta.ts, meta.thread_ts) {
(Some(channel), Some(ts), thread_ts) => ActionCard::Rewrite {
channel,
ts,
thread_ts,
},
(Some(channel), None, Some(thread_ts)) => ActionCard::Reply { channel, thread_ts },
_ => {
tracing::warn!("a view submission had nowhere to report; nothing ran");
return;
}
};
self.dispatch_action(act, &who, "frontdoor-modal", target);
}
}
async fn review_here_drafts(
slack: &Slack,
outbox_root: &std::path::Path,
channel: &str,
thread_ts: &str,
) {
let items = mecha_core::outbox::OutboxStore::open(outbox_root)
.ok()
.and_then(|s| s.items().ok())
.unwrap_or_default();
let mut pending: Vec<_> = items
.into_iter()
.filter(|i| i.status == "pending")
.collect();
pending.sort_by(|a, b| a.id.cmp(&b.id));
if pending.is_empty() {
let _ = chat::post_message(
slack,
channel,
Some(thread_ts),
"Nothing is pending in the outbox.",
None,
)
.await;
return;
}
let total = pending.len();
for item in pending.iter().take(REVIEW_HERE_MAX) {
let _ = chat::post_message(
slack,
channel,
Some(thread_ts),
"A draft is waiting for review.",
Some(draft_offer_blocks(item)),
)
.await;
}
if let Some(note) = review_here_note(total, "mecha outbox review") {
let _ = chat::post_message(slack, channel, Some(thread_ts), ¬e, None).await;
}
}
async fn review_here_requests(slack: &Slack, channel: &str, thread_ts: &str) {
let records = mecha_core::frontdoor::Frontdoor::open_default()
.and_then(|f| f.records())
.unwrap_or_default();
let waiting = super::frontdoor::waiting(&records);
if waiting.is_empty() {
let _ = chat::post_message(
slack,
channel,
Some(thread_ts),
"Nothing is waiting at the front door.",
None,
)
.await;
return;
}
let total = waiting.len();
for record in waiting.iter().take(REVIEW_HERE_MAX) {
let _ = chat::post_message(
slack,
channel,
Some(thread_ts),
&format!("Request {} is waiting.", record.seq),
Some(super::frontdoor::request_card(record)),
)
.await;
}
if let Some(note) = review_here_note(total, "mecha frontdoor list") {
let _ = chat::post_message(slack, channel, Some(thread_ts), ¬e, None).await;
}
}
enum ActionCard {
Rewrite {
channel: String,
ts: String,
thread_ts: Option<String>,
},
Reply {
channel: String,
thread_ts: String,
},
}
fn outcome_delivery(dispatch_landed: bool, slow: bool) -> (bool, bool) {
(dispatch_landed, !dispatch_landed || slow)
}
fn failed_auto_release_needs_card(surface: &str, item_status: Option<&str>) -> bool {
surface == "review-auto" && item_status == Some("pending")
}
fn draft_offer_blocks(item: &mecha_core::outbox::OutboxItem) -> Vec<serde_json::Value> {
let armed = item.taint.private && item.taint.untrusted;
let heading = if armed {
format!(
"*⚠️ Draft — written with the trifecta armed*\n`{}`",
item.tool
)
} else {
format!("*Draft*\n`{}`", item.tool)
};
vec![
blocks::section(&heading),
blocks::section(&format!("```\n{}\n```", truncate_for_slack(&item.summary))),
blocks::context(&format!("id `{}` · nothing has been sent", item.id)),
blocks::actions(vec![
blocks::button(actions::ids::OUTBOX_SEND, "Send", &item.id, Some("primary")),
blocks::button(
actions::ids::OUTBOX_REJECT,
"Reject",
&item.id,
Some("danger"),
),
]),
]
}
const REVIEW_HERE_MAX: usize = 8;
fn review_here_note(total: usize, terminal_cmd: &str) -> Option<String> {
let rest = total.saturating_sub(REVIEW_HERE_MAX);
(rest > 0)
.then(|| format!("{rest} more not shown here — the rest at the terminal: `{terminal_cmd}`"))
}
fn tainted_confirm_blocks(id: &str, args: &str) -> Vec<serde_json::Value> {
let shown = truncate_for_slack(args);
let cut = shown != args;
let mut card = vec![
blocks::section(
"*⚠️ This draft was written while the trifecta was armed.*\nPrivate data and \
untrusted content were both in the conversation that produced it. The full \
arguments are below — read them before sending.",
),
blocks::section(&format!("```\n{shown}\n```")),
];
if cut {
card.push(blocks::section(&format!(
"*The arguments are longer than this card can show* ({} of {} characters). \
A draft you cannot fully read is not released from here: run \
`mecha outbox show {id}` in a terminal to read all of it and release it \
there.",
shown.chars().count(),
args.chars().count()
)));
card.push(blocks::actions(vec![blocks::button(
actions::ids::OUTBOX_REJECT,
"Reject",
id,
Some("danger"),
)]));
} else {
card.push(blocks::actions(vec![
blocks::button(
actions::ids::OUTBOX_SEND_CONFIRM,
"Send anyway",
&actions::confirm_value(id, args),
Some("danger"),
),
blocks::button(actions::ids::OUTBOX_REJECT, "Reject", id, None),
]));
}
card
}
fn tainted_confirm_card(
id: &str,
item: Option<&mecha_core::outbox::OutboxItem>,
) -> (String, Vec<serde_json::Value>) {
match item {
Some(item) => {
let args = serde_json::to_string_pretty(&item.args).unwrap_or_default();
(
"This draft needs a second look before it is sent.".to_string(),
tainted_confirm_blocks(id, &args),
)
}
None => (
format!("Draft `{id}` could not be read — nothing is offered for sending."),
vec![
blocks::section(&format!(
"*Draft `{id}` could not be read from the outbox.*\nNothing is \
offered for sending on bytes nobody can see. Read and release it \
at a terminal: `mecha outbox show {id}`."
)),
blocks::actions(vec![blocks::button(
actions::ids::OUTBOX_REJECT,
"Reject",
id,
Some("danger"),
)]),
],
),
}
}
#[derive(Debug, PartialEq, Eq)]
enum ConfirmPress {
Send { id: String },
Recard { id: String },
Unreadable { id: String },
}
fn confirm_press(value: &str, item: Option<&mecha_core::outbox::OutboxItem>) -> ConfirmPress {
let (id, fingerprint) = match actions::parse_confirm_value(value) {
Some(parsed) => parsed,
None => {
return ConfirmPress::Recard {
id: value.to_string(),
}
}
};
let Some(item) = item else {
return ConfirmPress::Unreadable { id: id.to_string() };
};
let args = serde_json::to_string_pretty(&item.args).unwrap_or_default();
if truncate_for_slack(&args) != args {
return ConfirmPress::Recard { id: id.to_string() };
}
match fingerprint {
Some(fp) if fp == actions::fingerprint(&args) => ConfirmPress::Send { id: id.to_string() },
_ => ConfirmPress::Recard { id: id.to_string() },
}
}
fn releases_without_card(
setting: Option<&review::Setting>,
tainted: bool,
finished_clean: bool,
) -> bool {
setting.is_some_and(|s| crate::review_policy::auto_releases(s.mode, tainted, finished_clean))
}
fn safe_filename(raw: &str) -> String {
let base = raw.rsplit(['/', '\\']).next().unwrap_or(raw);
let cleaned: String = base
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '.' || c == '-' || c == '_' {
c
} else {
'-'
}
})
.collect();
let cleaned = cleaned.trim_matches(['.', '-']).to_string();
if !cleaned.chars().any(|c| c.is_ascii_alphanumeric()) {
return "attachment".into();
}
cleaned.chars().take(120).collect()
}
fn workspace_snapshot(root: &std::path::Path) -> HashMap<PathBuf, (u64, std::time::SystemTime)> {
fn walk(
dir: &std::path::Path,
depth: usize,
out: &mut HashMap<PathBuf, (u64, std::time::SystemTime)>,
) {
if depth == 0 {
return;
}
let Ok(entries) = std::fs::read_dir(dir) else {
return;
};
for entry in entries.filter_map(std::result::Result::ok) {
let path = entry.path();
let name = entry.file_name().to_string_lossy().to_string();
if name.starts_with('.') || name == "inbox" {
continue;
}
match entry.metadata() {
Ok(meta) if meta.is_dir() => walk(&path, depth - 1, out),
Ok(meta) => {
let stamp = meta.modified().unwrap_or(std::time::UNIX_EPOCH);
out.insert(path, (meta.len(), stamp));
}
Err(_) => {}
}
}
}
let mut out = HashMap::new();
walk(root, 6, &mut out);
out
}
fn pending_outbox_ids(root: &std::path::Path) -> std::collections::HashSet<String> {
let Ok(store) = mecha_core::outbox::OutboxStore::open(root) else {
return Default::default();
};
store
.items()
.map(|items| {
items
.into_iter()
.filter(|i| i.status == "pending")
.map(|i| i.id)
.collect()
})
.unwrap_or_default()
}
fn truncate_for_slack(text: &str) -> String {
mecha_slack::blocks::truncate(text, 2_600)
}
fn controls_blocks(key: &str, mode: &str) -> Vec<serde_json::Value> {
vec![
blocks::context(&format!("mode `{mode}`")),
blocks::actions(vec![
blocks::button("slack_stop", "Stop", key, Some("danger")),
blocks::button("slack_mode", "Mode", key, None),
]),
]
}
fn producer_root() -> Result<std::path::PathBuf> {
let root = mecha_core::work::producer_dir("slack")?;
std::fs::create_dir_all(&root).with_context(|| format!("creating {}", root.display()))?;
Ok(root)
}
fn thread_workspace(key: &str) -> Result<std::path::PathBuf> {
let root = mecha_core::work::producer_dir("slack")?.join(key);
if root.exists() && !root.is_dir() {
bail!("{} exists and is not a directory", root.display());
}
std::fs::create_dir_all(&root).with_context(|| format!("creating {}", root.display()))?;
mecha_core::work::ensure_outside_mecha_home(&root)?;
Ok(root)
}
#[cfg(test)]
mod tests {
use super::{
confirm_press, draft_offer_blocks, failed_auto_release_needs_card, outcome_delivery,
releases_without_card, review_here_note, safe_filename, tainted_confirm_blocks,
tainted_confirm_card, ConfirmPress, REVIEW_HERE_MAX,
};
use crate::review_policy::ReviewMode;
use crate::slack::actions;
use crate::slack::review::Setting;
fn draft(id: &str, tainted: bool, summary: &str) -> mecha_core::outbox::OutboxItem {
mecha_core::outbox::OutboxItem {
id: id.into(),
status: "pending".into(),
tool: "mail__send".into(),
kind: mecha_core::outbox::OutboxKind::Message,
args_before: serde_json::json!({}),
args: serde_json::json!({}),
summary: summary.into(),
session_id: None,
workspace: None,
taint: mecha_core::agent::Taint {
private: tainted,
untrusted: tainted,
},
created_at: "2026-08-14T00:00:00Z".into(),
resolved_at: None,
reason: None,
error: None,
}
}
#[test]
fn a_tainted_draft_in_a_review_here_batch_keeps_the_two_step_and_truncated_stays_reject_only() {
let batch = [
draft("a-1", false, "hello"),
draft("a-2", true, "the tainted one"),
];
let cards: Vec<String> = batch
.iter()
.map(|i| serde_json::to_string(&draft_offer_blocks(i)).unwrap())
.collect();
for card in &cards {
assert!(card.contains(actions::ids::OUTBOX_SEND), "{card}");
assert!(card.contains(actions::ids::OUTBOX_REJECT), "{card}");
assert!(
!card.contains(actions::ids::OUTBOX_SEND_CONFIRM),
"the confirm verb only ever appears on the second step: {card}"
);
}
assert!(cards[1].contains("trifecta armed"), "{}", cards[1]);
assert!(!cards[0].contains("trifecta armed"), "{}", cards[0]);
let second_step =
serde_json::to_string(&tainted_confirm_blocks("a-2", &"x".repeat(10_000))).unwrap();
assert!(
!second_step.contains(actions::ids::OUTBOX_SEND_CONFIRM),
"{second_step}"
);
assert!(
second_step.contains(actions::ids::OUTBOX_REJECT),
"{second_step}"
);
}
#[test]
fn review_auto_releases_nothing_after_an_early_stopped_run() {
let auto = Setting {
mode: ReviewMode::Auto,
set_by: "U_OWNER".into(),
};
assert!(
!releases_without_card(Some(&auto), false, false),
"an interrupted run's untainted draft cards, never auto-releases"
);
assert!(
releases_without_card(Some(&auto), false, true),
"a clean finish still releases the untainted draft"
);
assert!(
!releases_without_card(Some(&auto), true, true),
"tainted never releases, clean finish or not"
);
assert!(
!releases_without_card(None, false, true),
"no setting means now: card everything"
);
for mode in [ReviewMode::Now, ReviewMode::Later] {
let s = Setting {
mode,
set_by: "U_OWNER".into(),
};
assert!(!releases_without_card(Some(&s), false, true));
}
}
#[test]
fn an_unreadable_item_gets_an_error_card_and_never_a_send() {
let (text, card) = tainted_confirm_card("abc-123", None);
let json = serde_json::to_string(&card).unwrap();
assert!(
!json.contains(actions::ids::OUTBOX_SEND_CONFIRM),
"no Send on bytes nobody can see: {json}"
);
assert!(!json.contains("\"slack_outbox_send\""), "{json}");
assert!(json.contains(actions::ids::OUTBOX_REJECT), "{json}");
assert!(
json.contains("could not be read"),
"the error is stated: {json}"
);
assert!(json.contains("mecha outbox show abc-123"), "{json}");
assert!(text.contains("could not be read"), "{text}");
let item = draft("abc-123", true, "hello");
let (_, card) = tainted_confirm_card("abc-123", Some(&item));
let json = serde_json::to_string(&card).unwrap();
assert!(json.contains(actions::ids::OUTBOX_SEND_CONFIRM), "{json}");
}
#[test]
fn a_confirm_press_on_a_draft_edited_after_composition_recards_instead_of_sending() {
let mut item = draft("abc-123", true, "the tainted one");
item.args = serde_json::json!({"to": "a@x.org", "body": "as reviewed"});
let shown = serde_json::to_string_pretty(&item.args).unwrap();
let value = actions::confirm_value("abc-123", &shown);
assert_eq!(
confirm_press(&value, Some(&item)),
ConfirmPress::Send {
id: "abc-123".into()
}
);
item.args = serde_json::json!({"to": "attacker@evil.example", "body": "as reviewed"});
assert_eq!(
confirm_press(&value, Some(&item)),
ConfirmPress::Recard {
id: "abc-123".into()
},
"changed bytes must never send"
);
assert_eq!(
confirm_press(&value, None),
ConfirmPress::Unreadable {
id: "abc-123".into()
}
);
assert_eq!(
confirm_press("abc-123", Some(&item)),
ConfirmPress::Recard {
id: "abc-123".into()
}
);
item.args = serde_json::json!({"body": "x".repeat(10_000)});
let long = serde_json::to_string_pretty(&item.args).unwrap();
let long_value = actions::confirm_value("abc-123", &long);
assert_eq!(
confirm_press(&long_value, Some(&item)),
ConfirmPress::Recard {
id: "abc-123".into()
}
);
}
#[test]
fn a_tainted_confirm_cards_send_button_carries_the_shown_bytes_fingerprint() {
let args = "{\n \"to\": \"a@x.org\"\n}";
let json = serde_json::to_string(&tainted_confirm_blocks("abc-123", args)).unwrap();
assert!(
json.contains(&actions::confirm_value("abc-123", args)),
"the value is id#fingerprint: {json}"
);
assert!(json.contains("\"abc-123\""), "{json}");
}
#[test]
fn an_outcome_lands_exactly_once_plus_a_fresh_reply_when_slow_or_lost() {
assert_eq!(outcome_delivery(true, false), (true, false));
assert_eq!(
outcome_delivery(true, true),
(true, true),
"a slow rewrite earns the notifying reply too"
);
assert_eq!(outcome_delivery(false, false), (false, true));
assert_eq!(
outcome_delivery(false, true),
(false, true),
"never two fresh posts — that was the double-post"
);
}
#[test]
fn a_failed_auto_release_gets_the_draft_carded_for_the_phone() {
assert!(failed_auto_release_needs_card(
"review-auto",
Some("pending")
));
assert!(
!failed_auto_release_needs_card("review-auto", Some("sent")),
"a release that landed needs no card"
);
assert!(
!failed_auto_release_needs_card("draft-card", Some("pending")),
"a card-tapped release already has its card"
);
assert!(
!failed_auto_release_needs_card("review-auto", None),
"an unreadable item is reported, not guessed into a card"
);
}
#[test]
fn a_non_owners_view_submission_is_gated_on_the_signed_user_before_anything_is_read() {
use mecha_slack::binding::{self, Binding};
let bound = Binding {
team_id: "T1".into(),
enterprise_id: None,
owners: vec!["U_OWNER".into()],
bound_at: chrono::Utc::now(),
};
let stranger = binding::check(Some(&bound), Some("U_STRANGER"), Some("T1"));
assert!(
!stranger.is_allowed(),
"a stranger's submission is refused before the view is read"
);
assert!(actions::Action::from_submission(
actions::ids::FRONTDOOR_CLOSE_SUBMIT,
"5",
"a perfectly plausible reason"
)
.is_some());
let owner = binding::check(Some(&bound), Some("U_OWNER"), Some("T1"));
assert!(owner.is_allowed());
}
#[test]
fn a_review_here_batch_past_the_cap_names_the_terminal_for_the_rest() {
assert_eq!(
review_here_note(REVIEW_HERE_MAX, "mecha outbox review"),
None
);
let note =
review_here_note(REVIEW_HERE_MAX + 3, "mecha outbox review").expect("the cut says so");
assert!(note.contains("3 more"), "{note}");
assert!(note.contains("mecha outbox review"), "{note}");
let front = review_here_note(20, "mecha frontdoor list").unwrap();
assert!(front.contains("mecha frontdoor list"), "{front}");
}
#[test]
fn a_truncated_tainted_confirm_card_offers_reject_and_never_send() {
let long_args = "x".repeat(10_000);
let card = serde_json::to_string(&tainted_confirm_blocks("abc-123", &long_args)).unwrap();
assert!(
!card.contains("slack_outbox_send_confirm"),
"no Send on a draft the card could not fully show: {card}"
);
assert!(
card.contains("slack_outbox_reject"),
"declining needs no completeness: {card}"
);
assert!(card.contains("mecha outbox show abc-123"), "{card}");
assert!(card.contains("characters"), "{card}");
}
#[test]
fn a_tainted_confirm_card_that_fits_keeps_the_informed_send() {
let card =
serde_json::to_string(&tainted_confirm_blocks("abc-123", "{\"to\": \"a@x.org\"}"))
.unwrap();
assert!(card.contains("slack_outbox_send_confirm"), "{card}");
assert!(card.contains("slack_outbox_reject"), "{card}");
}
#[test]
fn a_review_mode_is_never_written_to_the_thread_record() {
let record = super::super::threads::ThreadRecord {
key: "D1-1".into(),
channel_id: "D1".into(),
thread_ts: "1.0".into(),
state: super::super::threads::ThreadState::Idle,
session_id: None,
workspace: None,
mode: "ask".into(),
last_seen_ts: None,
run: None,
stream_ts: None,
controls_ts: None,
updated_at: chrono::Utc::now(),
};
let json = serde_json::to_string(&record).unwrap();
assert!(
!json.contains("review"),
"a persisted review mode would outlive the attention that set it: {json}"
);
}
#[test]
fn approval_ids_never_repeat_even_as_entries_are_removed() {
let mut seq = 0u64;
let mut pending: std::collections::HashSet<String> = Default::default();
let mut minted = Vec::new();
for round in 0..50 {
seq += 1;
let id = format!("D1-1.0-{seq}");
assert!(pending.insert(id.clone()), "id {id} was reused");
minted.push(id.clone());
if round % 2 == 1 {
if let Some(old) = minted.first().cloned() {
pending.remove(&old);
minted.remove(0);
}
}
}
assert_eq!(seq, 50);
}
#[test]
fn a_snapshot_ignores_what_the_user_sent_and_what_is_hidden() {
let dir = std::env::temp_dir().join(format!(
"mecha-slack-snap-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(dir.join("inbox")).unwrap();
std::fs::create_dir_all(dir.join("sub")).unwrap();
std::fs::write(dir.join("made.txt"), "x").unwrap();
std::fs::write(dir.join("sub/also.txt"), "y").unwrap();
std::fs::write(dir.join("inbox/sent-by-user.png"), "z").unwrap();
std::fs::write(dir.join(".hidden"), "h").unwrap();
let snap = super::workspace_snapshot(&dir);
let names: Vec<String> = snap
.keys()
.filter_map(|p| p.file_name().map(|n| n.to_string_lossy().to_string()))
.collect();
assert!(names.contains(&"made.txt".to_string()), "{names:?}");
assert!(
names.contains(&"also.txt".to_string()),
"recurses: {names:?}"
);
assert!(
!names.contains(&"sent-by-user.png".to_string()),
"{names:?}"
);
assert!(!names.contains(&".hidden".to_string()), "{names:?}");
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_draft_summary_fits_inside_a_section_block_with_its_fence() {
let long = super::truncate_for_slack(&"x".repeat(10_000));
let fenced = format!("```\n{long}\n```");
assert!(
fenced.chars().count() < mecha_slack::blocks::limits::SECTION_TEXT,
"{} chars",
fenced.chars().count()
);
assert!(long.contains("truncated"), "and the cut says so");
}
#[test]
fn an_attachment_name_cannot_climb_out_of_the_inbox() {
for hostile in [
"../../.mecha/slack/credentials.json",
"..\\..\\windows",
"/etc/passwd",
"..",
"...",
] {
let safe = safe_filename(hostile);
assert!(!safe.contains('/'), "{hostile} -> {safe}");
assert!(!safe.contains('\\'), "{hostile} -> {safe}");
assert!(!safe.starts_with('.'), "{hostile} -> {safe}");
assert!(!safe.is_empty(), "{hostile} -> {safe}");
}
}
#[test]
fn an_ordinary_name_survives_intact() {
assert_eq!(safe_filename("screenshot.png"), "screenshot.png");
assert_eq!(safe_filename("notes-2026_08.md"), "notes-2026_08.md");
}
#[test]
fn a_nameless_attachment_still_lands_somewhere() {
assert_eq!(safe_filename(""), "attachment");
assert_eq!(safe_filename("???"), "attachment");
}
#[test]
fn a_threads_jail_sits_under_one_producer_so_retention_reaches_it() {
let path = mecha_core::work::producer_dir("slack")
.unwrap()
.join("D1-1");
let parent = path.parent().unwrap();
assert!(
parent.ends_with("slack"),
"every thread must be an entry inside one producer: {}",
path.display()
);
}
}