diff --git a/docs/concepts/system-prompt.md b/docs/concepts/system-prompt.md index 117e1bed895d..0fc4b388ea9c 100644 --- a/docs/concepts/system-prompt.md +++ b/docs/concepts/system-prompt.md @@ -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 diff --git a/src/agents/media-generation-task-status-shared.test.ts b/src/agents/media-generation-task-status-shared.test.ts index 36fcdf46a6cf..350dc75dae46 100644 --- a/src/agents/media-generation-task-status-shared.test.ts +++ b/src/agents/media-generation-task-status-shared.test.ts @@ -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>(), { + 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 { }; } -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(); + 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(); + 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(); + 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", diff --git a/src/agents/media-generation-task-status-shared.ts b/src/agents/media-generation-task-status-shared.ts index 6574beda9e18..45597d271391 100644 --- a/src/agents/media-generation-task-status-shared.ts +++ b/src/agents/media-generation-task-status-shared.ts @@ -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[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; diff --git a/src/agents/media-generation-task-status.cold.test.ts b/src/agents/media-generation-task-status.cold.test.ts new file mode 100644 index 000000000000..ea76530a38d1 --- /dev/null +++ b/src/agents/media-generation-task-status.cold.test.ts @@ -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); + }); +}); diff --git a/src/agents/media-generation-task-status.test.ts b/src/agents/media-generation-task-status.test.ts index fc337e5089ca..07e80857f977 100644 --- a/src/agents/media-generation-task-status.test.ts +++ b/src/agents/media-generation-task-status.test.ts @@ -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>, diff --git a/src/agents/media-generation-task-status.ts b/src/agents/media-generation-task-status.ts index 4b65fe6f417b..14565e3d8352 100644 --- a/src/agents/media-generation-task-status.ts +++ b/src/agents/media-generation-task-status.ts @@ -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, diff --git a/src/config/io.runtime.ts b/src/config/io.runtime.ts index d9e259a91599..c03323ec91a3 100644 --- a/src/config/io.runtime.ts +++ b/src/config/io.runtime.ts @@ -138,18 +138,20 @@ export function getRuntimeConfig(options?: { return loadConfig(options); } +type RuntimeConfigAsyncReader = (() => Promise) & { 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; +}): RuntimeConfigAsyncReader; export function captureRuntimeConfigAsyncReader(options?: { assertCurrent?: () => void; capture?: false; -}): () => Promise; +}): RuntimeConfigAsyncReader; export function captureRuntimeConfigAsyncReader( options: { assertCurrent?: () => void; capture?: boolean } = {}, -): () => Promise { +): RuntimeConfigAsyncReader { const sourceEnv = process.env; const cwd = tryProcessCwd(); const readSelectors = () => @@ -186,7 +188,7 @@ export function captureRuntimeConfigAsyncReader( }, }); let pending: Promise | 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: { diff --git a/test/vitest/vitest.database-worker-core-paths.mjs b/test/vitest/vitest.database-worker-core-paths.mjs index a34e35f89299..6c17b7a1c2e8 100644 --- a/test/vitest/vitest.database-worker-core-paths.mjs +++ b/test/vitest/vitest.database-worker-core-paths.mjs @@ -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",