diff --git a/src/node/services/taskService.taskLaunchFormalRepro.test.ts b/src/node/services/taskService.taskLaunchFormalRepro.test.ts index 88e07d74e64..92b66f1f7ef 100644 --- a/src/node/services/taskService.taskLaunchFormalRepro.test.ts +++ b/src/node/services/taskService.taskLaunchFormalRepro.test.ts @@ -13,6 +13,7 @@ * * Run: bun test ./src/node/services/taskService.taskLaunchFormalRepro.test.ts */ +import { EventEmitter } from "events"; import { afterEach, beforeEach, describe, expect, mock, spyOn, test } from "bun:test"; import { Err, Ok, type Result } from "@/common/types/result"; @@ -22,6 +23,10 @@ import { createMuxMessage } from "@/common/types/message"; import type { Config } from "@/node/config"; import type { InitStateManager } from "@/node/services/initStateManager"; import { UnsanitizedTaskCheckoutError } from "@/node/services/unsanitizedTaskCheckout"; +import type { SendMessageOptions } from "@/common/orpc/types"; +import type { AgentSession } from "@/node/services/agentSession"; +import { createAgentSessionHarness } from "@/node/services/agentSession.testHarness"; +import type { TurnCompletion } from "@/node/services/streamManager"; import * as runtimeFactory from "@/node/runtime/runtimeFactory"; import * as forkOrchestrator from "@/node/services/utils/forkOrchestrator"; import { WorkspaceBusyError, workspaceUseLeasesFor } from "@/node/services/workspaceUseLeases"; @@ -849,6 +854,195 @@ describe("task launch: formal-model counterexamples (formal/task-launch)", () => expect(sent.slice(1).filter((m) => m.includes(BRIEF)).length).toBe(1); }); + // #5544: the reactivation that prepends the kept brief is a send of its own. It can make its + // row durable and still return Err, as the launch's send can (path A). The brief must then be + // recognized on that row too, or the next reawakening prepends it again. + // `otherBackend`: another backend's turn runs the send that wrote the reawakening's row. While + // that turn runs, a Stop there can still roll the row back ("turn"). "rolledBack": it rolled + // the row back and ended right after this backend's lookup saw the row. + for (const otherBackend of ["none", "turn", "rolledBack"] as const) { + test(`a reawakening whose send accepted the kept brief and then failed does not send it again (#5544)${otherBackend === "none" ? "" : `, unless another backend's turn may roll the row back (${otherBackend})`}`, async () => { + const sent: string[] = []; + const internals: Array = []; + const box: { history?: Awaited>["historyService"] } = {}; + const s = await setUp({ + send: async (_workspaceId, message, _options, internal) => { + sent.push(message); + internals.push(internal); + // The launch's send fails before any row: the brief is kept for a reawakening. + if (sent.length === 1) { + return Err( + createUnknownSendMessageError(WORKSPACE_STOP_IN_PROGRESS_SEND_BLOCKED_MESSAGE) + ); + } + if (sent.length > 2) return Ok(undefined); + // The first reawakening: its row becomes durable as AgentSession publishes it (with the + // send's ids and digests), then the send fails. + const identities = internal?.sendIdentities ?? []; + if (box.history == null) throw new Error("history is set before the reawakening"); + const appended = await box.history.appendToHistory( + CHILD, + createMuxMessage( + "reawakening", + "user", + message, + identities.length > 0 + ? { + sendIds: identities.map((identity) => identity.id), + sendDigests: Object.fromEntries( + identities.map((identity) => [identity.id, identity.digest]) + ), + } + : {} + ) + ); + expect(appended.success).toBe(true); + return Err(createUnknownSendMessageError("the stream failed to start")); + }, + }); + box.history = s.historyService; + await spawn(s.taskService); + await s.launched; + await s.launchFailureRecorded; + expect(await s.briefsInHistory()).toBe(0); + + const first = await s.taskService.sendMessageToDescendantAgentTask( + ROOT, + CHILD, + "Keep going", + "tool-end" + ); + expect(first).toMatchObject({ success: false }); + expect(await s.briefsInHistory()).toBe(1); + + const turn = + otherBackend === "turn" + ? await workspaceUseLeasesFor(await createTestConfig(rootDir)).hold(CHILD, "turn") + : undefined; + if (otherBackend === "rolledBack") { + const leases = workspaceUseLeasesFor(s.config); + const realIsHeld = leases.isHeld.bind(leases); + spyOn(leases, "isHeld").mockImplementation(async (id, kind) => { + if (kind !== "turn") return realIsHeld(id, kind); + // The other backend's Stop rolled the row back, then its turn ended. + const deleted = await s.historyService.deleteMessages(CHILD, ["reawakening"]); + expect(deleted.success).toBe(true); + return false; + }); + } + let second: Awaited>; + try { + second = await s.taskService.sendMessageToDescendantAgentTask( + ROOT, + CHILD, + "Again", + "tool-end" + ); + } finally { + await turn?.release(); + } + + expect(second).toMatchObject({ success: true }); + expect(sent.length).toBe(3); + expect(sent[2]).toContain("Again"); + if (otherBackend !== "none") { + // The brief stays and is sent again: never lost. + expect(sent[2]).toContain(BRIEF); + return; + } + // Target assertion. + expect(sent[2]).not.toContain(BRIEF); + expect(findWorkspaceInConfig(s.config, CHILD)?.taskPrompt).toBeUndefined(); + // The reawakening's brief is its own row, as the launch's: on-send compaction would fold it + // into a follow-up dispatched later without its id. + expect(internals[1]?.skipOnSendCompaction).toBe(true); + }); + } + + // #5544 follow-up: the same path through a real AgentSession over the real HistoryService. The + // user's Stop lands after the launch's brief row is on disk and before the session accepts the + // turn, so the session rolls the row back. The kept brief must reach history once, through the + // reawakening. + test("real session: a Stop that rolls back the launch's brief row leaves the brief to the reawakening, once", async () => { + const box: { session?: AgentSession; stop?: Promise } = {}; + const s = await setUp({ + // A thin WorkspaceHost: the real session accepts or refuses each send. + send: (_workspaceId, message, options, internal) => { + if (box.session == null) throw new Error("the session exists before the launch"); + return box.session.sendMessage(message, options as SendMessageOptions, { + acceptanceOrigin: internal?.acceptanceOrigin ?? "automatic", + agentInitiated: internal?.agentInitiated, + sendIdentities: internal?.sendIdentities, + skipOnSendCompaction: internal?.skipOnSendCompaction, + admissionStale: internal?.admissionStale, + startStreamInBackground: true, + onAccepted: internal?.onAccepted, + onCanceled: internal?.onCanceled, + onAcceptedPreStreamFailure: internal?.onAcceptedPreStreamFailure, + }); + }, + }); + const aiEmitter = new EventEmitter(); + const harness = await createAgentSessionHarness({ + workspaceId: CHILD, + config: s.config, + historyService: s.historyService, + aiEmitter, + aiServiceOverrides: { + streamMessage: mock(() => + Promise.resolve( + Ok({ + messageId: "assistant-1", + completion: Promise.resolve({ status: "completed" } as TurnCompletion), + }) + ) + ), + }, + }); + box.session = harness.session; + try { + const internals = s.taskService as unknown as { + isWorkspaceStopInProgress: (id: string) => boolean; + }; + const realPublish = s.historyService.acceptCompactionReplacement.bind(s.historyService); + spyOn(s.historyService, "acceptCompactionReplacement").mockImplementation( + async (...args) => { + const published = await realPublish(...args); + if (args[0] === CHILD && box.stop == null) { + // The brief's row is on disk: the user's Stop lands now. + box.stop = s.taskService.stopDescendantAgentTask(ROOT, CHILD); + await waitUntil( + () => internals.isWorkspaceStopInProgress(CHILD), + "the Stop to latch the child" + ); + } + return published; + } + ); + + await spawn(s.taskService); + await s.launched; + await s.launchFailureRecorded; + await box.stop; + expect(findWorkspaceInConfig(s.config, CHILD)?.taskStatus).toBe("interrupted"); + // The session rolled the launch's row back: the brief is kept for the reawakening. + expect(await s.briefsInHistory()).toBe(0); + expect(findWorkspaceInConfig(s.config, CHILD)?.taskPrompt).toBe(BRIEF); + + const reawakened = await s.taskService.sendMessageToDescendantAgentTask( + ROOT, + CHILD, + "Keep going", + "tool-end" + ); + + expect(reawakened).toMatchObject({ success: true }); + expect(await s.briefsInHistory()).toBe(1); + } finally { + await harness.session.dispose(); + } + }); + test("control: a launch whose send succeeded does not resend the brief when a Stop and a message reawaken the child", async () => { const { s, sent } = await reawakenAfterLaunch("accept"); await waitUntil( diff --git a/src/node/services/taskService.ts b/src/node/services/taskService.ts index 7313c3f217e..681bd7b2a04 100644 --- a/src/node/services/taskService.ts +++ b/src/node/services/taskService.ts @@ -574,7 +574,11 @@ function mintTaskBriefSendId(): string { return `${MINTED_SEND_ID_PREFIX}${randomUUID()}`; } -/** The brief's send identity: the launch stamps it, and a lookup checks the same id. */ +/** + * The brief's send identity: the launch stamps it, and a lookup checks the same id. The digest is + * the brief's also when the row carries the brief followed by guidance (a reactivation that + * prepends a kept brief, #5544): the id stands for the brief that row delivers. + */ function taskBriefSendIdentity(prompt: string, sendId: string): SendIdentity { return { id: sendId, digest: computeSendDigest({ message: prompt.trim() }) }; } @@ -4922,6 +4926,8 @@ export class TaskService implements AgentTaskIntegration { // intent across interrupts, including repeated interrupts after the status is no longer queued. if (previousStatus !== "queued" && !persistedQueuedPrompt) { workspace.taskPrompt = undefined; + // The brief's send id means nothing without the brief (#5544). + workspace.taskPromptSendId = undefined; } return "interrupted"; } @@ -6207,10 +6213,10 @@ export class TaskService implements AgentTaskIntegration { * Before a reawakening prepends a kept taskPrompt: drop it when a history row already carries * its brief's send id (the launch's send accepted it, then failed or was stopped before * `running`). Rows without a brief send id keep the older behavior: the kept prompt is sent. - * Only while no launch of the task is in flight on any backend (its "launch" use lease) and - * no Stop is in progress here: a launch send in flight can roll its row back after this lookup - * saw it (a Stop on the launching backend), which would lose the brief. With the lease held, - * the kept prompt stays and is sent again, as before brief send ids. + * Only while no launch or turn of the task is in flight on any backend (its "launch" or "turn" + * use lease) and no Stop is in progress here: a send in flight can roll its row back after this + * lookup saw it (a Stop on the backend that runs it), which would lose the brief. With a lease + * held, the kept prompt stays and is sent again, as before brief send ids. */ private async dropKeptTaskPromptAlreadyInHistory(taskId: string): Promise { const workspace = findWorkspaceEntry(this.config.loadConfigOrDefault(), taskId)?.workspace; @@ -6224,6 +6230,13 @@ export class TaskService implements AgentTaskIntegration { // after this check. if (await workspaceUseLeasesFor(this.config).isHeld(taskId, "launch")) return; if (!(await this.isTaskBriefInHistory(taskId, prompt, sendId))) return; + // A reawakening's send carries a brief id too (#5544), and it can run on another backend, where + // a Stop can still roll its row back. Its session publishes its "turn" lease before the row + // is written and releases it once the send has settled, so: the row was seen, then no turn is + // live, then the row is still there. Only then is it permanent. A turn started after the + // config read above carries a newer id, which the edit below refuses. + if (await workspaceUseLeasesFor(this.config).isHeld(taskId, "turn")) return; + if (!(await this.isTaskBriefInHistory(taskId, prompt, sendId))) return; await this.editWorkspaceEntry( taskId, (ws) => { @@ -9358,7 +9371,12 @@ export class TaskService implements AgentTaskIntegration { private async reactivateInactiveAgentTask(params: { ancestorWorkspaceId: string; taskId: string; - buildPrompt: (refreshed: { workspace: WorkspaceConfigEntry }) => string; + message: string; + /** + * A stopped queued child keeps its only copy of the initial brief in taskPrompt: the message + * follows that brief in the reactivation prompt. + */ + prependKeptBrief?: true; queueDispatchMode: TaskMessageQueueDispatchMode; preTurnMessages?: MuxMessage[]; sendMessage?: WorkspaceTurnHost["sendMessage"]; @@ -9460,6 +9478,15 @@ export class TaskService implements AgentTaskIntegration { // whatever createWorkspaceTurn returns or throws (P3: never roll an id back) — a refused // reactivation leaves an owned, unsettled attempt that a later Stop settles. const previousAttemptId = refreshedEntry.workspace.taskAttemptId; + // A prepended brief gets a fresh brief send id, stored with the attempt commit before the + // send (#5544). The send can make its row durable and still return Err, which keeps + // taskPrompt; the launch's id is on no row, so only this id lets the next reawakening's + // lookup (dropKeptTaskPromptAlreadyInHistory) recognize the row and not prepend it again. + const keptBrief = + params.prependKeptBrief === true + ? coerceNonEmptyString(refreshedEntry.workspace.taskPrompt) + : undefined; + const briefSendId = keptBrief != null ? mintTaskBriefSendId() : undefined; // The terminated attempt's goal stays paused across reactivation (also covers a set_goal that // raced the termination); a failed pause leaves taskGoalPauseOwed fencing goal turns. await this.settleChildGoalPause(taskId, { force: true }); @@ -9512,6 +9539,10 @@ export class TaskService implements AgentTaskIntegration { return; } ws.taskAttemptId = reactivationAttemptId; + // Only the brief this reactivation prepends: a rewrite since the read is its writer's. + if (briefSendId != null && ws.taskPrompt === refreshedEntry.workspace.taskPrompt) { + ws.taskPromptSendId = briefSendId; + } committedProven = lineage.proven && ws.taskAttemptUnproven !== true; if (!committedProven) ws.taskAttemptUnproven = true; published = true; @@ -9564,7 +9595,7 @@ export class TaskService implements AgentTaskIntegration { try { execution = await this.getWorkspaceTurnManager().createWorkspaceTurn({ ownerWorkspaceId: ancestorWorkspaceId, - prompt: params.buildPrompt(refreshedEntry), + prompt: keptBrief != null ? `${keptBrief}\n\n${params.message}` : params.message, title: coerceNonEmptyString(refreshedEntry.workspace.title) ?? coerceNonEmptyString(refreshedEntry.workspace.name) ?? @@ -9578,6 +9609,13 @@ export class TaskService implements AgentTaskIntegration { attentionPolicy: "notify_on_terminal", ...(params.sendMessage != null ? { sendMessage: params.sendMessage } : {}), ...(agentTaskAi != null ? { agentTaskAi } : {}), + ...(keptBrief != null && briefSendId != null + ? { + sendIdentities: [ + { ...taskBriefSendIdentity(keptBrief, briefSendId), unpublished: true }, + ], + } + : {}), }); } finally { if (this.reawakeningsInFlight.get(taskId) === reactivationAttemptId) { @@ -9697,7 +9735,7 @@ export class TaskService implements AgentTaskIntegration { const result = await this.reactivateInactiveAgentTask({ ancestorWorkspaceId: parentWorkspaceId, taskId: workspaceId, - buildPrompt: () => prompt, + message: prompt, queueDispatchMode: "tool-end", sendMessage: send, // No Stop fence (unlike task_send_message, L1): the wake is the child's own monitor @@ -9826,7 +9864,7 @@ export class TaskService implements AgentTaskIntegration { return this.workspaceEventLocks.withLock(taskId, async () => this.withTaskTreeLifecycleLock(taskId, async () => { // A brief history already holds is not prepended again on reawakening (U4); the - // reactivation's buildPrompt reads the row after this. Awaited before the row read + // reactivation reads the kept brief from the row after this. Awaited before the row read // below, so the inactive decision and reactivateInactiveAgentTask see the same row. await this.dropKeptTaskPromptAlreadyInHistory(taskId); const cfg = this.config.loadConfigOrDefault(); @@ -9865,14 +9903,8 @@ export class TaskService implements AgentTaskIntegration { return this.reactivateInactiveAgentTask({ ancestorWorkspaceId, taskId, - // A stopped queued child keeps its only copy of the initial brief in taskPrompt; - // the guidance follows that brief in the reactivation prompt. - buildPrompt: (refreshed) => { - const preservedQueuedPrompt = coerceNonEmptyString(refreshed.workspace.taskPrompt); - return preservedQueuedPrompt - ? `${preservedQueuedPrompt}\n\n${labeledMessage}` - : labeledMessage; - }, + message: labeledMessage, + prependKeptBrief: true, queueDispatchMode, preTurnMessages: options?.preTurnMessages, ...(sender === "ancestor" ? { aiRefresh: { prepared: preparedReawakenAi } } : {}), diff --git a/src/node/services/workspaceTurnManager.ts b/src/node/services/workspaceTurnManager.ts index eca2d17a217..eb5c53402ef 100644 --- a/src/node/services/workspaceTurnManager.ts +++ b/src/node/services/workspaceTurnManager.ts @@ -29,6 +29,7 @@ import { } from "@/node/services/taskWorkspaceSeam"; import type { HistoryService } from "@/node/services/historyService"; import type { HistoryControlRow } from "@/node/services/historyScanner"; +import type { SendIdentity } from "@/node/services/sendIds"; import { isPlainObject } from "@/common/utils/isPlainObject"; import type { InitStateManager } from "@/node/services/initStateManager"; import { @@ -348,6 +349,12 @@ export interface WorkspaceTurnCreateArgs { * handle lifecycle around the send stays this manager's. */ sendMessage?: WorkspaceTurnHost["sendMessage"]; + /** + * Internal-only: send ids the row that accepts the prompt carries (a sub-agent reactivation + * that prepends its kept initial brief, #5544). The send skips on-send compaction, which would + * fold the prompt into a follow-up dispatched later without these ids. + */ + sendIdentities?: SendIdentity[]; } export interface WorkspaceTurnCreateResult { @@ -2117,6 +2124,9 @@ export class WorkspaceTurnManager { acceptanceOrigin: "automatic", startStreamInBackground: true, requireIdle: !queuedForExistingWorkspace, + ...(args.sendIdentities != null + ? { sendIdentities: args.sendIdentities, skipOnSendCompaction: true } + : {}), onCanceled: async (reason) => { const current = await this.taskHandleStore.getWorkspaceTurn(ownerWorkspaceId, handleId); if (