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

推荐订阅源

宝玉的分享
宝玉的分享
小众软件
小众软件
J
Java Code Geeks
I
InfoQ
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
腾讯CDC
L
LangChain Blog
博客园 - 司徒正美
量子位
Y
Y Combinator Blog
C
Check Point Blog
T
Tailwind CSS Blog
D
DataBreaches.Net
Blog — PlanetScale
Blog — PlanetScale
N
Netflix TechBlog - Medium
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
F
Fortinet All Blogs
云风的 BLOG
云风的 BLOG
A
About on SuperTechFans
B
Blog RSS Feed
酷 壳 – CoolShell
酷 壳 – CoolShell
大猫的无限游戏
大猫的无限游戏
V
V2EX
阮一峰的网络日志
阮一峰的网络日志

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
fix(telegram): drain outbound queue after polling reconne...
steipete · 2026-05-16 · via Recent Commits to openclaw:main

@@ -9,6 +9,7 @@ const createTelegramBotMock = vi.hoisted(() => vi.fn());

99

const isRecoverableTelegramNetworkErrorMock = vi.hoisted(() => vi.fn(() => true));

1010

const computeBackoffMock = vi.hoisted(() => vi.fn(() => 0));

1111

const sleepWithAbortMock = vi.hoisted(() => vi.fn(async () => undefined));

12+

const drainPendingDeliveriesMock = vi.hoisted(() => vi.fn(async (_opts: unknown) => undefined));

12131314

