diff --git a/packages/pilegram/src/queue.test.ts b/packages/pilegram/src/queue.test.ts new file mode 100644 index 0000000..7354741 --- /dev/null +++ b/packages/pilegram/src/queue.test.ts @@ -0,0 +1,20 @@ +import { expect, test } from "bun:test"; +import { SerialQueue } from "./queue.ts"; + +test("runs tasks in FIFO order after a failed task", async () => { + const queue = new SerialQueue(); + const events: string[] = []; + + const first = queue.enqueue(async () => { + events.push("first"); + throw new Error("expected"); + }); + const second = queue.enqueue(async () => { + events.push("second"); + }); + + await expect(first).rejects.toThrow("expected"); + await second; + await queue.flush(); + expect(events).toEqual(["first", "second"]); +}); diff --git a/packages/pilegram/src/queue.ts b/packages/pilegram/src/queue.ts new file mode 100644 index 0000000..5539fc5 --- /dev/null +++ b/packages/pilegram/src/queue.ts @@ -0,0 +1,23 @@ +/** + * A failure-tolerant FIFO promise queue. + * + * Each submitted task starts after every earlier task has settled. Individual + * task failures are returned to their caller but never poison later work. + */ +export class SerialQueue { + private tail: Promise = Promise.resolve(); + + enqueue(task: () => Promise): Promise { + const run = this.tail.then(task); + this.tail = run.then( + () => undefined, + () => undefined, + ); + return run; + } + + /** Wait for all work submitted before this call. */ + async flush(): Promise { + await this.tail; + } +} diff --git a/packages/pilegram/src/renderer.ts b/packages/pilegram/src/renderer.ts index 36be578..02f5fd8 100644 --- a/packages/pilegram/src/renderer.ts +++ b/packages/pilegram/src/renderer.ts @@ -155,7 +155,7 @@ export class Renderer { } /** Finalize on agent_settled. `finalText` is the authoritative answer. */ - onSettled(finalText: string | undefined) { + onSettled(finalText: string | undefined): Promise { const answer = finalText && finalText.trim() !== "" ? finalText : this.acc; const counts = new Map(this.toolCounts); const elapsedMs = this.turnStartAt ? Date.now() - this.turnStartAt : 0; @@ -175,14 +175,12 @@ export class Renderer { if (this.voiceMode) { // The Session sends this answer as a voice note; don't leave a provisional // text message behind while it does so. - void this.deletePreview(preview); this.log.info("finalize: voice-only (text not persisted)"); - return; + return this.deletePreview(preview); } if (answer.trim() === "") { - void this.deletePreview(preview); this.log.info("finalize: empty (preview deleted)"); - return; + return this.deletePreview(preview); } // Strip bidi-override / zero-width chars so a prompt-injected answer can't @@ -209,7 +207,7 @@ export class Renderer { html: !!extra, preview: !!preview, }); - void this.finalizePreview(preview, chunks, extra); + return this.finalizePreview(preview, chunks, extra); } /** diff --git a/packages/pilegram/src/session.ts b/packages/pilegram/src/session.ts index 1a7e7d5..df2696a 100644 --- a/packages/pilegram/src/session.ts +++ b/packages/pilegram/src/session.ts @@ -23,6 +23,7 @@ import { join } from "node:path"; import type { MessageLog } from "./context.ts"; import { errFields, log as rootLog } from "./log.ts"; import type { ImageContent } from "./media.ts"; +import { SerialQueue } from "./queue.ts"; import { Renderer } from "./renderer.ts"; import type { Route } from "./route.ts"; import { routeKey } from "./route.ts"; @@ -64,6 +65,8 @@ export interface SessionOptions { export class Session { private busy = false; + /** FIFO barrier between completed turns and the next turn's Telegram writes. */ + private readonly finalizations = new SerialQueue(); private voiceMode = false; private spokeThisTurn = false; // set if the agent sent a voice note via tg_send_voice this turn private readonly unsubscribe: () => void; @@ -186,22 +189,33 @@ export class Session { break; case "agent_settled": { const finalText = this.agent.getLastAssistantText(); - this.renderer.onSettled(finalText); + const voiceMode = this.voiceMode; + const spokeThisTurn = this.spokeThisTurn; this.busy = false; - // A voice-only turn's text is spoken, never rendered to Telegram — don't - // record it as the last-rendered answer, or reconcile would suppress the - // legitimate text repost if we crash before the voice note is sent. - this.onFinalized?.(this.voiceMode ? undefined : finalText); - // Voice mode: speak the answer as a voice note — unless the agent already - // sent one itself via tg_send_voice, which would double up. - if ( - this.voiceMode && - this.voice && - !this.spokeThisTurn && - finalText && - finalText.trim() !== "" - ) - void this.speak(finalText); + // Claim final Telegram writes in FIFO order before another turn starts. + // A fast following prompt waits on this queue instead of overtaking this + // turn's preview replacement or voice-note delivery. + void this.finalizations + .enqueue(async () => { + await this.renderer.onSettled(finalText); + // A voice-only turn's text is spoken, never rendered to Telegram — don't + // record it as the last-rendered answer, or reconcile would suppress the + // legitimate text repost if we crash before the voice note is sent. + this.onFinalized?.(voiceMode ? undefined : finalText); + // Voice mode: speak the answer as a voice note — unless the agent already + // sent one itself via tg_send_voice, which would double up. + if ( + voiceMode && + this.voice && + !spokeThisTurn && + finalText && + finalText.trim() !== "" + ) + await this.speak(finalText); + }) + .catch((e) => + this.log.error("turn finalization failed", errFields(e)), + ); break; } default: @@ -220,7 +234,7 @@ export class Session { async handlePrompt( text: string, opts?: { images?: ImageContent[]; messageId?: number; speak?: boolean }, - ) { + ): Promise { const images = opts?.images; if (opts?.messageId !== undefined) { if (this.turn) this.turn.messageId = opts.messageId; // for tg_react @@ -232,6 +246,9 @@ export class Session { await this.agent.steer(text, images); return; } + // `agent_settled` precedes its final preview edit. Drain all finalization + // work before this fresh turn can enqueue a draft or tool output. + await this.finalizations.flush(); this.busy = true; this.voiceMode = opts?.speak ?? false; // reply modality matches the input this.renderer.setVoiceMode(this.voiceMode); @@ -244,6 +261,8 @@ export class Session { this.renderer.onError(e); }) .finally(() => { + // A normal turn becomes idle at agent_settled. For failures that never + // settle, release the session here. this.busy = false; }); } diff --git a/packages/pilegram/src/writer.ts b/packages/pilegram/src/writer.ts index 77b9e8e..5a0f120 100644 --- a/packages/pilegram/src/writer.ts +++ b/packages/pilegram/src/writer.ts @@ -23,6 +23,7 @@ import type { } from "grammy/types"; import type { Route } from "./route.ts"; import { errFields, type Fields, log as rootLog } from "./log.ts"; +import { SerialQueue } from "./queue.ts"; const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); @@ -32,7 +33,7 @@ const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); const MAX_429_RETRIES = 8; export class Writer { - private tail: Promise = Promise.resolve(); + private readonly queue = new SerialQueue(); private readonly log: ReturnType; constructor( @@ -52,13 +53,7 @@ export class Writer { /** Serialize an op onto the route's queue, retrying on 429. */ private enqueue(label: string, op: () => Promise): Promise { - const run = this.tail.then(() => this.execWithRetry(label, op)); - // Keep the chain alive regardless of individual failures. - this.tail = run.then( - () => undefined, - () => undefined, - ); - return run; + return this.queue.enqueue(() => this.execWithRetry(label, op)); } private async execWithRetry( @@ -219,7 +214,7 @@ export class Writer { /** Resolve once everything enqueued so far has been sent (ordering barrier). */ async flush(): Promise { - await this.tail.catch(() => {}); + await this.queue.flush(); } /** Send a chat action ("typing", "record_voice", …). Best-effort, not queued