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
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion docs/guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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;
Expand Down
2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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"
},
Expand Down
71 changes: 70 additions & 1 deletion src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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())) {
Expand Down Expand Up @@ -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<void> {
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);
Expand All @@ -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);
Expand Down
216 changes: 216 additions & 0 deletions src/index.wiring.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, unknown> }[]; restore: () => void } {
const calls: { method: string; body: Record<string, unknown> }[] = [];
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<void> {
const prompt = '<telegram-message from_id="42" chat_id="42" chat_type="private" thread_id="9">hi</telegram-message>';
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"));
}
});
});
Loading
Loading