Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,17 @@ The version is carried in `pyproject.toml` and `src/assistant/__init__.py`
and reminders to finish older pending work before starting another task.

### Fixed
- Ping about an idle workspace only when it needs you: the agent's last lines
ask you something, a permission or plan prompt is up, or the session stopped
on an API error. Claude's idle alert fires about a minute after every turn,
and about 75% of recent turns ended with a status update instead.
- Match the watcher's patterns ("CI red", "PR opened", "tests passed") against
what the agent said this turn, not 50 lines of scrollback, so old tool output
no longer triggers them. Name workspaces cmux left as "Terminal" by Claude's own
session title, or the folder and branch.
- Mute the "stranded", "awaiting-review", and "emit-card" word lists by default
(168 pings in two weeks, mostly on status prose), and stop a muted pattern
from hiding a real one in the same message.
- Stop closing a healthy warm comms session whenever cmux is slow to answer.
A refused or timed-out check now leaves the session alone; only cmux saying
the workspace doesn't exist, or a successful workspace list without it,
Expand Down
183 changes: 153 additions & 30 deletions bin/cmux-watcher.py
Original file line number Diff line number Diff line change
Expand Up @@ -98,19 +98,24 @@
"patterns": [
{"id": "pr-opened", "regex": r"PR #\d+ (opened|created|shipped)",
"signal": "work_complete", "priority": "high"},
# The three needs-input word lists below matched status prose far more
# than real asks ("Awaiting the adversary, code-review…", "blocked" in a
# PR review): 168 pings in two weeks. A real ask now arrives as an
# AskUserQuestion or an idle alert whose last line asks something, and
# the pulse's Observer handles stalled workspaces itself.
{"id": "awaiting-review",
"regex": r"(awaiting|waiting for|needs).{0,30}(review|approval|your input)",
"signal": "needs_input", "priority": "high"},
"signal": "needs_input", "priority": "high", "suppress": True},
{"id": "ci-green",
"regex": r"(all CI (checks )?green|CI (is )?green|✓.*CI|Jenkins.*SUCCESS)",
"signal": "work_complete", "priority": "medium", "suppress": True},
{"id": "ci-red",
"regex": r"(CI (is )?red|CI fail|Jenkins.*FAIL|OURS.*failure)",
"signal": "needs_input", "priority": "high"},
{"id": "emit-card", "regex": r"(needs_user|emit.card|awaiting your)",
"signal": "needs_input", "priority": "high"},
"signal": "needs_input", "priority": "high", "suppress": True},
{"id": "stranded", "regex": r"(stranded|stuck|blocked|timed out|API error)",
"signal": "needs_input", "priority": "medium"},
"signal": "needs_input", "priority": "medium", "suppress": True},
{"id": "tests-pass", "regex": r"\d+ (tests?|specs?) (passing|passed|green)",
"signal": "work_complete", "priority": "low"},
{"id": "committed", "regex": r"\[main [a-f0-9]{7}\]",
Expand Down Expand Up @@ -452,9 +457,30 @@ def pending_question(records: list[dict]) -> str | None:
return None


def _is_prompt(rec: dict) -> bool:
"""A user turn someone typed, as opposed to a tool result or hook output."""
if agent_session.record_role(rec) != "user" or rec.get("isMeta"):
return False
msg = rec.get("message")
content = msg.get("content") if isinstance(msg, dict) else None
if isinstance(content, str):
return not content.startswith("<")
blocks = _content_blocks(rec)
return bool(blocks) and not any(b.get("type") == "tool_result" for b in blocks)


def current_turn(records: list[dict]) -> list[dict]:
"""The records after the newest typed prompt: this turn only, so an earlier
turn's words can't be judged again."""
for i in range(len(records) - 1, -1, -1):
if _is_prompt(records[i]):
return records[i + 1:]
return records


def last_assistant_text(records: list[dict]) -> str | None:
"""The newest non-empty text block the agent wrote."""
for rec in reversed(records):
"""The newest non-empty text block the agent wrote in the current turn."""
for rec in reversed(current_turn(records)):
if agent_session.record_role(rec) != "assistant":
continue
for block in reversed(_content_blocks(rec)):
Expand All @@ -464,6 +490,31 @@ def last_assistant_text(records: list[dict]) -> str | None:
return None


def pending_tool(records: list[dict]) -> str | None:
"""Name of a tool call in the current turn that has no result yet. At an
idle alert every call has its result, so a pending one means Claude is
showing a prompt — a permission request or plan approval."""
answered = {b.get("tool_use_id") for rec in records for b in _content_blocks(rec)
if b.get("type") == "tool_result"}
for rec in reversed(current_turn(records)):
if agent_session.record_role(rec) != "assistant":
continue
for block in _content_blocks(rec):
if block.get("type") == "tool_use" and block.get("id") not in answered:
name = block.get("name")
return name if isinstance(name, str) and name else "a tool"
return None


def ended_on_api_error(records: list[dict]) -> bool:
"""True if the agent's newest message is an API error Claude wrote into the
transcript (an expired login, an overloaded API): the session has stopped."""
for rec in reversed(current_turn(records)):
if agent_session.record_role(rec) == "assistant":
return rec.get("isApiErrorMessage") is True
return False


def trim_words(text: str, limit: int = LAST_MESSAGE_CHARS) -> str:
"""Collapse whitespace and cut to `limit` chars on a word boundary."""
flat = " ".join(text.split())
Expand All @@ -472,21 +523,76 @@ def trim_words(text: str, limit: int = LAST_MESSAGE_CHARS) -> str:
return flat[:limit].rsplit(" ", 1)[0] + "…"


def read_last_message(cwd: str | None, session_id: str | None, *,
question: bool) -> str | None:
"""What the agent last said, for the Slack ping: the pending question when
`question` (an AskUserQuestion event — the hook payload redacts its text),
else the last assistant text. Returns None on any failure so a missing or
malformed transcript never blocks the drop."""
# Words that ask the user for something. Most turns end in a status update
# ("Standing by for the gate agents…") that needs nobody; the ones that need
# the user ask a question or ask them to act ("say the word", "reply
# `cleanup`", "please run /login").
_ASKS_RE = re.compile(
r"\?|\b(let me know|your call|want me to|should i|shall i|say the word|awaiting your"
r"|your (approval|go-ahead|review|decision|sign-off)"
r"|please (run|paste|approve|confirm|reply))\b|\breply `|/login\b", re.I)
# An ask can sit a line or two above a closing note, so the last few lines count.
ASK_TAIL_LINES = 3


# A link's query string (`?config=squirrel`) isn't a question.
_URL_RE = re.compile(r"\b[a-z][a-z0-9+.-]*://\S+", re.I)


def asks_user(text: str | None) -> bool:
"""True if the message's last few non-empty lines ask the user something."""
lines = [ln.strip() for ln in (text or "").splitlines() if ln.strip()]
return any(_ASKS_RE.search(_URL_RE.sub("", ln)) for ln in lines[-ASK_TAIL_LINES:])


def session_title(records: list[dict]) -> str | None:
"""A name for the session when its workspace title is cmux's default:
Claude's own session title, else the working folder and git branch."""
for rec in reversed(records):
title = rec.get("aiTitle") if rec.get("type") == "ai-title" else None
if isinstance(title, str) and title.strip():
return title.strip()
for rec in reversed(records):
cwd = rec.get("cwd")
if isinstance(cwd, str) and cwd:
branch = rec.get("gitBranch")
name = Path(cwd).name or cwd
return f"{name} ({branch})" if isinstance(branch, str) and branch not in ("", "HEAD") else name
return None


def prompt_message(session: dict) -> str | None:
"""What an idle-alert ping should quote: the open question or prompt when
one is up, else the agent's last words."""
tool = session.get("pending_tool")
last_text = session.get("last_text")
if tool == "AskUserQuestion" and session.get("question"):
return session["question"]
if tool == "ExitPlanMode":
return "Waiting for you to approve its plan."
if tool and not asks_user(last_text):
return f"Waiting for your OK to use {tool}."
return last_text


def read_session(cwd: str | None, session_id: str | None) -> dict | None:
"""What the session's transcript says right now: the pending question (an
AskUserQuestion's text — the hook payload redacts it), the agent's last
message this turn, any tool call still waiting on a prompt, whether the
turn ended on an API error, and a title. None on any failure so a missing
or malformed transcript never blocks the drop."""
try:
path = transcript_path(cwd, session_id)
if path is None:
return None
records = tail_records(path)
text = pending_question(records) if question else last_assistant_text(records)
return {"question": pending_question(records),
"last_text": last_assistant_text(records),
"pending_tool": pending_tool(records),
"api_error": ended_on_api_error(records),
"title": session_title(records)}
except Exception: # noqa: BLE001 — transcripts are external; the ping must still go
return None
return trim_words(text) if text else None


# ─── event classification (pure) ──────────────────────────────────────────────
Expand Down Expand Up @@ -625,10 +731,10 @@ def cooled_down(self, ws_key: str, signal_type: str) -> bool:

def handle_event(evt: dict, bank: PatternBank, state: WatcherState,
resolver: WsRefResolver, *, screen_reader=read_screen,
message_reader=read_last_message) -> dict | None:
session_reader=read_session) -> dict | None:
"""Process one parsed event end-to-end. Returns the dropped item dict (for
tests/logging) or None when nothing was dropped. `screen_reader` and
`message_reader` are injectable so tests don't shell out to cmux or read
`session_reader` are injectable so tests don't shell out to cmux or read
real transcripts."""
cls = classify_event(evt)
if cls is None:
Expand All @@ -638,42 +744,59 @@ def handle_event(evt: dict, bank: PatternBank, state: WatcherState,

workspace_id = cls["workspace_id"]
ws_ref = resolver.resolve(workspace_id)
ws_title = resolver.title(workspace_id)
ws_key = ws_ref or workspace_id or "unknown"
session = session_reader(cls["cwd"], cls["session_id"]) or {}
ws_title = resolver.title(workspace_id) or session.get("title")
last_text = session.get("last_text")

if cls["signal"] == "needs_input":
# Always a signal — the agent is blocked on the user. Cooldown only.
pattern_matched = cls["event_name"].split(".")[-1] # Notification / AskUserQuestion
asking = pattern_matched == "AskUserQuestion"
# cmux sends every Claude alert as this one Notification event, with no
# message: a permission prompt looks the same as the idle "waiting for
# your input" alert that fires about a minute after every turn. Ping
# when a prompt is up (a tool call still waiting), the session died on
# an API error, or the agent's last words ask the user something; an
# unreadable transcript still pings, as before.
needs_user = (not session or session.get("pending_tool")
or session.get("api_error") or asks_user(last_text))
if not asking and not needs_user:
log(f"skip idle needs_input ws={ws_ref or workspace_id}: last message asks nothing")
return None
if not state.cooled_down(ws_key, "needs_input"):
return None
snippet = last_lines(screen_reader(workspace_id or ws_ref or ""))
pattern_matched = cls["event_name"].split(".")[-1] # Notification / AskUserQuestion
last_message = message_reader(cls["cwd"], cls["session_id"],
question=pattern_matched == "AskUserQuestion")
message = session.get("question") if asking else prompt_message(session)
item = drop_inbox_item(ws_ref, "needs_input", pattern_matched, snippet,
ws_title=ws_title, last_message=last_message)
ws_title=ws_title,
last_message=trim_words(message) if message else None)
record_fired(pattern_matched, ws_ref, "needs_input")
log(f"drop needs_input ws={ws_ref or workspace_id} via={pattern_matched} → {item.name}")
return {"path": str(item), "signal_type": "needs_input",
"pattern_matched": pattern_matched, "ws_ref": ws_ref}

