chahinebrini 6fa94116a4 fix(hub): offene Asks sichtbar machen + Erinnerungsflut stoppen (TSK-0406)
Am 09.08. stellte sich heraus, dass fuenf Asks auf pending standen, zwei
seit einer Woche. Vier Agenten waren blockiert — darunter einer, der einen
fertigen Windows-Build fuer einen externen Tester nicht ausliefern durfte.
Der Architekt hat die ganze Zeit das Board gepollt und nichts davon gesehen.
Aufgefallen ist es nur, weil ein Agent den Stau selbst benannt hat.

TEIL 1 — SICHTBARKEIT (Umsetzung: codex)
Asks sind ein eigener, BLOCKIERENDER Kanal, den keine der Oberflaechen
erwaehnte, die ein Architekt regelmaessig ansieht:
- core/status.ts nennt jetzt Anzahl, IDs und Alter des aeltesten Asks.
  Das Alter ist der eigentliche Alarm: "einer offen" ist normal,
  "seit sechs Tagen offen" ist ein Notfall.
- /health liefert pendingAsks + oldestAskAgeSec — der billigste Statuscheck
  des Architekten; was dort nicht steht, existiert fuer ihn nicht.
- Der Watchdog kannte Asks gar nicht. Er meldete brav "silent for 16m",
  also das Symptom, waehrend die Ursache unsichtbar blieb. Jetzt alarmiert
  er bei ueberfaelligen Asks mit steigender Dringlichkeit, und die
  irrefuehrende "silent"-Meldung wird bei einem offenen Ask ersetzt durch
  "wartet seit Xh auf ASK-NNNN — bitte antworten, nicht neu starten".
  Genau diese Verwechslung fuehrte dazu, dass wartende Agenten fuer tot
  gehalten und ihre Sessions neu gestartet wurden.
- watch --await-ask weckt den Architekten bei neuen Asks (mit --new-only).
- Das Board zeigt offene Asks mit Alter und wartendem Agenten.

TEIL 2 — DIE FLUT (Umsetzung: claude)
Der Watchdog erinnerte einen nicht beanspruchten Task unbegrenzt im
Basisintervall weiter — auch an Agenten, von denen er WUSSTE, dass sie
nicht im work-Loop sind. Ergebnis: ueber 1500 ungelesene Nachrichten fuer
einen einzigen Agenten, dessen Loop daran nicht mehr anlief. Ein Alarm,
der den Empfaenger handlungsunfaehig macht, ist schlimmer als kein Alarm.
- Ist der Assignee nachweislich aus dem Loop ausgestiegen (nicht: nie
  gesehen), wird KEINE Erinnerung mehr in sein Postfach geschrieben.
  Eine Nachricht ist ein dauerhaftes Artefakt; in ein Postfach zu
  schreiben, das niemand liest, baut nur den Rueckstau auf, der spaeter
  seinen eigenen Loop blockiert. Stattdessen genau EIN Hinweis an den
  Architekten — ein Session-Neustart ist eine menschliche Entscheidung.
- Das ephemere SSE-Reemit bleibt in beiden Zweigen: es kostet nichts und
  erreicht womoeglich einen Agenten, der gerade neu verbindet.
- Fuer erreichbare Agenten waechst der Abstand exponentiell (Basis x 2^n,
  gedeckelt bei 8x) statt konstant zu bleiben.
- Der Zaehler wird zurueckgesetzt, sobald der Task beansprucht wird, damit
  ein spaeteres Reopen nicht in einem halb stummgeschalteten Zustand
  startet.

Der bestehende Test "does not spam" pruefte, dass nach erneutem Ablauf der
BASIS-Schwelle eine zweite Erinnerung kommt — also genau das Verhalten,
das die Flut erzeugt hat. Er ist auf den Backoff angepasst; seine Absicht
(Anti-Spam) gilt jetzt ueber die Lebensdauer eines Tasks statt nur ueber
ein Cooldown-Fenster.

316 Tests gruen, tsc --noEmit sauber.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-09 12:29:09 +02:00

336 lines
13 KiB
TypeScript