vi.mock("@grammyjs/runner", () => ({

1415

run: runMock,

@@ -22,6 +23,10 @@ vi.mock("./network-errors.js", () => ({

2223

isRecoverableTelegramNetworkError: isRecoverableTelegramNetworkErrorMock,

2324

}));

242526+

vi.mock("openclaw/plugin-sdk/delivery-queue-runtime", () => ({

27+

drainPendingDeliveries: drainPendingDeliveriesMock,

28+

}));

29+2530

vi.mock("./api-logging.js", () => ({

2631

withTelegramApiErrorLogging: async ({ fn }: { fn: () => Promise<unknown> }) => await fn(),

2732

}));

@@ -54,6 +59,18 @@ type TelegramApiMiddleware = (

5459

method: string,

5560

payload: unknown,

5661

) => Promise<unknown>;

62+

type DrainPendingDeliveriesCall = {

63+

drainKey: string;

64+

logLabel: string;

65+

selectEntry: (

66+

entry: {

67+

channel: string;

68+

accountId?: string;

69+

lastError?: string;

70+

},

71+

now: number,

72+

) => { match: boolean; bypassBackoff: boolean };

73+

};

5774

type AsyncVoidFn = () => Promise<void>;

5875

type MockCallSource = { mock: { calls: Array<Array<unknown>> } };

5976

@@ -164,6 +181,14 @@ function expectTelegramBotTransportSequence(firstTransport: unknown, secondTrans

164181

expect(createTelegramBotMock.mock.calls.at(1)?.[0]?.telegramTransport).toBe(secondTransport);

165182

}

166183184+

function expectDrainPendingDeliveriesCall(index = 0): DrainPendingDeliveriesCall {

185+

const call = drainPendingDeliveriesMock.mock.calls[index]?.[0];

186+

if (!call || typeof call !== "object") {

187+

throw new Error(`Expected drainPendingDeliveries call ${index}`);

188+

}

189+

return call as DrainPendingDeliveriesCall;

190+

}

191+167192

function makeTelegramTransport() {

168193

return {

169194

fetch: globalThis.fetch,

@@ -292,6 +317,7 @@ describe("TelegramPollingSession", () => {

292317

isRecoverableTelegramNetworkErrorMock.mockReset().mockReturnValue(true);

293318

computeBackoffMock.mockReset().mockReturnValue(0);

294319

sleepWithAbortMock.mockReset().mockResolvedValue(undefined);

320+

drainPendingDeliveriesMock.mockReset().mockResolvedValue(undefined);

295321

});

296322297323

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

@@ -454,6 +480,58 @@ describe("TelegramPollingSession", () => {

454480

}

455481

});

456482483+

it("drains Telegram delivery queue after isolated ingress reports poll success", async () => {

484+

const abort = new AbortController();

485+

const init = vi.fn(async () => undefined);

486+

const bot = {

487+

api: {

488+

deleteWebhook: vi.fn(async () => true),

489+

config: { use: vi.fn() },

490+

},

491+

init,

492+

handleUpdate: vi.fn(async () => undefined),

493+

stop: vi.fn(async () => undefined),

494+

};

495+

createTelegramBotMock.mockReturnValueOnce(bot);

496+

let onMessage:

497+

| ((message: { type: "poll-success"; finishedAt: number; count: number }) => void)

498+

| undefined;

499+

let stopWorker: (() => void) | undefined;

500+

const workerDone = new Promise<void>((resolve) => {

501+

stopWorker = resolve;

502+

});

503+

const createWorker = vi.fn(() => ({

504+

onMessage: vi.fn((handler) => {

505+

onMessage = handler;

506+

return () => undefined;

507+

}),

508+

stop: vi.fn(async () => {

509+

stopWorker?.();

510+

}),

511+

task: vi.fn(async () => {

512+

await workerDone;

513+

}),

514+

}));

515+516+

const session = createPollingSession({

517+

abortSignal: abort.signal,

518+

isolatedIngress: {

519+

enabled: true,

520+

createWorker,

521+

drainIntervalMs: 10,

522+

},

523+

});

524+525+

const runPromise = session.runUntilAbort();

526+

await vi.waitFor(() => expect(init).toHaveBeenCalledTimes(1));

527+

onMessage?.({ type: "poll-success", finishedAt: Date.now(), count: 0 });

528+529+

await vi.waitFor(() => expect(drainPendingDeliveriesMock).toHaveBeenCalledTimes(1));

530+531+

abort.abort();

532+

await runPromise;

533+

});

534+457535

it("lets isolated ingress drain interleave different Telegram topic lanes", async () => {

458536

const abort = new AbortController();

459537

const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-spool-"));

@@ -1175,6 +1253,87 @@ describe("TelegramPollingSession", () => {

11751253

});

11761254

});

117712551256+

it("drains Telegram delivery queue after getUpdates confirms polling reconnect", async () => {

1257+

const abort = new AbortController();

1258+

const botStop = vi.fn(async () => undefined);

1259+

const runnerStop = vi.fn(async () => undefined);

1260+

const getApiMiddleware = mockBotCapturingApiMiddleware(botStop);

1261+

const resolveFirstTask = mockLongRunningPollingCycle(runnerStop);

1262+1263+

const session = createPollingSession({

1264+

abortSignal: abort.signal,

1265+

});

1266+1267+

const runPromise = session.runUntilAbort();

1268+

const apiMiddleware = await waitForApiMiddleware(getApiMiddleware);

1269+

await apiMiddleware(

1270+

vi.fn(async () => []),

1271+

"getUpdates",

1272+

{ offset: 1 },

1273+

);

1274+1275+

await vi.waitFor(() => expect(drainPendingDeliveriesMock).toHaveBeenCalledTimes(1));

1276+

const drain = expectDrainPendingDeliveriesCall();

1277+

expect(drain.drainKey).toBe("telegram:default");

1278+

expect(drain.logLabel).toBe("Telegram reconnect drain");

1279+

expect(drain.selectEntry({ channel: "telegram" }, Date.now())).toEqual({

1280+

match: true,

1281+

bypassBackoff: false,

1282+

});

1283+

expect(

1284+

drain.selectEntry(

1285+

{

1286+

channel: "telegram",

1287+

accountId: "default",

1288+

lastError: "Network request for 'sendMessage' failed!",

1289+

},

1290+

Date.now(),

1291+

),

1292+

).toEqual({

1293+

match: true,

1294+

bypassBackoff: false,

1295+

});

1296+

expect(drain.selectEntry({ channel: "telegram", accountId: "alerts" }, Date.now()).match).toBe(

1297+

false,

1298+

);

1299+

expect(drain.selectEntry({ channel: "whatsapp" }, Date.now()).match).toBe(false);

1300+1301+

abort.abort();

1302+

resolveFirstTask();

1303+

await runPromise;

1304+

});

1305+1306+

it("drains Telegram delivery queue after each getUpdates success", async () => {

1307+

const abort = new AbortController();

1308+

const botStop = vi.fn(async () => undefined);

1309+

const runnerStop = vi.fn(async () => undefined);

1310+

const getApiMiddleware = mockBotCapturingApiMiddleware(botStop);

1311+

const resolveFirstTask = mockLongRunningPollingCycle(runnerStop);

1312+1313+

const session = createPollingSession({

1314+

abortSignal: abort.signal,

1315+

});

1316+1317+

const runPromise = session.runUntilAbort();

1318+

const apiMiddleware = await waitForApiMiddleware(getApiMiddleware);

1319+

await apiMiddleware(

1320+

vi.fn(async () => []),

1321+

"getUpdates",

1322+

{ offset: 1 },

1323+

);

1324+

await apiMiddleware(

1325+

vi.fn(async () => []),

1326+

"getUpdates",

1327+

{ offset: 2 },

1328+

);

1329+1330+

await vi.waitFor(() => expect(drainPendingDeliveriesMock).toHaveBeenCalledTimes(2));

1331+1332+

abort.abort();

1333+

resolveFirstTask();

1334+

await runPromise;

1335+

});

1336+11781337

it("keeps polling marked connected across recoverable restart cycles", async () => {

11791338

const abort = new AbortController();

11801339

const recoverableError = new Error("recoverable polling error");