# turn_end: read the screen and pattern-match. Drop only on a non-muted hit.
# turn_end: pattern-match what the agent just said, not 50 lines of screen,
# where old scrollback and tool output matched words like "blocked" or
# "API error" that had nothing to do with this turn. The screen is only the
# fallback when the transcript can't be read.
screen = screen_reader(workspace_id or ws_ref or "")
if not screen:
return None # dead/headless workspace or read failure — nothing to judge
hits = bank.match(screen)
text = last_text or screen
if not text:
return None # nothing said this turn and no screen to read — nothing to judge
# A suppressed pattern doesn't count, so it can't hide a lower-priority
# real hit (a message saying both "awaiting review" and "CI is red").
hits = [h for h in bank.match(text) if not h.get("suppress")]
if not hits:
return None # plain turn-end with nothing notable — the noise floor
top = hits[0]
if top.get("suppress"):
return None
pat_signal = top.get("signal")
signal_type = pat_signal if pat_signal in ("needs_input", "work_complete") else "pattern_match"
if not state.cooled_down(ws_key, signal_type):
return None
snippet = last_lines(screen)
last_message = message_reader(cls["cwd"], cls["session_id"], question=False)
snippet = last_lines(screen) or (trim_words(last_text) if last_text else "")
item = drop_inbox_item(ws_ref, signal_type, top.get("id", ""), snippet,
ws_title=ws_title, last_message=last_message)
ws_title=ws_title,
last_message=trim_words(last_text) if last_text else None)
record_fired(top.get("id", ""), ws_ref, signal_type)
log(f"drop {signal_type} ws={ws_ref or workspace_id} pattern={top.get('id')} → {item.name}")
return {"path": str(item), "signal_type": signal_type,
Expand Down
3 changes: 2 additions & 1 deletion docs/assistant-comms-onboarding.md
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,8 @@ Replies come from a **warm cmux Claude session** (Sonnet, scoped `--add-dir`) th
|---|---|
| Assistant appends a verified action to its ledger | `I asked a workspace to merge its PR.` with the refs in a trailing italic line |
| Same, but evidence is `screen_read` (Assistant rejects this as weak) | `Heads up: I only confirmed this by reading the screen, which isn't reliable proof.` |
| A workspace asks you a question / needs input / finishes work | `*<workspace title>* is asking you: <question>` or `*<workspace title>* needs your input.` with the agent's last message quoted |
| A workspace asks you a question, shows a permission or plan prompt, stops on an API error, or goes idle on a message that asks you something | `*<workspace title>* is asking you: <question>` or `*<workspace title>* needs your input.` with the agent's last message quoted. Idle workspaces whose last message is a status update don't ping. |
| The agent's own message says CI is red, a PR opened, or tests passed | `*<workspace title>* needs your input.` or `*<workspace title>* looks done.` with the message quoted |
| Assistant heartbeat stale (>20 min) or status flips to `frozen`/`stale_world`/`respawn-requested` | `*Assistant's main loop has stopped* — no run for 25m …`, then `*Assistant's main loop is running again* after 3h.` |

Messages are Slack `mrkdwn`.
Expand Down
Loading
Loading