Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
194 changes: 194 additions & 0 deletions src/node/services/taskService.taskLaunchFormalRepro.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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";
Expand Down Expand Up @@ -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<SendMessageInternalOptions | undefined> = [];
const box: { history?: Awaited<ReturnType<typeof setUp>>["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<ReturnType<typeof s.taskService.sendMessageToDescendantAgentTask>>;
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<unknown> } = {};
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(
Expand Down
66 changes: 49 additions & 17 deletions src/node/services/taskService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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() }) };
}
Expand Down Expand Up @@ -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";
}
Expand Down Expand Up @@ -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<void> {
const workspace = findWorkspaceEntry(this.config.loadConfigOrDefault(), taskId)?.workspace;
Expand All @@ -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;
Comment thread
ThomasK33 marked this conversation as resolved.
if (!(await this.isTaskBriefInHistory(taskId, prompt, sendId))) return;
await this.editWorkspaceEntry(
taskId,
(ws) => {
Expand Down Expand Up @@ -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"];
Expand Down Expand Up @@ -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 });
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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) ??
Expand All @@ -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) {
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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 } } : {}),
Expand Down
Loading
Loading