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

推荐订阅源

MongoDB | Blog
MongoDB | Blog
B
Blog
Y
Y Combinator Blog
大猫的无限游戏
大猫的无限游戏
aimingoo的专栏
aimingoo的专栏
B
Blog RSS Feed
博客园 - Franky
V
V2EX
IT之家
IT之家
WordPress大学
WordPress大学
博客园 - 三生石上(FineUI控件)
J
Java Code Geeks
F
Fortinet All Blogs
I
InfoQ
云风的 BLOG
云风的 BLOG
腾讯CDC
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
月光博客
月光博客
让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
N
Netflix TechBlog - Medium
宝玉的分享
宝玉的分享
钛媒体:引领未来商业与生活新知
钛媒体:引领未来商业与生活新知
P
Proofpoint News Feed
Microsoft Security Blog
Microsoft Security Blog

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): centralize update offset tracking · opencl...
steipete · 2026-04-28 · via Recent Commits to openclaw:main

@@ -0,0 +1,176 @@

1+

import { describe, expect, it, vi } from "vitest";

2+

import {

3+

createTelegramUpdateTracker,

4+

type TelegramUpdateTrackerState,

5+

} from "./bot-update-tracker.js";

6+

import type { TelegramUpdateKeyContext } from "./bot-updates.js";

7+8+

const updateCtx = (updateId: number): TelegramUpdateKeyContext => ({

9+

update: { update_id: updateId },

10+

});

11+12+

async function flushTrackerMicrotasks() {

13+

await Promise.resolve();

14+

await Promise.resolve();

15+

}

16+17+

function deferred() {

18+

let resolve!: () => void;

19+

const promise = new Promise<void>((resolvePromise) => {

20+

resolve = resolvePromise;

21+

});

22+

return { promise, resolve };

23+

}

24+25+

