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
20 changes: 20 additions & 0 deletions packages/pilegram/src/queue.test.ts
Original file line number Diff line number Diff line change
@@ -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"]);
});
23 changes: 23 additions & 0 deletions packages/pilegram/src/queue.ts
Original file line number Diff line number Diff line change
@@ -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<void> = Promise.resolve();

enqueue<T>(task: () => Promise<T>): Promise<T> {
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<void> {
await this.tail;
}
}
10 changes: 4 additions & 6 deletions packages/pilegram/src/renderer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void> {
const answer = finalText && finalText.trim() !== "" ? finalText : this.acc;
const counts = new Map(this.toolCounts);
const elapsedMs = this.turnStartAt ? Date.now() - this.turnStartAt : 0;
Expand All @@ -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
Expand All @@ -209,7 +207,7 @@ export class Renderer {
html: !!extra,
preview: !!preview,
});
void this.finalizePreview(preview, chunks, extra);
return this.finalizePreview(preview, chunks, extra);
}

/**
Expand Down
51 changes: 35 additions & 16 deletions packages/pilegram/src/session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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:
Expand All @@ -220,7 +234,7 @@ export class Session {
async handlePrompt(
text: string,
opts?: { images?: ImageContent[]; messageId?: number; speak?: boolean },
) {
): Promise<void> {
const images = opts?.images;
if (opts?.messageId !== undefined) {
if (this.turn) this.turn.messageId = opts.messageId; // for tg_react
Expand All @@ -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);
Expand All @@ -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;
});
}
Expand Down
13 changes: 4 additions & 9 deletions packages/pilegram/src/writer.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>((r) => setTimeout(r, ms));

Expand All @@ -32,7 +33,7 @@ const sleep = (ms: number) => new Promise<void>((r) => setTimeout(r, ms));
const MAX_429_RETRIES = 8;

export class Writer {
private tail: Promise<unknown> = Promise.resolve();
private readonly queue = new SerialQueue();
private readonly log: ReturnType<typeof rootLog.child>;

constructor(
Expand All @@ -52,13 +53,7 @@ export class Writer {

/** Serialize an op onto the route's queue, retrying on 429. */
private enqueue<T>(label: string, op: () => Promise<T>): Promise<T> {
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<T>(
Expand Down Expand Up @@ -219,7 +214,7 @@ export class Writer {

/** Resolve once everything enqueued so far has been sent (ordering barrier). */
async flush(): Promise<void> {
await this.tail.catch(() => {});
await this.queue.flush();
}

/** Send a chat action ("typing", "record_voice", …). Best-effort, not queued
Expand Down