Skip to content

Commit b0679d1

Browse files
committed
refactor(channels): store inbound queues in SQLite
1 parent 80b7f56 commit b0679d1

24 files changed

Lines changed: 2046 additions & 444 deletions

AGENTS.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,7 @@ Skills own workflows; root owns hard policy and routing.
7474
- Core runtime consumes only current canonical shapes/config/data. Legacy or retired shapes normalize only in doctor/migration code before runtime; no runtime shims, aliases, or fallback readers.
7575
- State/storage migrations are database-first. Runtime reads/writes the canonical store only. Old file stores, sidecars, aliases, and fallback readers belong in `openclaw doctor --fix` migration code only, never steady-state runtime.
7676
- Storage default: SQLite only. Do not add JSON/JSONL/TXT/sidecar files for OpenClaw-owned runtime state, caches, queues, registries, indexes, cursors, checkpoints, or plugin scratch data.
77+
- SQLite runtime access uses Kysely helpers, not raw SQL statement strings, except schema DDL, migrations, low-level DB bootstrap, or narrowly justified SQLite primitives.
7778
- Use the shared state DB (`state/openclaw.sqlite`) for global runtime state and plugin KV data. Use the per-agent DB (`agents/<agentId>/agent/openclaw-agent.sqlite`) for agent-scoped state/cache. Use a dedicated SQLite DB only when schema, volume, or lifecycle clearly does not fit those stores.
7879
- Legacy state/cache files are migration debt. When touching code that reads/writes them, prefer moving the data into SQLite or calling out the refactor follow-up; do not add parallel file paths.
7980
- File storage must be a named product artifact: import/export, user attachment, log, backup, or external tool contract. If it is app state or cache, it belongs in SQLite.
Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,2 +1,2 @@
1-
eadfc9b897a05664735f8e2abcb70cb3f33c19427c20802bf8b035520b7a2ea1 plugin-sdk-api-baseline.json
2-
8e10e093068d73b9ac50d3f265bf7d892652b0392c677be4e332248499cf7ed0 plugin-sdk-api-baseline.jsonl
1+
19bdf1196ec771a00777a16fd1e9c3662b8fd788a81034e705c41a74ee79c7ec plugin-sdk-api-baseline.json
2+
43feff80c90adad0f821d1f1e184a9bff1e93d81e6d53a26a26fd9e2972be759 plugin-sdk-api-baseline.jsonl

extensions/telegram/src/polling-session.test.ts

Lines changed: 187 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,16 @@ import os from "node:os";
33
import path from "node:path";
44
import type { ChannelAccountSnapshot } from "openclaw/plugin-sdk/channel-contract";
55
import { MAX_TIMER_TIMEOUT_MS } from "openclaw/plugin-sdk/number-runtime";
6-
import { beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
6+
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
7+
import { createChannelIngressQueue } from "../../../src/channels/message/ingress-queue.js";
8+
import { executeSqliteQuerySync, getNodeSqliteKysely } from "../../../src/infra/kysely-sync.js";
9+
import type { DB as OpenClawStateKyselyDatabase } from "../../../src/state/openclaw-state-db.generated.js";
10+
import {
11+
closeOpenClawStateDatabaseForTest,
12+
openOpenClawStateDatabase,
13+
} from "../../../src/state/openclaw-state-db.js";
14+
import { clearTelegramRuntime, setTelegramRuntime } from "./runtime.js";
15+
import type { TelegramRuntime } from "./runtime.types.js";
716
import type { TelegramIngressWorkerMessage } from "./telegram-ingress-worker.js";
817

918
const runMock = vi.hoisted(() => vi.fn());
@@ -96,9 +105,21 @@ type WorkerPollErrorListener = (message: {
96105
type WorkerMessageListener = (message: TelegramIngressWorkerMessage) => void;
97106
type AsyncVoidFn = () => Promise<void>;
98107
type MockCallSource = { mock: { calls: Array<Array<unknown>> } };
108+
type TelegramPollingTestDatabase = Pick<OpenClawStateKyselyDatabase, "channel_ingress_events">;
99109

100110
const POLLING_TEST_WATCHDOG_INTERVAL_MS = 30_000;
101111

112+
function installTelegramIngressQueueRuntime(resolveStateDir: () => string): void {
113+
setTelegramRuntime({
114+
state: {
115+
resolveStateDir,
116+
openChannelIngressQueue: (
117+
options?: Omit<Parameters<typeof createChannelIngressQueue>[0], "channelId">,
118+
) => createChannelIngressQueue({ ...(options ?? {}), channelId: "telegram" }),
119+
},
120+
} as TelegramRuntime);
121+
}
122+
102123
function mockObjectArg(
103124
source: MockCallSource,
104125
label: string,
@@ -411,24 +432,67 @@ async function pendingUpdateIds(spoolDir: string, limit: number | "all" = 100):
411432
return (await listTelegramSpooledUpdates({ spoolDir, limit })).map((update) => update.updateId);
412433
}
413434

414-
async function failedUpdateIds(spoolDir: string): Promise<number[]> {
415-
const entries = await fs.readdir(spoolDir).catch((err) => {
416-
if ((err as { code?: string }).code === "ENOENT") {
417-
return [];
418-
}
419-
throw err;
435+
function normalizeTelegramTestAccountId(spoolDir: string): string {
436+
const trimmed = path.basename(spoolDir).trim();
437+
return trimmed ? trimmed.replace(/[^a-z0-9._-]+/gi, "_") : "default";
438+
}
439+
440+
function telegramTestQueueName(spoolDir: string): string {
441+
return JSON.stringify(["telegram", normalizeTelegramTestAccountId(spoolDir)]);
442+
}
443+
444+
function openTelegramSpoolTestKysely(spoolDir: string) {
445+
const database = openOpenClawStateDatabase({
446+
env: { ...process.env, OPENCLAW_STATE_DIR: spoolDir },
420447
});
421-
return entries
422-
.filter((entry) => entry.endsWith(".json.failed"))
423-
.map((entry) => Number(entry.slice(0, 16)))
424-
.toSorted((a, b) => a - b);
448+
return {
449+
database,
450+
kysely: getNodeSqliteKysely<TelegramPollingTestDatabase>(database.db),
451+
};
452+
}
453+
454+
async function failedUpdateIds(spoolDir: string): Promise<number[]> {
455+
const { database, kysely } = openTelegramSpoolTestKysely(spoolDir);
456+
const rows = executeSqliteQuerySync(
457+
database.db,
458+
kysely
459+
.selectFrom("channel_ingress_events")
460+
.select("event_id")
461+
.where("queue_name", "=", telegramTestQueueName(spoolDir))
462+
.where("status", "=", "failed")
463+
.orderBy("event_id", "asc"),
464+
).rows;
465+
return rows.map((row) => Number(row.event_id));
466+
}
467+
468+
async function adoptClaimOwner(params: {
469+
spoolDir: string;
470+
updateId: number;
471+
ownerId: string;
472+
claimedAt: number;
473+
}): Promise<void> {
474+
const { database, kysely } = openTelegramSpoolTestKysely(params.spoolDir);
475+
executeSqliteQuerySync(
476+
database.db,
477+
kysely
478+
.updateTable("channel_ingress_events")
479+
.set({
480+
claim_owner: params.ownerId,
481+
claimed_at: params.claimedAt,
482+
updated_at: params.claimedAt,
483+
})
484+
.where("queue_name", "=", telegramTestQueueName(params.spoolDir))
485+
.where("event_id", "=", String(params.updateId).padStart(16, "0"))
486+
.where("status", "=", "claimed"),
487+
);
425488
}
426489

427490
async function withTempSpool<T>(fn: (spoolDir: string) => Promise<T>): Promise<T> {
428491
const spoolDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
429492
try {
430493
return await fn(spoolDir);
431494
} finally {
495+
closeOpenClawStateDatabaseForTest();
432496
await fs.rm(spoolDir, { recursive: true, force: true });
433497
}
434498
}
@@ -526,6 +590,14 @@ describe("TelegramPollingSession", () => {
526590
sleepWithAbortMock.mockReset().mockResolvedValue(undefined);
527591
drainPendingDeliveriesMock.mockReset().mockResolvedValue(undefined);
528592
resetTelegramReplyFenceForTests();
593+
installTelegramIngressQueueRuntime(() =>
594+
path.join(os.tmpdir(), "openclaw-telegram-test-state"),
595+
);
596+
});
597+
598+
afterEach(() => {
599+
clearTelegramRuntime();
600+
closeOpenClawStateDatabaseForTest();
529601
});
530602

531603
it("uses backoff helpers for recoverable polling retries", async () => {
@@ -667,7 +739,14 @@ describe("TelegramPollingSession", () => {
667739

668740
const runPromise = session.runUntilAbort();
669741
await vi.waitFor(() => expect(handleUpdate).toHaveBeenCalledTimes(1));
670-
await vi.waitFor(async () => expect(await fs.readdir(tempDir)).toEqual([]));
742+
await vi.waitFor(async () => expect(await pendingUpdateIds(tempDir, "all")).toEqual([]));
743+
await vi.waitFor(async () =>
744+
expect(
745+
await listTelegramSpooledUpdateClaims({
746+
spoolDir: tempDir,
747+
}),
748+
).toEqual([]),
749+
);
671750
abort.abort();
672751
await runPromise;
673752

@@ -686,6 +765,76 @@ describe("TelegramPollingSession", () => {
686765
expect(init).toHaveBeenCalledBefore(handleUpdate);
687766
expect(handleUpdate).toHaveBeenCalledWith({ update_id: 42, message: { text: "hello" } });
688767
} finally {
768+
abort.abort();
769+
await fs.rm(tempDir, { recursive: true, force: true });
770+
}
771+
});
772+
773+
it("writes isolated worker updates through the main runtime queue", async () => {
774+
const abort = new AbortController();
775+
const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
776+
const handleUpdate = vi.fn(async () => undefined);
777+
const bot = {
778+
api: {
779+
deleteWebhook: vi.fn(async () => true),
780+
config: { use: vi.fn() },
781+
},
782+
init: vi.fn(async () => undefined),
783+
handleUpdate,
784+
stop: vi.fn(async () => undefined),
785+
};
786+
createTelegramBotMock.mockReturnValueOnce(bot);
787+
let onMessage: WorkerMessageListener | undefined;
788+
let stopWorker: (() => void) | undefined;
789+
const workerDone = new Promise<void>((resolve) => {
790+
stopWorker = resolve;
791+
});
792+
const ackSpooledUpdate = vi.fn();
793+
const createWorker = vi.fn(() => ({
794+
onMessage: vi.fn((listener: WorkerMessageListener) => {
795+
onMessage = listener;
796+
return () => undefined;
797+
}),
798+
ackSpooledUpdate,
799+
stop: vi.fn(async () => {
800+
stopWorker?.();
801+
}),
802+
task: vi.fn(async () => {
803+
await workerDone;
804+
}),
805+
}));
806+
807+
try {
808+
const session = createPollingSession({
809+
abortSignal: abort.signal,
810+
isolatedIngress: {
811+
enabled: true,
812+
spoolDir: tempDir,
813+
createWorker,
814+
drainIntervalMs: 10,
815+
},
816+
});
817+
818+
const runPromise = session.runUntilAbort();
819+
await vi.waitFor(() => expect(onMessage).toBeDefined());
820+
onMessage?.({
821+
type: "update",
822+
requestId: "write-1",
823+
update: { update_id: 42, message: { text: "hello" } },
824+
queued: 1,
825+
});
826+
827+
await vi.waitFor(() =>
828+
expect(ackSpooledUpdate).toHaveBeenCalledWith("write-1", { ok: true, updateId: 42 }),
829+
);
830+
await vi.waitFor(() =>
831+
expect(handleUpdate).toHaveBeenCalledWith({ update_id: 42, message: { text: "hello" } }),
832+
);
833+
await vi.waitFor(async () => expect(await pendingUpdateIds(tempDir, "all")).toEqual([]));
834+
abort.abort();
835+
await runPromise;
836+
} finally {
837+
abort.abort();
689838
await fs.rm(tempDir, { recursive: true, force: true });
690839
}
691840
});
@@ -735,7 +884,14 @@ describe("TelegramPollingSession", () => {
735884

736885
const runPromise = session.runUntilAbort();
737886
await vi.waitFor(() => expect(handleUpdate).toHaveBeenCalledTimes(1));
738-
await vi.waitFor(async () => expect(await fs.readdir(tempDir)).toEqual([]));
887+
await vi.waitFor(async () => expect(await pendingUpdateIds(tempDir, "all")).toEqual([]));
888+
await vi.waitFor(async () =>
889+
expect(
890+
await listTelegramSpooledUpdateClaims({
891+
spoolDir: tempDir,
892+
}),
893+
).toEqual([]),
894+
);
739895
abort.abort();
740896
await runPromise;
741897

@@ -1128,7 +1284,7 @@ describe("TelegramPollingSession", () => {
11281284
await runPromise;
11291285
expect(events).toEqual(["handled:42", "handled:44"]);
11301286
expect(await pendingUpdateIds(tempDir)).toEqual([43]);
1131-
expect((await fs.readdir(tempDir)).toSorted()).toEqual(["0000000000000043.json"]);
1287+
expect(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).toEqual([]);
11321288
stopWorker();
11331289
});
11341290
});
@@ -1189,21 +1345,12 @@ describe("TelegramPollingSession", () => {
11891345
if (!claimed) {
11901346
throw new Error("Expected claimed update");
11911347
}
1192-
await fs.writeFile(
1193-
claimed.path,
1194-
`${JSON.stringify({
1195-
version: 1,
1196-
updateId: 42,
1197-
receivedAt: interrupted.receivedAt,
1198-
update: interruptedUpdate,
1199-
claim: {
1200-
processId: "other-process",
1201-
processPid: process.pid,
1202-
claimedAt: Date.now(),
1203-
},
1204-
})}\n`,
1205-
{ mode: 0o600 },
1206-
);
1348+
await adoptClaimOwner({
1349+
spoolDir: tempDir,
1350+
updateId: 42,
1351+
ownerId: `${process.pid}:other-process`,
1352+
claimedAt: Date.now(),
1353+
});
12071354

12081355
const recovered = await recoverStaleTelegramSpooledUpdateClaims({
12091356
spoolDir: tempDir,
@@ -1213,10 +1360,11 @@ describe("TelegramPollingSession", () => {
12131360

12141361
expect(recovered).toBe(0);
12151362
expect(await pendingUpdateIds(tempDir)).toEqual([43]);
1216-
expect((await fs.readdir(tempDir)).toSorted()).toEqual([
1217-
"0000000000000042.json.processing",
1218-
"0000000000000043.json",
1219-
]);
1363+
expect(
1364+
(await listTelegramSpooledUpdateClaims({ spoolDir: tempDir })).map(
1365+
(claim) => claim.updateId,
1366+
),
1367+
).toEqual([42]);
12201368
});
12211369
});
12221370

@@ -2360,7 +2508,7 @@ describe("TelegramPollingSession", () => {
23602508
}
23612509
});
23622510

2363-
it("keeps a timed-out lane guarded when its failed tombstone cannot be written", async () => {
2511+
it("keeps a timed-out lane guarded when its failed state cannot be written", async () => {
23642512
vi.useFakeTimers({ shouldAdvanceTime: true });
23652513
const abort = new AbortController();
23662514
const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));
@@ -2371,15 +2519,10 @@ describe("TelegramPollingSession", () => {
23712519
const regularTurnDone = new Promise<void>((resolve) => {
23722520
releaseRegularTurn = resolve;
23732521
});
2374-
const originalWriteFile = fs.writeFile.bind(fs);
2375-
const writeFileSpy = vi
2376-
.spyOn(fs, "writeFile")
2377-
.mockImplementation(async (...args: Parameters<typeof fs.writeFile>) => {
2378-
if (typeof args[0] === "string" && args[0].includes(".json.failed.")) {
2379-
throw new Error("disk full");
2380-
}
2381-
return await originalWriteFile(...args);
2382-
});
2522+
const spoolModule = await import("./telegram-ingress-spool.js");
2523+
const failSpy = vi
2524+
.spyOn(spoolModule, "failTelegramSpooledUpdateClaim")
2525+
.mockRejectedValueOnce(new Error("disk full"));
23832526
createTelegramBotMock.mockReturnValueOnce({
23842527
api: {
23852528
deleteWebhook: vi.fn(async () => true),
@@ -2460,7 +2603,7 @@ describe("TelegramPollingSession", () => {
24602603
await vi.advanceTimersByTimeAsync(20_000);
24612604
await runPromise;
24622605
} finally {
2463-
writeFileSpy.mockRestore();
2606+
failSpy.mockRestore();
24642607
releaseRegularTurn?.();
24652608
abort.abort();
24662609
stopWorker?.();

0 commit comments

Comments
 (0)