diff --git a/README.md b/README.md index 170ec05..dd7eae6 100644 --- a/README.md +++ b/README.md @@ -8,7 +8,7 @@ commands like `/spawn`, `/sessions`, `/cleanup`, and `/status`. ## What you need -- [omp](https://github.com/can1357/oh-my-pi) 17.0.0 or newer +- [omp](https://github.com/can1357/oh-my-pi) 18.1.16 or newer - [Bun](https://bun.sh/) 1.3 or newer - A Telegram bot from [@BotFather](https://t.me/BotFather) - [herdr](https://herdr.dev/) for `/spawn`, `/sessions`, and stale-topic auto-resume diff --git a/docs/guide.md b/docs/guide.md index 4157717..692c28e 100644 --- a/docs/guide.md +++ b/docs/guide.md @@ -32,7 +32,7 @@ same poll lock and provides the identical routing behavior. ## Requirements -- omp ≥ 17.0.0, Bun ≥ 1.3 +- omp ≥ 18.1.16, Bun ≥ 1.3 - A Telegram bot token from [@BotFather](https://t.me/BotFather) - [herdr](https://herdr.dev/) for `/spawn` and `/sessions` (chat bridging works without it) @@ -470,6 +470,7 @@ Set it up: sessions you aren't actively watching. Turn it off explicitly. - **off** — nothing mirrors; asks stay on this terminal. +- **Run failures** — a Telegram-initiated run that fails at the provider (quota exhausted, 429, etc.) no longer dies silently: the first scheduled retry, the model fallback switch, and the terminal failure (with a hint at `/model` and `/sessions`) land in the same topic the reply would have used. A recovered retry reports back too. `explicit` streaming mode opts out, like it does for all automatic egress. - Only **locally-started** runs mirror — Telegram-initiated runs already stream their reply back, so they never double-notify. - With **topics mode** on, a session's prompts/pings land in its own topic; diff --git a/package.json b/package.json index 2352501..cdd3402 100644 --- a/package.json +++ b/package.json @@ -48,7 +48,7 @@ ] }, "devDependencies": { - "@oh-my-pi/pi-coding-agent": "^17.0.0", + "@oh-my-pi/pi-coding-agent": "^18.1.16", "typescript": "^5.7.0", "@types/node": "^22.0.0" }, diff --git a/src/index.ts b/src/index.ts index 42e9cbf..726e189 100644 --- a/src/index.ts +++ b/src/index.ts @@ -45,7 +45,7 @@ import { type BridgeHost, clearOwnerBotCommands, ensureControlTopic as ensureBri import { SpawnController, agentNameForSession, findSessionSpace, listControlSpaces, sendCommandMessage } from "./control"; import { daemonAlive, daemonDisableReason, ensureDaemon, readDaemonState } from "./daemon"; import { INBOX_MAX_FILE_BYTES, pruneInbox, storeInboxFile } from "./inbox"; -import { Outbound, finalAssistantText } from "./outbound"; +import { Outbound, finalAssistantText, lastRunError } from "./outbound"; import { PROMPT_SUPERSEDED, type PromptOption, type PromptQuestion, type PromptTarget, TelegramPromptController, formatPromptResult } from "./prompts"; import { DM_ROUTE_KEY, @@ -604,6 +604,10 @@ export default function telegramExtension(pi: ExtensionAPI): void { let activePromptTarget: PromptTarget | undefined; let savedPromptTools: string[] | undefined; let compacting = false; + /** Terminal run failure already announced to the active chats this run. */ + let runFailAnnounced = false; + /** Retry saga already announced this run (first auto_retry_start). */ + let retryNoticeSent = false; const poller = new Poller(); const outbound = new Outbound(() => access, log); const spawnController = new SpawnController({ @@ -2435,6 +2439,8 @@ export default function telegramExtension(pi: ExtensionAPI): void { ctx.ui.notify("telegram: away off — you're back at the terminal, runs stay on-screen (run /away to re-arm)", "info"); }); pi.on("before_agent_start", async (event, ctx) => { + runFailAnnounced = false; + retryNoticeSent = false; const telegramTarget = parseTelegramPromptTarget(event.prompt); const a = loadAccess(warn); if (isTaskSubagent(ctx.hasUI, pi.getActiveTools())) { @@ -2618,8 +2624,62 @@ export default function telegramExtension(pi: ExtensionAPI): void { lastCtx = ctx; await outbound.onTurnEnd(e.message); }); + /** First line of a provider error, single-spaced, capped for chat. */ + function shortError(message: string): string { + return message.replace(/\s+/g, " ").trim().slice(0, 300); + } + + function fmtDelay(ms: number): string { + if (!Number.isFinite(ms) || ms < 0) return "a bit"; + if (ms < 1000) return Math.round(ms) + "ms"; + const s = Math.round(ms / 1000); + if (s < 60) return s + "s"; + const m = Math.floor(s / 60); + if (m < 60) return m + "m"; + return Math.floor(m / 60) + "h" + (m % 60 === 0 ? "" : " " + (m % 60) + "m"); + } + + /** + * Run-status notice to the chats with a live Telegram turn. Skipped for task + * subagents and explicit mode — which opts out of all automatic egress, like + * the idle notify below. announce() self-gates on Telegram turns. + */ + async function announceRunNotice(ctx: ExtensionContext, text: string): Promise { + if (isTaskSubagent(ctx.hasUI, pi.getActiveTools())) return; + lastCtx = ctx; + if (effectiveStreaming(loadAccess(warn)) === "explicit") return; + await outbound.announce(text); + } + + pi.on("auto_retry_start", async (e, ctx) => { + if (retryNoticeSent) return; + retryNoticeSent = true; + await announceRunNotice( + ctx, + "⚠️ request failed: " + shortError(e.errorMessage) + " — retrying (" + e.attempt + "/" + e.maxAttempts + ", in ~" + fmtDelay(e.delayMs) + ")…", + ); + }); + pi.on("auto_retry_end", async (e, ctx) => { + if (e.success) { + if (!retryNoticeSent) return; + retryNoticeSent = false; + await announceRunNotice(ctx, "✅ request succeeded after " + e.attempt + (e.attempt === 1 ? " retry." : " retries.")); + return; + } + runFailAnnounced = true; + await announceRunNotice( + ctx, + "⚠️ run failed after " + e.attempt + (e.attempt === 1 ? " attempt" : " attempts") + ": " + shortError(e.finalError ?? "unknown error") + ". Use /model to switch, /sessions for state.", + ); + }); + pi.on("retry_fallback_applied", async (e, ctx) => { + await announceRunNotice(ctx, "🔀 model fallback: " + e.from + " → " + e.to + " (" + e.role + ")."); + }); pi.on("agent_end", async (e, ctx) => { if (isTaskSubagent(ctx.hasUI, pi.getActiveTools())) return; + // Retry/compaction continuations still own their chats and prompt tools. + // Only terminal settlement may tear those down or send an idle notice. + if (e.willContinue) return; lastCtx = ctx; for (const pending of pendingApprovals.values()) { clearTimeout(pending.timer); @@ -2628,6 +2688,15 @@ export default function telegramExtension(pi: ExtensionAPI): void { blockedPings.clear(); const wasActive = outbound.isActive(); const finalText = finalAssistantText(e.messages); + const terminalError = lastRunError(e.messages); + if (terminalError && !runFailAnnounced) { + runFailAnnounced = true; + const errText = + terminalError.status != null + ? terminalError.status + " " + shortError(terminalError.message) + : shortError(terminalError.message); + await announceRunNotice(ctx, "⚠️ run failed: " + errText + ". Use /model to switch, /sessions for state."); + } await outbound.onAgentEnd(finalText); await restorePromptTools(); await captureOwnSpace(ctx); diff --git a/src/index.wiring.test.ts b/src/index.wiring.test.ts index 61762df..2ff88a3 100644 --- a/src/index.wiring.test.ts +++ b/src/index.wiring.test.ts @@ -1216,3 +1216,219 @@ describe("telegram_ask execute (dual-surface)", () => { expect(res.content[0].text).toContain("no surface available"); }); }); + +describe("run failure notices", () => { + const noticeCtx = (sessionId: string) => ({ + hasUI: true, + ui: { notify() {} }, + isIdle: () => true, + sessionManager: { getSessionId: () => sessionId, getSessionFile: () => `/tmp/${sessionId}.jsonl` }, + }); + + function captureSend(): { calls: { method: string; body: Record }[]; restore: () => void } { + const calls: { method: string; body: Record }[] = []; + const previousFetch = globalThis.fetch; + globalThis.fetch = (async (input, init) => { + const method = String(input).split("/").pop()!; + if (method === "sendMessage") calls.push({ method, body: JSON.parse(String(init?.body ?? "{}")) }); + return new Response(JSON.stringify({ ok: true, result: { message_id: 60 } }), { status: 200 }); + }) as typeof fetch; + return { calls, restore: () => { globalThis.fetch = previousFetch; } }; + } + + async function activateTelegramTurn(h: Harness): Promise { + const prompt = 'hi'; + await h.handlers.get("before_agent_start")?.[0]?.( + { type: "before_agent_start", prompt, systemPrompt: [] }, + noticeCtx("session-1"), + ); + await h.handlers.get("message_start")?.[0]?.( + { type: "message_start", message: { role: "user", content: prompt } }, + noticeCtx("session-1"), + ); + } + + const failedRun = (message: string) => [ + { role: "assistant", content: [], stopReason: "error", errorStatus: 429, errorMessage: message }, + ]; + + test("a retry saga narrates its first retry and the recovery", async () => { + writeAccess({ enabled: true, allowFrom: ["42"] }); + const h = harness(["ask"]); + await startBridge(h); + const { calls, restore } = captureSend(); + try { + await activateTelegramTurn(h); + const retryStart = h.handlers.get("auto_retry_start")?.[0]; + await retryStart?.( + { type: "auto_retry_start", attempt: 1, maxAttempts: 3, delayMs: 62000, errorMessage: "429 boom" }, + noticeCtx("session-1"), + ); + expect(calls.length).toBe(1); + expect(String(calls[0].body.text)).toContain("429 boom"); + await h.handlers.get("agent_end")?.[0]?.( + { type: "agent_end", willContinue: true, messages: failedRun("429 boom") }, + noticeCtx("session-1"), + ); + await retryStart?.( + { type: "auto_retry_start", attempt: 2, maxAttempts: 3, delayMs: 120000, errorMessage: "429 boom" }, + noticeCtx("session-1"), + ); + expect(calls.length).toBe(1); + await h.handlers.get("auto_retry_end")?.[0]?.( + { type: "auto_retry_end", success: true, attempt: 2 }, + noticeCtx("session-1"), + ); + expect(calls.length).toBe(2); + const answer = { role: "assistant", content: [{ type: "text", text: "Recovered answer" }], stopReason: "stop" }; + await h.handlers.get("turn_end")?.[0]?.({ type: "turn_end", message: answer }, noticeCtx("session-1")); + await h.handlers.get("agent_end")?.[0]?.({ type: "agent_end", messages: [...failedRun("429 boom"), answer] }, noticeCtx("session-1")); + expect(calls).toHaveLength(3); + expect(calls[2].body.text).toBe("Recovered answer"); + expect(calls.every((call) => call.body.chat_id === "42" && call.body.message_thread_id === 9)).toBe(true); + } finally { + restore(); + await h.handlers.get("session_shutdown")?.[0]?.({ type: "session_shutdown" }, noticeCtx("session-1")); + } + }); + + test("exhausted retries announce the failure once, agent_end stays silent after", async () => { + writeAccess({ enabled: true, allowFrom: ["42"] }); + const h = harness(["ask"]); + await startBridge(h); + const { calls, restore } = captureSend(); + try { + await activateTelegramTurn(h); + await h.handlers.get("auto_retry_start")?.[0]?.( + { type: "auto_retry_start", attempt: 1, maxAttempts: 3, delayMs: 1000, errorMessage: "429 temporary" }, + noticeCtx("session-1"), + ); + await h.handlers.get("agent_end")?.[0]?.( + { type: "agent_end", willContinue: true, messages: failedRun("429 temporary") }, + noticeCtx("session-1"), + ); + await h.handlers.get("auto_retry_end")?.[0]?.( + { type: "auto_retry_end", success: false, attempt: 3, finalError: "429 no more" }, + noticeCtx("session-1"), + ); + expect(calls).toHaveLength(2); + expect(String(calls[1].body.text)).toContain("429 no more"); + await h.handlers.get("agent_end")?.[0]?.( + { type: "agent_end", messages: failedRun("429 no more") }, + noticeCtx("session-1"), + ); + expect(calls).toHaveLength(2); + expect(calls.every((call) => call.body.chat_id === "42" && call.body.message_thread_id === 9)).toBe(true); + expect(calls.every((call) => call.body.parse_mode === undefined)).toBe(true); + } finally { + restore(); + await h.handlers.get("session_shutdown")?.[0]?.({ type: "session_shutdown" }, noticeCtx("session-1")); + } + }); + + test("agent_end without a retry cycle announces the terminal error with status", async () => { + writeAccess({ enabled: true, allowFrom: ["42"] }); + const h = harness(["ask"]); + await startBridge(h); + const { calls, restore } = captureSend(); + try { + await activateTelegramTurn(h); + await h.handlers.get("agent_end")?.[0]?.( + { type: "agent_end", messages: failedRun("429 boom") }, + noticeCtx("session-1"), + ); + expect(calls.length).toBe(1); + expect(String(calls[0].body.text)).toContain("429 boom"); + } finally { + restore(); + await h.handlers.get("session_shutdown")?.[0]?.({ type: "session_shutdown" }, noticeCtx("session-1")); + } + }); + + test("continuations preserve the Telegram prompt surface until terminal settlement", async () => { + writeAccess({ enabled: true, allowFrom: ["42"] }); + const h = harness(["ask", "read"]); + await startBridge(h); + const { calls, restore } = captureSend(); + try { + await activateTelegramTurn(h); + await h.handlers.get("agent_end")?.[0]?.( + { type: "agent_end", willContinue: true, messages: failedRun("429 boom") }, + noticeCtx("session-1"), + ); + expect(calls).toEqual([]); + expect(h.active()).toContain("telegram_ask"); + expect(h.active()).not.toContain("ask"); + await h.handlers.get("agent_end")?.[0]?.( + { type: "agent_end", messages: failedRun("429 boom") }, + noticeCtx("session-1"), + ); + expect(calls).toHaveLength(1); + expect(h.active()).toContain("ask"); + } finally { + restore(); + await h.handlers.get("session_shutdown")?.[0]?.({ type: "session_shutdown" }, noticeCtx("session-1")); + } + }); + + test("a healthy run does not announce a historical provider failure", async () => { + writeAccess({ enabled: true, allowFrom: ["42"] }); + const h = harness(["ask"]); + await startBridge(h); + const { calls, restore } = captureSend(); + try { + await activateTelegramTurn(h); + const answer = { role: "assistant", content: [{ type: "text", text: "Healthy answer" }], stopReason: "stop" }; + await h.handlers.get("turn_end")?.[0]?.({ type: "turn_end", message: answer }, noticeCtx("session-1")); + await h.handlers.get("agent_end")?.[0]?.( + { type: "agent_end", messages: [...failedRun("old failure"), { role: "user", content: "new request" }, answer] }, + noticeCtx("session-1"), + ); + expect(calls.map((call) => call.body.text)).toEqual(["Healthy answer"]); + } finally { + restore(); + await h.handlers.get("session_shutdown")?.[0]?.({ type: "session_shutdown" }, noticeCtx("session-1")); + } + }); + + test("explicit mode opts out of failure notices", async () => { + writeAccess({ enabled: true, allowFrom: ["42"], streaming: "explicit" }); + const h = harness(["ask"]); + await startBridge(h); + const { calls, restore } = captureSend(); + try { + await activateTelegramTurn(h); + await h.handlers.get("auto_retry_end")?.[0]?.( + { type: "auto_retry_end", success: false, attempt: 1, finalError: "429 boom" }, + noticeCtx("session-1"), + ); + await h.handlers.get("agent_end")?.[0]?.( + { type: "agent_end", messages: failedRun("429 boom") }, + noticeCtx("session-1"), + ); + expect(calls).toEqual([]); + } finally { + restore(); + await h.handlers.get("session_shutdown")?.[0]?.({ type: "session_shutdown" }, noticeCtx("session-1")); + } + }); + + test("a model fallback is announced to the active chat", async () => { + writeAccess({ enabled: true, allowFrom: ["42"] }); + const h = harness(["ask"]); + await startBridge(h); + const { calls, restore } = captureSend(); + try { + await activateTelegramTurn(h); + await h.handlers.get("retry_fallback_applied")?.[0]?.( + { type: "retry_fallback_applied", from: "a/model", to: "b/model", role: "default" }, + noticeCtx("session-1"), + ); + expect(calls.length).toBe(1); + expect(String(calls[0].body.text)).toContain("a/model"); + } finally { + restore(); + await h.handlers.get("session_shutdown")?.[0]?.({ type: "session_shutdown" }, noticeCtx("session-1")); + } + }); +}); diff --git a/src/outbound.test.ts b/src/outbound.test.ts index 5412aec..634ff65 100644 --- a/src/outbound.test.ts +++ b/src/outbound.test.ts @@ -1,6 +1,6 @@ import { afterEach, test, expect, describe, setSystemTime } from "bun:test"; import { type Access, defaultAccess } from "./access"; -import { Outbound, assistantText, finalAssistantText } from "./outbound"; +import { Outbound, assistantText, finalAssistantText, lastRunError } from "./outbound"; import { mdToMarkdownV2 } from "./markdown"; const assistant = (text: string): unknown => ({ role: "assistant", content: [{ type: "text", text }] }); @@ -55,6 +55,57 @@ describe("finalAssistantText", () => { }); }); +describe("lastRunError", () => { + const failed = (message?: unknown): unknown => ({ role: "assistant", content: [], stopReason: "error", errorMessage: message, errorStatus: 429 }); + + test("returns the latest failure with its status", () => { + expect(lastRunError([assistant("ok"), failed("429 boom")])).toEqual({ message: "429 boom", status: 429 }); + }); + + test("stops at the latest assistant result or user boundary", () => { + expect(lastRunError([failed("old"), assistant("healthy")])).toBeUndefined(); + expect(lastRunError([failed("old"), { role: "user", content: "new request" }])).toBeUndefined(); + expect(lastRunError([failed("old"), failed(" ")])).toBeUndefined(); + expect(lastRunError([failed("old"), { role: "assistant", stopReason: "aborted" }])).toBeUndefined(); + expect(lastRunError([failed("old"), { role: "user", content: "new request" }, failed("current")])).toEqual({ message: "current", status: 429 }); + }); + + test("ignores aborts, textless failures, and non-assistant messages", () => { + expect(lastRunError([{ role: "assistant", content: [], stopReason: "aborted", errorMessage: "stop" }])).toBeUndefined(); + expect(lastRunError([failed(" ")])).toBeUndefined(); + expect(lastRunError([failed(undefined)])).toBeUndefined(); + expect(lastRunError([{ role: "user", content: "hi" }])).toBeUndefined(); + expect(lastRunError([])).toBeUndefined(); + }); +}); + +describe("Outbound.announce", () => { + test("notifies every active chat in plain text, silent with no token or no turn", async () => { + const calls: Array<{ method: string; payload: Record }> = []; + globalThis.fetch = (async (url, init) => { + const method = String(url).split("/").pop()!; + const payload = JSON.parse(String(init?.body)) as Record; + calls.push({ method, payload }); + return new Response(JSON.stringify({ ok: true, result: { message_id: 31 } }), { status: 200 }); + }) as typeof fetch; + + const outbound = new Outbound(() => ({ ...defaultAccess(), allowFrom: ["42"], richMessages: "on" })); + await outbound.announce("no token"); + expect(calls).toEqual([]); + outbound.setToken("secret"); + await outbound.announce("no turn"); + expect(calls).toEqual([]); + outbound.markActive("42", 9); + outbound.markActive("43"); + await outbound.announce("run failed: boom_bam"); + expect(calls.filter((c) => c.method === "sendMessage").map((c) => c.payload)).toEqual([ + { chat_id: "42", text: "run failed: boom_bam", message_thread_id: 9 }, + { chat_id: "43", text: "run failed: boom_bam" }, + ]); + outbound.shutdown(); + }); +}); + describe("Outbound Telegram delivery", () => { test("falls back to plain text when Telegram rejects MarkdownV2", async () => { const calls: Array<{ method: string; payload: Record }> = []; diff --git a/src/outbound.ts b/src/outbound.ts index a8e2b93..e0dbaf9 100644 --- a/src/outbound.ts +++ b/src/outbound.ts @@ -83,6 +83,25 @@ export function finalAssistantText(messages: readonly unknown[]): string { return ""; } +export interface RunError { + message: string; + status?: number; +} + +/** Failure of the current terminal assistant result, never an older turn. */ +export function lastRunError(messages: readonly unknown[]): RunError | undefined { + for (let i = messages.length - 1; i >= 0; i--) { + const m = messages[i]; + if (!m || typeof m !== "object") continue; + const r = m as { role?: unknown; stopReason?: unknown; errorMessage?: unknown; errorStatus?: unknown }; + if (r.role === "user") return undefined; + if (r.role !== "assistant") continue; + if (r.stopReason !== "error" || typeof r.errorMessage !== "string" || r.errorMessage.trim().length === 0) return undefined; + return { message: r.errorMessage, status: typeof r.errorStatus === "number" ? r.errorStatus : undefined }; + } + return undefined; +} + export class Outbound { #token = ""; readonly #getAccess: () => Access; @@ -127,6 +146,31 @@ export class Outbound { return this.#lastTarget; } + /** Chats with a live Telegram turn (normal replies route here). */ + activeTargets(): Array<{ chatId: string; threadId?: number }> { + const targets: Array<{ chatId: string; threadId?: number }> = []; + for (const key of this.#active) { + const st = this.#chats.get(key); + if (st) targets.push({ chatId: st.chatId, threadId: st.threadId }); + } + return targets; + } + + /** + * Run-status notice (retry/failure/recovery) to every chat with a live + * Telegram turn. Plain text: provider errors carry MarkdownV2-hostile + * characters. Gated to Telegram turns by construction — #active only fills + * from Telegram inbound. + */ + async announce(text: string): Promise { + if (!this.#token || this.#active.size === 0) return; + for (const t of this.activeTargets()) { + await this.send(t.chatId, text, { threadId: t.threadId, format: "text" }).catch((err) => + this.#log?.warn("[telegram] announce failed " + t.chatId + ": " + String(err)), + ); + } + } + /** Mark a chat (optionally a forum topic) as an active inbound source; starts typing. */ markActive(chatId: string, threadId?: number): void { this.#lastTarget = { chatId, threadId };