fix(heartbeat): Control UI completion replies leak to an explicit heartbeat channel (#161006)

* fix(heartbeat): keep WebChat exec completions out of the heartbeat target

A background command finishing in a WebChat session wakes a heartbeat turn.
With an explicit heartbeat target (e.g. discord #heartbeat, cadence 0m), the
reply was sent to that channel instead of publishing into the session that
ran the command. When the pending queue is only exec completions and a
WebChat projection exists, deliver by projection only; target:none is still
respected and ordinary heartbeat output keeps its configured target.

* fix(heartbeat): document session-owned exec completion routing and cover mixed batches

Address review: the override uses the existing internal-projection eligibility
(not WebChat only), so say so in code and docs; document that a failed session
write retries on a later wake with no channel fallback; assert a genuine poll
writes nothing into the session; add a mixed-batch test that keeps the explicit
target.

* fix(heartbeat): keep restart continuations for Control UI sessions out of the heartbeat target

A turn interrupted by a Gateway restart resumes through a restart-sentinel
heartbeat wake. With an explicit heartbeat target, that continuation reply was
sent to the heartbeat channel. On a restart-sentinel wake, events carrying the
sentinel's task:restart-sentinel: context key now count as session-owned
alongside exec completions: they publish into the owning session through the
internal projection and join the publication idempotency key. The sentinel reads
the prefix from the shared constant.

* fix(heartbeat): retain restart occurrences until session publication commits

* test(heartbeat): name captured occurrence invariants

* test(heartbeat): split system-event admission coverage

* chore: remove stale Discord max-lines baseline entry

---------

Co-authored-by: Marvin <marvin@local>
Co-authored-by: Patrick Erichsen <patrick.a.erichsen@gmail.com>
This commit is contained in:
Marvinthebored 2026-10-03 06:04:07 +08:00 • committed by GitHub
parent 2d1011a9a4
commit f0308b72d7
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
15 changed files with 589 additions and 63 deletions

View file

@ -63,7 +63,6 @@ extensions/discord/src/monitor/model-picker.test.ts
extensions/discord/src/monitor/model-picker.view.ts
extensions/discord/src/monitor/native-command.model-picker.test.ts
extensions/discord/src/monitor/native-command.plugin-dispatch.test.ts
extensions/discord/src/monitor/native-command.ts
extensions/discord/src/outbound-adapter.test.ts
extensions/discord/src/send.sends-basic-channel-messages.test.ts
extensions/fal/image-generation-provider.ts

View file

@ -342,7 +342,8 @@ Heartbeat configuration is strict: only the fields listed above are accepted. Ac
<AccordionGroup>
<Accordion title="Session and target routing">
- Heartbeats run in the agent's main session by default (`agent:<id>:main`), or `global` when `session.scope = "global"`. Set `session` to override to a specific channel session (Discord/WhatsApp/etc.).
- `session` only affects the run context. Delivery is controlled by `target` and `to`.
- `session` only affects the run context. Delivery is controlled by `target` and `to`, except for session-owned events in an internal session (see below).
- A wake whose pending events are all session-owned (background exec completions, or the continuation of a turn interrupted by a Gateway restart) in an internal session (Control UI/WebChat, or another operator-owned session without an external route) publishes the reply into that session's transcript instead of the `target`/`to` channel. `target: "none"` still suppresses it. If the session write fails, the event stays queued for a later wake and does not fall back to the channel. Batches that also contain other events use `target`/`to` as usual.
- The default `owner` target chooses an explicitly configured owner identity. It reuses the exact account/thread only when the session's last route is a direct chat to that owner.
- A wake that carries a channel and recipient uses that named origin before owner discovery. This event destination can be a group because it is explicit, not inferred.
- To deliver to a specific channel/recipient, set a channel `target` plus `to`. `target: "last"` is an explicit opt-in to the last external conversation, including groups.

View file

@ -170,6 +170,7 @@ export async function prepareReplyRunAdmission(context: PreparedReplyRunContext)
// A heartbeat may consume only its prepared generic selection, never
// dedicated reminders or arrivals that were not part of this turn.
events: context.isHeartbeat ? (eventContext?.events ?? []) : undefined,
deferredEventIds: context.isHeartbeat ? eventContext?.deferredEventIds : undefined,
});
if (eventsBlock) {
drainedSystemEventBlocks.push(eventsBlock);

View file

@ -21,7 +21,6 @@ import {
claimAgentRunDelegatedAuthority,
releaseAgentRunDelegatedAuthority,
} from "../../infra/agent-run-registry.js";
import { withSystemEventOwner } from "../../infra/system-event-ownership.js";
import {
enqueueSystemEvent,
enqueueSystemEventEntry,
@ -43,6 +42,7 @@ import {
} from "./get-reply-run-helpers.js";
import { runPreparedReply } from "./get-reply-run.js";
import { registerPendingRequesterAuthorityCases } from "./get-reply-run.requester-authority.test-support.js";
import { registerSystemEventAdmissionCases } from "./get-reply-run.system-event-admission.test-support.js";
import {
baseParams,
createInboundBody,
@ -2883,51 +2883,6 @@ describe("runPreparedReply media-only handling", () => {
expect(call.followupRun.run.skillWorkshopProposalRevision).not.toBe(proposalRevision);
});
it("admits only system events visible to the prepared agent", async () => {
const actualSystemEvents = await vi.importActual<typeof import("./session-system-events.js")>(
"./session-system-events.js",
);
vi.mocked(drainFormattedSystemEvents).mockImplementationOnce(
actualSystemEvents.drainFormattedSystemEvents,
);
enqueueSystemEvent(
"Alpha hook finished",
withSystemEventOwner({ sessionKey: "global" }, "alpha"),
);
enqueueSystemEvent(
"Beta hook finished",
withSystemEventOwner({ sessionKey: "global" }, "beta"),
);
enqueueSystemEvent("Alpha follow-up", withSystemEventOwner({ sessionKey: "global" }, "alpha"));
await runPreparedReply(
baseParams({
agentId: "alpha",
sessionKey: "global",
opts: withReplySystemEventContext(
{ isHeartbeat: true },
{ sessionKey: "global", events: peekSystemEventEntries("agent:alpha:global") },
),
}),
);
const call = requireRunReplyAgentCall();
const context = call.followupRun.currentInboundContext;
expect(call.followupRun.prompt).toBe("[User sent media without caption]");
for (const event of ["Alpha hook finished", "Alpha follow-up"]) {
expect(context?.text).toContain(event);
expect(context?.fragments).toContainEqual({
kind: "conversation-data",
text: expect.stringContaining(event),
});
expect(call.followupRun.transcriptPrompt).not.toContain(event);
}
expect(call.followupRun.prompt).not.toContain("Beta hook finished");
expect(context?.text).not.toContain("Beta hook finished");
expect(JSON.stringify(context?.fragments)).not.toContain("Beta hook finished");
expect(peekSystemEventEntries("agent:beta:global").map((event) => event.text)).toEqual([
"Beta hook finished",
]);
});
registerSystemEventAdmissionCases({ runPrepared, requireRunReplyAgentCall });
});
/* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */

View file

@ -0,0 +1,94 @@
import { expect, it, vi } from "vitest";
import { withSystemEventOwner } from "../../infra/system-event-ownership.js";
import { enqueueSystemEvent, peekSystemEventEntries } from "../../infra/system-events.js";
import type { runReplyAgent } from "./agent-runner.runtime.js";
import type { runPreparedReply } from "./get-reply-run.js";
import { drainFormattedSystemEvents } from "./session-system-events.js";
import { withReplySystemEventContext } from "./system-event-session-key.js";
export function registerSystemEventAdmissionCases({
runPrepared,
requireRunReplyAgentCall,
}: {
runPrepared: (
overrides?: Partial<Parameters<typeof runPreparedReply>[0]>,
) => ReturnType<typeof runPreparedReply>;
requireRunReplyAgentCall: () => Parameters<typeof runReplyAgent>[0];
}): void {
it("keeps delivery-owned restart occurrences queued through production reply admission", async () => {
const actualSystemEvents = await vi.importActual<typeof import("./session-system-events.js")>(
"./session-system-events.js",
);
vi.mocked(drainFormattedSystemEvents).mockImplementationOnce(
actualSystemEvents.drainFormattedSystemEvents,
);
const sessionKey = "agent:main:restart-admission-proof";
enqueueSystemEvent("Restart continuation retained for delivery", {
sessionKey,
contextKey: "task:restart-sentinel:admission-proof",
});
const captured = peekSystemEventEntries(sessionKey);
const deferredEventIds = captured.map((event) => event.id!).filter(Boolean);
await runPrepared({
agentId: "main",
sessionKey,
opts: withReplySystemEventContext(
{ isHeartbeat: true },
{
sessionKey,
events: captured,
deferredEventIds,
},
),
});
expect(requireRunReplyAgentCall().followupRun.currentInboundContext?.text).toContain(
"Restart continuation retained for delivery",
);
expect(peekSystemEventEntries(sessionKey).map((event) => event.id)).toEqual(deferredEventIds);
});
it("admits only system events visible to the prepared agent", async () => {
const actualSystemEvents = await vi.importActual<typeof import("./session-system-events.js")>(
"./session-system-events.js",
);
vi.mocked(drainFormattedSystemEvents).mockImplementationOnce(
actualSystemEvents.drainFormattedSystemEvents,
);
enqueueSystemEvent(
"Alpha hook finished",
withSystemEventOwner({ sessionKey: "global" }, "alpha"),
);
enqueueSystemEvent(
"Beta hook finished",
withSystemEventOwner({ sessionKey: "global" }, "beta"),
);
enqueueSystemEvent("Alpha follow-up", withSystemEventOwner({ sessionKey: "global" }, "alpha"));
await runPrepared({
agentId: "alpha",
sessionKey: "global",
opts: withReplySystemEventContext(
{ isHeartbeat: true },
{ sessionKey: "global", events: peekSystemEventEntries("agent:alpha:global") },
),
});
const call = requireRunReplyAgentCall();
const context = call.followupRun.currentInboundContext;
expect(call.followupRun.prompt).toBe("[User sent media without caption]");
for (const event of ["Alpha hook finished", "Alpha follow-up"]) {
expect(context?.text).toContain(event);
expect(context?.fragments).toContainEqual({
kind: "conversation-data",
text: expect.stringContaining(event),
});
expect(call.followupRun.transcriptPrompt).not.toContain(event);
}
expect(call.followupRun.prompt).not.toContain("Beta hook finished");
expect(context?.text).not.toContain("Beta hook finished");
expect(JSON.stringify(context?.fragments)).not.toContain("Beta hook finished");
expect(peekSystemEventEntries("agent:beta:global").map((event) => event.text)).toEqual([
"Beta hook finished",
]);
});
}

View file

@ -96,6 +96,7 @@ export async function drainFormattedSystemEvents(params: {
isMainSession: boolean;
isNewSession: boolean;
events?: readonly SystemEvent[];
deferredEventIds?: readonly string[];
}): Promise<string | undefined> {
const systemLines: string[] = [];
const queueKey = resolveSystemEventQueueKey(params.sessionKey, params.agentId);
@ -106,6 +107,7 @@ export async function drainFormattedSystemEvents(params: {
(params.events ?? peekSystemEventEntries(queueKey)).filter(
(event) => !isExecCompletionEvent(event.text),
),
{ deferredEventIds: params.deferredEventIds },
);
const sessionStateNotices = queued.flatMap((event) => {
const targetSessionKey = event.contextKey

View file

@ -5,6 +5,8 @@ const REPLY_SYSTEM_EVENT_CONTEXT = Symbol("openclaw.reply.systemEventContext");
type ReplySystemEventContext = {
sessionKey: string;
events?: readonly SystemEvent[];
/** Captured occurrences whose delivery owner, not prompt admission, settles them. */
deferredEventIds?: readonly string[];
};
/** Carry the queue and its optional prepared selection through internal option spreads. */

View file

@ -19,6 +19,7 @@ import {
} from "../infra/delivery-queue-state-context.js";
import { formatErrorMessage, toErrorObject } from "../infra/errors.js";
import type { GatewayScheduler } from "../infra/gateway-scheduler.js";
import { RESTART_CONTINUATION_CONTEXT_PREFIX } from "../infra/heartbeat-events-filter.js";
import { requestHeartbeat } from "../infra/heartbeat-wake.js";
import {
clearRestartSentinelIfRevision,
@ -117,7 +118,7 @@ function enqueueRestartSentinelWake(
const eventOptions = {
sessionKey,
// Recovered work keeps its ordinary turn budget when delivered by heartbeat.
contextKey: `task:restart-sentinel:${entry.id}`,
contextKey: `${RESTART_CONTINUATION_CONTEXT_PREFIX}${entry.id}`,
...(deliveryContext ? { deliveryContext } : {}),
};
enqueueSystemEvent(message, withSystemEventOwner(eventOptions, agentId));

View file

@ -275,10 +275,10 @@ async function prepareHeartbeatDispatchReply(
accountId: delivery.accountId,
});
if (consume && preflight.shouldInspectPendingEvents) {
consumeSelectedSystemEventEntries(
resolveSystemEventQueueKey(sessionKey, agentId),
prepared.inspectedSystemEventsToConsume,
);
consumeSelectedSystemEventEntries(resolveSystemEventQueueKey(sessionKey, agentId), [
...prepared.inspectedSystemEventsToConsume,
...prepared.deferredGenericEvents,
]);
if (prepared.hasExecCompletion && prepared.hasCronEvents) {
// Coalesced waiters share this turn, but exec and cron retain separate prompt/delivery policy.
requestHeartbeat({
@ -579,7 +579,12 @@ export async function deliverHeartbeatDispatch(
if (!internalProjection || policy.projectTarget === false) {
return { visibleReplySent: false };
}
const occurrenceIds = policy.prepared.inspectedSystemEventsToConsume.map((event) => event.id);
// Restart continuations are admitted as generic prompt text, so their queue
// identities join the publication key alongside inspected completions.
const occurrenceIds = [
...policy.prepared.inspectedSystemEventsToConsume,
...policy.prepared.deferredGenericEvents,
].map((event) => event.id);
if (!occurrenceIds.every((id): id is string => typeof id === "string" && id.length > 0)) {
policy.deliveryReason = "exec completion occurrence identity unavailable";
return { visibleReplySent: false };

View file

@ -165,6 +165,14 @@ function isHeartbeatNoiseEvent(evt: string): boolean {
);
}
/** Context-key prefix the restart sentinel gives a continuation queued for one session. */
export const RESTART_CONTINUATION_CONTEXT_PREFIX = "task:restart-sentinel:";
/** A restart continuation event resumes a specific session's interrupted turn. */
export function isRestartContinuationEvent(event: { contextKey?: string | null }): boolean {
return event.contextKey?.startsWith(RESTART_CONTINUATION_CONTEXT_PREFIX) ?? false;
}
export function isExecCompletionEvent(evt: string): boolean {
const trimmed = evt.trimStart();
const normalized = normalizeLowercaseStringOrEmpty(trimmed);

View file

@ -38,7 +38,7 @@ import { formatErrorMessage } from "./errors.js";
import { isWithinActiveHours } from "./heartbeat-active-hours.js";
import { tryResolveAmbientHeartbeatAgentId } from "./heartbeat-agent-resolution.js";
import { resolveHeartbeatForWake, type HeartbeatConfig } from "./heartbeat-config.js";
import { isExecCompletionEvent } from "./heartbeat-events-filter.js";
import { isExecCompletionEvent, isRestartContinuationEvent } from "./heartbeat-events-filter.js";
import { emitHeartbeatEvent } from "./heartbeat-events.js";
import { heartbeatLog as log } from "./heartbeat-log.js";
import { shouldUseHeartbeatResponseToolPrompt } from "./heartbeat-runner-config.js";
@ -363,10 +363,15 @@ export async function prepareHeartbeatRunStage(wake: ReadyHeartbeatWake) {
const projectionSessionKey = run.kind === "isolated" ? run.baseSessionKey : sessionKey;
// Capture the client-owned generation before routing can await. The inspected
// completion queue owns publication eligibility, not the coalesced wake source.
// Session-owned work: a background command completion, or the continuation of a turn
// interrupted by a Gateway restart. Both answer the session that owns them.
const isSessionOwnedEvent = (event: (typeof preflight.pendingEventEntries)[number]) =>
isExecCompletionEvent(event.text) ||
(wake.wakeSource === "restart-sentinel" && isRestartContinuationEvent(event));
const projectionCandidate =
scheduledTasks.length === 0 &&
preflight.shouldInspectPendingEvents &&
preflight.pendingEventEntries.some((event) => isExecCompletionEvent(event.text)) &&
preflight.pendingEventEntries.some(isSessionOwnedEvent) &&
!preflight.session.suppressOriginatingContext &&
!isInternalSessionEffectsKey(projectionSessionKey) &&
conversationEntry?.delivery?.kind === "internal" &&
@ -387,7 +392,7 @@ export async function prepareHeartbeatRunStage(wake: ReadyHeartbeatWake) {
// a new session ID (empty transcript) each run, avoiding the cost of
// sending the full conversation history (~100K tokens) to the LLM.
// Delivery routing uses the selected conversation, not the fresh execution row.
const delivery = await resolveHeartbeatDeliveryTargetWithSessionRoute({
const resolvedDelivery = await resolveHeartbeatDeliveryTargetWithSessionRoute({
cfg,
agentId,
entry: conversationEntry,
@ -403,7 +408,21 @@ export async function prepareHeartbeatRunStage(wake: ReadyHeartbeatWake) {
// an explicit target that never resolves to a route also reports `target-none`.
// Gate here so neither the relay prompt nor the session publication path can
// see a projection target the resolver already declined to deliver to.
const internalProjection = delivery.reason === "target-none" ? undefined : projectionCandidate;
const internalProjection =
resolvedDelivery.reason === "target-none" ? undefined : projectionCandidate;
// Session-owned work in an internal session (the same eligibility as the routeless
// projection above: Control UI/WebChat and other operator-owned internal sessions)
// answers in that session. The heartbeat target is for heartbeat output and must not
// capture it. Mixed batches keep the resolved route for their other events. If the
// session write fails, the events stay queued for a later wake; there is no channel
// fallback.
const sessionOwnedCompletion =
internalProjection !== undefined &&
preflight.pendingEventEntries.length > 0 &&
preflight.pendingEventEntries.every(isSessionOwnedEvent);
const delivery: typeof resolvedDelivery = sessionOwnedCompletion
? { ...resolvedDelivery, channel: "none", to: undefined, reason: "session-owned-completion" }
: resolvedDelivery;
// Routeless ambient polls are pure model burn, but only they may skip:
// triggered wakes (hook/manual/cron/exec), polls with queued events, and
// scheduled-task wakes must still run to process their payloads even when
@ -594,6 +613,11 @@ export async function prepareHeartbeatRunStage(wake: ReadyHeartbeatWake) {
(wake.wakeSource === undefined ||
wake.wakeSource === "interval" ||
wake.wakeSource === "manual")),
// Session publication owns restart custody through commit, not prompt admission.
deferredGenericEvents:
internalProjection && delivery.channel === "none"
? heartbeatRunPrompt.genericEvents.filter(isRestartContinuationEvent)
: [],
} as const;
}

View file

@ -163,6 +163,9 @@ export async function runHeartbeatOnce(opts: HeartbeatRunOptions): Promise<Heart
{
sessionKey: prepared.inspectsRunQueue ? prepared.sessionKey : runSessionKey,
events: prepared.inspectsRunQueue ? prepared.genericEvents : [],
deferredEventIds: prepared.deferredGenericEvents
.map((event) => event.id)
.filter((id): id is string => typeof id === "string"),
},
),
dispatcherOptions: {

View file

@ -0,0 +1,373 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { createHeartbeatToolResponsePayload } from "../auto-reply/heartbeat-tool-response.js";
import { drainFormattedSystemEvents } from "../auto-reply/reply/session-system-events.js";
import { getReplySystemEventContext } from "../auto-reply/reply/system-event-session-key.js";
import type { OpenClawConfig } from "../config/config.js";
import { loadTranscriptEvents } from "../config/sessions/session-accessor.js";
import { readTranscriptEventMessage } from "../config/sessions/session-accessor.sqlite-read.js";
import { setTestEnvValue } from "../test-utils/env.js";
import { resetHeartbeatEventsForTest } from "./heartbeat-events.js";
import { runHeartbeatOnce } from "./heartbeat-runner.js";
import {
readSessionStoreForTest,
seedMainSessionStore,
setupTelegramHeartbeatPluginRuntimeForTests,
withTempHeartbeatSandbox,
} from "./heartbeat-runner.test-utils.js";
import * as sessionPublication from "./heartbeat-session-publication.js";
import {
enqueueSystemEvent,
peekSystemEventEntries,
resetSystemEventsForTest,
} from "./system-events.js";
// A command started from an internal (WebChat) session belongs to that session. An
// explicit heartbeat target (a channel for heartbeat chatter, cadence disabled) must
// not capture the session-owned completion reply, and ordinary heartbeat output and
// mixed batches must keep using that target.
describe("exec completion from a WebChat session with an explicit heartbeat target", () => {
beforeEach(() => setupTelegramHeartbeatPluginRuntimeForTests());
afterEach(() => {
vi.restoreAllMocks();
resetSystemEventsForTest();
resetHeartbeatEventsForTest();
});
it("publishes into the WebChat session and does not send to the heartbeat channel", async () => {
await withTempHeartbeatSandbox(async ({ tmpDir, storePath }) => {
setTestEnvValue("OPENCLAW_STATE_DIR", tmpDir);
const marker = "WEBCHAT_EXEC_COMPLETION_STAYS_HOME";
const cfg: OpenClawConfig = {
agents: {
defaults: {
workspace: tmpDir,
heartbeat: { every: "0m", target: "telegram", to: "-100999000111" },
},
},
channels: { telegram: { allowFrom: ["*"] } },
messages: { visibleReplies: "message_tool" },
session: { store: storePath },
};
const sessionKey = await seedMainSessionStore(storePath, cfg, {
lastChannel: "webchat",
lastProvider: "",
lastTo: "",
sessionId: "webchat-exec-session",
lifecycleRevision: "webchat-exec-generation",
createdVia: "operator",
});
enqueueSystemEvent("Exec completed (bg-cmd, code 0) :: " + marker, { sessionKey });
const sendTelegram = vi
.fn()
.mockResolvedValue({ messageId: "leaked", chatId: "-100999000111" });
const reply = vi.fn().mockResolvedValue(
createHeartbeatToolResponsePayload({
outcome: "done",
notify: true,
summary: "private",
notificationText: marker,
}),
);
const result = await runHeartbeatOnce({
cfg,
agentId: "main",
source: "exec-event",
intent: "event",
reason: "exec-event",
deps: { getReplyFromConfig: reply, telegram: sendTelegram },
});
expect(result.status).toBe("ran");
expect(reply).toHaveBeenCalledOnce();
expect(
sendTelegram,
"session-owned completion leaked to the heartbeat channel",
).not.toHaveBeenCalled();
const entry = readSessionStoreForTest(storePath)[sessionKey];
const events = await loadTranscriptEvents({
agentId: "main",
sessionKey,
sessionId: entry!.sessionId!,
storePath,
});
const published = events
.map(readTranscriptEventMessage)
.filter((m) => m?.role === "assistant" && JSON.stringify(m.content).includes(marker));
expect(published).toHaveLength(1);
});
});
it("still sends a genuine heartbeat poll to the explicit heartbeat channel", async () => {
await withTempHeartbeatSandbox(async ({ tmpDir, storePath }) => {
setTestEnvValue("OPENCLAW_STATE_DIR", tmpDir);
const cfg: OpenClawConfig = {
agents: {
defaults: {
workspace: tmpDir,
heartbeat: { every: "5m", target: "telegram", to: "-100999000111" },
},
},
channels: { telegram: { allowFrom: ["*"] } },
session: { store: storePath },
};
const sessionKey = await seedMainSessionStore(storePath, cfg, {
lastChannel: "webchat",
lastProvider: "",
lastTo: "",
createdVia: "operator",
});
const sendTelegram = vi.fn().mockResolvedValue({ messageId: "hb", chatId: "-100999000111" });
const reply = vi.fn().mockResolvedValue({ text: "Heartbeat alert: disk 95%" });
await runHeartbeatOnce({
cfg,
agentId: "main",
source: "manual",
intent: "immediate",
reason: "wake",
deps: { getReplyFromConfig: reply, telegram: sendTelegram },
});
expect(sendTelegram).toHaveBeenCalledTimes(1);
expect(await publishedAssistantTexts(storePath, sessionKey, "Heartbeat alert")).toHaveLength(
0,
);
});
});
it("keeps the explicit target for a batch that mixes an exec completion with another event", async () => {
await withTempHeartbeatSandbox(async ({ tmpDir, storePath }) => {
setTestEnvValue("OPENCLAW_STATE_DIR", tmpDir);
const marker = "MIXED_BATCH_USES_HEARTBEAT_TARGET";
const cfg: OpenClawConfig = {
agents: {
defaults: {
workspace: tmpDir,
heartbeat: { every: "0m", target: "telegram", to: "-100999000111" },
},
},
channels: { telegram: { allowFrom: ["*"] } },
session: { store: storePath },
};
const sessionKey = await seedMainSessionStore(storePath, cfg, {
lastChannel: "webchat",
lastProvider: "",
lastTo: "",
createdVia: "operator",
});
enqueueSystemEvent("Exec completed (bg-cmd, code 0) :: done", { sessionKey });
enqueueSystemEvent("Reminder: rotate the backup disk", { sessionKey });
const sendTelegram = vi
.fn()
.mockResolvedValue({ messageId: "mixed", chatId: "-100999000111" });
const reply = vi.fn().mockResolvedValue({ text: marker });
await runHeartbeatOnce({
cfg,
agentId: "main",
source: "exec-event",
intent: "event",
reason: "exec-event",
deps: { getReplyFromConfig: reply, telegram: sendTelegram },
});
expect(sendTelegram).toHaveBeenCalledTimes(1);
expect(JSON.stringify(sendTelegram.mock.calls[0])).toContain(marker);
});
});
it("publishes a restart continuation into the WebChat session, not the heartbeat channel", async () => {
await withTempHeartbeatSandbox(async ({ tmpDir, storePath }) => {
setTestEnvValue("OPENCLAW_STATE_DIR", tmpDir);
const marker = "RESTART_CONTINUATION_STAYS_HOME";
const cfg: OpenClawConfig = {
agents: {
defaults: {
workspace: tmpDir,
heartbeat: { every: "0m", target: "telegram", to: "-100999000111" },
},
},
channels: { telegram: { allowFrom: ["*"] } },
messages: { visibleReplies: "message_tool" },
session: { store: storePath },
};
const sessionKey = await seedMainSessionStore(storePath, cfg, {
lastChannel: "webchat",
lastProvider: "",
lastTo: "",
sessionId: "webchat-restart-session",
lifecycleRevision: "webchat-restart-generation",
createdVia: "operator",
});
enqueueSystemEvent("Gateway restarted. Continue the interrupted turn.", {
sessionKey,
contextKey: "task:restart-sentinel:queue-1",
});
const sendTelegram = vi
.fn()
.mockResolvedValue({ messageId: "leaked", chatId: "-100999000111" });
const reply = vi.fn().mockResolvedValue(
createHeartbeatToolResponsePayload({
outcome: "done",
notify: true,
summary: "private",
notificationText: marker,
}),
);
const result = await runHeartbeatOnce({
cfg,
agentId: "main",
sessionKey,
source: "restart-sentinel",
intent: "immediate",
reason: "wake",
deps: { getReplyFromConfig: reply, telegram: sendTelegram },
});
expect(result.status).toBe("ran");
expect(reply).toHaveBeenCalledOnce();
expect(
sendTelegram,
"restart continuation leaked to the heartbeat channel",
).not.toHaveBeenCalled();
expect(await publishedAssistantTexts(storePath, sessionKey, marker)).toHaveLength(1);
});
});
it.each(["rejected write", "thrown write", "failed model"] as const)(
"retains a restart occurrence after %s and settles only the successful retry",
async (failure) => {
await withTempHeartbeatSandbox(async ({ tmpDir, storePath }) => {
setTestEnvValue("OPENCLAW_STATE_DIR", tmpDir);
const marker = "RESTART_RETRY_COMMITTED";
const cfg: OpenClawConfig = {
agents: {
defaults: {
workspace: tmpDir,
heartbeat: { every: "0m", target: "telegram", to: "-100999000111" },
},
},
channels: { telegram: { allowFrom: ["*"] } },
session: { store: storePath },
};
const sessionKey = await seedMainSessionStore(storePath, cfg, {
lastChannel: "webchat",
lastProvider: "",
lastTo: "",
sessionId: "restart-retry-session",
lifecycleRevision: "restart-retry-generation",
createdVia: "operator",
});
const continuation = "Gateway restarted. Continue the interrupted turn.";
enqueueSystemEvent(continuation, {
sessionKey,
contextKey: "task:restart-sentinel:retry-1",
});
const captured = peekSystemEventEntries(sessionKey);
const sendTelegram = vi
.fn()
.mockResolvedValue({ messageId: "leaked", chatId: "-100999000111" });
const publish = vi.spyOn(sessionPublication, "publishHeartbeatSessionReply");
if (failure === "rejected write") {
publish.mockResolvedValueOnce({ ok: false, reason: "injected write rejection" });
} else if (failure === "thrown write") {
publish.mockRejectedValueOnce(new Error("injected write failure"));
}
const reply = vi.fn().mockImplementation(async (_ctx, options) => {
const context = getReplySystemEventContext(options);
// Exercise real admission: previous tests injected a reply without formatting
// generic events, which hid the pre-publication drain.
const block = await drainFormattedSystemEvents({
cfg,
agentId: "main",
sessionKey,
isMainSession: false,
isNewSession: false,
events: context?.events ?? [],
deferredEventIds: context?.deferredEventIds,
});
expect(block).toContain(continuation);
if (failure === "failed model" && reply.mock.calls.length === 1) {
throw new Error("injected model failure after admission");
}
return createHeartbeatToolResponsePayload({
outcome: "done",
notify: true,
summary: "private",
notificationText: marker,
});
});
const run = () =>
runHeartbeatOnce({
cfg,
agentId: "main",
sessionKey,
source: "restart-sentinel",
intent: "immediate",
reason: "wake",
deps: { getReplyFromConfig: reply, telegram: sendTelegram },
});
await run();
expect(publish).toHaveBeenCalledTimes(failure === "failed model" ? 0 : 1);
expect(peekSystemEventEntries(sessionKey).map((event) => event.id)).toEqual(
captured.map((event) => event.id),
);
expect(await publishedAssistantTexts(storePath, sessionKey, marker)).toHaveLength(0);
expect(sendTelegram).not.toHaveBeenCalled();
await run();
expect(reply).toHaveBeenCalledTimes(2);
expect(publish).toHaveBeenCalledTimes(failure === "failed model" ? 1 : 2);
expect(await publishedAssistantTexts(storePath, sessionKey, marker)).toHaveLength(1);
expect(peekSystemEventEntries(sessionKey)).toEqual([]);
expect(sendTelegram).not.toHaveBeenCalled();
});
},
);
it("keeps the explicit target when a restart wake carries an untagged event", async () => {
await withTempHeartbeatSandbox(async ({ tmpDir, storePath }) => {
setTestEnvValue("OPENCLAW_STATE_DIR", tmpDir);
const marker = "UNTAGGED_WAKE_USES_HEARTBEAT_TARGET";
const cfg: OpenClawConfig = {
agents: {
defaults: {
workspace: tmpDir,
heartbeat: { every: "0m", target: "telegram", to: "-100999000111" },
},
},
channels: { telegram: { allowFrom: ["*"] } },
session: { store: storePath },
};
const sessionKey = await seedMainSessionStore(storePath, cfg, {
lastChannel: "webchat",
lastProvider: "",
lastTo: "",
createdVia: "operator",
});
enqueueSystemEvent("Reminder: rotate the backup disk", { sessionKey });
const sendTelegram = vi.fn().mockResolvedValue({ messageId: "hb", chatId: "-100999000111" });
const reply = vi.fn().mockResolvedValue({ text: marker });
await runHeartbeatOnce({
cfg,
agentId: "main",
sessionKey,
source: "restart-sentinel",
intent: "immediate",
reason: "wake",
deps: { getReplyFromConfig: reply, telegram: sendTelegram },
});
expect(sendTelegram).toHaveBeenCalledTimes(1);
expect(JSON.stringify(sendTelegram.mock.calls[0])).toContain(marker);
});
});
});
async function publishedAssistantTexts(storePath: string, sessionKey: string, needle: string) {
const entry = readSessionStoreForTest(storePath)[sessionKey];
if (!entry?.sessionId) {
return [];
}
const events = await loadTranscriptEvents({
agentId: "main",
sessionKey,
sessionId: entry.sessionId,
storePath,
});
return events
.map(readTranscriptEventMessage)
.filter((m) => m?.role === "assistant" && JSON.stringify(m.content).includes(needle));
}

View file

@ -33,6 +33,57 @@ import {
type SystemEvent,
} from "./system-events.js";
describe("delivery-owned system event selection", () => {
beforeEach(() => resetSystemEventsForTest());
afterEach(() => resetSystemEventsForTest());
it("formats retained occurrences in captured order and leaves late arrivals untouched", async () => {
const sessionKey = "agent:main:deferred-order-proof";
enqueueSystemEvent("Restart first", { sessionKey, contextKey: "task:restart-sentinel:first" });
enqueueSystemEvent("Newer instruction second", { sessionKey });
const captured = peekSystemEventEntries(sessionKey);
const retainedId = expectDefined(captured[0]?.id, "captured restart occurrence ID");
enqueueSystemEvent("Late arrival third", { sessionKey });
const text = await drainFormattedSystemEvents({
cfg: {},
agentId: "main",
sessionKey,
isMainSession: false,
isNewSession: false,
events: captured,
deferredEventIds: [retainedId],
});
expect(text?.indexOf("Restart first")).toBeLessThan(text!.indexOf("Newer instruction second"));
expect(text).not.toContain("Late arrival third");
expect(peekSystemEvents(sessionKey)).toEqual(["Restart first", "Late arrival third"]);
consumeSelectedSystemEventEntries(sessionKey, [captured[0]!]);
expect(peekSystemEvents(sessionKey)).toEqual(["Late arrival third"]);
});
it("does not re-admit a captured retained occurrence consumed by another turn", async () => {
const sessionKey = "agent:main:deferred-membership-proof";
enqueueSystemEvent("Already handled restart", {
sessionKey,
contextKey: "task:restart-sentinel:handled",
});
const captured = peekSystemEventEntries(sessionKey);
const retainedId = expectDefined(captured[0]?.id, "captured restart occurrence ID");
consumeSelectedSystemEventEntries(sessionKey, captured);
enqueueSystemEvent("Replacement arrival", { sessionKey });
const text = await drainFormattedSystemEvents({
cfg: {},
agentId: "main",
sessionKey,
isMainSession: false,
isNewSession: false,
events: captured,
deferredEventIds: [retainedId],
});
expect(text).toBeUndefined();
expect(peekSystemEvents(sessionKey)).toEqual(["Replacement arrival"]);
});
});
type SystemEventsModule = typeof import("./system-events.js");
const systemEventsModuleUrl = new URL("./system-events.ts", import.meta.url).href;

View file

@ -258,25 +258,32 @@ function resetQueueState(key: string, entry: SessionQueue) {
export function consumeSelectedSystemEventEntries(
sessionKey: string,
consumedEntries: readonly SystemEvent[],
options?: { deferredEventIds?: readonly string[] },
): SystemEvent[] {
const key = requireSessionKey(sessionKey);
const entry = queues.get(key);
if (!entry || entry.queue.length === 0 || consumedEntries.length === 0) {
return [];
}
const removed: SystemEvent[] = [];
// Prompt admission can defer captured occurrences to a delivery owner. Selection
// still resolves against the live queue, in captured order, never late arrivals.
const deferredIds = new Set(options?.deferredEventIds);
const selected: SystemEvent[] = [];
for (const consumed of consumedEntries) {
const index = entry.queue.findIndex((event) => matchesConsumedSystemEvent(event, consumed));
if (index === -1) {
continue;
}
const [event] = entry.queue.splice(index, 1);
const event = entry.queue[index];
if (event) {
removed.push(cloneSystemEvent(event));
if (!event.id || !deferredIds.has(event.id)) {
entry.queue.splice(index, 1);
}
selected.push(cloneSystemEvent(event));
}
}
resetQueueState(key, entry);
return removed;
return selected;
}
export function drainSystemEvents(sessionKey: string): string[] {