describe("createTelegramUpdateTracker", () => {

26+

it("persists accepted offsets before earlier pending updates complete", async () => {

27+

const onAcceptedUpdateId = vi.fn();

28+

const tracker = createTelegramUpdateTracker({

29+

initialUpdateId: 100,

30+

onAcceptedUpdateId,

31+

});

32+33+

const update101 = tracker.beginUpdate(updateCtx(101));

34+

if (!update101.accepted) {

35+

throw new Error("expected update 101 to be accepted");

36+

}

37+

await flushTrackerMicrotasks();

38+

expect(onAcceptedUpdateId).toHaveBeenCalledWith(101);

39+40+

const update102 = tracker.beginUpdate(updateCtx(102));

41+

if (!update102.accepted) {

42+

throw new Error("expected update 102 to be accepted");

43+

}

44+

tracker.finishUpdate(update102.update, { completed: true });

45+

await flushTrackerMicrotasks();

46+47+

expect(onAcceptedUpdateId.mock.calls.map((call) => Number(call[0]))).toEqual([101, 102]);

48+

expect(tracker.getState()).toMatchObject({

49+

highestAcceptedUpdateId: 102,

50+

highestPersistedAcceptedUpdateId: 102,

51+

highestCompletedUpdateId: 102,

52+

safeCompletedUpdateId: 100,

53+

pendingUpdateIds: [101],

54+

failedUpdateIds: [],

55+

} satisfies Partial<TelegramUpdateTrackerState>);

56+57+

tracker.finishUpdate(update101.update, { completed: true });

58+

expect(tracker.getState()).toMatchObject({

59+

highestCompletedUpdateId: 102,

60+

safeCompletedUpdateId: 102,

61+

pendingUpdateIds: [],

62+

} satisfies Partial<TelegramUpdateTrackerState>);

63+

});

64+65+

it("skips restart replays once the accepted offset is restored", async () => {

66+

const onAcceptedUpdateId = vi.fn();

67+

const firstProcess = createTelegramUpdateTracker({

68+

initialUpdateId: 100,

69+

onAcceptedUpdateId,

70+

});

71+72+

const accepted = firstProcess.beginUpdate(updateCtx(101));

73+

expect(accepted.accepted).toBe(true);

74+

await flushTrackerMicrotasks();

75+76+

const restartedProcess = createTelegramUpdateTracker({

77+

initialUpdateId: Number(onAcceptedUpdateId.mock.calls.at(-1)?.[0]),

78+

});

79+80+

expect(restartedProcess.beginUpdate(updateCtx(101))).toEqual({

81+

accepted: false,

82+

reason: "accepted-watermark",

83+

});

84+

});

85+86+

it("serializes and coalesces accepted offset persistence", async () => {

87+

const firstWrite = deferred();

88+

const secondWrite = deferred();

89+

const writes: number[] = [];

90+

const onAcceptedUpdateId = vi.fn((updateId: number) => {

91+

writes.push(updateId);

92+

if (updateId === 101) {

93+

return firstWrite.promise;

94+

}

95+

return secondWrite.promise;

96+

});

97+

const tracker = createTelegramUpdateTracker({

98+

initialUpdateId: 100,

99+

onAcceptedUpdateId,

100+

});

101+102+

const update101 = tracker.beginUpdate(updateCtx(101));

103+

const update102 = tracker.beginUpdate(updateCtx(102));

104+

const update103 = tracker.beginUpdate(updateCtx(103));

105+

expect(update101.accepted).toBe(true);

106+

expect(update102.accepted).toBe(true);

107+

expect(update103.accepted).toBe(true);

108+109+

await flushTrackerMicrotasks();

110+

expect(writes).toEqual([101]);

111+

expect(tracker.getState()).toMatchObject({

112+

highestAcceptedUpdateId: 103,

113+

highestPersistedAcceptedUpdateId: 100,

114+

} satisfies Partial<TelegramUpdateTrackerState>);

115+116+

firstWrite.resolve();

117+

await flushTrackerMicrotasks();

118+

expect(writes).toEqual([101, 103]);

119+

expect(onAcceptedUpdateId).not.toHaveBeenCalledWith(102);

120+121+

secondWrite.resolve();

122+

await flushTrackerMicrotasks();

123+

expect(tracker.getState()).toMatchObject({

124+

highestPersistedAcceptedUpdateId: 103,

125+

} satisfies Partial<TelegramUpdateTrackerState>);

126+

});

127+128+

it("keeps failed accepted updates retryable in the same process", () => {

129+

const tracker = createTelegramUpdateTracker({ initialUpdateId: 200 });

130+

const first = tracker.beginUpdate(updateCtx(201));

131+

if (!first.accepted) {

132+

throw new Error("expected first update to be accepted");

133+

}

134+

tracker.finishUpdate(first.update, { completed: false });

135+136+

expect(tracker.getState()).toMatchObject({

137+

highestAcceptedUpdateId: 201,

138+

highestCompletedUpdateId: 200,

139+

safeCompletedUpdateId: 200,

140+

failedUpdateIds: [201],

141+

} satisfies Partial<TelegramUpdateTrackerState>);

142+143+

const retry = tracker.beginUpdate(updateCtx(201));

144+

if (!retry.accepted) {

145+

throw new Error("expected failed update retry to be accepted");

146+

}

147+

tracker.finishUpdate(retry.update, { completed: true });

148+149+

expect(tracker.getState()).toMatchObject({

150+

highestAcceptedUpdateId: 201,

151+

highestCompletedUpdateId: 201,

152+

safeCompletedUpdateId: 201,

153+

failedUpdateIds: [],

154+

} satisfies Partial<TelegramUpdateTrackerState>);

155+

expect(tracker.beginUpdate(updateCtx(201))).toEqual({

156+

accepted: false,

157+

reason: "accepted-watermark",

158+

});

159+

});

160+161+

it("dedupes handler dispatch separately from the accepted watermark", () => {

162+

const onSkip = vi.fn();

163+

const tracker = createTelegramUpdateTracker({ initialUpdateId: 300, onSkip });

164+

const accepted = tracker.beginUpdate(updateCtx(301));

165+

if (!accepted.accepted) {

166+

throw new Error("expected update to be accepted");

167+

}

168+169+

expect(tracker.shouldSkipHandlerDispatch(updateCtx(301))).toBe(false);

170+

expect(tracker.shouldSkipHandlerDispatch(updateCtx(301))).toBe(true);

171+

expect(onSkip).toHaveBeenCalledWith("update:301");

172+173+

tracker.finishUpdate(accepted.update, { completed: true });

174+

expect(tracker.shouldSkipHandlerDispatch(updateCtx(301))).toBe(true);

175+

});

176+

});