fix(media): avoid blocking cold task lookups (#148583)

* fix(agents): prepare media requester config asynchronously

Use the existing immutable captured-config reader for legacy media requester resolution. Preserve recorded requester fast paths and reread current tasks after config preparation, with task-owner checks before selection.

The real cold image status and duplicate-action regression failed on the original synchronous fallback with 506 calling-thread SQLite prepare calls, then passed with six zero method counters through async close. All 66 scoped media/config cases pass, and independent P2 review is clean. Remaining selected static checks and compiled runtime proof are in progress before publication.

* test(agents): use lifecycle replacement factory in fixture

* fix(media): skip config for terminal-only status lookups

Preserve the active-status and prompt fast path when legacy rows are already terminal. Duplicate guards still resolve terminal requester ownership so persisted success and failure replace cached running starts.

Two terminal-only controls failed the preceding candidate before this correction. All 43 media tests pass with completed and failed duplicate coverage; the staged full candidate passed independent P2 review.

* test(tasks): prepare worker before publication race guards
This commit is contained in:
Peter Steinberger 2026-09-21 11:32:53 -07:00 • committed by GitHub
parent 6bac9f2e95
commit a31a2aba5c
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 563 additions and 24 deletions

View file

@ -51,6 +51,8 @@ The prompt is compact, with fixed sections:
Large stable content (including **Project Context** and static **Memory Recall** instructions) stays above the internal prompt cache boundary. Volatile per-turn sections (**UI Presentation**, Control UI embed guidance, **Messaging**, **Collapsible Details**, **Voice**, **Group Chat Context**, **Reactions**, **Runtime**, **Project Memory** facts, channel-specific ACP hints, delegation/orchestration mode, and the current elevated level) are appended below that boundary so local backends with prefix caches can reuse the stable workspace prefix across channel turns. Exec, subagent, and media facts use the later Runtime Context carrier to preserve the conversation-history prefix too; their capability-based instructions stay in the system prompt. The boundary is internal transport metadata: every section remains system-prompt guidance for CLI backends. Tool descriptions should avoid embedding current channel names when the accepted schema already carries that runtime detail.
Media task facts include only enabled media tools and tasks belonging to the current requester. Restored tasks without a recorded requester use the configured session owner; completed tasks are omitted.
Tooling also carries long-running-work guidance:
- use cron for future follow-up (`check back later`, reminders, recurring work) instead of `exec` sleep loops, `yieldMs` delay tricks, or repeated `process` polling

View file

@ -1,22 +1,54 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../test/helpers/promise.js";
import type { CapturedRuntimeConfigRead } from "../config/runtime-config-capture-state.js";
import type { OpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.types.js";
import type { TaskRecord } from "../tasks/task-registry.types.js";
import {
buildActiveMediaGenerationTaskPromptContext,
createMediaGenerationTaskStatusOwner,
MEDIA_GENERATION_DELIVERING_COMPLETION_PROGRESS,
recordRecentMediaGenerationTaskStartForSession,
} from "./media-generation-task-status-shared.js";
import { resetRecentMediaGenerationDuplicateGuardsForTests } from "./media-generation-task-status-shared.test-support.js";
import { buildMediaTaskRuntimeContext } from "./media-generation-task-status.js";
const taskRuntimeInternalMocks = vi.hoisted(() => ({
listFreshTasksForOwnerKey: vi.fn(),
}));
const configMocks = vi.hoisted(() => ({
getRuntimeConfig: vi.fn(),
readConfig: Object.assign(vi.fn<() => Promise<CapturedRuntimeConfigRead>>(), {
assertCurrent: vi.fn(),
}),
captureRuntimeConfigAsyncReader: vi.fn(),
}));
const ownerMocks = vi.hoisted(() => ({
assertCurrent: vi.fn(),
context: {
admission: {
databasePath: "/synthetic/media/state.sqlite",
identity: { key: "media-test", canonicalPath: "/synthetic/media/state.sqlite" },
assertCurrent: vi.fn(),
},
environment: { OPENCLAW_STATE_DIR: "/synthetic/media" },
coordinatorRuntime: { directory: "/synthetic/coordinator", keepAlive: false },
} satisfies OpenClawStateWorkerContext,
}));
vi.mock("../tasks/runtime-internal.js", () => taskRuntimeInternalMocks);
vi.mock("../config/config.js", () => configMocks);
vi.mock("../config/io.runtime.js", () => ({
captureRuntimeConfigAsyncReader: configMocks.captureRuntimeConfigAsyncReader,
}));
vi.mock("../state/openclaw-state-worker-context.js", () => ({
captureOpenClawStateWorkerContext: () => ownerMocks.context,
}));
vi.mock("../tasks/task-registry-state.js", () => ({
assertTaskRegistryOwnerCurrent: ownerMocks.assertCurrent,
}));
vi.mock("../tasks/task-registry.store.js", () => ({
getTaskRegistryStore: () => ({}),
}));
const videoTaskStatusOwner = createMediaGenerationTaskStatusOwner({
taskKind: "video_generation",
@ -48,17 +80,25 @@ function makeTask(overrides: Partial<TaskRecord> = {}): TaskRecord {
};
}
beforeEach(() => {
resetRecentMediaGenerationDuplicateGuardsForTests();
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockReset();
configMocks.getRuntimeConfig.mockReset().mockReturnValue({
const capturedConfig: CapturedRuntimeConfigRead = {
config: {
session: { scope: "global", store: "/tmp/shared-sessions.sqlite" },
agents: {
ownership: "explicit",
defaults: { sessionStore: { agentId: "ops" } },
entries: { ops: {}, research: {} },
},
});
},
env: {},
};
beforeEach(() => {
resetRecentMediaGenerationDuplicateGuardsForTests();
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockReset();
configMocks.readConfig.mockReset().mockResolvedValue(capturedConfig);
configMocks.readConfig.assertCurrent.mockReset();
configMocks.captureRuntimeConfigAsyncReader.mockReset().mockReturnValue(configMocks.readConfig);
ownerMocks.assertCurrent.mockReset();
});
describe("media generation delivery-phase prompt guard", () => {
@ -156,8 +196,261 @@ describe("media generation delivery-phase prompt guard", () => {
await videoTaskStatusOwner.findActiveTaskForSession("global", { agentId: "ops" }),
).toEqual(task);
expect(await videoTaskStatusOwner.listActiveTasksForSession("global", "research")).toEqual([]);
configMocks.readConfig.mockResolvedValue({
...capturedConfig,
config: {
...capturedConfig.config,
agents: { ...capturedConfig.config.agents, entries: { research: {} } },
},
});
expect(await videoTaskStatusOwner.listActiveTasksForSession("global", "research")).toEqual([]);
});
it.each([
{ ownerKey: "global", requesterAgentId: "ops" },
{ ownerKey: "agent:ops:main", requesterAgentId: undefined },
])("uses recorded requester identity without loading config for $ownerKey", async (identity) => {
const task = makeTask({ ...identity, agentId: "research" });
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockResolvedValue([
makeTask({ taskId: "unrelated", taskKind: "image_generation" }),
task,
]);
expect(await videoTaskStatusOwner.listActiveTasksForSession(identity.ownerKey, "ops")).toEqual([
task,
]);
expect(configMocks.readConfig).not.toHaveBeenCalled();
});
it("keeps known requesters visible when legacy config cannot be prepared", async () => {
const known = makeTask({ taskId: "known", ownerKey: "global", requesterAgentId: "ops" });
const legacy = makeTask({ taskId: "legacy", ownerKey: "global", agentId: "ops" });
const updated = { ...known, progressSummary: "Rendering final frames" };
taskRuntimeInternalMocks.listFreshTasksForOwnerKey
.mockResolvedValueOnce([legacy, known])
.mockResolvedValue([legacy, updated]);
configMocks.readConfig.mockRejectedValue(new Error("config unavailable"));
expect(await videoTaskStatusOwner.listActiveTasksForSession("global", "ops")).toEqual([
updated,
]);
});
it("includes restored legacy media in the shared prompt only for its requester", async () => {
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockResolvedValue([
makeTask({ ownerKey: "global", agentId: "research" }),
makeTask({
taskId: "music-1",
ownerKey: "global",
taskKind: "music_generation",
sourceId: "music_generate",
}),
]);
const params = {
capabilityToolNames: new Set(["video_generate", "music_generate"]),
sessionKey: "global",
};
expect(await buildMediaTaskRuntimeContext({ ...params, agentId: "ops" })).toBe(
'## Media Generation Tasks\n- tool=music_generate; task=music-1; status=running\n- tool=video_generate; task=task-1; status=running; provider_json="byteplus"',
);
expect(configMocks.readConfig).toHaveBeenCalledTimes(1);
expect(await buildMediaTaskRuntimeContext({ ...params, agentId: "research" })).toBe(
"## Media Generation Tasks\n- tool=music_generate; none\n- tool=video_generate; none",
);
});
it.each(["succeeded", "failed"] as const)(
"skips requester config for terminal-only %s status and prompt lookups",
async (status) => {
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockResolvedValue([
makeTask({ ownerKey: "global", status }),
]);
expect(await videoTaskStatusOwner.listActiveTasksForSession("global", "ops")).toEqual([]);
expect(
await buildMediaTaskRuntimeContext({
capabilityToolNames: new Set(["video_generate"]),
sessionKey: "global",
agentId: "ops",
}),
).toBe("## Media Generation Tasks\n- tool=video_generate; none");
expect(configMocks.readConfig).not.toHaveBeenCalled();
},
);
it.each(["succeeded", "failed"] as const)(
"resolves a persisted legacy %s task before applying a cached duplicate guard",
async (status) => {
const task = makeTask({ ownerKey: "global", status });
recordRecentMediaGenerationTaskStartForSession({
sessionKey: "global",
agentId: "ops",
taskKind: "video_generation",
sourcePrefix: "video_generate",
taskId: task.taskId,
runId: task.runId,
taskLabel: task.task,
requestKey: "same-request",
progressSummary: "Generating video",
});
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockResolvedValue([task]);
const duplicate = await videoTaskStatusOwner.findDuplicateGuardTaskForSession("global", {
agentId: "ops",
requestKey: "same-request",
});
expect(duplicate).toEqual(status === "succeeded" ? task : undefined);
},
);
it.each(
(["active", "duplicate", "prompt"] as const).flatMap((lookup) =>
(["completed", "deleted"] as const).flatMap((change) =>
(["resolves", "rejects"] as const).map((completion) => ({ lookup, change, completion })),
),
),
)(
"refreshes $lookup selection after a task is $change while config $completion",
async ({ lookup, change, completion }) => {
const started = createDeferred();
const config = createDeferred<CapturedRuntimeConfigRead>();
configMocks.readConfig.mockImplementation(() => {
started.resolve();
return config.promise;
});
const legacy = makeTask({ taskId: "legacy", ownerKey: "global" });
const known = makeTask({ taskId: "known", ownerKey: "global", requesterAgentId: "ops" });
let records = [legacy, known];
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockImplementation(async () => records);
const pending =
lookup === "active"
? videoTaskStatusOwner.listActiveTasksForSession("global", "ops")
: lookup === "prompt"
? buildMediaTaskRuntimeContext({
capabilityToolNames: new Set(["video_generate"]),
sessionKey: "global",
agentId: "ops",
})
: videoTaskStatusOwner.findDuplicateGuardTaskForSession("global", { agentId: "ops" });
await started.promise;
records =
change === "deleted"
? []
: records.map((task) => Object.assign({}, task, { status: "succeeded" as const }));
if (completion === "resolves") {
config.resolve(capturedConfig);
} else {
config.reject(new Error("config unavailable"));
}
expect(await pending).toEqual(
lookup === "active"
? []
: lookup === "prompt"
? "## Media Generation Tasks\n- tool=video_generate; none"
: undefined,
);
},
);
it.each(["resolves", "rejects"])(
"propagates retired task ownership when pending config %s",
async (completion) => {
const started = createDeferred();
const config = createDeferred<CapturedRuntimeConfigRead>();
configMocks.readConfig.mockImplementation(() => {
started.resolve();
return config.promise;
});
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockResolvedValue([
makeTask({ ownerKey: "global" }),
]);
const pending = videoTaskStatusOwner.findDuplicateGuardTaskForSession("global", {
agentId: "ops",
});
await started.promise;
const retired = new Error("task owner retired");
ownerMocks.assertCurrent.mockImplementation(() => {
throw retired;
});
const rejected = expect(pending).rejects.toBe(retired);
if (completion === "resolves") {
config.resolve(capturedConfig);
} else {
config.reject(new Error("config unavailable"));
}
await rejected;
},
);
it("keeps a recent start added while requester config is being prepared", async () => {
const started = createDeferred();
const config = createDeferred<CapturedRuntimeConfigRead>();
configMocks.readConfig.mockImplementation(() => {
started.resolve();
return config.promise;
});
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockResolvedValue([
makeTask({ ownerKey: "global", task: "different prompt" }),
]);
const pending = videoTaskStatusOwner.findDuplicateGuardTaskForSession("global", {
agentId: "ops",
prompt: "new request",
requestKey: "new-request-key",
});
await started.promise;
recordRecentMediaGenerationTaskStartForSession({
sessionKey: "global",
agentId: "ops",
taskKind: "video_generation",
sourcePrefix: "video_generate",
taskId: "recent-start",
taskLabel: "new request",
requestKey: "new-request-key",
progressSummary: "Generating video",
});
config.resolve(capturedConfig);
expect(await pending).toMatchObject({ taskId: "recent-start", status: "running" });
});
it.each(["active", "duplicate", "prompt"] as const)(
"rejects %s selection after ownership retires between preparation and continuation",
async (lookup) => {
let settled = false;
let queued = false;
let retired = false;
const retirement = new Error("task owner retired after preparation");
taskRuntimeInternalMocks.listFreshTasksForOwnerKey.mockImplementation(async () => {
settled = true;
return [makeTask({ ownerKey: "global", requesterAgentId: "ops" })];
});
ownerMocks.assertCurrent.mockImplementation(() => {
if (retired) {
throw retirement;
}
if (settled && !queued) {
queued = true;
queueMicrotask(() => {
retired = true;
});
}
});
const pending =
lookup === "active"
? videoTaskStatusOwner.listActiveTasksForSession("global", "ops")
: lookup === "prompt"
? buildMediaTaskRuntimeContext({
capabilityToolNames: new Set(["video_generate"]),
sessionKey: "global",
agentId: "ops",
})
: videoTaskStatusOwner.findDuplicateGuardTaskForSession("global", { agentId: "ops" });
await expect(pending).rejects.toBe(retirement);
},
);
it("blocks the same prompt while allowing a distinct prompt", async () => {
const task = makeTask({
task: "generate clip 01",

View file

@ -11,9 +11,13 @@ import {
normalizeOptionalString,
} from "@openclaw/normalization-core/string-coerce";
import { truncateUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
import { getRuntimeConfig } from "../config/config.js";
import { captureRuntimeConfigAsyncReader } from "../config/io.runtime.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { parseAgentSessionKey } from "../routing/session-key.js";
import { captureOpenClawStateWorkerContext } from "../state/openclaw-state-worker-context.js";
import { listFreshTasksForOwnerKey } from "../tasks/runtime-internal.js";
import { assertTaskRegistryOwnerCurrent } from "../tasks/task-registry-state.js";
import { getTaskRegistryStore } from "../tasks/task-registry.store.js";
import type { TaskRecord } from "../tasks/task-registry.types.js";
import { resolveSessionAgentId } from "./agent-scope.js";
import { sanitizeForPromptLiteral } from "./sanitize-for-prompt.js";
@ -93,7 +97,10 @@ function mediaGenerationTaskLabelMatches(task: TaskRecord, taskLabel: string): b
return normalizeOptionalString(task.task) === taskLabel;
}
function resolveMediaGenerationTaskRequesterAgentId(task: TaskRecord): string | undefined {
function resolveMediaGenerationTaskRequesterAgentId(
task: TaskRecord,
config?: OpenClawConfig,
): string | undefined {
const explicit = normalizeOptionalString(task.requesterAgentId);
if (explicit) {
return explicit;
@ -103,16 +110,68 @@ function resolveMediaGenerationTaskRequesterAgentId(task: TaskRecord): string |
if (parsed) {
return parsed;
}
if (!ownerKey) {
if (!ownerKey || !config) {
return undefined;
}
try {
return resolveSessionAgentId({ config: getRuntimeConfig(), sessionKey: ownerKey });
return resolveSessionAgentId({ config, sessionKey: ownerKey });
} catch {
return undefined;
}
}
export async function prepareMediaGenerationTaskLookup(params: {
sessionKey: string;
agentId?: string;
taskIdentities: readonly { taskKind: string; sourcePrefix: string }[];
includeTerminalTasks?: boolean;
}) {
const context = captureOpenClawStateWorkerContext();
const store = getTaskRegistryStore();
const assertTaskCurrent = () => assertTaskRegistryOwnerCurrent(context, store);
const readConfig = params.agentId
? captureRuntimeConfigAsyncReader({ assertCurrent: assertTaskCurrent, capture: true })
: undefined;
let configReadStarted = false;
const assertCurrent = () => {
assertTaskCurrent();
if (configReadStarted) {
readConfig?.assertCurrent();
}
};
let tasks = await listFreshTasksForOwnerKey(params.sessionKey);
assertCurrent();
let config: OpenClawConfig | undefined;
if (
readConfig &&
tasks.some(
(task) =>
task.runtime === "cli" &&
task.scopeKind === "session" &&
(params.includeTerminalTasks || isTaskStillBlockingDuplicateGuard(task)) &&
params.taskIdentities.some(
({ taskKind, sourcePrefix }) =>
task.taskKind === taskKind &&
(!sourcePrefix || mediaGenerationSourceMatches(task, sourcePrefix)),
) &&
Boolean(normalizeOptionalString(task.ownerKey ?? task.requesterSessionKey)) &&
!resolveMediaGenerationTaskRequesterAgentId(task),
)
) {
configReadStarted = true;
try {
({ config } = await readConfig());
} catch {
// Unreadable config keeps legacy requester ownership unresolved.
}
assertCurrent();
// Config preparation may outlive a task's completion or deletion.
tasks = await listFreshTasksForOwnerKey(params.sessionKey);
}
assertCurrent();
return { tasks, config, assertCurrent };
}
function isTaskStillBlockingDuplicateGuard(task: TaskRecord): boolean {
return task.status === "queued" || task.status === "running";
}
@ -151,6 +210,7 @@ function recentMediaGenerationTaskStartMatches(
function findPersistedTaskForRecentMediaGenerationStart(params: {
tasks: readonly TaskRecord[];
config?: OpenClawConfig;
agentId?: string;
cachedTask: TaskRecord;
taskKind: string;
@ -162,7 +222,8 @@ function findPersistedTaskForRecentMediaGenerationStart(params: {
task.scopeKind !== "session" ||
task.taskKind !== params.taskKind ||
!mediaGenerationSourceMatches(task, params.sourcePrefix) ||
(params.agentId && resolveMediaGenerationTaskRequesterAgentId(task) !== params.agentId)
(params.agentId &&
resolveMediaGenerationTaskRequesterAgentId(task, params.config) !== params.agentId)
) {
return false;
}
@ -240,6 +301,7 @@ export function recordRecentMediaGenerationTaskStartForSession(params: {
/** Finds a recent started media task from memory or persisted task state. */
function findRecentStartedMediaGenerationTaskForSession(params: {
tasks: readonly TaskRecord[];
config?: OpenClawConfig;
sessionKey?: string;
agentId?: string;
taskKind: string;
@ -269,6 +331,7 @@ function findRecentStartedMediaGenerationTaskForSession(params: {
const task = entry.task;
const persistedTask = findPersistedTaskForRecentMediaGenerationStart({
agentId: params.agentId,
config: params.config,
tasks: params.tasks,
cachedTask: task,
taskKind: params.taskKind,
@ -361,12 +424,19 @@ async function listActiveMediaGenerationTasksForSession(params: {
if (!sessionKey) {
return [];
}
return selectActiveMediaGenerationTasks(params, await listFreshTasksForOwnerKey(sessionKey));
const lookup = await prepareMediaGenerationTaskLookup({
sessionKey,
agentId: params.agentId,
taskIdentities: [params],
});
lookup.assertCurrent();
return selectActiveMediaGenerationTasks(params, lookup.tasks, lookup.config);
}
function selectActiveMediaGenerationTasks(
params: Parameters<typeof listActiveMediaGenerationTasksForSession>[0],
tasks: readonly TaskRecord[],
config?: OpenClawConfig,
): TaskRecord[] {
const taskLabel = normalizeOptionalString(params.taskLabel);
const sourcePrefix = normalizeOptionalString(params.sourcePrefix);
@ -379,7 +449,10 @@ function selectActiveMediaGenerationTasks(
) {
return false;
}
if (params.agentId && resolveMediaGenerationTaskRequesterAgentId(task) !== params.agentId) {
if (
params.agentId &&
resolveMediaGenerationTaskRequesterAgentId(task, config) !== params.agentId
) {
return false;
}
if (sourcePrefix && !mediaGenerationSourceMatches(task, sourcePrefix)) {
@ -416,10 +489,16 @@ async function findDuplicateGuardMediaGenerationTaskForSession(params: {
if (!sessionKey) {
return undefined;
}
const tasks = await listFreshTasksForOwnerKey(sessionKey);
const lookup = await prepareMediaGenerationTaskLookup({
sessionKey,
agentId: params.agentId,
taskIdentities: [params],
includeTerminalTasks: true,
});
lookup.assertCurrent();
return (
findRecentStartedMediaGenerationTaskForSession({ ...params, sessionKey, tasks }) ??
selectActiveMediaGenerationTasks(params, tasks)[0]
findRecentStartedMediaGenerationTaskForSession({ ...params, sessionKey, ...lookup }) ??
selectActiveMediaGenerationTasks(params, lookup.tasks, lookup.config)[0]
);
}
@ -509,6 +588,7 @@ function buildMediaGenerationTaskStatusListText(params: {
/** Builds bounded current-turn facts without instructions or elapsed-time fields. */
export function buildActiveMediaGenerationTaskPromptContext(params: {
tasks: readonly TaskRecord[];
config?: OpenClawConfig;
agentId?: string;
taskKind: string;
sourcePrefix: string;
@ -516,6 +596,7 @@ export function buildActiveMediaGenerationTaskPromptContext(params: {
const tasks = selectActiveMediaGenerationTasks(
{ ...params, excludeDeliveringCompletion: true },
params.tasks,
params.config,
);
if (tasks.length === 0) {
return undefined;

View file

@ -0,0 +1,148 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { getRuntimeConfigSnapshot, setRuntimeConfigSnapshot } from "../config/runtime-snapshot.js";
import { createPluginCache, withPluginCache } from "../plugins/plugin-cache.js";
import { closeOpenClawStateDatabaseAsync } from "../state/openclaw-state-db.js";
import * as taskRuntime from "../tasks/runtime-internal.js";
import { upsertTaskWithDeliveryStateToSqlite } from "../tasks/task-registry.store.sqlite.js";
import {
resetTaskRegistryForTests,
resetTaskFlowRegistryForTests,
} from "../tasks/task-runtime.test-helpers.js";
import { observeMainThreadSql } from "../test-utils/main-thread-sql-spies.js";
import {
createOpenClawTestState,
type OpenClawTestState,
} from "../test-utils/openclaw-test-state.js";
import {
buildMediaTaskRuntimeContext,
IMAGE_GENERATION_TASK_KIND,
} from "./media-generation-task-status.js";
import {
createImageGenerateDuplicateGuardResult,
createImageGenerateStatusActionResult,
} from "./tools/image-generate-tool.actions.js";
let state: OpenClawTestState;
beforeEach(async () => {
state = await createOpenClawTestState({
layout: "state-only",
prefix: "media-task-status-cold-",
});
resetTaskRegistryForTests({ persist: false });
resetTaskFlowRegistryForTests({ persist: false });
await state.writeConfig({
gateway: { mode: "local" },
session: { scope: "global", store: state.statePath("legacy-sessions.sqlite") },
agents: {
ownership: "explicit",
defaults: { sessionStore: { agentId: "ops" } },
entries: { ops: {}, research: {} },
},
});
vi.spyOn(process, "cwd").mockReturnValue(state.workspaceDir);
upsertTaskWithDeliveryStateToSqlite({
task: {
taskId: "legacy-media",
runtime: "cli",
requesterSessionKey: "global",
ownerKey: "global",
scopeKind: "session",
task: "Synthetic restore",
status: "running",
deliveryStatus: "not_applicable",
notifyPolicy: "silent",
createdAt: 10,
runId: "legacy-media-run",
agentId: "research",
taskKind: IMAGE_GENERATION_TASK_KIND,
sourceId: "image_generate:synthetic",
},
});
await closeOpenClawStateDatabaseAsync();
});
afterEach(async () => {
resetTaskRegistryForTests({ persist: false });
resetTaskFlowRegistryForTests({ persist: false });
vi.restoreAllMocks();
await state.cleanup();
});
describe("cold media generation task status", () => {
it("resolves legacy requester identity from cold config without parent SQLite through close", async () => {
const mainSql = observeMainThreadSql();
expect(getRuntimeConfigSnapshot()).toBeNull();
expect(
await buildMediaTaskRuntimeContext({
capabilityToolNames: new Set(["image_generate"]),
sessionKey: "global",
agentId: "ops",
}),
).toBe(
'## Media Generation Tasks\n- tool=image_generate; task=legacy-media; status=running; provider_json="synthetic"',
);
await withPluginCache(createPluginCache(), async () => {
expect((await createImageGenerateStatusActionResult("global", "ops")).details).toMatchObject({
active: true,
task: { taskId: "legacy-media" },
});
expect(
(await createImageGenerateStatusActionResult("global", "research")).details,
).toMatchObject({
active: false,
});
expect(
(await createImageGenerateDuplicateGuardResult("global", { agentId: "ops" }))?.details,
).toMatchObject({ task: { taskId: "legacy-media" } });
});
await closeOpenClawStateDatabaseAsync();
mainSql.expectIdle();
});
it.each([1, 2])("rejects a changed config source during task read %i", async (readToChange) => {
const listTasks = taskRuntime.listFreshTasksForOwnerKey;
let reads = 0;
vi.spyOn(taskRuntime, "listFreshTasksForOwnerKey").mockImplementation(async (ownerKey) => {
const tasks = await listTasks(ownerKey);
if (++reads === readToChange) {
process.env.OPENCLAW_CONFIG_PATH = state.statePath("replacement.json");
}
return tasks;
});
await expect(
buildMediaTaskRuntimeContext({
capabilityToolNames: new Set(["image_generate"]),
sessionKey: "global",
agentId: "ops",
}),
).rejects.toThrow("Runtime config source changed");
expect(reads).toBe(readToChange);
});
it("retains its captured requester across an ordinary same-source config reload", async () => {
const listTasks = taskRuntime.listFreshTasksForOwnerKey;
let reads = 0;
vi.spyOn(taskRuntime, "listFreshTasksForOwnerKey").mockImplementation(async (ownerKey) => {
const tasks = await listTasks(ownerKey);
if (++reads === 2) {
setRuntimeConfigSnapshot({
agents: {
ownership: "explicit",
defaults: { sessionStore: { agentId: "research" } },
entries: { ops: {}, research: {} },
},
});
}
return tasks;
});
expect(
await buildMediaTaskRuntimeContext({
capabilityToolNames: new Set(["image_generate"]),
sessionKey: "global",
agentId: "ops",
}),
).toContain("task=legacy-media; status=running");
expect(reads).toBe(2);
});
});

View file

@ -28,6 +28,9 @@ const taskRuntimeInternalMocks = vi.hoisted(() => {
});
vi.mock("../tasks/runtime-internal.js", () => taskRuntimeInternalMocks);
vi.mock("../tasks/task-registry-state.js", () => ({
assertTaskRegistryOwnerCurrent: vi.fn(),
}));
function expectActiveImageGenerationTask(
task: Awaited<ReturnType<typeof findDuplicateGuardImageGenerationTaskForSession>>,

View file

@ -5,10 +5,10 @@
* source id, duplicate-guard timing, and prompt/status wording.
*/
import { normalizeOptionalString } from "@openclaw/normalization-core/string-coerce";
import { listFreshTasksForOwnerKey } from "../tasks/runtime-internal.js";
import {
buildActiveMediaGenerationTaskPromptContext,
createMediaGenerationTaskStatusOwner,
prepareMediaGenerationTaskLookup,
} from "./media-generation-task-status-shared.js";
export const IMAGE_GENERATION_TASK_KIND = "image_generation";
@ -91,11 +91,19 @@ export async function buildMediaTaskRuntimeContext(params: {
return undefined;
}
const sessionKey = normalizeOptionalString(params.sessionKey);
const tasks = sessionKey ? await listFreshTasksForOwnerKey(sessionKey) : [];
const lookup = sessionKey
? await prepareMediaGenerationTaskLookup({
sessionKey,
agentId: params.agentId,
taskIdentities: enabled.map(([sourcePrefix, taskKind]) => ({ sourcePrefix, taskKind })),
})
: undefined;
lookup?.assertCurrent();
const facts = enabled.map(
([tool, taskKind]) =>
buildActiveMediaGenerationTaskPromptContext({
tasks,
tasks: lookup?.tasks ?? [],
config: lookup?.config,
agentId: params.agentId,
taskKind,
sourcePrefix: tool,

View file

@ -138,18 +138,20 @@ export function getRuntimeConfig(options?: {
return loadConfig(options);
}
type RuntimeConfigAsyncReader<T> = (() => Promise<T>) & { assertCurrent: () => void };
/** Capture the config source before a task read, and load only if its owner needs config facts. */
export function captureRuntimeConfigAsyncReader(options: {
assertCurrent?: () => void;
capture: true;
}): () => Promise<CapturedRuntimeConfigRead>;
}): RuntimeConfigAsyncReader<CapturedRuntimeConfigRead>;
export function captureRuntimeConfigAsyncReader(options?: {
assertCurrent?: () => void;
capture?: false;
}): () => Promise<OpenClawConfig>;
}): RuntimeConfigAsyncReader<OpenClawConfig>;
export function captureRuntimeConfigAsyncReader(
options: { assertCurrent?: () => void; capture?: boolean } = {},
): () => Promise<OpenClawConfig | CapturedRuntimeConfigRead> {
): RuntimeConfigAsyncReader<OpenClawConfig | CapturedRuntimeConfigRead> {
const sourceEnv = process.env;
const cwd = tryProcessCwd();
const readSelectors = () =>
@ -186,7 +188,7 @@ export function captureRuntimeConfigAsyncReader(
},
});
let pending: Promise<OpenClawConfig | CapturedRuntimeConfigRead> | undefined;
return () => {
const read = () => {
assertCurrent();
const loadFresh = async (assertPinned: () => void) => {
try {
@ -211,6 +213,7 @@ export function captureRuntimeConfigAsyncReader(
? loadPinnedRuntimeConfigAsync(loadFresh, { assertCurrent, capture: true })
: loadPinnedRuntimeConfigAsync(loadFresh, { assertCurrent }));
};
return Object.assign(read, { assertCurrent });
}
function createCurrentConfigReader(params: {

View file

@ -343,6 +343,7 @@ export const databaseWorkerCoreTestFiles = [
"src/tasks/task-registry.restore-ownership.test.ts",
"src/agents/agent-harness-completion-delivery.test.ts",
"src/agents/openclaw-tools.subagents.scope.test.ts",
"src/agents/media-generation-task-status.cold.test.ts",
"src/agents/tools/media-generate-tool.donor-resources.test.ts",
"src/agents/tools/media-generate-tool.resources.test.ts",
"src/agents/embedded-agent-runner/context-engine-maintenance.lifecycle.test.ts",