fix: show native background commands in Tasks (#158038)

* fix: show native background commands in Tasks

* fix: preserve task cancellation context and terminal cleanup

Retain the admitting plugin registry for native command cancellation while preserving the current Gateway caller. Release terminal ownership after failed publication so existing maintenance can reconcile it. Keep cancellation errors scoped to their task details and visible in the list.

* fix: retain gateway authority for task cancellation

* test: expose doctor isolation failure reports

* fix: distinguish failed and stopped command tasks
This commit is contained in:
Peter Steinberger 2026-09-25 07:21:28 -07:00 • committed by GitHub
parent 22fbcfca99
commit bdbe64d019
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
24 changed files with 1322 additions and 81 deletions

View file

@ -402,6 +402,16 @@ continuation records that result without rewriting the earlier turn's snapshot.
The existing unknown-outcome audit diagnostic remains; cancellation and a
command with no confirmed live owner retain their failure handling.
These retained commands also appear in **Tasks**, where you can follow completion
or stop an individual command. The task follows the native process after the
foreground turn ends; its final status does not rewrite the earlier tool row.
A known nonzero exit reports **Command failed**, even if a Stop request races
with completion. A confirmed Stop with no native exit result reports **Command
stopped**; this records the acknowledged request without attributing the exit to
a particular signal. Task updates do not automatically start another model
turn. If the native connection is lost before completion is confirmed, the task
reports an unknown outcome instead of success.
Stopping an active Codex run interrupts its turn. With the OpenClaw sandbox
exec-server, cleanup stops the concrete processes admitted by that turn and
preserves independent background work in the same reused thread. Each process

View file

@ -98,6 +98,33 @@ store, identity, binding, and live authority, removes only the exact upstream
link, then invokes backend cleanup. Queue selection, native protocol/policy,
and resource cleanup remain with the backend; core owns host session lifecycle.
## Background command tasks
Official harnesses can use `createAgentHarnessCommandTask` from the existing
private `openclaw/plugin-sdk/agent-harness-task-runtime` entrypoint to expose a
native command in Tasks after its foreground turn ends. Pass the host-issued
task scope and retain the original native connection and source authority. The
helper creates a worker-persisted CLI task and binds cancellation to that exact
task run; it does not take custody of the native process.
The cancellation callback receives `assertTaskCurrent`; call it after awaited
preparation and immediately before stopping work, alongside the retained source
and concrete command checks. Publish the native terminal outcome with `finish`.
It returns `"published"` after terminal publication or `"retired"` when the original
task was replaced. Retirement releases the old binding without changing its
successor; both results let the harness release its native observation leases.
A successful stop requires the original task to settle as cancelled; natural
completion racing Stop remains success. Failed publication retains the run owner;
the harness must either own a subsequent settlement attempt or release the binding
so normal task recovery can reconcile the row. A one-shot terminal notification
must not leave a finished command holding live ownership indefinitely. Release
the binding when the native owner closes and cannot publish an outcome. Restored
rows do not recreate native process authority.
Command previews use the shared redacted exec formatter, and Incognito content
stays private. These tasks are silent: recording completion does not schedule a
new model turn.
## Subagent task history
Native subagents can expose the shared task transcript view through the optional

View file

@ -278,7 +278,7 @@ Run-error banners offer **Refresh** to reload the conversation without resending
- File paths recognized in chat messages read as their basename with a small glyph for the file type in front — a Markdown page, a `package.json` manifest, a TypeScript source, a `.tsx` component, a config or data file, a shell script, and an image each get their own mark, and anything else falls back to a plain document. When two links in the same message share a basename, each keeps just enough of its trailing path to stay distinct. The full path stays on the link: it is what the tooltip shows, what opens in the file panel, and what the message's **Copy** action returns, since copy hands back the original Markdown. Labels you write yourself in a `[label](path)` link are never rewritten. The glyph is drawn from the bundled icon set, never fetched from the network, and is decorative only: it is not read by screen readers and is not part of copied text. Text that is not a recognizable path — anything carrying spaces, parentheses, a `#` fragment, or a `?` query — stays plain prose.
- Clicking a file reference in chat, a file path in an expanded read/edit/write tool card, or a file row in **Files** opens its own filename tab in the shared side-panel header. Reopening the same file selects its existing tab and rereads its content when there is no unsaved draft. Selecting a filename tab keeps its current preview; unsaved drafts are never replaced by a file reopen. The folder action returns to the file browser without closing previews. The last opened file stays highlighted in both session and project lists, including after refreshing the file list. A pending listing cannot clear a newer file selection or replace results and errors for a different folder or search. If the folder being browsed becomes unavailable, **Files** keeps its parent-folder action so you can continue browsing without reloading or changing sessions. Session file labels show the filename and enough parent folders to distinguish matching names; hovering or copying a path keeps the full path. Closing a filename tab closes only its preview, never the underlying file. Closing or replacing a file preview cancels a delayed copy fallback; an already issued native clipboard write may still finish. Open previews are scoped to the current session, agent, and connection, and are not persisted across reconnects. HTML files open a sandboxed **Preview**, with **Source** in the same filename tab. Other UTF-8 text files use a CodeMirror-based code view with syntax highlighting, line numbers, jump-to-line, in-file search, copy actions, and an open-in-external-editor menu. The code view has a **Word wrap** toolbar toggle, including in HTML **Source** view. Wrapping starts off; the browser remembers your choice across files and reloads without changing file contents. Search follows the displayed line numbers for LF, CRLF, and CR line endings; editing preserves the original line endings, including when pasted text uses different line endings. Read-only previews, including files with mixed line endings, do not create unsaved drafts or block interface reloads. Escape closes in-file search and returns keyboard focus to **Search in file** in the toolbar. AVIF, GIF, JPEG, PNG, and WebP images no larger than 256 KiB render inline; other binary files show metadata without lossy text decoding. When the Gateway advertises `sessions.files.set` to an `operator.admin` connection, the text panel adds an Edit mode with dirty tracking and Cmd/Ctrl-S save; unsaved drafts survive file, panel, and session navigation in the current browser tab until explicitly saved or discarded. Saves are compare-and-swap on a content hash returned by `sessions.files.get`: if the file changed on disk since it was loaded (for example because the agent kept working), the panel shows a conflict notice with Reload (take the latest content) and Overwrite (keep the local edit) actions. Writes go through the same fs-safe workspace guards as reads — path containment, symlink/hardlink rejection, and a 256 KiB UTF-8 cap — and only overwrite existing files; the editor never creates or deletes them. If the editor cannot load, use **Retry** or **View Raw Text**. A missing editor chunk after an update offers **Reload**, which waits for the Gateway to become reachable.
- Subagent runs appear in inline transcript activity rows, the chat **Tasks** tab, and the Tasks page. They have no sidebar row; opening a run in the main chat view is view-only. The composer identifies the parent session and offers **Open parent session** so you can continue the conversation there. Message input, reply actions, model and access pickers, microphone, and attachment controls are hidden. This does not change copy or fork availability; **Open parent session** takes you to the conversation where you can reply. **Stop** remains available when the Gateway reports an abortable run. Spawned persistent sessions (visible sessions in the session tree) are not subagents: a subagent run ends, a session does not, and you can always type in it.
- The **Tasks** tab lists the current agent's background tasks and subagents (`tasks.list` scoped by agent, kept live by `task` events): running work shows a live elapsed timer, tool-use count, the tool currently in use, and a stop control, while the collapsible finished section adds run durations. Inline subagent activity rows show ongoing status and progress without per-task edit counters; finished runs appear only in Tasks history. Task details retain each task’s cumulative edit-activity counter; the checkout chip above the composer shows the session checkout’s actual Git diff. Selecting a task from either a task row or an inline subagent activity row opens its live status and transcript inside **Tasks** without replacing the main conversation or the **Review** diff; tasks whose session is the current conversation show their prompt and output inspector there instead. Select **Back to tasks** to return to the list. Closing the Tasks tab clears inspection; switching tabs or minimizing the panel preserves it. Open **Tasks** with the title-bar activity toggle or the panel's **+** menu; the task snapshot loads eagerly, so the title-bar toggle carries a running-count badge without opening the tab first. The Tasks page remains the full cross-agent ledger.
- The **Tasks** tab lists the current agent's background tasks and subagents (`tasks.list` scoped by agent, kept live by `task` events): running work shows a live elapsed timer, tool-use count, the tool currently in use, and a stop control, while the collapsible finished section adds run durations. Inline subagent activity rows show ongoing status and progress without per-task edit counters; finished runs appear only in Tasks history. Task details retain each task’s cumulative edit-activity counter; the checkout chip above the composer shows the session checkout’s actual Git diff. Selecting a task from either a task row or an inline subagent activity row opens its live status and transcript inside **Tasks** without replacing the main conversation or the **Review** diff; tasks whose session is the current conversation show their prompt and output inspector there instead. Select **Back to tasks** to return to the list. A failed **Stop** remains visible on that task and in the list; opening another task does not show the failure there. Closing the Tasks tab clears inspection; switching tabs or minimizing the panel preserves it. Open **Tasks** with the title-bar activity toggle or the panel's **+** menu; the task snapshot loads eagerly, so the title-bar toggle carries a running-count badge without opening the tab first. The Tasks page remains the full cross-agent ledger.
- After a chat turn finishes, remaining background work appears as an inline task count followed by elapsed time. Hover or focus the count to preview tasks; select it to open **Tasks**. The status disappears when no active tasks remain or the Gateway disconnects.
- **Tasks** and **Review** retain their selections independently of each other and of file tabs. Reloading restores the selected task from current Tasks data; if that task is no longer available, Tasks says so instead of showing workspace Git changes. A pending file or artifact updates only its own open tab: it cannot select itself over a newer tab, reopen a closed preview, or return after you leave the chat page. Switching tabs or hiding the whole side panel preserves the pending preview without changing your chosen layout when it finishes. Text attachments retain their Preview or View Raw Text mode while switching between open files. Background download-link refreshes keep an unchanged attachment's reader in place, including keyboard focus and code-block controls.
- Each task has a main view and a unified side panel. The task toolbar's **Swap** button exchanges the main view and active side-panel tab; its tooltip names both views, for example **Swap Chat and Dashboard**. Chat, Dashboard, Browser, Terminal, Files, Tasks, and Review can all be main. Other side-panel tabs remain available. **Focus** in the main pane header gives that view the full task area; **Restore split** brings the side panel back. Swapping or focusing preserves live content and drafts. Closing the whole side panel hides it without changing the main view, and the browser remembers each task's arrangement.

View file

@ -0,0 +1,295 @@
import { randomUUID } from "node:crypto";
import { embeddedAgentLog, formatErrorMessage } from "openclaw/plugin-sdk/agent-harness-runtime";
import { createAgentHarnessCommandTask } from "openclaw/plugin-sdk/agent-harness-task-runtime";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import { protectCodexAppServerLiveThread } from "./client-runtime.js";
import { isCodexNotificationForTurn } from "./notification-correlation.js";
import { isJsonObject, type CodexAppServerRequestResult } from "./protocol.js";
import type { CodexAttemptResources } from "./run-attempt-resources.js";
import { retainSharedCodexAppServerClientIfCurrent } from "./shared-client.js";
import { isSameCodexAppServerThreadOwner } from "./thread-ownership.js";
import { waitForPromiseOrAbort } from "./timeout.js";
type CommandTask = Awaited<ReturnType<typeof createAgentHarnessCommandTask>>;
type Terminal = Parameters<CommandTask["finish"]>[0];
type Inventory = CodexAppServerRequestResult<"thread/backgroundTerminals/list">["data"];
type Entry = {
processId?: string;
task?: CommandTask;
terminal?: Terminal;
settlement?: Promise<void>;
cancellationAttempt?: Promise<void>;
done: ReturnType<typeof createDeferred<void>>;
cancellation: "idle" | "pending" | "confirmed";
nativeCompletion?: { exitCode: number | undefined };
};
/** Task projection retains native custody; it never becomes a second process registry. */
export function prepareCodexNativeCommandTasks(
resources: CodexAttemptResources,
turnId: string,
pending: ReadonlyMap<string, string | null>,
) {
const { connection } = resources.prompt.context.runtime;
const { params, bindingStore, bindingIdentity, appServer } = connection;
const scope = params.agentHarnessTaskRuntimeScope;
if (!scope) {
return undefined;
}
const source = params.hostCapabilities.retainSourceAuthority?.();
if (!source) {
return undefined;
}
const { client, thread } = resources.state;
const custody = resources.nativeProcessAuthority;
const agentId = params.agentId;
const requestTimeoutMs = appServer.requestTimeoutMs;
const entries = new Map<string, Entry>(
[...pending.keys()].map((itemId) => [
itemId,
{
done: createDeferred<void>(),
cancellation: "idle",
},
]),
);
let admitting = true;
let closed = false;
let released = false;
const releaseClient = retainSharedCodexAppServerClientIfCurrent(client);
const releaseThread = protectCodexAppServerLiveThread(client, thread.threadId);
const assertCurrent = () => {
source.assertCurrent();
if (
closed ||
released ||
!isSameCodexAppServerThreadOwner(bindingStore.read(bindingIdentity), thread)
) {
throw new Error("Native command no longer belongs to this source");
}
};
const releaseIfFinished = () => {
if (released || admitting || entries.size > 0) {
return;
}
released = true;
unwatch();
unwatchClose();
source.signal?.removeEventListener("abort", sourceEnded);
releaseThread();
releaseClient?.();
source.release();
};
const settle = async (itemId: string) => {
const entry = entries.get(itemId);
if (!entry?.task || !entry.terminal || entry.cancellation === "pending") {
return;
}
const terminal =
entry.nativeCompletion &&
entry.cancellation === "confirmed" &&
entry.terminal.status === "failed" &&
// Codex uses -1 when no exit result is available; Stop can acknowledge an already-exited process.
(entry.nativeCompletion.exitCode === undefined || entry.nativeCompletion.exitCode === -1)
? {
...entry.terminal,
status: "cancelled" as const,
error: "Stop confirmed; native exit result unavailable.",
terminalSummary: "Command stopped",
}
: entry.terminal;
await (entry.settlement ??= entry.task
.finish(terminal)
.then(() => {
entries.delete(itemId);
entry.done.resolve();
releaseIfFinished();
})
.catch((error: unknown) => {
// Leave unconfirmed durable outcomes to canonical task recovery, not a live run claim.
entry.task?.release();
entries.delete(itemId);
entry.done.resolve();
releaseIfFinished();
throw error;
}));
};
const ownerEnded = async () => {
closed = true;
await Promise.all(
[...entries].map(async ([itemId, entry]) => {
entry.terminal ??= {
status: "failed",
endedAt: Date.now(),
error: "Native command owner closed before its outcome was collected.",
terminalSummary: "Command outcome unknown",
};
await settle(itemId);
}),
);
};
const report = (error: unknown) =>
embeddedAgentLog.warn("Native command task settlement failed", {
error: formatErrorMessage(error),
});
const sourceEnded = () => {
void ownerEnded().catch(report);
};
const unwatchClose = client.addCloseHandler(sourceEnded);
const unwatch = client.addNotificationHandler(async (notification) => {
if (
notification.method !== "item/completed" ||
!isCodexNotificationForTurn(notification.params, thread.threadId, turnId) ||
!isJsonObject(notification.params)
) {
return;
}
const item = notification.params.item;
if (!isJsonObject(item) || item.type !== "commandExecution" || typeof item.id !== "string") {
return;
}
const entry = entries.get(item.id);
if (
!entry ||
entry.nativeCompletion ||
entry.settlement ||
(entry.processId && typeof item.processId === "string" && item.processId !== entry.processId)
) {
return;
}
const exitCode = typeof item.exitCode === "number" ? item.exitCode : undefined;
const succeeded = item.status === "completed" && exitCode === 0;
entry.nativeCompletion = { exitCode };
// A collected native result supersedes an owner-close placeholder until settlement starts.
entry.terminal = {
status: succeeded ? "succeeded" : "failed",
endedAt: Date.now(),
terminalSummary: succeeded ? "Command completed" : "Command failed",
...(succeeded ? { clearError: true } : { error: "Native command failed." }),
...(exitCode !== undefined ? { detail: { exitCode } } : {}),
};
await settle(item.id);
});
source.signal?.addEventListener("abort", sourceEnded, { once: true });
return {
async retain(inventory: Inventory, retained: ReadonlyMap<string, string>) {
try {
for (const [itemId, entry] of entries) {
const processId = retained.get(itemId);
const command = inventory.find(
(item) => item.itemId === itemId && item.processId === processId,
);
if (!command || entry.terminal) {
entries.delete(itemId);
continue;
}
assertCurrent();
entry.processId = command.processId;
entry.task = await createAgentHarnessCommandTask({
scope,
runId: `codex-command:${randomUUID()}`,
taskKind: "codex-command",
command: command.command,
agentId,
startedAt: Date.now(),
assertCurrent,
cancel: async (_reason, assertTaskCurrent) => {
assertCurrent();
assertTaskCurrent();
return (entry.cancellationAttempt ??= (async () => {
const options = {
timeoutMs: requestTimeoutMs,
signal: AbortSignal.timeout(requestTimeoutMs),
};
if (!entry.terminal) {
if (
custody?.ownsCurrentCommand(client, {
threadId: thread.threadId,
turnId,
itemId,
})
) {
entry.cancellation = "pending";
try {
entry.cancellation = (await custody.cancelCommand(client, {
threadId: thread.threadId,
turnId,
itemId,
}))
? "confirmed"
: "idle";
} catch (error) {
entry.cancellation = "idle";
throw error;
} finally {
await settle(itemId);
}
} else {
const live = await client.request(
"thread/backgroundTerminals/list",
{ threadId: thread.threadId },
options,
);
assertCurrent();
assertTaskCurrent();
if (!entry.terminal) {
if (
!live.data.some(
(item) => item.itemId === itemId && item.processId === entry.processId,
)
) {
throw new Error("Native command identity is no longer current");
}
entry.cancellation = "pending";
// The native API controls thread-owned handles, never OS PIDs. It has no expected-item CAS.
try {
const response = await client.request(
"thread/backgroundTerminals/terminate",
{ threadId: thread.threadId, processId: command.processId },
options,
);
entry.cancellation = response.terminated ? "confirmed" : "idle";
if (!response.terminated) {
throw new Error("Native command termination was not confirmed");
}
} catch (error) {
entry.cancellation = "idle";
throw error;
} finally {
await settle(itemId);
}
}
}
}
await settle(itemId);
if (!(await waitForPromiseOrAbort(entry.done.promise, options.signal))) {
options.signal.throwIfAborted();
}
})().finally(() => {
entry.cancellationAttempt = undefined;
}));
},
});
await settle(itemId);
}
} finally {
admitting = false;
for (const [itemId, entry] of entries) {
if (!entry.task) {
entries.delete(itemId);
}
}
releaseIfFinished();
}
},
async closeAdmission() {
admitting = false;
for (const [itemId, entry] of entries) {
if (!entry.task) {
entries.delete(itemId);
}
}
releaseIfFinished();
},
};
}

View file

@ -28,7 +28,7 @@ const assertActive = () => {};
const metadata = { threadId: command.threadId, toolCallId: command.itemId };
describe("native process custody", () => {
it("rechecks foreground permission before a pending spawn without closing background custody", () => {
it("rechecks foreground permission before a pending spawn without closing background custody", async () => {
const client = createClientHarness();
const origin = source();
let active = true;
@ -38,13 +38,29 @@ describe("native process custody", () => {
throw new Error("foreground admission closed");
}
});
const process = getCodexNativeProcessClient(client.client).claim(metadata, async () => {});
const stop = vi.fn(async () => {});
const process = getCodexNativeProcessClient(client.client).claim(metadata, stop);
const sibling = { ...command, itemId: "sibling" };
origin.owner.admit(client.client, sibling, assertActive);
const stopSibling = vi.fn(async () => {});
const siblingProcess = getCodexNativeProcessClient(client.client).claim(
{ ...metadata, toolCallId: sibling.itemId },
stopSibling,
);
try {
active = false;
expect(() => process.assertAdmission()).toThrow("foreground admission closed");
expect(() => process.assertCurrent()).not.toThrow();
expect(await origin.owner.cancelCommand(client.client, { ...command, turnId: "stale" })).toBe(
false,
);
expect(stop).not.toHaveBeenCalled();
expect(await origin.owner.cancelCommand(client.client, command)).toBe(true);
expect(stop).toHaveBeenCalledOnce();
expect(stopSibling).not.toHaveBeenCalled();
} finally {
process.settle();
siblingProcess.settle();
origin.owner.release();
client.client.close();
}

View file

@ -6,7 +6,7 @@ import {
readCodexNotificationThreadId,
readCodexNotificationTurnId,
} from "./notification-correlation.js";
import { isJsonObject } from "./protocol.js";
import { isJsonObject, type CodexAppServerRequestResult } from "./protocol.js";
import { retainSharedCodexAppServerClientIfCurrent } from "./shared-client.js";
type RetainedSource = NonNullable<
@ -57,7 +57,10 @@ export async function readCodexRetainedBackgroundCommands(params: {
assertCurrent: () => void;
signal: AbortSignal;
timeoutMs: number;
}): Promise<() => ReadonlyMap<string, string>> {
}): Promise<{
inventory: CodexAppServerRequestResult<"thread/backgroundTerminals/list">["data"];
readCurrent: () => ReadonlyMap<string, string>;
}> {
params.assertCurrent();
const { data } = await params.client.request(
"thread/backgroundTerminals/list",
@ -68,26 +71,29 @@ export async function readCodexRetainedBackgroundCommands(params: {
params.assertCurrent();
// Consumption follows a second notification drain. Recheck source custody then,
// so revocation during that await cannot turn an orphan into retained work.
return () => {
params.signal.throwIfAborted();
params.assertCurrent();
const retained = new Map<string, string>();
for (const { itemId, processId } of data) {
if (
params.commands.has(itemId) &&
// Approval starts omit the process ID; the native inventory supplies it.
(params.commands.get(itemId) === null || params.commands.get(itemId) === processId) &&
(!params.authority ||
params.authority.ownsCurrentCommand(params.client, {
threadId: params.threadId,
turnId: params.turnId,
itemId,
}))
) {
retained.set(itemId, processId);
return {
inventory: data,
readCurrent: () => {
params.signal.throwIfAborted();
params.assertCurrent();
const retained = new Map<string, string>();
for (const { itemId, processId } of data) {
if (
params.commands.has(itemId) &&
// Approval starts omit the process ID; the native inventory supplies it.
(params.commands.get(itemId) === null || params.commands.get(itemId) === processId) &&
(!params.authority ||
params.authority.ownsCurrentCommand(params.client, {
threadId: params.threadId,
turnId: params.turnId,
itemId,
}))
) {
retained.set(itemId, processId);
}
}
}
return retained;
return retained;
},
};
}
@ -376,6 +382,23 @@ export class CodexNativeProcessAuthority {
);
}
async cancelCommand(client: CodexAppServerClient, receipt: NativeCommand): Promise<boolean> {
this.assertCurrent();
const command = [...this.commands].find(
(candidate) =>
candidate.client === clients.get(client) &&
candidate.threadId === receipt.threadId &&
candidate.turnId === receipt.turnId &&
candidate.itemId === receipt.itemId &&
candidate.processes.size > 0,
);
if (!command) {
return false;
}
await this.terminate([command]);
return true;
}
cancelClient(client: CodexNativeProcessClient): void {
void this.terminate([...this.commands].filter((command) => command.client === client)).catch(
this.onCleanupFailure,

View file

@ -692,7 +692,9 @@ type CodexAppServerRequestParamsOverride = {
};
type CodexAppServerRequestResultMap = {
"thread/backgroundTerminals/list": { data: { itemId: string; processId: string }[] };
"thread/backgroundTerminals/list": {
data: { itemId: string; processId: string; command: string; cwd: string }[];
};
"thread/backgroundTerminals/terminate": { terminated: boolean };
initialize: CodexInitializeResponse;
"account/rateLimits/read": JsonValue;

View file

@ -2,6 +2,7 @@ import { addAbortListener } from "node:events";
import { embeddedAgentLog, formatErrorMessage } from "openclaw/plugin-sdk/agent-harness-runtime";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import { TURN_FINALIZE_DRAIN_ABORT_GRACE_MS } from "./attempt-timeouts.js";
import { prepareCodexNativeCommandTasks } from "./native-command-tasks.js";
import { readCodexRetainedBackgroundCommands } from "./native-process-authority.js";
import type { CodexAttemptActiveTurn } from "./run-attempt-active-turn.js";
import type { CodexAttemptNotificationController } from "./run-attempt-notification-controller.js";
@ -60,8 +61,10 @@ export async function beginCodexAttemptSettlement(
const settlement = drainNotificationQueue().then(async () => {
const commands = activeProjector.getPendingNativeCommands();
if (commands.size > 0 && !params.oneShotCliRun && !state.pluginRuntimeRefreshStop) {
let tasks: ReturnType<typeof prepareCodexNativeCommandTasks>;
try {
const readCurrent = await readCodexRetainedBackgroundCommands({
tasks = prepareCodexNativeCommandTasks(resources, activeTurnId, commands);
const { readCurrent, inventory } = await readCodexRetainedBackgroundCommands({
client: resourceState.client,
threadId: resourceState.thread.threadId,
turnId: activeTurnId,
@ -79,11 +82,15 @@ export async function beginCodexAttemptSettlement(
return new Map();
}
};
await drainNotificationQueue();
await tasks?.retain(inventory, readCurrent());
} catch (error) {
embeddedAgentLog.debug("could not confirm retained native commands", {
threadId: resourceState.thread.threadId,
error: formatErrorMessage(error),
});
} finally {
await tasks?.closeAdmission();
}
// Native exit may arrive while the inventory RPC is pending.
await drainNotificationQueue();

View file

@ -1,5 +1,8 @@
import * as commandTaskRuntime from "openclaw/plugin-sdk/agent-harness-task-runtime";
import { createAgentHarnessTaskRuntime } from "openclaw/plugin-sdk/agent-harness-task-runtime";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import { describe, expect, it } from "vitest";
import { createAdmittedHostCapabilityTestFixture } from "openclaw/plugin-sdk/plugin-test-runtime";
import { describe, expect, it, vi } from "vitest";
import { itemNotification, turnCompleted } from "./protocol.test-helpers.js";
import {
createStartedThreadHarness,
@ -11,6 +14,288 @@ import {
setupRunAttemptTestHooks();
describe("native background command outcomes", () => {
it.each([
"complete",
"cancel",
"natural success",
"natural failure",
"natural failure during stop",
"source retired during stop",
"publication failure",
"refused stop",
"failed stop",
"concurrent stop",
"authority retired during inventory",
"changed handle",
"source retired",
"client closed",
] as const)(
"keeps a retained native command owned after its foreground turn ends (%s)",
async (scenario) => {
const finished = createDeferred<void>();
const terminateStarted = createDeferred<void>();
const terminateRelease = createDeferred<void>();
const releaseSource = vi.fn();
const releaseTask = vi.fn();
let cancellationCurrent = true;
let finishAttempts = 0;
let terminationAttempts = 0;
let stopAttempts: Promise<PromiseSettledResult<void>[]> | undefined;
let cancelOwner: (() => Promise<void>) | undefined;
let commandTaskId: string | undefined;
const createTask = commandTaskRuntime.createAgentHarnessCommandTask;
const admission = vi
.spyOn(commandTaskRuntime, "createAgentHarnessCommandTask")
.mockImplementation(async (input) => {
cancelOwner = () =>
input.cancel("Cancelled by operator.", () => {
if (!cancellationCurrent) {
throw new Error("Task cancellation authority retired");
}
});
const task = await createTask(input);
commandTaskId = task.task.taskId;
return {
...task,
release() {
releaseTask();
task.release();
},
async finish(terminal) {
try {
finishAttempts += 1;
if (scenario === "client closed" || scenario === "publication failure") {
throw new Error("Synthetic terminal publication failure");
}
return await task.finish(terminal);
} finally {
finished.resolve();
}
},
};
});
const params = createTestParams();
const host = await createAdmittedHostCapabilityTestFixture(params);
params.agentHarnessTaskRuntimeScope = host.agentHarnessTaskRuntimeScope;
const source = new AbortController();
params.hostCapabilities = {
...params.hostCapabilities,
retainSourceAuthority: () => ({
modelPolicyRequired: false,
assertCurrent: () => source.signal.throwIfAborted(),
signal: source.signal,
release: releaseSource,
bindModelExecution: () => ({
signal: source.signal,
assertCurrent: () => source.signal.throwIfAborted(),
release() {},
}),
}),
};
const command = {
type: "commandExecution",
id: "visible-command",
command: "python3 synthetic-worker.py",
cwd: "/workspace",
processId: "54527",
status: "inProgress",
commandActions: [],
aggregatedOutput: null,
exitCode: null,
durationMs: null,
};
let retained = false;
const terminal = (exitCode: number) =>
itemNotification("item/completed", {
...command,
status: exitCode === 0 ? "completed" : "failed",
exitCode,
aggregatedOutput: "synthetic outcome",
durationMs: 1,
});
const harness = createStartedThreadHarness(async (method, input) => {
if (method === "thread/backgroundTerminals/list") {
if (retained && scenario === "authority retired during inventory") {
cancellationCurrent = false;
}
return {
data: [
{
itemId:
retained && scenario === "changed handle" ? "successor-command" : command.id,
processId: "54527",
command: command.command,
cwd: "/workspace",
},
],
nextCursor: null,
};
}
if (method === "thread/backgroundTerminals/terminate") {
expect(input).toEqual({ threadId: "thread-1", processId: "54527" });
terminationAttempts += 1;
if (scenario === "failed stop" && terminationAttempts === 1) {
throw new Error("Synthetic transport failure");
}
if (scenario === "concurrent stop" || scenario === "client closed") {
if (terminationAttempts > 1) {
return { terminated: false };
}
terminateStarted.resolve();
await terminateRelease.promise;
}
if (scenario === "refused stop") {
return { terminated: false };
}
if (scenario === "source retired during stop") {
source.abort();
}
await harness.notify(
terminal(
scenario === "natural success"
? 0
: scenario === "natural failure during stop" ||
scenario === "source retired during stop"
? 7
: -1,
),
);
return { terminated: true };
}
return undefined;
});
const run = runCodexAppServerAttempt(params);
try {
await run.waitForTurnAccepted();
await harness.notify(itemNotification("item/started", command));
await harness.notify(turnCompleted({ id: "turn-1", status: "completed", items: [] }));
await run;
if (!host.agentHarnessTaskRuntimeScope) {
throw new Error("Expected an admitted task scope");
}
const tasks = createAgentHarnessTaskRuntime({
runtime: "cli",
scope: host.agentHarnessTaskRuntimeScope,
});
expect(tasks.listTaskRecords()).toContainEqual(
expect.objectContaining({
taskId: commandTaskId,
status: "running",
task: command.command,
ownerKey: params.sessionKey,
}),
);
retained = true;
if (!cancelOwner) {
throw new Error("Expected the registered native cancellation owner");
}
if (scenario === "publication failure") {
const releasedBeforeCompletion = releaseSource.mock.calls.length;
await expect(harness.notify(terminal(0))).rejects.toThrow(
"Synthetic terminal publication failure",
);
expect(releaseTask).toHaveBeenCalledOnce();
expect(releaseSource).toHaveBeenCalledTimes(releasedBeforeCompletion + 1);
expect(finishAttempts).toBe(1);
expect(terminationAttempts).toBe(0);
expect(source.signal.aborted).toBe(false);
// The failed publication has no confirmed durable terminal result.
expect(tasks.listTaskRecords()).toContainEqual(
expect.objectContaining({ taskId: commandTaskId, status: "running" }),
);
return;
}
if (scenario === "client closed") {
const releasedBeforeStop = releaseSource.mock.calls.length;
stopAttempts = Promise.allSettled([cancelOwner()]);
await terminateStarted.promise;
harness.close();
terminateRelease.resolve();
expect(await stopAttempts).toMatchObject([{ status: "rejected" }]);
expect(releaseTask).toHaveBeenCalledOnce();
expect(releaseSource).toHaveBeenCalledTimes(releasedBeforeStop + 1);
expect(tasks.listTaskRecords()).toContainEqual(
expect.objectContaining({ taskId: commandTaskId, status: "running" }),
);
await expect(cancelOwner()).rejects.toThrow();
return;
} else if (scenario === "source retired") {
source.abort();
await finished.promise;
await expect(cancelOwner()).rejects.toThrow();
} else if (scenario === "complete" || scenario === "natural failure") {
await harness.notify(terminal(scenario === "complete" ? 0 : 7));
} else if (scenario === "concurrent stop") {
stopAttempts = Promise.allSettled([cancelOwner(), cancelOwner()]);
await terminateStarted.promise;
terminateRelease.resolve();
const results = await stopAttempts;
expect(terminationAttempts).toBe(1);
expect(results).toMatchObject([{ status: "fulfilled" }, { status: "fulfilled" }]);
} else if (
[
"refused stop",
"failed stop",
"changed handle",
"authority retired during inventory",
].includes(scenario)
) {
await expect(cancelOwner()).rejects.toThrow();
expect(tasks.listTaskRecords()).toContainEqual(
expect.objectContaining({ taskId: commandTaskId, status: "running" }),
);
if (scenario === "failed stop") {
await cancelOwner();
} else {
await harness.notify(terminal(7));
}
} else {
await cancelOwner();
}
const expected =
scenario === "complete" || scenario === "natural success"
? "succeeded"
: scenario === "cancel" || scenario === "concurrent stop" || scenario === "failed stop"
? "cancelled"
: "failed";
expect(tasks.listTaskRecords()).toContainEqual(
expect.objectContaining({
taskId: commandTaskId,
status: expected,
task: command.command,
terminalSummary:
scenario === "source retired"
? "Command outcome unknown"
: {
succeeded: "Command completed",
failed: "Command failed",
cancelled: "Command stopped",
}[expected],
}),
);
if (
scenario === "changed handle" ||
scenario === "source retired" ||
scenario === "authority retired during inventory"
) {
expect(
harness.requests.some(
(request) => request.method === "thread/backgroundTerminals/terminate",
),
).toBe(false);
}
} finally {
terminateRelease.resolve();
await stopAttempts;
source.abort();
harness.close();
host.closeHost();
host.closeAdmission();
await Promise.allSettled([run]);
admission.mockRestore();
}
},
);
it.each([
["retained", "52627"],
["foreign item", "52627"],

View file

@ -178,6 +178,7 @@ describe("background exec task tracking", () => {
timedOut: false,
} satisfies ExecProcessOutcome,
status: "succeeded",
terminalSummary: "Command completed",
error: undefined,
},
{
@ -194,6 +195,7 @@ describe("background exec task tracking", () => {
reason: "secret output\nCommand timed out",
} satisfies ExecProcessOutcome,
status: "timed_out",
terminalSummary: "Command timed out",
error: "Command timed out",
},
{
@ -207,6 +209,7 @@ describe("background exec task tracking", () => {
timedOut: false,
} satisfies ExecProcessOutcome,
status: "failed",
terminalSummary: "Command failed",
error: "Command failed (exit code 17)",
},
{
@ -223,11 +226,12 @@ describe("background exec task tracking", () => {
reason: "secret output\nCommand aborted",
} satisfies ExecProcessOutcome,
status: "cancelled",
terminalSummary: "Command stopped",
error: "Cancelled by operator",
},
])(
"finalizes $label before wake without persisting process output",
async ({ outcome, status, error }) => {
async ({ outcome, status, terminalSummary, error }) => {
const finalize = vi.fn();
await finalizeBackgroundExecTask({
handle: {
@ -242,6 +246,7 @@ describe("background exec task tracking", () => {
expect(finalize).toHaveBeenCalledWith(
expect.objectContaining({
status,
terminalSummary,
...(error ? { error } : { clearError: true }),
}),
);

View file

@ -1,8 +1,9 @@
// Projects detached exec processes into the durable task ledger used by clients.
import { truncateWithMarker } from "@openclaw/normalization-core/utf16-slice";
import { stripAnsi } from "../../packages/terminal-core/src/ansi.js";
import { redactToolPayloadText } from "../logging/redact.js";
import { createSubsystemLogger } from "../logging/subsystem.js";
import {
backgroundCommandTaskContent,
backgroundCommandTaskSummary,
} from "../tasks/background-command-task-content.js";
import { BACKGROUND_EXEC_TASK_KIND } from "../tasks/background-exec-task-contract.js";
import type { DetachedTaskTerminalState } from "../tasks/detached-task-runtime-contract.js";
import { prepareRunningTaskRun } from "../tasks/detached-task-runtime.js";
@ -39,16 +40,7 @@ export function createBackgroundExecTask(params: {
return null;
};
try {
// Redact the complete command before compacting it so truncated secrets cannot escape masking.
const command = stripAnsi(redactToolPayloadText(params.command))
.replace(/\p{Cc}/gu, (control) => ("\r\n\t".includes(control) ? control : ""))
.trim();
const label =
truncateWithMarker(command.replace(/\s+/gu, " "), 120, {
marker: "…",
reserve: 1,
trimEnd: true,
}) || "CLI command";
const content = backgroundCommandTaskContent(params.command);
const prepared = prepareRunningTaskRun(
{
runtime: "cli",
@ -60,8 +52,7 @@ export function createBackgroundExecTask(params: {
agentId: params.agentId,
requesterAgentId: params.agentId,
runId,
label,
task: command || label,
...content,
notifyPolicy: "silent",
deliveryStatus: "not_applicable",
startedAt: params.startedAt,
@ -134,12 +125,7 @@ export function finalizeBackgroundExecTask(params: {
status,
endedAt,
lastEventAt: endedAt,
terminalSummary:
status === "succeeded"
? "Command completed"
: status === "failed"
? "Command failed"
: "Command stopped",
terminalSummary: backgroundCommandTaskSummary(status),
...(status === "succeeded" ? { clearError: true } : { error: execTaskError(params.outcome) }),
detail: {
exitCode: params.outcome.exitCode,

View file

@ -54,12 +54,11 @@ export async function verifyDoctorLintOAuthStateIsolation(
]);
const stdout = vi.spyOn(process.stdout, "write").mockImplementation(() => true);
try {
await expect(
runDoctorLintCli(runtime, {
json: true,
onlyIds: ["core/doctor/runtime-tool-schemas"],
}),
).resolves.toBe(0);
const exitCode = await runDoctorLintCli(runtime, {
json: true,
onlyIds: ["core/doctor/runtime-tool-schemas"],
});
expect(exitCode, String(stdout.mock.calls.at(-1)?.[0])).toBe(0);
const report = JSON.parse(String(stdout.mock.calls.at(-1)?.[0]));
expect(report).toMatchObject({
ok: true,

View file

@ -24,6 +24,7 @@ import {
listTaskRecordPage,
prepareTaskRegistryRead,
} from "../../tasks/runtime-internal.js";
import { withTaskCancellationContext } from "../../tasks/task-cancellation-context.js";
import type { TaskRecord, TaskStatus } from "../../tasks/task-registry.types.js";
import { readGatewayAccessRevision } from "../gateway-access-revision.js";
import { resolveRequestedSessionAgentId } from "../session-request-agent.js";
@ -31,6 +32,7 @@ import {
canAccessTaskRequesterSession,
prepareTaskSessionReadFilter,
} from "../task-session-access.js";
import { readGatewayRequestMutationAuthority } from "./session-mutation-guards.js";
import { taskHistoryHandler } from "./task-history.js";
import { mapTaskSummary } from "./task-summary.js";
import type { GatewayRequestHandler, GatewayRequestHandlers } from "./types.js";
@ -351,30 +353,63 @@ export const tasksHandlers: GatewayRequestHandlers = {
respond(true, { task: mapTaskSummary(task, { includePrompt: true }) });
},
"tasks.history": taskHistoryHandler,
"tasks.cancel": async ({ params, respond, context, client }) => {
"tasks.cancel": async (options) => {
const { params, respond, context, client, sessionMutationAuthorization } = options;
if (!assertValidParams(params, validateTasksCancelParams, "tasks.cancel", respond)) {
return;
}
const authority = readGatewayRequestMutationAuthority(options);
const assertRequestCurrent = () => {
authority.assertCurrent();
sessionMutationAuthorization?.assertCurrent();
};
assertRequestCurrent();
const taskId = params.taskId;
const reason = normalizeOptionalString(params.reason);
const { cancelDetachedTaskRunByIdCore } =
await import("../../tasks/task-executor-cancel.runtime.js");
assertRequestCurrent();
const cfg = context.getRuntimeConfig();
const task = getTaskById(taskId);
if (task && !canAccessTaskRequesterSession({ access: "write", cfg, client, task })) {
respond(true, { found: false, cancelled: false });
return;
}
const result = await cancelDetachedTaskRunByIdCore({
cfg,
taskId,
...(reason ? { reason } : {}),
});
const cancel = () =>
cancelDetachedTaskRunByIdCore({ cfg, taskId, ...(reason ? { reason } : {}) });
const result = task
? await withTaskCancellationContext(
() => {
assertRequestCurrent();
const current = getTaskById(taskId);
if (
!current ||
current.requesterSessionKey !== task.requesterSessionKey ||
current.requesterAgentId !== task.requesterAgentId ||
!canAccessTaskRequesterSession({
access: "write",
cfg: context.getRuntimeConfig(),
client,
task: current,
})
) {
throw new Error("Task cancellation authority changed.");
}
},
cancel,
{ selectedTask: task },
)
: await cancel();
assertRequestCurrent();
const responseTask = result.task && getTaskById(result.task.taskId);
respond(true, {
found: result.found,
cancelled: result.cancelled,
...(result.reason ? { reason: result.reason } : {}),
...(result.task ? { task: mapTaskSummary(result.task) } : {}),
...(responseTask &&
canAccessTaskRequesterSession({ cfg: context.getRuntimeConfig(), client, task: responseTask })
? { task: mapTaskSummary(responseTask) }
: {}),
});
},
"tasks.retry": createTaskRecoveryHandler("tasks.retry"),

View file

@ -1,20 +1,36 @@
import { expectDefined } from "@openclaw/normalization-core";
import { expect, it } from "vitest";
import { upsertSessionEntryCore } from "../config/sessions/session-accessor.js";
import type { GatewayOperatorRoleDefinition } from "../config/types.gateway.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { identifiedClient, runTaskHandler } from "../gateway/server-methods/tasks.test-helpers.js";
import { createPluginMetadataSnapshotFixture } from "../plugins/plugin-metadata.test-support.js";
import { createEmptyPluginRegistry } from "../plugins/registry-empty.js";
import { markPluginRegistryRetired } from "../plugins/registry-lifecycle.js";
import {
getPluginRuntimeGatewayRequestScope,
withPluginRuntimeGatewayRequestScope,
} from "../plugins/runtime/gateway-request-scope.js";
import { withPluginRuntimeGenerationScope } from "../plugins/runtime/generation-scope.js";
import { ensureProfileForEmail } from "../state/user-profiles.js";
import { createAgentHarnessTaskRuntimeScope } from "../tasks/agent-harness-task-runtime-scope.js";
import { getTaskFlowRegistryStore } from "../tasks/task-flow-registry.store.js";
import {
getTaskActivitySnapshot,
recordTaskActivityEvent,
} from "../tasks/task-registry-activity.js";
import { updateTask } from "../tasks/task-registry-mutation.js";
import { getTaskById } from "../tasks/task-registry.js";
import { getTaskRegistryStore } from "../tasks/task-registry.store.js";
import {
resetTaskRegistryForTests,
withTaskRegistryTempDir,
} from "../tasks/task-registry.test-support.js";
import { getTaskRunOwner } from "../tasks/task-run-owner.js";
import {
captureAgentHarnessTaskAssignment,
createAgentHarnessTaskRuntime,
createAgentHarnessCommandTask,
} from "./agent-harness-task-runtime.js";
it.each([false, true])(
@ -91,3 +107,203 @@ it.each([false, true])(
);
},
);
it.each([
"cancelled",
"succeeded",
"replacement",
"same-scope replacement",
"replacement before stop",
"revoked",
"incognito",
"incognito failed",
"registry retired",
"caller revoked",
"requester reassigned",
"session write revoked",
"session read revoked after stop",
] as const)("binds command cancellation to its original task outcome (%s)", async (scenario) => {
await withTaskRegistryTempDir(
async () => {
let current = true;
let stops = 0;
const agentRegistry = createEmptyPluginRegistry();
const requestSignal = new AbortController().signal;
const sessionPolicyCase =
scenario === "session write revoked" || scenario === "session read revoked after stop";
const role: GatewayOperatorRoleDefinition = {
sessions: { others: "write" },
agents: "*",
scopes: ["operator.read", "operator.write"],
};
const config: OpenClawConfig = sessionPolicyCase
? {
gateway: {
roles: { default: "task-operator", definitions: { "task-operator": role } },
},
}
: {};
const requestClient = sessionPolicyCase
? identifiedClient(role.scopes, ensureProfileForEmail("viewer@example.com").id)
: identifiedClient(["operator.admin"]);
let replacement: ReturnType<typeof updateTask> | undefined;
const incognito = scenario === "incognito" || scenario === "incognito failed";
const ownerKey = incognito ? "agent:main:dashboard:incognito-command" : "agent:main:command";
if (sessionPolicyCase) {
await upsertSessionEntryCore(
{ agentId: "main", sessionKey: ownerKey },
{
sessionId: "command-session",
updatedAt: 1,
createdActor: { type: "human", source: "profile", id: "owner@example.com" },
visibility: "shared",
},
);
}
const command = await withPluginRuntimeGenerationScope(
{ metadataSnapshot: createPluginMetadataSnapshotFixture(), pluginRegistry: agentRegistry },
() =>
createAgentHarnessCommandTask({
scope: createAgentHarnessTaskRuntimeScope({ requesterSessionKey: ownerKey }),
runId: "native-command:original",
taskKind: "test-native-command",
command: "SYNTHETIC_TASK_CONTENT",
startedAt: Date.now(),
assertCurrent() {
if (!current) {
throw new Error("source retired");
}
},
async cancel(_reason, assertTaskCurrent) {
expect(getPluginRuntimeGatewayRequestScope()?.signal).toBe(requestSignal);
expect(getPluginRuntimeGatewayRequestScope()?.client).toBe(requestClient);
if (
scenario === "caller revoked" ||
scenario === "requester reassigned" ||
scenario === "session write revoked"
) {
await Promise.resolve();
if (scenario === "caller revoked") {
requestClient.invalidated = true;
} else if (scenario === "requester reassigned") {
updateTask(command.task.taskId, {
requesterSessionKey: "agent:main:another-session",
});
} else {
role.sessions = { others: "view" };
}
assertTaskCurrent();
}
stops += 1;
if (scenario === "replacement") {
replacement = updateTask(command.task.taskId, {
runId: "native-command:successor",
status: "cancelled",
});
} else if (scenario === "same-scope replacement") {
replacement = updateTask(command.task.taskId, { taskKind: "successor-command" });
assertTaskCurrent();
} else {
await command.finish({
status:
scenario === "succeeded"
? "succeeded"
: scenario === "incognito failed"
? "failed"
: "cancelled",
endedAt: Date.now(),
...(scenario === "incognito failed"
? {
terminalSummary: "SYNTHETIC_TASK_CONTENT",
error: "SYNTHETIC_TASK_CONTENT",
}
: {}),
});
if (scenario === "session read revoked after stop") {
role.sessions = { others: "none" };
}
}
},
}),
);
try {
if (scenario === "registry retired") {
markPluginRegistryRetired(agentRegistry);
}
if (scenario === "revoked") {
current = false;
}
if (scenario === "replacement before stop") {
replacement = updateTask(command.task.taskId, { taskKind: "successor-command" });
}
const request = withPluginRuntimeGatewayRequestScope(
{
pluginRegistry: createEmptyPluginRegistry(),
client: requestClient,
signal: requestSignal,
isWebchatConnect: () => true,
},
() =>
runTaskHandler("tasks.cancel", { taskId: command.task.taskId }, config, requestClient),
);
if (scenario === "caller revoked") {
await expect(request).rejects.toThrow("Gateway requester authority changed");
} else {
const result = await request;
expect(result.payload?.cancelled).toBe(
scenario === "cancelled" ||
scenario === "incognito" ||
scenario === "session read revoked after stop",
);
if (scenario === "session read revoked after stop") {
expect(result.payload).toMatchObject({ found: true, cancelled: true });
expect(result.payload?.task).toBeUndefined();
}
}
expect(stops).toBe(
scenario === "revoked" ||
scenario === "replacement before stop" ||
scenario === "registry retired" ||
scenario === "caller revoked" ||
scenario === "requester reassigned" ||
scenario === "session write revoked"
? 0
: 1,
);
if (incognito) {
const task = getTaskById(command.task.taskId);
expect(task).toMatchObject({
status: scenario === "incognito failed" ? "failed" : "cancelled",
terminalSummary: scenario === "incognito failed" ? "Command failed" : "Command stopped",
...(scenario === "incognito failed" ? { error: "Incognito task error." } : {}),
});
expect(JSON.stringify(task)).not.toContain("SYNTHETIC_TASK_CONTENT");
}
if (scenario === "replacement") {
expect(replacement).toMatchObject({
runId: "native-command:successor",
status: "cancelled",
});
await expect(command.finish({ status: "failed", endedAt: Date.now() })).resolves.toBe(
"retired",
);
expect(getTaskById(command.task.taskId)?.runId).toBe("native-command:successor");
}
if (scenario === "same-scope replacement" || scenario === "replacement before stop") {
expect(replacement).toMatchObject({ taskKind: "successor-command", status: "running" });
const outcome = await command
.finish({ status: "failed", endedAt: Date.now() })
.catch(() => "rejected");
const successor = expectDefined(getTaskById(command.task.taskId), "replacement task");
expect(successor).toMatchObject({ taskKind: "successor-command", status: "running" });
expect(outcome).toBe("retired");
expect(getTaskRunOwner(successor)).toBeUndefined();
}
} finally {
command.release();
markPluginRegistryRetired(agentRegistry);
}
},
{ durableStore: true },
);
});

View file

@ -56,6 +56,8 @@ import {
} from "../tasks/task-registry-records.js";
import type { TaskPersistenceReceipt, TaskRunTransition } from "../tasks/task-registry.types.js";
export { createAgentHarnessCommandTask } from "../tasks/agent-harness-command-task.js";
export { createAgentHarnessTaskEventSink } from "../tasks/agent-harness-completion-custody.js";
export type { AgentHarnessCompletionCustody };
export {

View file

@ -0,0 +1,181 @@
import { err, ok, type Result } from "@openclaw/normalization-core/result";
import { requireActivePluginRegistry } from "../plugins/runtime.js";
import { withPluginRuntimeRegistryScope } from "../plugins/runtime/gateway-request-scope.js";
import { isIncognitoSessionKey } from "../shared/incognito-session-key.js";
import {
assertAgentHarnessTaskRuntimeScope,
type AgentHarnessTaskRuntimeScope,
} from "./agent-harness-task-runtime-scope.js";
import {
backgroundCommandTaskContent,
backgroundCommandTaskSummary,
} from "./background-command-task-content.js";
import {
DetachedTaskAssignmentUnsupportedError,
type DetachedTaskTerminalState,
} from "./detached-task-runtime-contract.js";
import { captureDetachedTaskRuntimeOwner } from "./detached-task-runtime-state.js";
import { prepareRunningTaskRun } from "./detached-task-runtime.js";
import { captureTaskCancellationControl } from "./task-cancellation-context.js";
import { getTaskById } from "./task-registry-query.js";
import {
captureTaskPersistenceReceipt,
matchesTaskPersistenceReceipt,
} from "./task-registry-records.js";
import type { TaskPersistenceReceipt, TaskRecord } from "./task-registry.types.js";
import { getTaskRunOwner } from "./task-run-owner.js";
import type { TaskRunOwnerBinding } from "./task-run-owner.types.js";
/** Bind a harness-owned command to the existing ledger without taking over its process. */
export async function createAgentHarnessCommandTask(params: {
scope: AgentHarnessTaskRuntimeScope;
runId: string;
taskKind: string;
command: string;
agentId?: string;
startedAt: number;
assertCurrent: () => void;
/** Recheck current task authority after awaits and immediately before stopping native work. */
cancel: (reason: string, assertTaskCurrent: () => void) => Promise<void>;
}) {
const scope = assertAgentHarnessTaskRuntimeScope(params.scope);
const registry = requireActivePluginRegistry();
const runtime = captureDetachedTaskRuntimeOwner();
// Legacy custom runtimes cannot bind a receipt-owned cancellation capability.
if (runtime.runtime) {
throw new DetachedTaskAssignmentUnsupportedError();
}
const assertCurrent = () => {
runtime.assertCurrent();
params.assertCurrent();
};
assertCurrent();
const incognito = isIncognitoSessionKey(scope.requesterSessionKey);
const prepared = prepareRunningTaskRun(
{
runtime: "cli",
taskKind: params.taskKind,
requesterSessionKey: scope.requesterSessionKey,
ownerKey: scope.requesterSessionKey,
scopeKind: "session",
agentId: params.agentId,
requesterAgentId: params.agentId,
runId: params.runId,
...backgroundCommandTaskContent(incognito ? "Incognito task" : params.command),
notifyPolicy: "silent",
deliveryStatus: "not_applicable",
startedAt: params.startedAt,
lastEventAt: params.startedAt,
},
assertCurrent,
);
if (prepared.kind !== "receipt") {
throw new DetachedTaskAssignmentUnsupportedError();
}
const receipt = await prepared.create();
if (!receipt) {
throw new Error("Native command task persistence failed");
}
let binding: TaskRunOwnerBinding | undefined;
let expectedTask: TaskPersistenceReceipt | undefined;
let settlement: Promise<"published" | "retired"> | undefined;
const originalTask = () => {
const current = getTaskById(receipt.task.taskId);
if (
!expectedTask ||
!current ||
!matchesTaskPersistenceReceipt(current, expectedTask) ||
!binding ||
getTaskRunOwner(current) !== binding.owner
) {
binding?.release();
return undefined;
}
return current;
};
const finish = (terminal: DetachedTaskTerminalState) =>
(settlement ??= (async (): Promise<"published" | "retired"> => {
if (!originalTask()) {
return "retired";
}
await receipt.finalizeActive(
{
...terminal,
...(incognito
? {
terminalSummary: backgroundCommandTaskSummary(terminal.status),
...(terminal.error ? { error: "Incognito task error." } : {}),
}
: {}),
},
(task) =>
Boolean(
expectedTask &&
matchesTaskPersistenceReceipt(task, expectedTask) &&
binding &&
getTaskRunOwner(task) === binding.owner,
),
);
const current = originalTask();
if (!current) {
return "retired";
}
if (current.status !== terminal.status) {
throw new Error(
"Native command terminal publication did not settle for its original task.",
);
}
binding?.release();
return "published";
})().catch((error: unknown) => {
settlement = undefined;
throw error;
}));
try {
binding = await receipt.bindRunOwner(
(reason) =>
// Keep the admitting registry without reviving its old Gateway request authority.
withPluginRuntimeRegistryScope(registry, async (): Promise<Result<TaskRecord, string>> => {
try {
const control = captureTaskCancellationControl();
const assertTaskCurrent = () => {
assertCurrent();
if (!originalTask()) {
throw new Error("Native command no longer belongs to its original task.");
}
control?.assertCurrent();
};
assertTaskCurrent();
await params.cancel(reason, assertTaskCurrent);
const current = getTaskById(receipt.task.taskId);
return expectedTask &&
current &&
matchesTaskPersistenceReceipt(current, expectedTask) &&
current.status === "cancelled"
? ok(current)
: err("Native command did not settle as cancelled.");
} catch {
return err("Native command could not be cancelled by its current owner.");
}
}),
assertCurrent,
);
const current = getTaskById(receipt.task.taskId);
if (!current || getTaskRunOwner(current) !== binding.owner) {
throw new Error("Native command task owner was replaced during admission.");
}
expectedTask = captureTaskPersistenceReceipt(current);
return { task: receipt.task, finish, release: () => binding?.release() };
} catch (error) {
binding?.release();
await receipt.settleUnstarted(
{
status: "failed",
endedAt: Date.now(),
error: "Native command task admission ended.",
},
(task) => !getTaskRunOwner(task),
);
throw error;
}
}

View file

@ -0,0 +1,27 @@
import { truncateWithMarker } from "@openclaw/normalization-core/utf16-slice";
import { stripAnsi } from "../../packages/terminal-core/src/ansi.js";
import { redactToolPayloadText } from "../logging/redact.js";
import type { DetachedTaskTerminalState } from "./detached-task-runtime-contract.js";
export function backgroundCommandTaskSummary(status: DetachedTaskTerminalState["status"]) {
return {
succeeded: "Command completed",
failed: "Command failed",
cancelled: "Command stopped",
timed_out: "Command timed out",
}[status];
}
/** Redact before truncation so a shortened secret cannot escape masking. */
export function backgroundCommandTaskContent(input: string) {
const command = stripAnsi(redactToolPayloadText(input))
.replace(/\p{Cc}/gu, (control) => ("\r\n\t".includes(control) ? control : ""))
.trim();
const label =
truncateWithMarker(command.replace(/\s+/gu, " "), 120, {
marker: "…",
reserve: 1,
trimEnd: true,
}) || "CLI command";
return { label, task: command || label };
}

View file

@ -9,6 +9,7 @@ import { readTaskBackingInstance } from "./task-backing-records.js";
import { getTaskActivitySnapshot } from "./task-registry-activity.js";
import { resolveTaskAgentId } from "./task-registry-records.js";
import { isTerminalTaskStatus, type TaskRecord } from "./task-registry.types.js";
import { getTaskRunOwner } from "./task-run-owner.js";
import { sanitizeTaskStatusText, TASK_STATUS_DETAIL_MAX_CHARS } from "./task-status.js";
function sanitizeOptionalTaskText(value: unknown): string | undefined {
@ -19,6 +20,9 @@ function observeCliExecution(task: TaskRecord): "queued" | "running" | undefined
if (task.runtime !== "cli" || !task.runId) {
return undefined;
}
if (getTaskRunOwner(task)) {
return "running";
}
const context = getAgentRunContext(task.runId);
const sessionKey =
task.childSessionKey ?? (task.scopeKind === "session" ? task.ownerKey : undefined);

View file

@ -4,6 +4,7 @@ import type {
readSessionBackingFactsInWorker,
SessionBackingFact,
} from "../config/sessions/session-accessor.js";
import { getAgentRunContext } from "../infra/agent-run-registry.js";
import type { parseAgentSessionKey } from "../routing/session-key.js";
import {
deriveSessionChatTypeFromKey,
@ -11,6 +12,7 @@ import {
} from "../sessions/session-chat-type-shared.js";
import { sessionChanges } from "../sessions/session-row-changes.js";
import type { TaskRecord } from "./task-registry.types.js";
import { getTaskRunOwner } from "./task-run-owner.js";
export type BackingSessionRuntime = {
readSessionBackingFacts: typeof readSessionBackingFacts;
@ -134,3 +136,21 @@ export function findTaskSessionEntry(
// An unprepared or concurrently changed key is unknown, never evidence of death.
return entries.get(target.sessionKey);
}
export function hasActiveCliRun(task: TaskRecord): boolean {
if (getTaskRunOwner(task)) {
return true;
}
const candidateRunIds = [task.sourceId, task.runId];
for (const candidate of candidateRunIds) {
const runId = candidate?.trim();
if (runId && getAgentRunContext(runId)) {
return true;
}
}
return false;
}
export function hasCliRunIdentity(task: TaskRecord): boolean {
return [task.sourceId, task.runId].some((candidate) => Boolean(candidate?.trim()));
}

View file

@ -20,6 +20,7 @@ import {
} from "../test-utils/openclaw-test-state.js";
import { holdStateDatabaseCoordinator } from "../test-utils/state-database-contention.js";
import { getDetachedTaskLifecycleRuntime } from "./detached-task-runtime.js";
import { getTaskExecutionObservation } from "./task-execution-observation.js";
import { createRunningTaskRunCoreWithReceiptAsync } from "./task-executor-create.async.js";
import { readResidentTaskFlow } from "./task-flow-registry.js";
import { loadTaskFlowRegistryStateFromSqliteReadOnly } from "./task-flow-registry.store.sqlite.js";
@ -41,6 +42,7 @@ import {
createTaskFixture,
reloadTaskRegistryFromStoreAsync,
} from "./task-registry.test-support.js";
import { bindTaskRunOwner } from "./task-run-owner.js";
import {
resetDetachedTaskLifecycleRuntimeForTests,
resetTaskFlowRegistryForTests,
@ -72,6 +74,44 @@ afterEach(async () => {
});
describe("task maintenance session metadata", () => {
it("retains a CLI task until its live run owner releases it", async () => {
await withMaintenanceState("openclaw-task-maintenance-run-owner-", async () => {
resetTaskRegistryForTests({ persist: false });
configureTaskRegistryMaintenance({ runtimeAuthoritative: true });
const task = createTaskFixture("cli", {
runId: "retained-native-command",
task: "Background command after its foreground turn",
notifyPolicy: "silent",
lastEventAt: Date.now() - 40 * 60_000,
});
const release = bindTaskRunOwner(task, async () => ({
ok: false,
error: "No cancellation requested in this scenario.",
}));
try {
expect(reconcileInspectableTasks()).toContainEqual(
expect.objectContaining({ taskId: task.taskId, status: "running" }),
);
expect(getTaskExecutionObservation(task)).toEqual({ state: "running" });
expect(getTaskExecutionObservation({ ...task, runId: "replacement-command" })).toEqual({
state: "unknown",
});
expect((await runTaskRegistryMaintenance()).reconciled).toBe(0);
expect(getTaskById(task.taskId)?.status).toBe("running");
release();
expect(getTaskExecutionObservation(task)).toEqual({ state: "unknown" });
expect((await runTaskRegistryMaintenance()).reconciled).toBe(1);
expect(getTaskById(task.taskId)).toMatchObject({
status: "lost",
error: "backing session missing",
});
} finally {
release();
}
});
});
it.each(["publication", "coordinator hold"] as const)(
"retains task payloads without synchronous refreshes during %s",
async (boundary) => {

View file

@ -69,6 +69,8 @@ import { createTaskMaintenanceScheduler } from "./task-registry-maintenance-sche
import {
createBackingSessionLookupContext,
findTaskSessionEntry,
hasActiveCliRun,
hasCliRunIdentity,
prepareBackingSessionFacts,
observeBackingSessionFacts,
resolveSessionChatType,
@ -273,21 +275,6 @@ function resolveDurableCronTaskRecovery(
};
}
function hasActiveCliRun(task: TaskRecord): boolean {
const candidateRunIds = [task.sourceId, task.runId];
for (const candidate of candidateRunIds) {
const runId = candidate?.trim();
if (runId && getAgentRunContext(runId)) {
return true;
}
}
return false;
}
function hasCliRunIdentity(task: TaskRecord): boolean {
return [task.sourceId, task.runId].some((candidate) => Boolean(candidate?.trim()));
}
function hasBackingSession(task: TaskRecord, context: BackingSessionLookupContext): boolean {
const hasProcessLocalLiveness =
task.runtime === "cron" || task.runtime === "cli" || task.runtime === "acp";

View file

@ -637,6 +637,7 @@ export const databaseWorkerCoreTestFiles = [
"src/logging/diagnostic-session-context.test.ts",
"src/logging/diagnostic-stuck-session-recovery.runtime.test.ts",
"src/memory/memory-artifact-provenance.test.ts",
"src/plugin-sdk/agent-harness-task-runtime.persistence.test.ts",
"src/plugin-sdk/memory-host-core.test.ts",
"src/plugin-sdk/memory-host-event-export.test.ts",
"src/plugin-sdk/memory-host-events.test.ts",

View file

@ -205,8 +205,8 @@ describe("chat pane embedded panels", () => {
},
);
it.each(["refusal", "rejection"] as const)(
"keeps a cancellation %s visible in task detail and after Back to the list",
it.each(["refusal", "rejection", "late rejection"] as const)(
"keeps a cancellation %s scoped to its task detail and visible in the list",
async (outcome) => {
const { mount, rails, renderPanels, state, task } = createReviewFixture({
status: "running",
@ -219,11 +219,36 @@ describe("chat pane embedded panels", () => {
rails().backgroundTasks.onOpenTaskDetail?.(task);
await renderPanels();
await renderPanels();
const completedTask = {
...task,
id: "completed-task",
taskId: "completed-task",
status: "completed" as const,
title: "Another completed task",
};
state.backgroundTasksState!.tasks!.push(completedTask);
state.backgroundTasksState!.taskDetails.set(completedTask.id, completedTask);
const openCompletedTask = async () => {
rails().backgroundTasks.onOpenTaskList?.();
await renderPanels();
rails().backgroundTasks.onToggleFinished();
await renderPanels();
[...mount.querySelectorAll<HTMLButtonElement>(".chat-tasks-rail__task-open")]
.find((button) => button.textContent?.includes(completedTask.title))!
.click();
await renderPanels();
expect(mount.querySelector("[data-task-detail-panel] .sidebar-title")?.textContent).toBe(
completedTask.title,
);
};
const message =
outcome === "refusal" ? "Task cannot be cancelled" : "Cancellation unavailable";
const pending = createDeferred<unknown>();
const request = vi.spyOn(state.client!, "request");
if (outcome === "refusal") {
request.mockResolvedValueOnce({ found: true, cancelled: false, reason: message });
} else if (outcome === "late rejection") {
request.mockReturnValueOnce(pending.promise);
} else {
request.mockRejectedValueOnce(new Error(message));
}
@ -233,7 +258,22 @@ describe("chat pane embedded panels", () => {
)!
.click();
expect(request).toHaveBeenCalledExactlyOnceWith("tasks.cancel", { taskId: task.id });
await vi.waitFor(() => expect(state.backgroundTasksState?.error).toBe(message));
if (outcome === "late rejection") {
await openCompletedTask();
pending.reject(new Error(message));
}
await request.mock.results[0]!.value.catch(() => undefined);
expect(state.backgroundTasksState?.error).toBe(message);
if (outcome !== "late rejection") {
await renderPanels();
expect(
mount.querySelector('[data-task-detail-panel] [role="alert"]')?.textContent,
).toContain(message);
await openCompletedTask();
}
await renderPanels();
expect.soft(mount.querySelector('[data-task-detail-panel] [role="alert"]')).toBeNull();
rails().backgroundTasks.onOpenTaskDetail?.(task);
await renderPanels();
expect
.soft(mount.querySelector('[data-task-detail-panel] [role="alert"]')?.textContent)
@ -248,6 +288,28 @@ describe("chat pane embedded panels", () => {
);
expect(state.backgroundTasksState?.error).toBe(message);
expect(request).toHaveBeenCalledOnce();
if (outcome === "refusal") {
const refreshed = createDeferred();
const requestUpdate = state.requestUpdate;
vi.spyOn(state, "requestUpdate").mockImplementation(() => {
requestUpdate?.();
if (!state.backgroundTasksState?.loading) {
refreshed.resolve();
}
});
request.mockRejectedValue(new Error("Task list unavailable"));
rails().backgroundTasks.onRefresh();
await refreshed.promise;
await renderPanels();
expect(mount.querySelector('.chat-tasks-rail [role="alert"]')?.textContent).toContain(
"Task list unavailable",
);
rails().backgroundTasks.onOpenTaskDetail?.(completedTask);
await renderPanels();
expect(
mount.querySelector('[data-task-detail-panel] [role="alert"]')?.textContent,
).toContain("Task list unavailable");
}
},
);

View file

@ -60,6 +60,7 @@ type BackgroundTasksState = BackgroundTaskObservations & {
connectionClient: GatewayBrowserClient | null;
connectionEpoch: number | undefined;
error: string | null;
errorTaskId?: string;
explicitReadRequested: boolean;
deferredRetryAttempt?: number;
finishedCollapsed: boolean;
@ -244,6 +245,7 @@ function loadBackgroundTasks(
state.pendingTaskEvents = buffer;
state.loading = true;
state.error = null;
delete state.errorTaskId;
state.pendingReload = false;
host.requestUpdate?.();
void (async () => {
@ -283,6 +285,7 @@ function loadBackgroundTasks(
current.tasks = replayTaskEvents([], buffer.events);
}
current.error = formatUiError(error, t("tasksPage.loadFailed"));
delete current.errorTaskId;
}
} finally {
const current = getBackgroundTasksState(host);
@ -299,6 +302,7 @@ function loadBackgroundTasks(
current.loadedClient = null;
// The superseded read's error must not block its queued replacement on resume.
current.error = null;
delete current.errorTaskId;
}
}
host.requestUpdate?.();
@ -376,6 +380,7 @@ export function handleBackgroundTasksEvent(
state.tasks = null;
state.loadedClient = null;
state.error = null;
delete state.errorTaskId;
host.requestUpdate?.();
return;
}
@ -508,6 +513,7 @@ async function cancelBackgroundTask(
}
state.cancellingTaskIds = new Set([...state.cancellingTaskIds, taskId]);
state.error = null;
delete state.errorTaskId;
host.requestUpdate?.();
try {
const payload = await client.request("tasks.cancel", { taskId });
@ -530,10 +536,12 @@ async function cancelBackgroundTask(
if (!result?.cancelled) {
const reason = result?.reason;
state.error = reason ? formatUiError(reason) : t("tasksPage.cancelFailed");
state.errorTaskId = taskId;
}
} catch (error) {
if (getBackgroundTasksState(host) === state) {
state.error = formatUiError(error, t("tasksPage.cancelFailed"));
state.errorTaskId = taskId;
}
} finally {
if (getBackgroundTasksState(host) === state) {
@ -611,7 +619,10 @@ export function createBackgroundTasksProps(
// tasks.cancel needs operator.write; read-only operators get no button.
canCancel: host.connected && hasOperatorWriteAccess(host.hello?.auth ?? null),
loading: state.loading,
error: state.error,
error:
state.errorTaskId && opts.selectedTaskId && state.errorTaskId !== opts.selectedTaskId
? null
: state.error,
tasks: state.tasks,
activeCount: state.tasks?.filter(isActiveTask).length ?? 0,
subagentActivity,