惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

MyScale Blog
MyScale Blog
人人都是产品经理
人人都是产品经理
云风的 BLOG
云风的 BLOG
小众软件
小众软件
F
Fortinet All Blogs
爱范儿
爱范儿
WordPress大学
WordPress大学
N
Netflix TechBlog - Medium
Recent Announcements
Recent Announcements
Google DeepMind News
Google DeepMind News
C
Check Point Blog
博客园 - 聂微东
D
Docker
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
aimingoo的专栏
aimingoo的专栏
Vercel News
Vercel News
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
A
About on SuperTechFans
博客园 - 【当耐特】
Microsoft Azure Blog
Microsoft Azure Blog
B
Blog
宝玉的分享
宝玉的分享
Jina AI
Jina AI
H
Hackread – Cybersecurity News, Data Breaches, AI and More

Recent Commits to openclaw:main

test: merge chat side-result checks · openclaw/openclaw@ddd2c2a test: merge cron history checks · openclaw/openclaw@f7eb746 test: merge responsive navigation shell checks · openclaw/openclaw@c2e4b47 docs(changelog): add codex oauth fixes · openclaw/openclaw@628e6cd test: merge navigation routing cases · openclaw/openclaw@5d8cecb Tests: mock channel registry bundled fallback · openclaw/openclaw@2b08233 Secrets: avoid broad web search discovery for single plugin config · openclaw/openclaw@a464f59 test: merge config view browser checks · openclaw/openclaw@20cf511 fix(status): align oauth health with runtime · openclaw/openclaw@eed7116 feat: add macOS screen snapshots for monitor preview (#67954) thanks … · openclaw/openclaw@f377db1 fix: report shared auth scopes in hello-ok (#67810) thanks @BunsDev · openclaw/openclaw@0b6c39b Auto-reply: avoid eager bundled route fallback · openclaw/openclaw@3ea1bf4 Tests: narrow session binding contract setup · openclaw/openclaw@54e4e16 fix(macOS): enable undo/redo in webchat composer text input (#34962) · openclaw/openclaw@00951dc Tests: speed up channel setup promotion · openclaw/openclaw@82b529a Docs: refresh agent instructions · openclaw/openclaw@5775fe2 fix(auth): serialize OAuth refresh across agents to fix #26322 (#67876) · openclaw/openclaw@8e79080 test: allow ollama public surface boundary test · openclaw/openclaw@7d4f1a6 Docs: add test performance guardrails · openclaw/openclaw@89706d3 Tests: restore context-engine usage proof · openclaw/openclaw@e4c4f95 Tests: slim context engine runtime coverage · openclaw/openclaw@74c198f ci: retry failed custom checkouts · openclaw/openclaw@0ee5baf test: trim duplicate provider auth onboarding cases · openclaw/openclaw@1ffc02e matrix: fix sessions_spawn --thread subagent session spawning (#67643) · openclaw/openclaw@1ce2596 test: reduce auth choice fixture churn · openclaw/openclaw@857b9cd test: mock health status config boundaries · openclaw/openclaw@9d5ab4a test: mock onboard config io boundary · openclaw/openclaw@299694d test: mock legacy state plugin boundaries · openclaw/openclaw@2713089 test: mock channel install boundaries · openclaw/openclaw@b945248 test: mock doctor preview channel boundaries · openclaw/openclaw@b1a3ad4
refactor(channels): store inbound queues in SQLite · open...
steipete · 2026-06-01 · via Recent Commits to openclaw:main

@@ -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";

817918

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">;

99109100110

const POLLING_TEST_WATCHDOG_INTERVAL_MS = 30_000;

101111112+

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

}

413434414-

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

}

426489427490

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

});

530602531603

it("uses backoff helpers for recoverable polling retries", async () => {

@@ -667,7 +739,14 @@ describe("TelegramPollingSession", () => {

667739668740

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", () => {

735884736885

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+

});

1207135412081355

const recovered = await recoverStaleTelegramSpooledUpdateClaims({

12091356

spoolDir: tempDir,

@@ -1213,10 +1360,11 @@ describe("TelegramPollingSession", () => {

1213136012141361

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

});

236225102363-

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?.();