mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-04 02:00:10 +00:00
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:
parent
2d1011a9a4
commit
f0308b72d7
15 changed files with 589 additions and 63 deletions
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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. */
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
]);
|
||||
});
|
||||
}
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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. */
|
||||
|
|
|
|||
|
|
@ -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));
|
||||
|
|
|
|||
|
|
@ -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 };
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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: {
|
||||
|
|
|
|||
373
src/infra/heartbeat-runner.webchat-exec-explicit-target.test.ts
Normal file
373
src/infra/heartbeat-runner.webchat-exec-explicit-target.test.ts
Normal 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));
|
||||
}
|
||||
|
|
@ -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;
|
||||
|
|
|
|||
|
|
@ -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[] {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue