fix: distinguish active subagents from retained task state (#152514)

Related: #101656

## What Problem This Solves

Subagent status could report active or queued work when only retained state remained, or borrow activity from a replacement execution.

## User Impact

Tasks, `/status`, and subagent controls distinguish current execution, an actual queue reservation, child or approval waits, and unavailable activity. Completed work and pending result delivery remain separate. Quiet active work stays active; a saved row alone no longer proves liveness.

## Why This Change Was Made

Consumers use the existing execution owners instead of inferring liveness from stored task status. Prepared observations retain exact task and generation identity. Integrating current main exposed an early stored-queued shortcut; subagents now consult their actual reservation owner while terminal tasks keep the existing fast path.

The status-transcript prerequisite is already merged in #152512. This PR adds no schema, storage, notification default, or completion-delivery owner.

## Evidence

- Registered Gateway Tasks reads, status consumers, task progress, and subagent observation checks passed. The reservation-removal and replacement cases failed before correcting the current-main integration, then passed with the owner-level fix.
- All 15 focused task and native-owner cases passed after that correction, preserving terminal fast paths and quiet live execution.
- A real isolated Gateway answered `/status` over WebSocket while parent and child model requests were held. Status history remained visible, the original run continued, and status did not become new model input. Only external model responses were mocked.
- Crash recovery was exercised through real Gateway processes and `tasks.get`, without inserting or mutating registry rows. After a hard crash, unchanged baseline observation code reported the retained child as running before recovery; this candidate reported unknown, then running after execution was reacquired. Both runs first spawned real held-provider children, and both fixtures cleaned up.
- The supported runtime build, core typecheck, focused type-aware lint, and whitespace checks passed.

This does not complete #101656. Retained-card presentation remains separate in #152517, and completion accounting in #150933 remains with its maintainer.

Co-authored-by: Ayaan Zaidi <hi@obviy.us>
This commit is contained in:
Ayaan Zaidi 2026-09-19 13:32:38 +05:30 • committed by GitHub
parent 36ac37eb41
commit a877d3612c
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
16 changed files with 951 additions and 112 deletions

View file

@ -141,7 +141,7 @@ stateDiagram-v2
| Status | What it means |
| ----------- | --------------------------------------------------------------------------- |
| `queued` | Created, waiting for the agent to start |
| `running` | Agent turn is actively executing |
| `running` | Started, with no terminal lifecycle outcome recorded |
| `succeeded` | Completed successfully |
| `failed` | Completed with an error |
| `timed_out` | Exceeded the configured timeout |
@ -150,6 +150,12 @@ stateDiagram-v2
Transitions happen automatically - agent run lifecycle events (start, end, error) update the task status; you do not manage it manually.
Stored task status and live execution observation are separate. A task can retain
`running` status while its execution observation is `waiting` or `unknown`.
`unknown` means current execution cannot be observed; it does not mean the task
failed or stopped. A retained subagent registration alone does not prove that
its execution is running.
Execution and result delivery are separate. A subagent task can remain
`succeeded` while its `deliveryStatus` is `session_queued` or `failed`. The
terminal outcome is `succeeded` after delivery and `blocked` when the work

View file

@ -24,6 +24,14 @@ by default). Use `sessions_history` for a bounded, safety-filtered recall
view from within an agent turn, or inspect the transcript path on disk for
the raw full transcript.
`/status` keeps full sub-agent counts but shows at most three current detail rows.
Each row distinguishes running (with a safe tool name when available), queued,
waiting for approval, input, children, agent messages, or external work,
and finished execution with settlement still pending. Pending children may have
finished but still owe completion delivery; they are not necessarily executing.
“Current activity unavailable” means no current execution is observable, not that
the task failed. These observations do not change retained active/done counts.
In the Control UI, subagent runs appear in inline transcript activity rows, the
chat **Tasks** tab, and the [Tasks page](/automation/tasks#control-ui). They do not
appear as sidebar rows or add an expand control to their parent. The parent's

View file

@ -1,10 +1,14 @@
/** Controller identity, authorization, and controlled-run read scope. */
import type { TaskSummary } from "../../../../packages/gateway-protocol/src/schema/tasks.js";
import type { OpenClawConfig } from "../../../config/types.openclaw.js";
import {
isSubagentSessionKey,
normalizeAgentId,
parseAgentSessionKey,
} from "../../../routing/session-key.js";
import { readTaskBackingInstance } from "../../../tasks/task-backing-records.js";
import { getTaskExecutionObservation } from "../../../tasks/task-execution-observation.js";
import { findTaskByRunId } from "../../../tasks/task-registry-query.js";
import { resolveSessionAgentId } from "../../agent-scope.js";
import { resolveSubagentRequesterAgentId } from "../../subagent-requester-owner.js";
import {
@ -12,7 +16,8 @@ import {
resolveMainSessionAlias,
} from "../../tools/sessions-helpers.js";
import { resolveStoredSubagentCapabilities } from "../spawn/subagent-capabilities.js";
import { subagentRuns } from "./subagent-registry-memory.js";
import { observeSubagentExecution } from "./subagent-execution-observation.js";
import { getSubagentRunsForRequesterSession, subagentRuns } from "./subagent-registry-memory.js";
import { buildSubagentRunReadIndexFromRuns } from "./subagent-registry-queries.js";
import { getLatestLiveSubagentRunByChildSessionKey } from "./subagent-registry-read.js";
import { getSubagentRunsSnapshotForRead } from "./subagent-registry-state.js";
@ -107,21 +112,25 @@ export function isSubagentRunVisibleToSession(
);
}
export type ControlledSubagentRunsReadContext = {
runs: SubagentRunRecord[];
countPendingDescendantRuns(rootSessionKey: string): number;
getExecutionObservation(entry: SubagentRunRecord): NonNullable<TaskSummary["execution"]>;
};
/** Builds one stable snapshot for controlled-run listing and descendant status reads. */
export function buildControlledSubagentRunsReadContext(
controllerSessionKey: string,
controllerAgentId?: string,
cfg?: OpenClawConfig,
): {
runs: SubagentRunRecord[];
countPendingDescendantRuns(rootSessionKey: string): number;
} {
): ControlledSubagentRunsReadContext {
const key = controllerSessionKey.trim();
const agentId = controllerAgentId ?? parseAgentSessionKey(key)?.agentId;
if (!key || !agentId) {
return {
runs: [],
countPendingDescendantRuns: () => 0,
getExecutionObservation: () => ({ state: "unknown" }),
};
}
@ -134,6 +143,32 @@ export function buildControlledSubagentRunsReadContext(
runs: sortSubagentRuns(filtered),
countPendingDescendantRuns: (rootSessionKey) =>
readIndex.countPendingDescendantRuns(rootSessionKey),
getExecutionObservation: (entry) => {
const taskRunId = entry.taskRunId ?? entry.runId;
const task = findTaskByRunId(taskRunId);
const backing = readTaskBackingInstance(task?.detail);
const requesterAgentId = resolveRunRequesterAgentId(entry, cfg);
// Child sessions and logical tasks survive successor runs; only the selected
// backing generation may contribute activity to this snapshot's detail row.
if (
task?.runtime === "subagent" &&
task.runId === taskRunId &&
task.childSessionKey === entry.childSessionKey &&
task.requesterSessionKey === entry.requesterSessionKey &&
requesterAgentId !== undefined &&
task.requesterAgentId === requesterAgentId &&
task.agentId ===
(parseAgentSessionKey(entry.childSessionKey)?.agentId ?? requesterAgentId) &&
backing?.runtime === "subagent" &&
backing.generation === entry.generation
) {
return getTaskExecutionObservation(task);
}
return observeSubagentExecution(
entry,
getSubagentRunsForRequesterSession(entry.childSessionKey),
);
},
};
}

View file

@ -1,4 +1,5 @@
import { afterEach, describe, expect, it } from "vitest";
import { claimAgentRunContext, releaseAgentRunContext } from "../../../infra/agent-run-registry.js";
import { getSubagentExecutionObservation } from "./subagent-execution-observation.js";
import { subagentRuns } from "./subagent-registry-memory.js";
import type { SubagentRunRecord } from "./subagent-registry.types.js";
@ -25,8 +26,22 @@ describe("subagent execution observation", () => {
const original = run("original");
subagentRuns.set(original.runId, original);
const target = { taskRunId: original.runId, childSessionKey: original.childSessionKey };
const claim = claimAgentRunContext(
original.runId,
{ sessionKey: original.childSessionKey },
{ trackOwner: true, ownsContext: true },
);
try {
expect(getSubagentExecutionObservation(target)).toEqual({
state: "running",
executionRunId: original.runId,
});
} finally {
releaseAgentRunContext(original.runId, claim);
}
original.execution = { status: "queued" };
expect(getSubagentExecutionObservation(target)).toEqual({
state: "running",
state: "unknown",
executionRunId: original.runId,
});

View file

@ -1,13 +1,18 @@
import {
getSubagentRunsForChildSession,
getSubagentRunsForRequesterSession,
subagentRuns,
} from "./subagent-registry-memory.js";
import type { SubagentRunRecord } from "./subagent-registry.types.js";
import {
compareSubagentRunGeneration,
recordLatestSubagentRun,
} from "./subagent-run-generation.js";
import { hasSubagentRunEnded, isRetainedUnendedSubagentRun } from "./subagent-run-liveness.js";
import {
hasSubagentRunEnded,
isSubagentRunLive,
isSubagentRunQueued,
} from "./subagent-run-liveness.js";
export type SubagentExecutionObservation = {
state: "queued" | "running" | "waiting" | "finished" | "unknown";
@ -74,10 +79,27 @@ export function observeSubagentExecution(
if (hasSubagentRunEnded(entry)) {
return { state: "finished" };
}
if (entry.execution.status === "interrupted" || !isRetainedUnendedSubagentRun(entry)) {
if (entry.execution.status === "interrupted") {
return { state: "unknown" };
}
return { state: entry.execution.status === "queued" ? "queued" : "running" };
// Snapshots must match the current registration before using live or queued ownership.
const current = subagentRuns.get(entry.runId);
if (
!current ||
current.childSessionKey !== entry.childSessionKey ||
current.requesterSessionKey !== entry.requesterSessionKey ||
(current.taskRunId ?? current.runId) !== (entry.taskRunId ?? entry.runId) ||
compareSubagentRunGeneration(current, entry) !== 0
) {
return { state: "unknown" };
}
if (isSubagentRunLive(current)) {
return { state: current.execution.status === "queued" ? "queued" : "running" };
}
if (isSubagentRunQueued(current)) {
return { state: "queued" };
}
return { state: "unknown" };
}
/** Observe only the current memory owner of this exact delegated task. */

View file

@ -1,6 +1,7 @@
/** Registry projections must agree with admitted execution and queue owners. */
import { afterEach, expect, it, vi } from "vitest";
import { createDeferred } from "../../../../test/helpers/promise.js";
import { buildSubagentsStatusLine } from "../../../auto-reply/reply/commands-status-subagents.js";
import { getRuntimeConfig } from "../../../config/config.js";
import { projectGatewaySessionRunState } from "../../../gateway/session-utils-display.js";
import {
@ -31,6 +32,7 @@ import {
removeQueuedSwarmRun,
reserveSwarmRun,
} from "../swarm/swarm-scheduler.js";
import { buildControlledSubagentRunsReadContext } from "./subagent-control-scope.js";
import { useSubagentControlFixture } from "./subagent-control.test-support.js";
import { buildSubagentList } from "./subagent-list.js";
import { subagentRuns } from "./subagent-registry-memory.js";
@ -200,8 +202,17 @@ it("retains an exact queued collector reservation without calling it executor-li
onStartFailure: () => true,
});
await Promise.resolve();
const prepared = buildControlledSubagentRunsReadContext(parent, "main", getRuntimeConfig());
expect(
buildSubagentList({
cfg: getRuntimeConfig(),
runs: prepared.runs,
recentMinutes: 30,
}).active.map((row) => ({ status: row.status, execution: row.execution.state })),
).toEqual([{ status: "queued", execution: "queued" }]);
now.mockReturnValue(olderThanCutoff);
expect(isSwarmRunWaitingForCapacity(entry.runId, entry)).toBe(true);
expect(prepared.getExecutionObservation(prepared.runs[0]!)).toMatchObject({ state: "queued" });
expect(isSubagentRunQueued(entry)).toBe(true);
expect(isSubagentRunQueued({ ...entry })).toBe(false);
expect(isSubagentRunLive(entry)).toBe(false);
@ -241,6 +252,7 @@ it("retains an exact queued collector reservation without calling it executor-li
expect(captured.countPendingDescendantRuns(parent)).toBe(0);
expect(captured.hasDescendantRunAwaitingSettle(parent)).toBe(false);
expect(isSubagentRunQueued(entry)).toBe(false);
expect(prepared.getExecutionObservation(prepared.runs[0]!)).toMatchObject({ state: "unknown" });
expect(countActiveRunsForSession(parent, { collect: true })).toBe(0);
expect(hasDescendantRunAwaitingSettle(parent)).toBe(false);
const released = buildSubagentSessionListReadIndex();
@ -285,6 +297,52 @@ it("does not retain an old run after its last claim releases preserved routing m
}
});
it("does not borrow a same-run-ID successor's live claim through a prepared observation", async () => {
vi.spyOn(Date, "now").mockReturnValue(start);
const original = await register("live-generation");
const prepared = buildControlledSubagentRunsReadContext(parent, "main", getRuntimeConfig());
const successor = await register(original.runId);
const claim = claimAgentRunContext(
successor.runId,
{ sessionKey: successor.childSessionKey },
{ trackOwner: true, ownsContext: true },
);
try {
expect(successor.generation).toBeGreaterThan(original.generation!);
expect(hasLiveAgentRunContext(successor.runId)).toBe(true);
expect(isSubagentRunLive(successor)).toBe(true);
const current = buildControlledSubagentRunsReadContext(parent, "main", getRuntimeConfig());
expect(prepared.getExecutionObservation(prepared.runs[0]!)).toMatchObject({ state: "unknown" });
expect(current.getExecutionObservation(current.runs[0]!)).toMatchObject({ state: "running" });
for (const [context, state] of [
[prepared, "unknown"],
[current, "running"],
] as const) {
expect(
buildSubagentList({
cfg: getRuntimeConfig(),
runs: context.runs,
recentMinutes: 30,
}).active.map((row) => ({ runId: row.runId, execution: row.execution.state })),
).toEqual([{ runId: successor.runId, execution: state }]);
const status = buildSubagentsStatusLine({ context, verboseEnabled: false });
expect(status).toContain("Subagents: 1 active");
if (state === "unknown") {
expect(status).toMatch(/unknown|unavailable/i);
expect(status).not.toMatch(/\brunning\b/i);
} else {
expect(status).toMatch(/\brunning\b/i);
}
}
expect(findTaskByRunId(successor.runId)?.status).toBe("running");
expect(countActiveRunsForSession(parent)).toBe(1);
expect(countPendingDescendantRuns(parent)).toBe(1);
} finally {
releaseAgentRunContext(successor.runId, claim);
}
});
it("does not transfer read retention across replaced queue owners or copied reservations", async () => {
const now = vi.spyOn(Date, "now").mockReturnValue(start);
const id = "queue-generation";
@ -297,6 +355,7 @@ it("does not transfer read retention across replaced queue owners or copied rese
});
expect(reserve()).toBe(true);
const original = await register(id, true);
const prepared = buildControlledSubagentRunsReadContext(parent, "main", getRuntimeConfig());
now.mockReturnValue(olderThanCutoff);
const snapshot = buildSubagentSessionListReadIndex();
const compact = snapshot.listDescendantRunsForRequester(parent)[0]!;
@ -316,6 +375,8 @@ it("does not transfer read retention across replaced queue owners or copied rese
expect(replacement.generation).toBeGreaterThan(original.generation!);
expect(isSubagentRunQueued(original)).toBe(false);
expect(isSubagentRunQueued(replacement)).toBe(false);
const replaced = buildControlledSubagentRunsReadContext(parent, "main", getRuntimeConfig());
expect(replaced.getExecutionObservation(replaced.runs[0]!)).toMatchObject({ state: "unknown" });
expect(snapshot.countActiveDescendantRuns(parent)).toBe(0);
expect(snapshot.hasDescendantRunAwaitingSettle(parent)).toBe(false);
expect(buildSubagentSessionListReadIndex().countPendingDescendantRuns(parent)).toBe(0);
@ -326,6 +387,7 @@ it("does not transfer read retention across replaced queue owners or copied rese
const successor = await register(id, true);
now.mockReturnValue(olderThanCutoff);
expect(isSubagentRunQueued(successor)).toBe(true);
expect(prepared.getExecutionObservation(prepared.runs[0]!)).toMatchObject({ state: "unknown" });
expect(buildSubagentSessionListReadIndex().countActiveDescendantRuns(parent)).toBe(1);
// A previously issued projection cannot borrow the new generation's owner.
expect(
@ -395,6 +457,12 @@ it("does not keep a recent orphan executor-live after its admitted owner closes"
expect(isSubagentRunQueued(entry)).toBe(false);
expect(findTaskByRunId(entry.runId)?.status).toBe("running");
expect.soft(isSubagentSessionRunActive(entry.childSessionKey)).toBe(false);
expect(countActiveRunsForSession(parent)).toBe(1);
expect(
buildSubagentList({ cfg: getRuntimeConfig(), runs: [entry], recentMinutes: 30 }).active.map(
(row) => ({ runId: row.runId, execution: row.execution.state }),
),
).toEqual([{ runId: entry.runId, execution: "unknown" }]);
});
it("retains durable suspended completion debt without reporting a live executor or awaiting automatic settlement", async () => {

View file

@ -0,0 +1,235 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { clearAgentHarnesses } from "../../agents/harness/registry.js";
import {
addSubagentRunForTests,
resetSubagentRegistryForTests,
} from "../../agents/subagents/registry/subagent-registry.test-helpers.js";
import { emitAgentEvent } from "../../infra/agent-events.js";
import { clearAgentRunContext, registerAgentRunContext } from "../../infra/agent-run-registry.js";
import { createSubagentTaskBackingDetail } from "../../tasks/task-backing-records.js";
import { createRunningTaskRunCore } from "../../tasks/task-executor.js";
import { resetTaskRegistryForTests } from "../../tasks/task-runtime.test-helpers.js";
import { buildStatusReplyForTest } from "./commands-status.test-support.js";
import { configureInMemoryTaskRegistryStoreForTests } from "./commands.test-harness.js";
vi.mock("../../status/status-plugin-health.runtime.js", () => ({
collectRuntimePluginHealthSnapshot: () => ({
plugins: [],
diagnostics: [],
contextEngineQuarantines: [],
runtimeToolQuarantines: [],
channelPluginFailures: [],
}),
}));
describe("buildStatusReply execution observations", () => {
beforeEach(() => {
clearAgentHarnesses();
resetSubagentRegistryForTests();
resetTaskRegistryForTests({ persist: false });
configureInMemoryTaskRegistryStoreForTests();
});
afterEach(() => {
clearAgentHarnesses();
resetSubagentRegistryForTests();
resetTaskRegistryForTests({ persist: false });
});
it("shows canonical successor tool activity resuming after approval without changing counts", async () => {
const runId = "status-observed-successor";
const taskRunId = "status-observed-original";
const childSessionKey = "agent:main:subagent:status-observed";
addSubagentRunForTests({
runId,
taskRunId,
generation: 2,
childSessionKey,
requesterSessionKey: "agent:main:main",
requesterDisplayKey: "main",
task: "observed worker",
cleanup: "keep",
createdAt: Date.now() - 60_000,
startedAt: Date.now() - 60_000,
});
createRunningTaskRunCore({
runtime: "subagent",
requesterSessionKey: "agent:main:main",
childSessionKey,
runId: taskRunId,
task: "observed worker",
detail: createSubagentTaskBackingDetail(2),
});
registerAgentRunContext(runId, { sessionKey: childSessionKey, projectSessionActive: true });
try {
emitAgentEvent({
runId,
stream: "tool",
data: { phase: "start", name: "read", toolCallId: "status-read" },
});
const running = await buildStatusReplyForTest({});
const runningDetail = running?.text
?.split("\n")
.find((line) => line.includes("• observed worker"));
expect(running?.text).toContain("Subagents: 1 active");
expect(running?.text).toContain("Tasks: 1 active · 1 total");
expect(runningDetail).toMatch(/running/i);
expect(runningDetail).toMatch(/\bread\b/);
emitAgentEvent({
runId,
stream: "execution",
data: { approval: { id: "status-approval", state: "pending" } },
});
const approval = await buildStatusReplyForTest({});
const approvalDetail = approval?.text
?.split("\n")
.find((line) => line.includes("• observed worker"));
expect(approvalDetail).toMatch(/wait.*approval/i);
expect(approvalDetail).not.toMatch(/\brunning\b/i);
expect(approval?.text).toContain("Tasks: 1 active · 1 total");
emitAgentEvent({
runId,
stream: "execution",
data: { approval: { id: "status-approval", state: "resolved" } },
});
emitAgentEvent({
runId,
stream: "tool",
data: { phase: "start", name: "read", toolCallId: "status-read" },
});
const resumed = await buildStatusReplyForTest({});
const resumedDetail = resumed?.text
?.split("\n")
.find((line) => line.includes("• observed worker"));
expect(resumedDetail).toMatch(/\brunning\b/i);
expect(resumedDetail).toMatch(/\bread\b/);
expect(resumedDetail).not.toMatch(/approval/i);
expect(resumed?.text).toContain("Subagents: 1 active");
expect(resumed?.text).toContain("Tasks: 1 active · 1 total");
} finally {
clearAgentRunContext(runId);
}
});
it("keeps recent ownerless rows and task counts while reporting unknown activity", async () => {
const runId = "status-ownerless";
const childSessionKey = "agent:main:subagent:status-ownerless";
addSubagentRunForTests({
runId,
generation: 1,
childSessionKey,
requesterSessionKey: "agent:main:main",
requesterDisplayKey: "main",
task: "retained worker",
cleanup: "keep",
createdAt: Date.now() - 60_000,
startedAt: Date.now() - 60_000,
});
createRunningTaskRunCore({
runtime: "subagent",
requesterSessionKey: "agent:main:main",
childSessionKey,
runId,
task: "retained worker",
detail: createSubagentTaskBackingDetail(1),
});
const reply = await buildStatusReplyForTest({});
const detail = reply?.text?.split("\n").find((line) => line.includes("• retained worker"));
expect(reply?.text).toContain("Subagents: 1 active");
expect(reply?.text).toContain("Tasks: 1 active · 1 total");
expect(detail).toMatch(/unknown|unavailable/i);
expect(detail).not.toMatch(/\b(running|queued)\b/i);
});
it.each(["task", "generation"] as const)(
"does not borrow approval activity from a different canonical %s in the same child session",
async (replacement) => {
const childSessionKey = "agent:main:subagent:status-replaced";
const now = Date.now();
addSubagentRunForTests({
runId: "status-previous",
generation: 1,
childSessionKey,
task: "previous worker",
createdAt: now - 2_000,
startedAt: now - 2_000,
});
createRunningTaskRunCore({
runtime: "subagent",
requesterSessionKey: "agent:main:main",
childSessionKey,
runId: "status-previous",
task: "previous worker",
detail: createSubagentTaskBackingDetail(1),
});
emitAgentEvent({
runId: "status-previous",
stream: "execution",
data: { approval: { id: "previous-approval", state: "pending" } },
});
addSubagentRunForTests({
runId: "status-current",
taskRunId: replacement === "generation" ? "status-previous" : "status-current",
generation: 2,
childSessionKey,
task: "replacement worker",
createdAt: now - 1_000,
startedAt: now - 1_000,
});
const reply = await buildStatusReplyForTest({});
const detail = reply?.text?.split("\n").find((line) => line.includes("• replacement worker"));
expect(reply?.text).toContain("Subagents: 1 active");
expect(detail).toMatch(/unknown|unavailable/i);
expect(detail).not.toMatch(/approval/i);
expect(reply?.text).not.toContain("• previous worker");
},
);
it("retains ended child delivery debt without calling the child active", async () => {
const parentKey = "agent:main:subagent:status-delivery-parent";
const now = Date.now();
addSubagentRunForTests({
runId: "status-delivery-parent",
childSessionKey: parentKey,
task: "delivery orchestrator",
createdAt: now - 120_000,
startedAt: now - 120_000,
endedAt: now - 60_000,
outcome: { status: "ok" },
});
addSubagentRunForTests({
runId: "status-delivery-child",
childSessionKey: `${parentKey}:subagent:child`,
requesterSessionKey: parentKey,
requesterDisplayKey: parentKey,
task: "completed child",
createdAt: now - 90_000,
startedAt: now - 90_000,
endedAt: now - 30_000,
outcome: { status: "ok" },
expectsCompletionMessage: true,
completion: {
required: true,
resultText: "private completed result",
capturedAt: now - 30_000,
},
delivery: { status: "pending" },
});
const reply = await buildStatusReplyForTest({});
const detail = reply?.text
?.split("\n")
.find((line) => line.includes("• delivery orchestrator"));
expect(reply?.text).toContain("Subagents: 1 active");
expect(detail).toMatch(/pending|delivery|settle/i);
expect(detail).not.toMatch(/child(?:ren)? active|\brunning\b/i);
expect(reply?.text).not.toContain("private completed result");
});
});

View file

@ -1,15 +1,45 @@
// Formats subagent status rows for the status command response.
import type { buildControlledSubagentRunsReadContext } from "../../agents/subagents/registry/subagent-control-scope.js";
import type { TaskSummary } from "../../../packages/gateway-protocol/src/schema/tasks.js";
import type { ControlledSubagentRunsReadContext } from "../../agents/subagents/registry/subagent-control-scope.js";
import {
hasSubagentRunEnded,
isRetainedUnendedSubagentRun,
} from "../../agents/subagents/registry/subagent-run-liveness.js";
import { formatDurationCompact } from "../../infra/format-time/format-duration.ts";
import { sanitizeTaskStatusText } from "../../tasks/task-status.js";
import { formatRunLabel } from "./subagents-utils.js";
function formatExecutionObservation(observation: NonNullable<TaskSummary["execution"]>): string {
switch (observation.state) {
case "running": {
const tool = sanitizeTaskStatusText(observation.currentTool?.name, { maxChars: 60 });
return tool ? `running ${tool}` : "running";
}
case "queued":
return "queued";
case "waiting":
switch (observation.wait?.kind) {
case "approval":
return "waiting for approval";
case "user_input":
return "waiting for input";
case "children":
return "waiting for child tasks";
case "agent_messages":
return "waiting for agent messages";
default:
return "waiting for external work";
}
case "finished":
return "finished · settlement pending";
default:
return "current activity unavailable";
}
}
/** Builds the compact status line from the controller's ordered snapshot and descendant index. */
export function buildSubagentsStatusLine(params: {
context: ReturnType<typeof buildControlledSubagentRunsReadContext>;
context: ControlledSubagentRunsReadContext;
verboseEnabled: boolean;
now?: number;
}): string | undefined {
@ -36,11 +66,12 @@ export function buildSubagentsStatusLine(params: {
);
const duration = formatDurationCompact(durationMs, { spaced: true }) ?? "0s";
const label = formatRunLabel(entry, { maxLength: 56 });
const executionText = formatExecutionObservation(context.getExecutionObservation(entry));
const descendantText =
pendingDescendants > 0
? ` · ${pendingDescendants} child${pendingDescendants === 1 ? "" : "ren"} active`
? ` · ${pendingDescendants} child${pendingDescendants === 1 ? "" : "ren"} pending`
: "";
detailLines.push(` • ${label} · ${duration}${descendantText}`);
detailLines.push(` • ${label} · ${duration} · ${executionText}${descendantText}`);
} else if (hasSubagentRunEnded(entry) && pendingDescendants === 0) {
done += 1;
}

View file

@ -0,0 +1,37 @@
import type { OpenClawConfig } from "../../config/config.js";
import { buildStatusReply } from "./commands-status.js";
import { baseCommandTestConfig, buildCommandTestParams } from "./commands.test-harness.js";
export async function buildStatusReplyForTest(params: {
sessionKey?: string;
agentId?: string;
cfg?: OpenClawConfig;
verbose?: boolean;
}) {
const cfg = params.cfg ?? baseCommandTestConfig;
const commandParams = buildCommandTestParams("/status", cfg);
const sessionKey = params.sessionKey ?? commandParams.sessionKey;
return await buildStatusReply({
cfg,
agentId: params.agentId,
command: commandParams.command,
sessionEntry: commandParams.sessionEntry,
sessionKey,
parentSessionKey: sessionKey,
sessionScope: commandParams.sessionScope,
storePath: commandParams.storePath,
provider: "anthropic",
model: "claude-opus-4-6",
contextTokens: 0,
resolvedThinkLevel: commandParams.resolvedThinkLevel,
resolvedFastMode: false,
resolvedVerboseLevel: params.verbose ? "on" : commandParams.resolvedVerboseLevel,
resolvedReasoningLevel: commandParams.resolvedReasoningLevel,
resolvedElevatedLevel: commandParams.resolvedElevatedLevel,
resolveDefaultThinkingLevel: commandParams.resolveDefaultThinkingLevel,
isGroup: commandParams.isGroup,
defaultGroupActivation: commandParams.defaultGroupActivation,
modelAuthOverride: "api-key",
activeModelAuthOverride: "api-key",
});
}

View file

@ -31,6 +31,7 @@ import {
import { resetTaskRegistryForTests } from "../../tasks/task-runtime.test-helpers.js";
import { withEnvAsync } from "../../test-utils/env.js";
import { buildStatusPluginsReply, buildStatusReply, buildStatusText } from "./commands-status.js";
import { buildStatusReplyForTest } from "./commands-status.test-support.js";
import {
baseCommandTestConfig,
buildCommandTestParams,
@ -147,40 +148,6 @@ function createStatusDisplayParams(
} satisfies Partial<StatusTextParams>;
}
async function buildStatusReplyForTest(params: {
sessionKey?: string;
agentId?: string;
cfg?: OpenClawConfig;
verbose?: boolean;
}) {
const cfg = params.cfg ?? baseCfg;
const commandParams = buildCommandTestParams("/status", cfg);
const sessionKey = params.sessionKey ?? commandParams.sessionKey;
return await buildStatusReply({
cfg,
agentId: params.agentId,
command: commandParams.command,
sessionEntry: commandParams.sessionEntry,
sessionKey,
parentSessionKey: sessionKey,
sessionScope: commandParams.sessionScope,
storePath: commandParams.storePath,
provider: "anthropic",
model: "claude-opus-4-6",
contextTokens: 0,
resolvedThinkLevel: commandParams.resolvedThinkLevel,
resolvedFastMode: false,
resolvedVerboseLevel: params.verbose ? "on" : commandParams.resolvedVerboseLevel,
resolvedReasoningLevel: commandParams.resolvedReasoningLevel,
resolvedElevatedLevel: commandParams.resolvedElevatedLevel,
resolveDefaultThinkingLevel: commandParams.resolveDefaultThinkingLevel,
isGroup: commandParams.isGroup,
defaultGroupActivation: commandParams.defaultGroupActivation,
modelAuthOverride: "api-key",
activeModelAuthOverride: "api-key",
});
}
function registerStatusCodexHarness(): void {
const codexProviders = new Set(["codex", "openai"]);
const harness: AgentHarness = {
@ -2336,38 +2303,12 @@ describe("buildStatusReply error handling", () => {
vi.restoreAllMocks();
});
async function runStatusReply() {
const commandParams = buildCommandTestParams("/status", baseCfg);
return await buildStatusReply({
cfg: baseCfg,
command: commandParams.command,
sessionEntry: commandParams.sessionEntry,
sessionKey: commandParams.sessionKey,
parentSessionKey: commandParams.sessionKey,
sessionScope: commandParams.sessionScope,
storePath: commandParams.storePath,
provider: "anthropic",
model: "claude-opus-4-6",
contextTokens: 0,
resolvedThinkLevel: commandParams.resolvedThinkLevel,
resolvedFastMode: false,
resolvedVerboseLevel: commandParams.resolvedVerboseLevel,
resolvedReasoningLevel: commandParams.resolvedReasoningLevel,
resolvedElevatedLevel: commandParams.resolvedElevatedLevel,
resolveDefaultThinkingLevel: commandParams.resolveDefaultThinkingLevel,
isGroup: commandParams.isGroup,
defaultGroupActivation: commandParams.defaultGroupActivation,
modelAuthOverride: "api-key",
activeModelAuthOverride: "api-key",
});
}
it("delivers a fixed generic reply and logs details when status rendering throws", async () => {
const logError = vi.spyOn(logger, "logError").mockImplementation(() => {});
vi.spyOn(statusText, "buildStatusReplyParts").mockRejectedValue(
new Error("Unexpected rendering error"),
);
const reply = await runStatusReply();
const reply = await buildStatusReplyForTest({});
// Exact object equality also pins that no stale presentation or internal
// error text reaches the channel; diagnostics belong to the log sink only.
@ -2385,7 +2326,7 @@ describe("buildStatusReply error handling", () => {
text: "plain status",
presentation,
});
const reply = await runStatusReply();
const reply = await buildStatusReplyForTest({});
expect(reply).toMatchObject({
text: "plain status",

View file

@ -212,20 +212,23 @@ describe("subagents status", () => {
});
}
expect(
buildSubagentsStatusLine({
context: buildControlledSubagentRunsReadContext("agent:main:main"),
verboseEnabled: true,
now,
}),
).toBe(
[
"🤖 Subagents: 4 active · 1 done",
" • first worker · 1s",
" • tie-b worker · 2s",
` • tie-a worker · 2s · ${children} child${children === 1 ? "" : "ren"} active`,
].join("\n"),
);
const text = buildSubagentsStatusLine({
context: buildControlledSubagentRunsReadContext("agent:main:main"),
verboseEnabled: true,
now,
});
const details = text?.split("\n").slice(1);
expect(text).toContain("Subagents: 4 active · 1 done");
expect(details).toEqual([
expect.stringContaining("first worker"),
expect.stringContaining("tie-b worker"),
expect.stringContaining("tie-a worker"),
]);
expect(details?.[0]).toContain("1s");
expect(details?.[1]).toContain("2s");
expect(details?.[2]).toContain("2s");
expect(details?.[2]).toMatch(new RegExp(`\\b${children} child`));
},
);
@ -244,13 +247,14 @@ describe("subagents status", () => {
createdAt: 1_000,
execution: { status: "running", startedAt: 1_000, endedAt },
};
expect(
buildSubagentsStatusLine({
context: { runs: [run], countPendingDescendantRuns: () => 0 },
verboseEnabled: false,
now: 5_000,
}),
).toBe(`🤖 Subagents: 1 active\n • active worker · ${duration}`);
addSubagentRunForTests(run);
const text = buildSubagentsStatusLine({
context: buildControlledSubagentRunsReadContext("agent:main:main"),
verboseEnabled: false,
now: 5_000,
});
expect(text).toContain("Subagents: 1 active");
expect(text).toContain(`active worker · ${duration}`);
});
});

View file

@ -1,11 +1,19 @@
import { afterEach, describe, expect, it } from "vitest";
import { useSubagentControlFixture } from "../../agents/subagents/registry/subagent-control.test-support.js";
import { subagentRuns } from "../../agents/subagents/registry/subagent-registry-memory.js";
import { registerSubagentRun } from "../../agents/subagents/registry/subagent-registry.js";
import { writeSubagentSessionEntry } from "../../agents/subagents/registry/subagent-registry.persistence.test-support.js";
import { emitAgentEvent } from "../../infra/agent-events.js";
import {
claimAgentRunContext,
releaseAgentRunContext,
resetAgentRunRegistryForTest,
} from "../../infra/agent-run-registry.js";
import { markTaskTerminalById } from "../../tasks/runtime-internal.js";
import {
findTaskByRunId,
getTaskById,
markTaskTerminalById,
} from "../../tasks/runtime-internal.js";
import { clearTaskActivity } from "../../tasks/task-registry-activity.js";
import { createTaskFixture } from "../../tasks/task-registry.test-support.js";
import {
@ -15,10 +23,126 @@ import {
} from "./tasks.fixture.test-support.js";
import { runTaskHandler } from "./tasks.test-helpers.js";
useTaskGatewayFixture();
afterEach(resetAgentRunRegistryForTest);
describe("registered subagent execution", () => {
const fixture = useSubagentControlFixture();
it.each(["tasks.get", "tasks.list"] as const)(
"%s separates retained tasks from current execution ownership",
async (method) => {
const runId = "retained-execution";
const childSessionKey = `agent:main:subagent:${runId}`;
await writeSubagentSessionEntry({
stateDir: fixture.stateDir,
agentId: "main",
sessionKey: childSessionKey,
defaultSessionId: `${runId}-session`,
lifecycleRevision: `${runId}-revision`,
});
registerSubagentRun({
runId,
childSessionKey,
requesterSessionKey: mainSessionTaskScope.requesterSessionKey,
requesterAgentId: "main",
requesterDisplayKey: "main",
task: "Inspect execution ownership",
cleanup: "keep",
expectsCompletionMessage: false,
});
const task = findTaskByRunId(runId)!;
const entry = subagentRuns.get(runId)!;
const readTask = async () => {
if (method === "tasks.get") {
return (await getTaskPayload(task.taskId)).payload?.task;
}
const { payload } = await runTaskHandler("tasks.list", {});
return payload?.tasks?.find((row) => row.id === task.taskId);
};
expect(await readTask()).toMatchObject({
id: task.taskId,
status: "running",
execution: { state: "unknown" },
});
const claim = claimAgentRunContext(
runId,
{ sessionKey: childSessionKey },
{ trackOwner: true, ownsContext: true },
);
try {
const sourceId = "owned-observer";
const executionId = "owned-turn";
emitAgentEvent({
runId,
stream: "execution",
data: { state: "running", sourceId, executionId },
});
expect(await readTask()).toMatchObject({
status: "running",
execution: { state: "running" },
});
// Losing the current observation does not release the run context or settle its task.
emitAgentEvent({
runId,
stream: "execution",
data: { state: "unknown", sourceId, invalidate: true },
});
expect(await readTask()).toMatchObject({
id: task.taskId,
status: "running",
execution: { state: "unknown" },
});
expect(getTaskById(task.taskId)?.status).toBe("running");
emitAgentEvent({
runId,
stream: "execution",
data: { state: "running", sourceId, executionId },
});
expect(await readTask()).toMatchObject({
status: "running",
execution: { state: "running" },
});
emitAgentEvent({
runId,
stream: "tool",
data: { phase: "start", name: "read", toolCallId: "owned-read" },
});
expect(await readTask()).toMatchObject({
status: "running",
execution: { state: "running", currentTool: { name: "read" } },
});
// The old task must not borrow a replacement registration's owner or activity.
subagentRuns.set(runId, { ...entry, generation: entry.generation! + 1 });
const replaced = await readTask();
expect(replaced).toMatchObject({
status: "running",
execution: { state: "unknown" },
});
expect(replaced?.execution).not.toHaveProperty("currentTool");
subagentRuns.set(runId, entry);
} finally {
releaseAgentRunContext(runId, claim);
}
const released = await readTask();
expect(released).toMatchObject({
id: task.taskId,
status: "running",
execution: { state: "unknown" },
});
expect(released?.execution).not.toHaveProperty("currentTool");
expect(getTaskById(task.taskId)?.status).toBe("running");
},
);
});
describe("tasks gateway execution and activity", () => {
useTaskGatewayFixture();
afterEach(resetAgentRunRegistryForTest);
it("reports a live CLI run before activity arrives and after transient activity is cleared", async () => {
const runId = "run-cli-owner";
const sessionKey = "agent:main:dashboard:cli-owner";

View file

@ -92,8 +92,8 @@ it.each(["agent:main:dashboard:stored", "global"])(
},
);
it("projects fixed task statuses without observing retained native executions", () => {
const statuses = ["queued", "succeeded", "failed", "timed_out", "cancelled", "lost"] as const;
it("projects terminal task statuses without observing retained native executions", () => {
const statuses = ["succeeded", "failed", "timed_out", "cancelled", "lost"] as const;
const rows = Array.from({ length: 1_000 }, (_, index) => {
const status = statuses[index % statuses.length]!;
const record = task(`fixed-${index}`, status);
@ -125,8 +125,7 @@ it("projects fixed task statuses without observing retained native executions",
expect(rows.map(({ record }) => getTaskExecutionObservation(record))).toEqual(
rows.map(({ record, timestamp }) => ({
state:
record.status === "lost" ? "unknown" : record.status === "queued" ? "queued" : "finished",
state: record.status === "lost" ? "unknown" : "finished",
...(timestamp !== undefined ? { lastActivityAt: timestamp } : {}),
})),
);
@ -137,6 +136,11 @@ it("projects fixed task statuses without observing retained native executions",
it("keeps running task observations current through generation replacement and deletion", () => {
const record = task("running-task", "running");
const original = registerRun(record);
claimAgentRunContext(
original.runId,
{ sessionKey: original.childSessionKey },
{ trackOwner: true, ownsContext: true },
);
recordTaskActivityEvent(record, {
runId: original.runId,
seq: 1,
@ -167,6 +171,11 @@ it("keeps running task observations current through generation replacement and d
successor.pauseReason = undefined;
successor.execution = { status: "running", startedAt: 30 };
claimAgentRunContext(
successor.runId,
{ sessionKey: successor.childSessionKey },
{ trackOwner: true, ownsContext: true },
);
recordTaskActivityEvent(record, {
runId: successor.runId,
seq: 1,

View file

@ -48,7 +48,7 @@ export function getTaskExecutionObservation(
? "unknown"
: isTerminalTaskStatus(task.status)
? "finished"
: task.status === "queued"
: task.status === "queued" && task.runtime !== "subagent"
? "queued"
: undefined;
if (fixedState) {
@ -101,7 +101,10 @@ export function getTaskExecutionObservation(
state: currentActivity?.executionState ?? observeCliExecution(task) ?? "unknown",
...(currentActivity?.executionWait ? { wait: currentActivity.executionWait } : {}),
};
if (execution.state === "running" && currentActivity?.executionWait) {
if (
execution.state === "running" &&
(currentActivity?.executionState || currentActivity?.executionWait)
) {
execution.state = currentActivity.executionState ?? "waiting";
execution.wait = currentActivity.executionWait;
}

View file

@ -5,7 +5,11 @@ import { subagentRuns } from "../agents/subagents/registry/subagent-registry-mem
import { settleRequesterTurnAfterSessionSpawns } from "../agents/subagents/registry/subagent-registry-requester-yield.js";
import type { SubagentRunRecord } from "../agents/subagents/registry/subagent-registry.types.js";
import { emitAgentEvent, resetAgentEventsForTest } from "../infra/agent-events.js";
import { registerAgentRunContext } from "../infra/agent-run-registry.js";
import {
claimAgentRunContext,
registerAgentRunContext,
releaseAgentRunContext,
} from "../infra/agent-run-registry.js";
import { peekSystemEvents, resetSystemEventsForTest } from "../infra/system-events.js";
import {
markGatewayRestartDraining,
@ -48,6 +52,7 @@ const origin = {
threadId: "test-thread",
};
const sendMessage = vi.fn<deliveryRuntime.TaskRegistryDeliveryRuntime["sendMessage"]>();
const runContextClaims = new Map<string, string>();
function child(name: string, notifyPolicy: TaskNotifyPolicy = "state_changes") {
const entry: SubagentRunRecord = {
@ -68,6 +73,15 @@ function child(name: string, notifyPolicy: TaskNotifyPolicy = "state_changes") {
expectsCompletionMessage: true,
};
subagentRuns.set(entry.runId, entry);
const claim = claimAgentRunContext(
entry.runId,
{ sessionKey: entry.childSessionKey },
{ trackOwner: true, ownsContext: true },
);
if (!claim) {
throw new Error("Expected child execution ownership");
}
runContextClaims.set(entry.runId, claim);
const task = createTaskRecord({
runtime: "subagent",
ownerKey: PARENT,
@ -87,7 +101,7 @@ function child(name: string, notifyPolicy: TaskNotifyPolicy = "state_changes") {
if (!task) {
throw new Error("Expected accepted task");
}
return { entry, task };
return { entry, task, claim };
}
function yieldParent() {
@ -150,6 +164,10 @@ afterEach(() => {
resetTaskRegistryForTests({ persist: false });
resetTaskFlowRegistryForTests({ persist: false });
resetTaskRegistryDeliveryRuntimeForTests();
for (const [runId, claim] of runContextClaims) {
releaseAgentRunContext(runId, claim);
}
runContextClaims.clear();
resetAgentEventsForTest({ preserveListeners: true });
resetSystemEventsForTest();
subagentRuns.clear();
@ -208,9 +226,8 @@ describe("yielded subagent progress delivery", () => {
gatewayOwnedDelivery: true,
mirror: { sessionKey: PARENT, agentId: "main" },
});
expect(message.content).toBe(
"Background work is still in progress:\n- First: running read; 1 tool call started.\n- Second: running read; 40 tool calls started.",
);
expect(message.content).toContain("First: running read; 1 tool call started.");
expect(message.content).toContain("Second: running read; 40 tool calls started.");
expect(message.content).not.toMatch(/Quiet|private-/);
expect(getTaskById(first.task.taskId)).toMatchObject({
status: "running",
@ -220,6 +237,23 @@ describe("yielded subagent progress delivery", () => {
expect(peekSystemEvents(PARENT)).toEqual([]);
});
it("does not report retained tool activity as running after execution ownership releases", async () => {
const item = child("Worker");
tool(item.entry);
yieldParent();
releaseAgentRunContext(item.entry.runId, item.claim);
await vi.advanceTimersByTimeAsync(15_000);
expect(sendMessage).toHaveBeenCalledOnce();
const message = sendMessage.mock.calls[0]![0];
expect(message.content).toContain("Worker: current activity unavailable");
expect(message.content).not.toMatch(/running read|private-/);
expect(getTaskById(item.task.taskId)).toMatchObject({
status: "running",
toolUseCount: 1,
deliveryStatus: "pending",
});
});
it.each(["done_only", "silent"] as const)(
"preserves %s notification policy across yield",
async (policy) => {

View file

@ -1,6 +1,11 @@
import { createServer, type IncomingMessage, type ServerResponse } from "node:http";
import type { AddressInfo } from "node:net";
import { afterEach, describe, expect, it, vi } from "vitest";
import type { ChatEvent } from "../packages/gateway-protocol/src/schema/logs-chat.js";
import type {
TasksGetResult,
TasksListResult,
} from "../packages/gateway-protocol/src/schema/tasks.js";
import { writeSubagentSessionEntry } from "../src/agents/subagents/registry/subagent-registry.persistence.test-support.js";
import {
loadSubagentRegistryFromSqlite,
@ -11,6 +16,7 @@ import { getSessionKysely } from "../src/config/sessions/session-accessor.sqlite
import type { OpenClawConfig } from "../src/config/types.openclaw.js";
import { connectGatewayClient, disconnectGatewayClient } from "../src/gateway/test-helpers.e2e.js";
import { executeSqliteQuerySync } from "../src/infra/kysely-sync.js";
import { extractFirstTextBlock } from "../src/shared/chat-message-content.js";
import { withOpenClawAgentDatabaseReadOnly } from "../src/state/openclaw-agent-db-readonly.js";
import { closeOpenClawStateDatabaseForTest } from "../src/state/openclaw-state-db.js";
import {
@ -21,7 +27,8 @@ import {
createOpenClawTestInstance,
type OpenClawTestInstance,
} from "./helpers/openclaw-test-instance.js";
import { createDeferred } from "./helpers/promise.js";
import { createDeferred, withTestTimeout } from "./helpers/promise.js";
import { runQaGatewayFixture } from "./helpers/qa-gateway-cleanup.js";
const TEST_TIMEOUT_MS = 180_000;
const MODEL_REF = "requester-owner/synthetic";
@ -62,6 +69,261 @@ afterEach(async () => {
});
describe("REQUESTER-OWNER requester agent id survives completion dispatch", () => {
it(
"answers status through WebSocket while parent and child provider requests are held",
{ timeout: TEST_TIMEOUT_MS },
async () => {
const parentGate = createDeferred();
const childGate = createDeferred();
const modelServer = await startProofModelServer({
yieldAfterSpawn: parentGate.promise,
childReply: childGate.promise,
});
modelServers.push(modelServer);
const instance = await createOpenClawTestInstance({
name: "requester-owner-busy-status",
config: createTestConfig(modelServer.url),
env: { OPENCLAW_SKIP_PROVIDERS: undefined, OPENCLAW_TEST_MINIMAL_GATEWAY: undefined },
});
instances.push(instance);
instance.state.applyEnv();
await instance.startGateway();
const sessionKey = `agent:${REQUESTER_AGENT_ID}:${REQUESTER_KEY}`;
const statusRunId = "busy-parent-status";
const statusReply = createDeferred<ChatEvent>();
const client = await connectGatewayClient({
url: instance.url,
token: instance.gatewayToken,
onEvent: (event) => {
if (event.event !== "chat") {
return;
}
const chat = event.payload as ChatEvent;
if (
chat.runId === statusRunId &&
(chat.state === "final" || chat.state === "error" || chat.state === "aborted")
) {
statusReply.resolve(chat);
}
},
});
let parent: Promise<unknown> | undefined;
let parentSettled = false;
try {
parent = client.request(
"agent",
{
sessionKey: REQUESTER_KEY,
agentId: REQUESTER_AGENT_ID,
idempotencyKey: "busy-status-parent-turn",
message: PARENT_PROMPT,
deliver: false,
},
{ expectFinal: true },
);
void parent.then(
() => {
parentSettled = true;
},
() => {
parentSettled = true;
},
);
await vi.waitFor(
() => {
expect(modelServer.requestCount(), instance.logs()).toBe(3);
expect(
modelServer
.bodies()
.filter(
(body) => body.includes(PARENT_PROMPT) && body.includes("function_call_output"),
),
).toHaveLength(1);
expect(
modelServer
.bodies()
.filter(
(body) => body.includes(CHILD_TASK) && !body.includes("function_call_output"),
),
).toHaveLength(1);
},
{ interval: 50, timeout: 60_000 },
);
const runs = [...loadSubagentRegistryFromSqlite().values()];
expect(runs, instance.logs()).toHaveLength(1);
const run = runs[0]!;
expect(run).toMatchObject({
requesterAgentId: REQUESTER_AGENT_ID,
requesterSessionKey: sessionKey,
requesterTurnRunId: "busy-status-parent-turn",
completionTarget: "parent",
});
expect(run.childSessionKey).toMatch(/^agent:beta:subagent:/);
const listed = await client.request<TasksListResult>("tasks.list", {
sessionKey,
agentId: REQUESTER_AGENT_ID,
});
const children = listed.tasks.filter((task) => task.runtime === "subagent");
expect(children).toHaveLength(1);
const child = children[0]!;
expect(child).toMatchObject({
runId: run.taskRunId ?? run.runId,
childSessionKey: run.childSessionKey,
agentId: REQUESTER_AGENT_ID,
sessionKey,
status: "running",
deliveryStatus: "pending",
execution: { state: "running" },
});
// A held HTTP response is not a tool call or an explicit execution wait.
expect(child.execution?.currentTool).toBeUndefined();
expect(child.execution?.wait).toBeUndefined();
const detail = await client.request<TasksGetResult>("tasks.get", { taskId: child.id });
expect(detail.task).toMatchObject({
id: child.id,
runId: child.runId,
sessionKey,
childSessionKey: run.childSessionKey,
prompt: CHILD_TASK,
status: "running",
execution: { state: "running" },
deliveryStatus: "pending",
});
expect(detail.task.execution?.currentTool).toBeUndefined();
expect(detail.task.execution?.wait).toBeUndefined();
expect(
await client.request<TasksListResult>("tasks.list", {
sessionKey: `agent:${OTHER_AGENT_ID}:${REQUESTER_KEY}`,
agentId: OTHER_AGENT_ID,
}),
).toEqual({ tasks: [] });
const requestsBeforeStatus = modelServer.requestCount();
expect(
await client.request("chat.send", {
sessionKey,
agentId: REQUESTER_AGENT_ID,
message: "/status",
idempotencyKey: statusRunId,
}),
).toMatchObject({ runId: statusRunId });
const reply = await withTestTimeout(
statusReply.promise,
30_000,
"status did not finish while the parent provider was held",
);
expect(reply, instance.logs()).toMatchObject({
state: "final",
runId: statusRunId,
sessionKey,
});
const statusText = "message" in reply ? extractFirstTextBlock(reply.message) : undefined;
const childLine = statusText
?.split("\n")
.find((line) => line.includes("requester-owner-child"));
expect(childLine).toMatch(/\brunning\b/);
expect(childLine).not.toMatch(/\b(waiting|approval|unknown|unavailable)\b/i);
expect(statusText).not.toContain(CHILD_MARKER);
expect(parentSettled).toBe(false);
expect(modelServer.requestCount()).toBe(requestsBeforeStatus);
childGate.resolve();
await vi.waitFor(
() => {
expect(loadSubagentRegistryFromSqlite().get(run.runId), instance.logs()).toMatchObject({
execution: { status: "terminal", outcome: { status: "ok" } },
completion: { resultText: CHILD_MARKER },
delivery: { status: "pending" },
});
},
{ interval: 50, timeout: 30_000 },
);
const finished = await client.request<TasksGetResult>("tasks.get", { taskId: child.id });
expect(finished.task).toMatchObject({
id: child.id,
runId: child.runId,
agentId: REQUESTER_AGENT_ID,
sessionKey,
childSessionKey: run.childSessionKey,
status: "completed",
execution: { state: "finished" },
deliveryStatus: "pending",
});
expect(finished.task.execution?.currentTool).toBeUndefined();
expect(finished.task.execution?.wait).toBeUndefined();
expect(parentSettled).toBe(false);
expect(modelServer.requestCount()).toBe(requestsBeforeStatus);
expect(modelServer.countRequestsContaining(CHILD_MARKER)).toBe(0);
parentGate.resolve();
expect(await parent, instance.logs()).toMatchObject({ status: "ok" });
await vi.waitFor(
async () => {
const delivered = await client.request<TasksGetResult>("tasks.get", {
taskId: child.id,
});
expect(delivered.task, instance.logs()).toMatchObject({
status: "completed",
execution: { state: "finished" },
deliveryStatus: "delivered",
});
},
{ interval: 50, timeout: 30_000 },
);
const history = await client.request<{
messages: Array<{ role?: string; content?: unknown }>;
}>("chat.history", { sessionKey, agentId: REQUESTER_AGENT_ID, limit: 30 });
const visible = history.messages.map((message) => ({
role: message.role,
text: extractFirstTextBlock(message),
}));
expect(
visible.filter(({ role, text }) => role === "user" && text === "/status"),
).toHaveLength(1);
expect(
visible.filter(({ role, text }) => role === "assistant" && text === statusText),
).toHaveLength(1);
const statusIndex = visible.findIndex(
({ role, text }) => role === "user" && text === "/status",
);
expect(visible[statusIndex + 1]).toEqual({ role: "assistant", text: statusText });
expect(
visible.filter(({ role, text }) => role === "assistant" && text === CHILD_MARKER),
).toHaveLength(1);
expect(JSON.stringify(history.messages)).not.toContain("This turn ended before a reply");
const laterModelText = modelServer
.bodies()
.slice(requestsBeforeStatus)
.flatMap((body) => {
const request = JSON.parse(body) as {
input: Array<{ role?: string; content?: unknown }>;
};
return request.input
.filter(({ role }) => role === "user" || role === "assistant")
.map(extractFirstTextBlock);
});
expect(laterModelText.some((text) => text?.includes(PARENT_PROMPT))).toBe(true);
expect(laterModelText).not.toContain("/status");
expect(laterModelText).not.toContain(statusText);
} finally {
childGate.resolve();
parentGate.resolve();
await runQaGatewayFixture(
async () => {
await withTestTimeout(
Promise.allSettled(parent ? [parent] : []),
30_000,
"parent request did not settle after releasing provider gates",
);
},
() => disconnectGatewayClient(client),
() => instance.stopGateway(),
() => closeOpenClawStateDatabaseForTest(),
);
}
},
);
it(
"delivers a private result once when the child finishes before the parent yields",
{ timeout: TEST_TIMEOUT_MS },
@ -405,6 +667,7 @@ describe("REQUESTER-OWNER requester agent id survives completion dispatch", () =
function createTestConfig(baseUrl: string): OpenClawConfig {
return {
logging: { file: "${OPENCLAW_STATE_DIR}/logs/requester-owner-e2e.log" },
plugins: { enabled: false },
agents: {
ownership: "explicit",
@ -494,6 +757,7 @@ function buildToolCallEvents(name: string, args: Record<string, unknown>): SseEv
async function startProofModelServer(options?: {
yieldAfterSpawn: Promise<void>;
childReply?: Promise<void>;
}): Promise<ProofModelServer> {
const requestBodies: string[] = [];
let parentCheckedChildren = false;
@ -546,6 +810,9 @@ async function startProofModelServer(options?: {
}
if (body.includes(CHILD_TASK) && !body.includes("function_call_output")) {
if (options?.childReply) {
await options.childReply;
}
writeOpenAiResponsesText(response, {
text: CHILD_MARKER,
responseId: `response-${++responseSequence}`,