Skip to content

Commit 1c9f62f

Browse files
VACIncBunsDev
andauthored
fix(gateway): restart sentinel wakes session after restart and preserves thread routing (#53940) thanks @VACInc
Co-authored-by: VACInc <[email protected]> Co-authored-by: Val Alexander <[email protected]>
1 parent 23a4932 commit 1c9f62f

9 files changed

Lines changed: 587 additions & 31 deletions

CHANGELOG.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ Docs: https://docs.openclaw.ai
2424

2525
### Fixes
2626

27+
- Gateway/restart sentinel: wake the interrupted agent session via heartbeat after restart instead of only sending a best-effort restart note, retry outbound delivery once on transient failure, and preserve explicit thread/topic routing through the wake path so replies land in the correct Telegram topic or Slack thread. (#53940) Thanks @VACInc.
2728
- Memory/builtin sqlite: cut redundant sync and status query churn by snapshotting file state once per source, reusing sync statements, and consolidating status aggregation reads, which reduces builtin memory overhead on sync/status/doctor-style paths. Thanks @vincentkoc.
2829
- ACP/direct chats: always deliver a terminal ACP result when final TTS does not yield audio, even if block text already streamed earlier, and skip redundant empty-text final synthesis. (#53692) Thanks @w-sss.
2930
- Doctor/image generation: seed migrated legacy Nano Banana Google provider config with the `/v1beta` API root and an empty model list so `openclaw doctor --fix` completes and the migrated native Google image path keeps hitting the correct endpoint. (#53757) Thanks @mahopan.

src/gateway/server-restart-sentinel.test.ts

Lines changed: 174 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
import { describe, expect, it, vi } from "vitest";
1+
import { beforeEach, describe, expect, it, vi } from "vitest";
22

33
const mocks = vi.hoisted(() => ({
44
resolveSessionAgentId: vi.fn(() => "agent-from-key"),
@@ -25,8 +25,13 @@ const mocks = vi.hoisted(() => ({
2525
})),
2626
normalizeChannelId: vi.fn((channel: string) => channel),
2727
resolveOutboundTarget: vi.fn(() => ({ ok: true as const, to: "+15550002" })),
28-
deliverOutboundPayloads: vi.fn(async () => []),
28+
deliverOutboundPayloads: vi.fn(async () => [{ channel: "whatsapp", messageId: "msg-1" }]),
29+
enqueueDelivery: vi.fn(async () => "queue-1"),
30+
ackDelivery: vi.fn(async () => {}),
31+
failDelivery: vi.fn(async () => {}),
2932
enqueueSystemEvent: vi.fn(),
33+
requestHeartbeatNow: vi.fn(),
34+
logWarn: vi.fn(),
3035
}));
3136

3237
vi.mock("../agents/agent-scope.js", () => ({
@@ -72,23 +77,187 @@ vi.mock("../infra/outbound/deliver.js", () => ({
7277
deliverOutboundPayloads: mocks.deliverOutboundPayloads,
7378
}));
7479

80+
vi.mock("../infra/outbound/delivery-queue.js", () => ({
81+
enqueueDelivery: mocks.enqueueDelivery,
82+
ackDelivery: mocks.ackDelivery,
83+
failDelivery: mocks.failDelivery,
84+
}));
85+
7586
vi.mock("../infra/system-events.js", () => ({
7687
enqueueSystemEvent: mocks.enqueueSystemEvent,
7788
}));
7889

90+
vi.mock("../infra/heartbeat-wake.js", () => ({
91+
requestHeartbeatNow: mocks.requestHeartbeatNow,
92+
}));
93+
94+
vi.mock("../logging/subsystem.js", () => ({
95+
createSubsystemLogger: vi.fn(() => ({
96+
warn: mocks.logWarn,
97+
})),
98+
}));
99+
79100
const { scheduleRestartSentinelWake } = await import("./server-restart-sentinel.js");
80101

81102
describe("scheduleRestartSentinelWake", () => {
82-
it("forwards session context to outbound delivery", async () => {
83-
await scheduleRestartSentinelWake({ deps: {} as never });
103+
beforeEach(() => {
104+
vi.useRealTimers();
105+
mocks.consumeRestartSentinel.mockResolvedValue({
106+
payload: {
107+
sessionKey: "agent:main:main",
108+
deliveryContext: {
109+
channel: "whatsapp",
110+
to: "+15550002",
111+
accountId: "acct-2",
112+
},
113+
},
114+
});
115+
mocks.deliverOutboundPayloads.mockReset();
116+
mocks.deliverOutboundPayloads.mockResolvedValue([{ channel: "whatsapp", messageId: "msg-1" }]);
117+
mocks.enqueueDelivery.mockReset();
118+
mocks.enqueueDelivery.mockResolvedValue("queue-1");
119+
mocks.ackDelivery.mockClear();
120+
mocks.failDelivery.mockClear();
121+
mocks.enqueueSystemEvent.mockClear();
122+
mocks.requestHeartbeatNow.mockClear();
123+
mocks.logWarn.mockClear();
124+
});
125+
126+
it("enqueues the sentinel note and wakes the session even when outbound delivery succeeds", async () => {
127+
const deps = {} as never;
128+
129+
await scheduleRestartSentinelWake({ deps });
84130

85131
expect(mocks.deliverOutboundPayloads).toHaveBeenCalledWith(
86132
expect.objectContaining({
87133
channel: "whatsapp",
88134
to: "+15550002",
89135
session: { key: "agent:main:main", agentId: "main" },
136+
deps,
137+
bestEffort: false,
138+
skipQueue: true,
139+
}),
140+
);
141+
expect(mocks.enqueueDelivery).toHaveBeenCalledWith(
142+
expect.objectContaining({
143+
channel: "whatsapp",
144+
to: "+15550002",
145+
payloads: [{ text: "restart message" }],
146+
bestEffort: false,
147+
}),
148+
);
149+
expect(mocks.ackDelivery).toHaveBeenCalledWith("queue-1");
150+
expect(mocks.failDelivery).not.toHaveBeenCalled();
151+
expect(mocks.enqueueSystemEvent).toHaveBeenCalledWith(
152+
"restart message",
153+
expect.objectContaining({
154+
sessionKey: "agent:main:main",
155+
}),
156+
);
157+
expect(mocks.requestHeartbeatNow).toHaveBeenCalledWith({
158+
reason: "wake",
159+
sessionKey: "agent:main:main",
160+
});
161+
expect(mocks.logWarn).not.toHaveBeenCalled();
162+
});
163+
164+
it("retries outbound delivery once and logs a warning without dropping the agent wake", async () => {
165+
vi.useFakeTimers();
166+
mocks.deliverOutboundPayloads
167+
.mockRejectedValueOnce(new Error("transport not ready"))
168+
.mockResolvedValueOnce([{ channel: "whatsapp", messageId: "msg-2" }]);
169+
170+
const wakePromise = scheduleRestartSentinelWake({ deps: {} as never });
171+
await vi.runAllTimersAsync();
172+
await wakePromise;
173+
174+
expect(mocks.enqueueDelivery).toHaveBeenCalledTimes(1);
175+
expect(mocks.deliverOutboundPayloads).toHaveBeenCalledTimes(2);
176+
expect(mocks.deliverOutboundPayloads).toHaveBeenNthCalledWith(
177+
1,
178+
expect.objectContaining({
179+
skipQueue: true,
90180
}),
91181
);
92-
expect(mocks.enqueueSystemEvent).not.toHaveBeenCalled();
182+
expect(mocks.deliverOutboundPayloads).toHaveBeenNthCalledWith(
183+
2,
184+
expect.objectContaining({
185+
skipQueue: true,
186+
}),
187+
);
188+
expect(mocks.ackDelivery).toHaveBeenCalledWith("queue-1");
189+
expect(mocks.failDelivery).not.toHaveBeenCalled();
190+
expect(mocks.enqueueSystemEvent).toHaveBeenCalledTimes(1);
191+
expect(mocks.requestHeartbeatNow).toHaveBeenCalledTimes(1);
192+
expect(mocks.logWarn).toHaveBeenCalledWith(
193+
expect.stringContaining("retrying in 750ms"),
194+
expect.objectContaining({
195+
channel: "whatsapp",
196+
to: "+15550002",
197+
sessionKey: "agent:main:main",
198+
attempt: 1,
199+
maxAttempts: 2,
200+
}),
201+
);
202+
});
203+
204+
it("keeps one queued restart notice when outbound retries are exhausted", async () => {
205+
vi.useFakeTimers();
206+
mocks.deliverOutboundPayloads
207+
.mockRejectedValueOnce(new Error("transport not ready"))
208+
.mockRejectedValueOnce(new Error("transport still not ready"));
209+
210+
const wakePromise = scheduleRestartSentinelWake({ deps: {} as never });
211+
await vi.runAllTimersAsync();
212+
await wakePromise;
213+
214+
expect(mocks.enqueueDelivery).toHaveBeenCalledTimes(1);
215+
expect(mocks.deliverOutboundPayloads).toHaveBeenCalledTimes(2);
216+
expect(mocks.ackDelivery).not.toHaveBeenCalled();
217+
expect(mocks.failDelivery).toHaveBeenCalledWith("queue-1", "transport still not ready");
218+
});
219+
220+
it("prefers top-level sentinel threadId for wake routing context", async () => {
221+
// Legacy or malformed sentinel JSON can still carry a nested threadId.
222+
mocks.consumeRestartSentinel.mockResolvedValue({
223+
payload: {
224+
sessionKey: "agent:main:main",
225+
deliveryContext: {
226+
channel: "whatsapp",
227+
to: "+15550002",
228+
accountId: "acct-2",
229+
threadId: "stale-thread",
230+
} as never,
231+
threadId: "fresh-thread",
232+
},
233+
} as Awaited<ReturnType<typeof mocks.consumeRestartSentinel>>);
234+
235+
await scheduleRestartSentinelWake({ deps: {} as never });
236+
237+
expect(mocks.enqueueSystemEvent).toHaveBeenCalledWith(
238+
"restart message",
239+
expect.objectContaining({
240+
sessionKey: "agent:main:main",
241+
deliveryContext: expect.objectContaining({
242+
threadId: "fresh-thread",
243+
}),
244+
}),
245+
);
246+
});
247+
248+
it("does not wake the main session when the sentinel has no sessionKey", async () => {
249+
mocks.consumeRestartSentinel.mockResolvedValue({
250+
payload: {
251+
message: "restart message",
252+
},
253+
} as unknown as Awaited<ReturnType<typeof mocks.consumeRestartSentinel>>);
254+
255+
await scheduleRestartSentinelWake({ deps: {} as never });
256+
257+
expect(mocks.enqueueSystemEvent).toHaveBeenCalledWith("restart message", {
258+
sessionKey: "agent:main:main",
259+
});
260+
expect(mocks.requestHeartbeatNow).not.toHaveBeenCalled();
261+
expect(mocks.deliverOutboundPayloads).not.toHaveBeenCalled();
93262
});
94263
});

0 commit comments

Comments
 (0)