/**
* `agenthub watch` — subscribe to live events from a running AgentHub server.
*
* Connects to GET /events (Server-Sent Events) and prints each change
* compactly to stdout. No new dependencies: uses Node's built-in fetch +
* ReadableStream reader.
*
* Flags:
* --once Exit 0 after the first event is received. Lets a
* turn-based agent use this as a blocking wait: run in
* background, get woken up when something changes.
* --role <role> Client-side filter: suppress events whose `role` field
* doesn't match. (The server also accepts ?role= for a
* server-side filter, reducing traffic; both can be used.)
*/
import type { AgentHubEvent } from '../../server/events.js';
import { messageRecipientAliases } from '../../core/services/messageService.js';
// Re-export so tests can import type + helpers from one place.
export type { AgentHubEvent } from '../../server/events.js';
/**
* Parse all complete SSE events from a text buffer.
*
* SSE wire format: "data: <json>\n\n" per event, ": \n\n" for keepalives.
* Frame separators are CRLF-tolerant: both "\n\n" and "\r\n\r\n" (and mixed
* line endings inside a frame) are accepted. An event with multiple `data:`
* lines is joined with "\n" before parsing, per the SSE spec.
* Anything that didn't end with a blank line is returned as `remaining`
* (to be prepended to the next chunk).
*/
export function parseSSEBuffer(buffer: string): { events: AgentHubEvent[]; remaining: string } {
const parts = buffer.split(/\r?\n\r?\n/);
const remaining = parts.pop() ?? ''; // last segment may be incomplete
const events: AgentHubEvent[] = [];
for (const part of parts) {
// A keepalive block looks like ":" — no data line.
const lines = part.split(/\r?\n/);
const eventName = lines.find((l) => l.startsWith('event: '))?.slice(7);
if (eventName && eventName !== 'message') continue;
const idLine = lines.find((l) => l.startsWith('id: '));
const seq = idLine ? Number.parseInt(idLine.slice(4), 10) : undefined;
const dataLines = lines.filter((l) => l.startsWith('data: ')).map((l) => l.slice(6));
if (dataLines.length === 0) continue;
try {
const event = JSON.parse(dataLines.join('\n')) as AgentHubEvent;
if (seq !== undefined && Number.isFinite(seq)) event.seq = seq;
events.push(event);
} catch {
// Ignore malformed JSON — should never happen in practice.
}
}
return { events, remaining };
}
/**
* Map an event to a human label + trailing detail. Keeps the wire-level
* type/action out of the user-facing line in favor of a verb the user reads
* at a glance ("Task received", "Task done", "Handoff", ...).
*/
function describeEvent(event: AgentHubEvent): { label: string; detail: string } {
switch (event.type) {
case 'task': {
if (event.action === 'created') {
const status = event.status ? `[${event.status}]` : '';
return { label: 'Task received', detail: [event.title, status].filter(Boolean).join(' ') };
}
switch (event.status) {
case 'done':
return { label: 'Task done', detail: event.assignedTo ? `by ${event.assignedTo}` : '' };
case 'in_progress':
return { label: 'Task claimed', detail: event.assignedTo ? `by ${event.assignedTo}` : '' };
case 'review':
return { label: 'Task review', detail: event.assignedTo ? `by ${event.assignedTo}` : '' };
case 'cancelled':
return { label: 'Task cancelled', detail: '' };
case 'open':
return { label: 'Task reopened', detail: '' };
default:
return { label: 'Task updated', detail: event.status ? `status=${event.status}` : '' };
}
}
case 'handoff':
return { label: 'Handoff', detail: [event.title, event.role ? `${event.role}` : ''].filter(Boolean).join(' ') };
case 'decision':
return { label: 'Decision', detail: event.title ?? '' };
case 'memory':
return { label: 'Memory', detail: event.title ?? '' };
case 'message':
// A read message surfaces to the SENDER as a delivery/read receipt.
if (event.action === 'updated' && event.status === 'read') {
return { label: '✓ Read', detail: event.assignedTo ? `read by ${event.assignedTo}` : 'read' };
}
return { label: 'Message', detail: event.title ?? '' };
case 'ask':
return { label: 'Ask', detail: event.title ?? '' };
default:
// 'agent' presence events are formatted in formatEvent() before reaching
// here; this fallback only satisfies exhaustiveness.
return { label: 'Event', detail: event.title ?? '' };
}
}
/**
* Format an event as a single branded log line, e.g.:
* AgentHub: Task received TSK-0019 codex: Magic-App Audit [open]
* AgentHub: Task done TSK-0014 by windows-claude
* AgentHub: Handoff HOF-0012 → implementer
*
* Every line carries the `AgentHub:` prefix so it's recognizable in any
* agent's console (Claude / Codex / Kimi) regardless of surrounding output.
*/
export function formatEvent(event: AgentHubEvent): string {
// Presence events read more naturally as "<agent> joined (<role>)" than the
// padded "<Label> <id>" form used for entities.
if (event.type === 'agent') {
const verb = event.action === 'left' ? 'left' : 'joined';
const role = event.role ? ` (${event.role})` : '';
return `AgentHub: ${event.id} ${verb}${role}`;
}
const { label, detail } = describeEvent(event);
const head = `AgentHub: ${label.padEnd(14)} ${event.id}`;
return detail ? `${head} ${detail}` : head;
}
/** Fetch any tasks currently in `review` (an implementer is awaiting a verdict). */
async function fetchReviewTasks(serverUrl: string): Promise<AgentHubEvent[]> {
try {
const res = await fetch(new URL('/tasks', serverUrl).toString());
if (!res.ok) return [];
const tasks = (await res.json()) as Array<{
id: string;
title?: string;
status?: string;
role?: string;
assignedTo?: string;
}>;
return tasks
.filter((t) => t.status === 'review')
.map((t) => ({ type: 'task', action: 'updated', id: t.id, title: t.title, status: 'review', role: t.role, assignedTo: t.assignedTo }));
} catch {
return [];
}
}
/** Fetch unread messages currently addressed to an agent or its role alias. */
async function fetchUnreadMessages(serverUrl: string, agent: string): Promise<AgentHubEvent[]> {
try {
const res = await fetch(new URL(`/messages?agent=${encodeURIComponent(agent)}&unread=1`, serverUrl).toString());
if (!res.ok) return [];
const messages = (await res.json()) as Array<{
id: string;
from?: string;
to?: string;
status?: string;
}>;
return messages.map((m) => ({
type: 'message',
action: 'created',
id: m.id,
title: `${m.from ?? ''}${m.to ?? ''}`,
status: m.status ?? 'unread',
assignedTo: m.to,
}));
} catch {
return [];
}
}
/** Fetch pending Asks for the architect/legacy claude recipient aliases. */
async function fetchPendingArchitectAsks(serverUrl: string): Promise<AgentHubEvent[]> {
try {
const res = await fetch(new URL('/asks?status=pending', serverUrl).toString());
if (!res.ok) return [];
const asks = (await res.json()) as Array<{ id: string; from?: string; to?: string; status?: string }>;
const recipients = messageRecipientAliases('architect');
return asks
.filter((a) => a.to && recipients.has(a.to.toLowerCase()))
.map((a) => ({
type: 'ask', action: 'created', id: a.id,
title: `${a.from ?? ''}${a.to ?? ''}`,
status: a.status ?? 'pending', assignedTo: a.to,
}));
} catch {
return [];
}
}
function isAskForArchitect(event: AgentHubEvent): boolean {
return event.type === 'ask'
&& event.action === 'created'
&& !!event.assignedTo
&& messageRecipientAliases('architect').has(event.assignedTo.toLowerCase());
}
function isMessageFor(event: AgentHubEvent, agent: string, ignoreFrom: string[] = []): boolean {
if (event.type !== 'message' || event.action !== 'created') return false;
const to = event.assignedTo;
if (!to || !messageRecipientAliases(agent).has(String(to).toLowerCase())) return false;
// Absender ausblenden, die den Notifier nur zumüllen. Konkreter Fall: der
// Watchdog schreibt dem Architekten alle paar Minuten Erinnerungen — ohne
// Filter weckt der Notifier ihn im Takt dieser Erinnerungen, obwohl nichts
// Neues passiert ist, und das Signal entwertet sich selbst.
const sender = String(event.from ?? (event.title ?? '').split('→')[0] ?? '').trim().toLowerCase();
return !ignoreFrom.some((ignored) => ignored.trim().toLowerCase() === sender);
}
/**
* Connect to the AgentHub server's SSE endpoint and stream events to stdout.
*
* Exits the process when:
* - `--once` is set and the first event arrives (exit 0).
* - `--await-review` is set and a task enters `review` (an implementer
* submitted) — including any task already in review on connect. Lets the
* architect run it in the background and be notified the moment a
* submission needs a verdict.
* - The server closes the stream (normal exit).
* - A connection error occurs (exit 1).
*/
export async function watchEvents(
serverUrl: string,
options: { once?: boolean; role?: string; awaitReview?: boolean; awaitMessage?: string; awaitAsk?: boolean; newOnly?: boolean; ignoreFrom?: string[] } = {},
): Promise<void> {
const url = new URL('/events', serverUrl);
// Pass role to the server for an additional server-side filter (saves
// bandwidth on high-volume setups, optional).
if (options.role) url.searchParams.set('role', options.role);
let response: Response;
let lastEventId: number | undefined;
try {
response = await fetch(url.toString(), {
headers: { Accept: 'text/event-stream' },
});
} catch (err) {
const msg = err instanceof Error ? err.message : String(err);
console.error(`Cannot reach AgentHub server at ${serverUrl}: ${msg}`);
process.exit(1);
return; // unreachable; satisfies TypeScript
}
if (!response.ok || !response.body) {
console.error(`AgentHub server returned ${response.status} for /events`);
process.exit(1);
return;
}
console.log(options.role ? `AgentHub: connected (${options.role})` : 'AgentHub: connected');
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = '';
// --await-review: surface a submission that's ALREADY pending on connect,
// so the architect isn't blind to reviews submitted before this watcher.
// --new-only suppresses this backlog check so the notifier can be re-armed
// without instantly exiting on a still-non-empty review queue (no spin) —
// it then fires only on the NEXT genuine task→review transition.
if (options.awaitReview && !options.newOnly) {
const pending = await fetchReviewTasks(serverUrl);
if (pending.length > 0) {
for (const ev of pending) console.log(formatEvent(ev));
await reader.cancel();
return;
}
}
if (options.awaitMessage && !options.newOnly) {
const pending = await fetchUnreadMessages(serverUrl, options.awaitMessage);
if (pending.length > 0) {
for (const ev of pending) console.log(formatEvent(ev));
await reader.cancel();
return;
}
}
if (options.awaitAsk && !options.newOnly) {
const pending = await fetchPendingArchitectAsks(serverUrl);
if (pending.length > 0) {
for (const ev of pending) console.log(formatEvent(ev));
await reader.cancel();
return;
}
}
while (true) {
let done: boolean;
let value: Uint8Array | undefined;
try {
({ done, value } = await reader.read());
} catch {
// Connection closed (e.g. AbortController or server shutdown).
break;
}
if (done) break;
if (value) {
buffer += decoder.decode(value, { stream: true });
}
const { events, remaining } = parseSSEBuffer(buffer);
buffer = remaining;
for (const event of events) {
if (event.seq !== undefined) lastEventId = event.seq;
// Client-side role filter: skip tasks that don't match the requested
// role. Non-task events (handoffs, decisions, memory) always print.
if (options.role && event.type === 'task' && event.role !== undefined && event.role !== options.role) {
continue;
}
console.log(formatEvent(event));
if (options.once) {
await reader.cancel();
return;
}
// --await-review: exit the moment an implementer submits (task → review),
// so the architect (running this in the background) gets notified.
if (options.awaitReview && event.type === 'task' && event.status === 'review') {
await reader.cancel();
return;
}
if (options.awaitMessage && isMessageFor(event, options.awaitMessage, options.ignoreFrom)) {
await reader.cancel();
return;
}
if (options.awaitAsk && isAskForArchitect(event)) {
await reader.cancel();
return;
}
}
}
}