diff --git a/packages/gateway-protocol/src/schema/task-suggestions.ts b/packages/gateway-protocol/src/schema/task-suggestions.ts index 2d84c17b0c87..c55d1dffb789 100644 --- a/packages/gateway-protocol/src/schema/task-suggestions.ts +++ b/packages/gateway-protocol/src/schema/task-suggestions.ts @@ -4,9 +4,9 @@ import { Type } from "typebox"; import { closedObject } from "./closed-object.js"; const TaskIdSchema = Type.String({ minLength: 1, maxLength: 128 }); -const TaskTitleSchema = Type.String({ minLength: 1, maxLength: 60 }); -const TaskPromptSchema = Type.String({ minLength: 1, maxLength: 32_768 }); -const TaskTldrSchema = Type.String({ minLength: 1, maxLength: 1_024 }); +const TaskTitleSchema = Type.String({ minLength: 1, maxLength: 60, pattern: "\\S" }); +const TaskPromptSchema = Type.String({ minLength: 1, maxLength: 32_768, pattern: "\\S" }); +const TaskTldrSchema = Type.String({ minLength: 1, maxLength: 1_024, pattern: "\\S" }); const TaskCwdSchema = Type.String({ minLength: 1, maxLength: 4_096 }); const TaskSessionKeySchema = Type.String({ minLength: 1, maxLength: 512 }); const TaskAgentIdSchema = Type.String({ minLength: 1, maxLength: 128 }); diff --git a/src/agents/tools/task-suggestion-tools.test.ts b/src/agents/tools/task-suggestion-tools.test.ts index f73f01e44a20..0f492a9b955b 100644 --- a/src/agents/tools/task-suggestion-tools.test.ts +++ b/src/agents/tools/task-suggestion-tools.test.ts @@ -43,6 +43,14 @@ describe("task suggestion tools", () => { expect(result?.content).toEqual([ { type: "text", text: JSON.stringify({ task_id: "task_123" }, null, 2) }, ]); + expect(spawnTask?.description).toContain("absolute path inside a git checkout"); + expect(spawnTask?.parameters).toMatchObject({ + properties: { + cwd: { + description: "Absolute path inside a git checkout; defaults to the current project.", + }, + }, + }); expect(spawnTask?.outputSchema).toBeDefined(); expect(Value.Check(spawnTask!.outputSchema!, result?.details)).toBe(true); expect(compactToolOutputHint(spawnTask?.outputSchema)).toBe("{ task_id: string }"); diff --git a/src/agents/tools/task-suggestion-tools.ts b/src/agents/tools/task-suggestion-tools.ts index 4022784a5ec9..571e13be6a70 100644 --- a/src/agents/tools/task-suggestion-tools.ts +++ b/src/agents/tools/task-suggestion-tools.ts @@ -33,7 +33,7 @@ const SpawnTaskToolSchema = Type.Object( Type.String({ minLength: 1, maxLength: 4_096, - description: "Absolute project directory; defaults to the current project.", + description: "Absolute path inside a git checkout; defaults to the current project.", }), ), }, @@ -75,7 +75,7 @@ export function createTaskSuggestionTools(params: { displaySummary: SPAWN_TASK_TOOL_DISPLAY_SUMMARY, description: [ "Suggest confirmed valuable out-of-scope follow-up: dead code, stale docs, missing coverage, verified TODO, security issue.", - "Operator suggestion only; does not start work.", + "Operator suggestion only; does not start work. cwd must be an absolute path inside a git checkout.", ].join(" "), parameters: SpawnTaskToolSchema, outputSchema: SpawnTaskOutputSchema, diff --git a/src/gateway/server-methods/task-suggestions.test.ts b/src/gateway/server-methods/task-suggestions.test.ts index 50b7dc529dd3..a36d2f236407 100644 --- a/src/gateway/server-methods/task-suggestions.test.ts +++ b/src/gateway/server-methods/task-suggestions.test.ts @@ -16,9 +16,10 @@ type Method = | "taskSuggestions.accept" | "taskSuggestions.dismiss"; -async function call(method: Method, params: Record) { +const GIT_CWD = process.cwd(); + +async function call(method: Method, params: Record, broadcast = vi.fn()) { const calls: Parameters[] = []; - const broadcast = vi.fn(); const respond: RespondFn = (...args) => { calls.push(args); }; @@ -55,10 +56,10 @@ afterEach(async () => { describe("task suggestion gateway methods", () => { it("creates, lists, and resolves an ephemeral suggestion", async () => { const created = await call("taskSuggestions.create", { - title: "Remove stale adapter", - prompt: "Delete src/example.ts and update its tests.", - tldr: "The adapter is unreachable and adds maintenance cost.", - cwd: "/repo", + title: " Remove stale adapter ", + prompt: " Delete src/example.ts and update its tests. ", + tldr: " The adapter is unreachable and adds maintenance cost. ", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); const payload = requirePayload(created) as { taskId: string }; @@ -67,7 +68,12 @@ describe("task suggestion gateway methods", () => { "task.suggestion", expect.objectContaining({ action: "created", - suggestion: expect.objectContaining({ agentId: "main" }), + suggestion: expect.objectContaining({ + agentId: "main", + title: "Remove stale adapter", + prompt: "Delete src/example.ts and update its tests.", + tldr: "The adapter is unreachable and adds maintenance cost.", + }), }), { dropIfSlow: true }, ); @@ -77,7 +83,15 @@ describe("task suggestion gateway methods", () => { agentId: "main", }); expect(listed.response?.[1]).toMatchObject({ - suggestions: [{ id: payload.taskId, cwd: "/repo" }], + suggestions: [ + { + id: payload.taskId, + cwd: GIT_CWD, + title: "Remove stale adapter", + prompt: "Delete src/example.ts and update its tests.", + tldr: "The adapter is unreachable and adds maintenance cost.", + }, + ], }); const resolved = await call("taskSuggestions.dismiss", { @@ -94,12 +108,12 @@ describe("task suggestion gateway methods", () => { expect(empty.response?.[1]).toEqual({ suggestions: [] }); }); - it("preserves accepted-session replay when successful admission expires a pending suggestion", async () => { + it("evicts accepted-session replay before an unseen pending suggestion", async () => { const created = await call("taskSuggestions.create", { title: "Remove stale adapter", prompt: "Delete src/example.ts and update its tests.", tldr: "The adapter is unreachable and adds maintenance cost.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", agentId: "main", }); @@ -114,7 +128,7 @@ describe("task suggestion gateway methods", () => { label: "Remove stale adapter", task: "Delete src/example.ts and update its tests.", worktree: true, - cwd: "/repo", + cwd: GIT_CWD, }); sessionKey = (params as { key: string }).key; expect(sessionKey).toMatch(/^agent:main:dashboard:/); @@ -122,13 +136,13 @@ describe("task suggestion gateway methods", () => { }); const first = await call("taskSuggestions.accept", { taskId }); - for (let index = 0; index < 100; index += 1) { + for (let index = 0; index < 99; index += 1) { requirePayload( await call("taskSuggestions.create", { title: `Pending follow up ${index}`, prompt: `Complete pending follow-up task ${index}.`, tldr: "The operator has not accepted this follow-up.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }), ); @@ -141,23 +155,28 @@ describe("task suggestion gateway methods", () => { expect(oldestPending).toBeDefined(); const admitted = await call("taskSuggestions.create", { title: "Latest follow up", - prompt: "Preserve the accepted task while admitting new work.", - tldr: "Evict only the oldest pending follow-up.", - cwd: "/repo", + prompt: "Prefer unseen pending work over accepted replay state.", + tldr: "Evict only the completed task replay.", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); expect(admitted.response?.[0]).toBe(true); - expect(admitted.broadcast).toHaveBeenNthCalledWith( - 1, + expect(admitted.broadcast).toHaveBeenCalledTimes(1); + expect(admitted.broadcast).toHaveBeenCalledWith( "task.suggestion", - { action: "resolved", taskId: oldestPending?.id, resolution: "expired" }, + expect.objectContaining({ action: "created" }), { dropIfSlow: true }, ); const retry = await call("taskSuggestions.accept", { taskId }); + const listed = await call("taskSuggestions.list", {}); expect(first.response?.[1]).toEqual({ taskId, key: sessionKey }); - expect(retry.response?.[1]).toEqual({ taskId, key: sessionKey }); + expect(retry.response?.[0]).toBe(false); + expect(retry.response?.[2]).toMatchObject({ code: "INVALID_REQUEST" }); + expect( + (requirePayload(listed) as { suggestions: Array<{ id: string }> }).suggestions, + ).toContainEqual(expect.objectContaining({ id: oldestPending?.id })); expect(createSession).toHaveBeenCalledTimes(1); expect(first.broadcast).toHaveBeenCalledWith( "task.suggestion", @@ -179,7 +198,7 @@ describe("task suggestion gateway methods", () => { title: `Accepted follow up ${index}`, prompt: `Complete accepted follow-up task ${index}.`, tldr: "This follow-up already created its managed task session.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); const taskId = (requirePayload(created) as { taskId: string }).taskId; @@ -192,7 +211,7 @@ describe("task suggestion gateway methods", () => { title: "Latest follow up", prompt: "Keep accepting new suggestions after earlier tasks completed.", tldr: "Accepted-session replay is bounded best-effort state.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); @@ -214,7 +233,7 @@ describe("task suggestion gateway methods", () => { title: "Add coverage", prompt: "Add the missing regression test.", tldr: "The edge case is untested.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); const taskId = (requirePayload(created) as { taskId: string }).taskId; @@ -244,7 +263,7 @@ describe("task suggestion gateway methods", () => { title: "Add coverage", prompt: "Add the missing regression test.", tldr: "The edge case is untested.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", agentId: "main", }); @@ -298,7 +317,7 @@ describe("task suggestion gateway methods", () => { title: "Add coverage", prompt: "Add the missing regression test.", tldr: "The edge case is untested.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", agentId: "main", }); @@ -329,7 +348,7 @@ describe("task suggestion gateway methods", () => { title: "Add coverage", prompt: "Add the missing regression test.", tldr: "The edge case is untested.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", agentId: "main", }); @@ -356,6 +375,47 @@ describe("task suggestion gateway methods", () => { expect(listed.response?.[1]).toEqual({ suggestions: [] }); }); + it("abandons an acceptance when rollback inspection throws", async () => { + const created = await call("taskSuggestions.create", { + title: "Add coverage", + prompt: "Add the missing regression test.", + tldr: "The edge case is untested.", + cwd: GIT_CWD, + sessionKey: "agent:main:main", + agentId: "main", + }); + const taskId = (requirePayload(created) as { taskId: string }).taskId; + vi.spyOn(sessionCreateHandlers, "sessions.create").mockRejectedValue( + new Error("initial dispatch failed"), + ); + vi.spyOn(sessionDeleteHandlers, "sessions.delete").mockImplementation(async ({ respond }) => { + respond(true, { ok: true, deleted: true }, undefined); + }); + vi.spyOn(managedWorktrees, "findLiveByOwner").mockImplementation(() => { + throw new Error("worktree registry unavailable"); + }); + const broadcast = vi.fn(); + + await expect(call("taskSuggestions.accept", { taskId }, broadcast)).rejects.toThrow( + "worktree registry unavailable", + ); + + expect(broadcast).toHaveBeenCalledTimes(1); + expect(broadcast).toHaveBeenCalledWith( + "task.suggestion", + { action: "resolved", taskId, resolution: "expired" }, + { dropIfSlow: true }, + ); + const listed = await call("taskSuggestions.list", {}); + expect(listed.response?.[1]).toEqual({ suggestions: [] }); + const retry = await call("taskSuggestions.accept", { taskId }); + expect(retry.response?.[0]).toBe(false); + expect(retry.response?.[2]).toMatchObject({ + code: "INVALID_REQUEST", + message: "task suggestion cannot be accepted: dismissed", + }); + }); + it("rejects a relative cwd before recording or broadcasting", async () => { const result = await call("taskSuggestions.create", { title: "Add coverage", @@ -370,12 +430,50 @@ describe("task suggestion gateway methods", () => { expect(result.broadcast).not.toHaveBeenCalled(); }); + it("rejects a cwd outside a git checkout before recording or broadcasting", async () => { + const result = await call("taskSuggestions.create", { + title: "Add coverage", + prompt: "Add the missing regression test.", + tldr: "The edge case is untested.", + cwd: "/not-a-git-checkout", + sessionKey: "agent:main:main", + }); + + expect(result.response?.[0]).toBe(false); + expect(result.response?.[2]).toMatchObject({ + code: "INVALID_REQUEST", + message: "task suggestion cwd must be inside a git checkout", + }); + expect(result.broadcast).not.toHaveBeenCalled(); + }); + + it.each(["title", "prompt", "tldr"] as const)( + "rejects whitespace-only %s before recording or broadcasting", + async (field) => { + const params = { + title: "Add coverage", + prompt: "Add the missing regression test.", + tldr: "The edge case is untested.", + cwd: GIT_CWD, + sessionKey: "agent:main:main", + }; + params[field] = " \n\t "; + + const result = await call("taskSuggestions.create", params); + + expect(result.response?.[0]).toBe(false); + expect(result.response?.[2]).toMatchObject({ code: "INVALID_REQUEST" }); + expect(result.response?.[2]?.message).toContain(field); + expect(result.broadcast).not.toHaveBeenCalled(); + }, + ); + it("rejects an agent that conflicts with the source session", async () => { const result = await call("taskSuggestions.create", { title: "Add coverage", prompt: "Add the missing regression test.", tldr: "The edge case is untested.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", agentId: "work", }); @@ -393,7 +491,7 @@ describe("task suggestion gateway methods", () => { title: "Add coverage", prompt: "x".repeat(32_769), tldr: "The edge case is untested.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); @@ -409,7 +507,7 @@ describe("task suggestion gateway methods", () => { title: `Follow up ${index}`, prompt: `${index}: ${"x".repeat(32_760)}`, tldr: "The follow-up remains useful.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); taskIds.push((requirePayload(created) as { taskId: string }).taskId); @@ -430,7 +528,7 @@ describe("task suggestion gateway methods", () => { title: `Follow up ${index}`, prompt: `Complete follow-up task ${index}.`, tldr: `Follow-up task ${index} remains useful.`, - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); requirePayload(created); @@ -445,7 +543,7 @@ describe("task suggestion gateway methods", () => { title: "Latest follow up", prompt: "Complete the latest follow-up task.", tldr: "The latest follow-up remains useful.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); @@ -472,7 +570,7 @@ describe("task suggestion gateway methods", () => { title: `Running follow up ${index}`, prompt: largePrompt, tldr: "This accepted task is still starting.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); const taskId = (requirePayload(created) as { taskId: string }).taskId; @@ -484,7 +582,7 @@ describe("task suggestion gateway methods", () => { title: "Preserve accepted task", prompt: "Keep its accepted result available for retries.", tldr: "A rejected admission must not discard completed results.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); const acceptedTaskId = (requirePayload(accepted) as { taskId: string }).taskId; @@ -495,7 +593,7 @@ describe("task suggestion gateway methods", () => { title: "Keep this follow up", prompt: "Do not discard this pending task.", tldr: "The operator has not accepted it yet.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); const pendingTaskId = (requirePayload(pending) as { taskId: string }).taskId; @@ -503,7 +601,7 @@ describe("task suggestion gateway methods", () => { title: "One oversized follow up", prompt: largePrompt, tldr: "This valid task cannot fit beside protected tasks.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); @@ -536,7 +634,7 @@ describe("task suggestion gateway methods", () => { title: `Follow up ${index}`, prompt: `Complete follow-up task ${index}.`, tldr: `Follow-up task ${index} remains useful.`, - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); const taskId = (requirePayload(created) as { taskId: string }).taskId; @@ -548,7 +646,7 @@ describe("task suggestion gateway methods", () => { title: "One too many", prompt: "Complete one more follow-up task.", tldr: "This follow-up can wait until capacity returns.", - cwd: "/repo", + cwd: GIT_CWD, sessionKey: "agent:main:main", }); diff --git a/src/gateway/server-methods/task-suggestions.ts b/src/gateway/server-methods/task-suggestions.ts index c86d28fcc9c8..aa8fe7e9d2bc 100644 --- a/src/gateway/server-methods/task-suggestions.ts +++ b/src/gateway/server-methods/task-suggestions.ts @@ -12,6 +12,7 @@ import { validateTaskSuggestionsListParams, } from "../../../packages/gateway-protocol/src/index.js"; import { resolveDefaultAgentId } from "../../agents/agent-scope.js"; +import { insideGitCheckout } from "../../agents/worktrees/git.js"; import { managedWorktrees } from "../../agents/worktrees/service.js"; import { formatErrorMessage } from "../../infra/errors.js"; import { normalizeAgentId, parseAgentSessionKey } from "../../routing/session-key.js"; @@ -43,6 +44,19 @@ type TaskSuggestionAcceptanceResult = const activeAcceptances = new Map>(); +function abandonSuggestedTaskAcceptance( + taskId: string, + options: GatewayRequestHandlerOptions, +): void { + if (abandonTaskSuggestionAcceptance(taskId)) { + options.context.broadcast( + "task.suggestion", + { action: "resolved", taskId, resolution: "expired" }, + { dropIfSlow: true }, + ); + } +} + async function rollbackSuggestedTaskSession(params: { key: string; agentId?: string; @@ -119,13 +133,7 @@ async function failSuggestedTaskSession(params: { } return { ok: false, error: params.error }; } - if (abandonTaskSuggestionAcceptance(params.taskId)) { - params.options.context.broadcast( - "task.suggestion", - { action: "resolved", taskId: params.taskId, resolution: "expired" }, - { dropIfSlow: true }, - ); - } + abandonSuggestedTaskAcceptance(params.taskId, params.options); return { ok: false, error: errorShape( @@ -258,6 +266,14 @@ export const taskSuggestionsHandlers: GatewayRequestHandlers = { ); return; } + if (!insideGitCheckout(params.cwd)) { + respond( + false, + undefined, + errorShape(ErrorCodes.INVALID_REQUEST, "task suggestion cwd must be inside a git checkout"), + ); + return; + } const sessionAgentId = parseAgentSessionKey(params.sessionKey)?.agentId; const requestedAgentId = params.agentId ? normalizeAgentId(params.agentId) : undefined; if ( @@ -342,6 +358,9 @@ export const taskSuggestionsHandlers: GatewayRequestHandlers = { taskId: params.taskId, suggestion: acceptance.suggestion, options, + }).catch((error: unknown) => { + abandonSuggestedTaskAcceptance(params.taskId, options); + throw error; }); activeAcceptances.set(params.taskId, pending); try { diff --git a/src/gateway/task-suggestion-registry.test.ts b/src/gateway/task-suggestion-registry.test.ts new file mode 100644 index 000000000000..48c0a495d081 --- /dev/null +++ b/src/gateway/task-suggestion-registry.test.ts @@ -0,0 +1,55 @@ +import { importFreshModule } from "openclaw/plugin-sdk/test-fixtures"; +import { describe, expect, it } from "vitest"; + +const BASE_SUGGESTION = { + title: "Follow up", + prompt: "Complete the follow-up task.", + tldr: "The follow-up remains useful.", + cwd: process.cwd(), + sessionKey: "agent:main:main", + agentId: "main", +}; + +describe("task suggestion registry", () => { + it("evicts accepted replay before pending work", async () => { + const { + beginTaskSuggestionAcceptance, + completeTaskSuggestionAcceptance, + createTaskSuggestion, + listTaskSuggestions, + } = await importFreshModule( + import.meta.url, + "./task-suggestion-registry.js?scope=eviction-priority", + ); + const accepted = createTaskSuggestion(BASE_SUGGESTION); + expect(accepted.status).toBe("created"); + if (accepted.status !== "created") { + throw new Error("expected accepted suggestion admission"); + } + expect(beginTaskSuggestionAcceptance(accepted.suggestion.id).status).toBe("claimed"); + completeTaskSuggestionAcceptance(accepted.suggestion.id, "agent:main:dashboard:accepted"); + + let oldestPendingTaskId = ""; + for (let index = 0; index < 99; index += 1) { + const pending = createTaskSuggestion({ + ...BASE_SUGGESTION, + title: `Pending follow up ${index}`, + }); + expect(pending.status).toBe("created"); + if (pending.status === "created" && index === 0) { + oldestPendingTaskId = pending.suggestion.id; + } + } + + const replacement = createTaskSuggestion({ + ...BASE_SUGGESTION, + title: "Latest follow up", + }); + + expect(replacement).toMatchObject({ status: "created", evictedPendingTaskIds: [] }); + expect(beginTaskSuggestionAcceptance(accepted.suggestion.id)).toEqual({ status: "missing" }); + expect(listTaskSuggestions({}).map((suggestion) => suggestion.id)).toContain( + oldestPendingTaskId, + ); + }); +}); diff --git a/src/gateway/task-suggestion-registry.ts b/src/gateway/task-suggestion-registry.ts index 0c7fff731722..d174f556e3c1 100644 --- a/src/gateway/task-suggestion-registry.ts +++ b/src/gateway/task-suggestion-registry.ts @@ -31,9 +31,9 @@ function planTaskSuggestionEvictions( let projectedCount = suggestions.size + 1; let projectedBytes = retainedSuggestionBytes + suggestionBytes + 1; const planned: Array<[string, TaskSuggestionRecord]> = []; - // Accepted replay is best effort: protect it behind pending work, but never - // let completed entries permanently prevent new suggestions from starting. - for (const status of ["dismissed", "pending", "accepted"] as const) { + // Accepted replay is best-effort and evicted before unseen pending work; + // completed entries must not displace suggestions awaiting operator action. + for (const status of ["dismissed", "accepted", "pending"] as const) { for (const [taskId, record] of suggestions) { if ( projectedCount <= MAX_TASK_SUGGESTIONS && @@ -61,9 +61,9 @@ export function createTaskSuggestion( ): CreateTaskSuggestionResult { const suggestion: TaskSuggestion = { id: `task_${randomUUID()}`, - title: params.title, - prompt: params.prompt, - tldr: params.tldr, + title: params.title.trim(), + prompt: params.prompt.trim(), + tldr: params.tldr.trim(), cwd: params.cwd, sessionKey: params.sessionKey, ...(params.agentId ? { agentId: params.agentId } : {}),