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
6 changes: 5 additions & 1 deletion src/node/services/agentResolution.ts
Original file line number Diff line number Diff line change
Expand Up @@ -277,9 +277,13 @@ export async function resolveAgentForStream(
// Missing, disabled and unverifiable agents fail instead of falling back to exec.
const failClosed = strictTopLevel || automaticGoalTurn;
const failResolution = (strictMessage: string, goalTurnDetail: string) => {
// A sub-agent resolves its task-pinned agentType, which the agent picker does not repin
// (#5452), so its recovery is restoring that agent, not selecting another one.
const errorMessage = strictTopLevel
? strictMessage
: `Selected agent '${requestedAgentId}' is unavailable: ${goalTurnDetail}. Automatic goal turns never fall back to exec; select an available agent and resume the goal.`;
: isSubagentWorkspace
? `This sub-agent's agent '${requestedAgentId}' is unavailable: ${goalTurnDetail}. Automatic goal turns never fall back to exec; restore or enable agent '${requestedAgentId}' before reactivating this task.`
: `Selected agent '${requestedAgentId}' is unavailable: ${goalTurnDetail}. Automatic goal turns never fall back to exec; select an available agent and resume the goal.`;
emitError(
createErrorEvent(workspaceId, {
messageId: createAssistantMessageId(),
Expand Down
8 changes: 7 additions & 1 deletion src/node/services/aiService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1469,14 +1469,20 @@ describe("AIService.streamMessage compaction boundary slicing", () => {
});

it("a sub-agent's goal turn does not fall through to a later candidate or exec", async () => {
const { result, harness } = await streamGoalTurn("child", "researcher", {
const { result, harness, errors } = await streamGoalTurn("child", "researcher", {
parentWorkspaceId: "parent-workspace",
agentType: "researcher",
agentId: "exec",
});
const topLevel = await streamGoalTurn("child-top-level", "researcher");

expect(result.success).toBe(false);
expect(harness.startStreamCalls).toHaveLength(0);
// #5452: the agent picker does not repin a child's agentType, so a child's refusal takes
// its own recovery branch rather than the top-level one for the same unavailable agent.
expect(errors).toHaveLength(1);
expect(topLevel.errors).toHaveLength(1);
expect(errors[0]).not.toBe(topLevel.errors[0]);
});
});

Expand Down
131 changes: 130 additions & 1 deletion src/node/services/taskService.childGoals.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -283,7 +283,11 @@ describe("TaskService child goals", () => {

await streamEnd(t.taskService, t.proseEnd("assistant-1"));

expect(refuse).toHaveBeenCalledWith(childId, expect.objectContaining({ agentId: "explore" }));
expect(refuse).toHaveBeenCalledWith(
childId,
expect.objectContaining({ agentId: "explore" }),
expect.any(Function)
);
expect(t.sends()).toEqual([]);
expect((await t.goals.getGoal(childId))?.status).toBe("paused");
expect(await t.parentReports()).toHaveLength(1);
Expand All @@ -305,6 +309,131 @@ describe("TaskService child goals", () => {
expect(await t.parentReports()).toHaveLength(1);
});

// #5452 item 1: the agent check awaits; a refusal that went stale meanwhile must neither show
// its chat error (the probe) nor pause the goal a newer attempt now runs.
test("a refusal whose attempt was replaced during the check leaves the goal alone", async () => {
let currentAtCheck: boolean | undefined;
let editChild: ((mutate: (workspace: WorkspaceConfigEntry) => void) => Promise<void>) | null =
null;
const refuse = mock(
async (_id: string, _options: unknown, isCurrent?: () => boolean | Promise<boolean>) => {
await editChild?.((workspace) => {
workspace.taskAttemptId = "att_00000000000000c2";
});
currentAtCheck = await isCurrent?.();
return "Selected agent 'explore' is unavailable: it is disabled";
}
);
const t = await setup({}, undefined, { refuseUnavailableGoalTurnAgent: refuse });
editChild = t.editChild;
await t.setChildGoal();

await streamEnd(t.taskService, t.proseEnd("assistant-1"));

expect(refuse).toHaveBeenCalledTimes(1);
expect((await t.goals.getGoal(childId))?.status).toBe("active");
expect(currentAtCheck).toBe(false);
expect(t.sends()).toEqual([]);
});

test("a refusal whose goal was paused and resumed during the check shows no chat error", async () => {
let currentAtCheck: boolean | undefined;
let resumedAtCheck = false;
let goals: WorkspaceGoalService | null = null;
const resumes: Array<Promise<unknown>> = [];
const refuse = mock(
async (_id: string, _options: unknown, isCurrent?: () => boolean | Promise<boolean>) => {
if (goals == null || resumes.length > 0) return null;
await goals.setGoal({ workspaceId: childId, status: "paused" });
// The resume's own continuation waits for this stream end's event lock.
resumes.push(goals.setGoal({ workspaceId: childId, status: "active" }));
for (
let i = 0;
i < 500 && (await goals.readGoalSerialized(childId))?.status !== "active";
i++
) {
await new Promise((resolve) => setTimeout(resolve, 2));
Comment thread
ThomasK33 marked this conversation as resolved.
}
resumedAtCheck = (await goals.readGoalSerialized(childId))?.status === "active";
currentAtCheck = await isCurrent?.();
return "Selected agent 'explore' is unavailable: it is disabled";
}
);
const t = await setup({}, undefined, { refuseUnavailableGoalTurnAgent: refuse });
goals = t.goals;
await t.setChildGoal();

await streamEnd(t.taskService, t.proseEnd("assistant-1"));
await Promise.all(resumes);

// The race really happened: the goal was active again (resumed) when the probe ran.
expect(resumedAtCheck).toBe(true);
expect(currentAtCheck).toBe(false);
});

// #5452: the pause itself is fenced under the goal file lock, so an attempt replaced after the
// last check but before the pause write does not get its goal paused by the stale refusal.
test("an attempt replaced just before the refusal's pause write keeps its goal", async () => {
const refuse = mock(() =>
Promise.resolve<string | null>("Selected agent 'explore' is unavailable: it is disabled")
);
const t = await setup({}, undefined, { refuseUnavailableGoalTurnAgent: refuse });
await t.setChildGoal();
const realPause = t.goals.pauseForUnavailableAgent.bind(t.goals);
const pause = spyOn(t.goals, "pauseForUnavailableAgent").mockImplementationOnce(
async (...args) => {
await t.editChild((workspace) => {
workspace.taskAttemptId = "att_00000000000000c2";
});
return realPause(...args);
}
);

await streamEnd(t.taskService, t.proseEnd("assistant-1"));

expect(pause).toHaveBeenCalledTimes(1);
expect((await t.goals.getGoal(childId))?.status).toBe("active");
});

// An unreadable goal during the staleness probe keeps the fail-closed path: the goal pauses
// and the report publishes (the stream end is never left unhandled).
test("a failed goal read in the staleness probe still pauses and reports", async () => {
let currentAtCheck: boolean | undefined;
let replacedAtCheck: boolean | undefined;
let editChild: ((mutate: (workspace: WorkspaceConfigEntry) => void) => Promise<void>) | null =
null;
let readGoal: (() => void) | null = null;
const refuse = mock(
async (_id: string, _options: unknown, isCurrent?: () => boolean | Promise<boolean>) => {
currentAtCheck = await isCurrent?.();
// The same failed read after the attempt was replaced: the attempt fence still holds.
readGoal?.();
await editChild?.((workspace) => {
workspace.taskAttemptId = "att_00000000000000c2";
});
replacedAtCheck = await isCurrent?.();
await editChild?.((workspace) => {
workspace.taskAttemptId = "att_00000000000000c1";
});
return "Selected agent 'explore' is unavailable: it is disabled";
}
);
const t = await setup({}, undefined, { refuseUnavailableGoalTurnAgent: refuse });
await t.setChildGoal();
editChild = t.editChild;
const read = spyOn(t.goals, "readGoalSerialized").mockRejectedValueOnce(new Error("EIO"));
readGoal = () => {
read.mockRejectedValueOnce(new Error("EIO"));
};

await streamEnd(t.taskService, t.proseEnd("assistant-1"));

expect(currentAtCheck).toBe(true);
expect(replacedAtCheck).toBe(false);
expect((await t.goals.getGoal(childId))?.status).toBe("paused");
expect(await t.parentReports()).toHaveLength(1);
});

test("a stream without final prose continues an active goal (agent_report is progress)", async () => {
const t = await setup();
await t.setChildGoal();
Expand Down
23 changes: 21 additions & 2 deletions src/node/services/taskService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18057,18 +18057,37 @@ export class TaskService implements AgentTaskIntegration {
if (!(await this.canChildAgentDriveGoal(workspaceId))) return "none";
// Fail closed when the pinned agent is unavailable (#5402): pause; the normal path applies.
const childEntry = findWorkspaceEntry(this.config.loadConfigOrDefault(), workspaceId);
// The check awaits (#5452): a dispatch that went stale meanwhile (its attempt replaced, or
// its goal replaced, settled or paused and resumed) shows no chat error and pauses nothing,
// since that goal may now belong to a newer turn. Its stream end then takes the normal path.
// Synchronous, so the pause can evaluate it under the goal file lock right before its write.
const stillCurrent = (live: GoalRecordV1 | null) =>
this.currentTaskAttemptId(workspaceId) === expectedAttemptId &&
live?.goalId === goal.goalId &&
live.status === goal.status &&
(live.lastUserActivationAtMs ?? null) === (goal.lastUserActivationAtMs ?? null);
// An unreadable goal counts as current while the attempt is: the refusal stands (fail
// closed) and the pause write, fenced again under its lock, owns any error.
const isCurrent = async () => {
try {
return stillCurrent(await goalService.readGoalSerialized(workspaceId));
} catch {
return this.currentTaskAttemptId(workspaceId) === expectedAttemptId;
}
};
const refusal =
childEntry == null
? null
: await this.workspaceService.refuseUnavailableGoalTurnAgent(
workspaceId,
buildTaskTurnSendOptions(childEntry.workspace)
buildTaskTurnSendOptions(childEntry.workspace),
isCurrent
);
if (refusal != null) {
// A failed write is safe to ignore: "none" leads to the report (or a report prompt, whose
// stream end re-arbitrates here), and the reported transition marks the pause owed
// (taskGoalPauseOwed fences goal turns) and settles it.
await goalService.pauseForUnavailableAgent(workspaceId, goal, refusal);
await goalService.pauseForUnavailableAgent(workspaceId, goal, refusal, stillCurrent);
return "none";
}
const goalAdmission = await goalService.buildGoalRedispatchAdmission(
Expand Down
2 changes: 1 addition & 1 deletion src/node/services/taskWorkspaceSeam.ts
Original file line number Diff line number Diff line change
Expand Up @@ -504,7 +504,7 @@ export interface WorkspaceTurnHost {
refuseUnavailableGoalTurnAgent(
workspaceId: string,
options: Pick<SendMessageOptions, "agentId" | "disableWorkspaceAgents">,
isCurrent?: () => boolean
isCurrent?: () => boolean | Promise<boolean>
): Promise<string | null>;
waitForPendingCompactionCompletionDecision(
workspaceId: string,
Expand Down
36 changes: 33 additions & 3 deletions src/node/services/workspaceGoalService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,13 @@ export interface SetGoalInput {
* it. Not part of the public oRPC schema.
*/
requireSelectedAgentId?: string | null;
/**
* Internal fence for automatic writes decided from an earlier read (#5452): checked against
* the current record under the goal file lock, right before the write is installed or
* persisted; false refuses the write as a goal_conflict. Must be synchronous. Not part of the
* public oRPC schema.
*/
stillCurrent?: ((current: GoalRecordV1 | null) => boolean) | null;
}

export type { GoalStreamOriginKind } from "./goalContinuationPolicy";
Expand Down Expand Up @@ -1942,12 +1949,18 @@ export class WorkspaceGoalService {
* existing matching reservation is already correct and left untouched) so
* the recovered stream's end cannot arm and fire a second wrap-up.
*/
async reserveBudgetWrapupForRedispatch(workspaceId: string, goalId: string): Promise<void> {
async reserveBudgetWrapupForRedispatch(
workspaceId: string,
goalId: string,
/** See SetGoalInput.stillCurrent. */
stillCurrent?: (current: GoalRecordV1 | null) => boolean
): Promise<void> {
assert(workspaceId.trim().length > 0, "reserveBudgetWrapupForRedispatch requires workspaceId");
assert(goalId.trim().length > 0, "reserveBudgetWrapupForRedispatch requires goalId");
await this.fileLocks.withLock(workspaceId, async () => {
const current = await this.readGoalFile(workspaceId);
if (
stillCurrent?.(current) === false ||
current?.goalId !== goalId ||
current.status !== "budget_limited" ||
current.budgetLimitOriginKind === "user" ||
Expand Down Expand Up @@ -2354,11 +2367,13 @@ export class WorkspaceGoalService {
* user resumes it after selecting an available agent); a budget-limited goal's one wrap-up is
* skipped, consumed as settleChildGoalPause does, so it stays budget_limited with nothing owed.
* False only when the write failed; a refused transition (the goal changed meanwhile) is settled.
* `stillCurrent` fences the write under the goal file lock (see SetGoalInput.stillCurrent).
*/
async pauseForUnavailableAgent(
workspaceId: string,
goal: GoalRecordV1,
reason: string
reason: string,
stillCurrent?: (current: GoalRecordV1 | null) => boolean
): Promise<boolean> {
try {
if (goal.status === "active") {
Expand All @@ -2367,9 +2382,10 @@ export class WorkspaceGoalService {
status: "paused",
initiator: "auto",
expectedGoalId: goal.goalId,
...(stillCurrent != null ? { stillCurrent } : {}),
});
} else if (goal.status === "budget_limited") {
await this.reserveBudgetWrapupForRedispatch(workspaceId, goal.goalId);
await this.reserveBudgetWrapupForRedispatch(workspaceId, goal.goalId, stillCurrent);
}
} catch (error) {
log.warn("WorkspaceGoalService: could not settle a goal whose agent is unavailable", {
Expand Down Expand Up @@ -2767,6 +2783,18 @@ export class WorkspaceGoalService {
return { type: "goal_conflict", expectedGoalId, actualGoalId };
}

private conflictForStillCurrent(
current: GoalRecordV1 | null,
input: SetGoalInput
): GoalSetError | null {
if (input.stillCurrent == null || input.stillCurrent(current)) return null;
return {
type: "goal_conflict",
expectedGoalId: input.expectedGoalId ?? null,
actualGoalId: current?.goalId ?? null,
};
}

private conflictForReplacementGuard(
current: GoalRecordV1 | null,
replacementGuard: SetGoalReplacementGuard | null | undefined
Expand Down Expand Up @@ -3262,6 +3290,7 @@ export class WorkspaceGoalService {
const current = await this.readGoalFile(input.workspaceId);
const conflict =
this.conflictForExpectedGoalId(current, input.expectedGoalId) ??
this.conflictForStillCurrent(current, input) ??
this.conflictForReplacementGuard(current, input.replacementGuard);
if (conflict) {
return Err(conflict);
Expand Down Expand Up @@ -3626,6 +3655,7 @@ export class WorkspaceGoalService {
const current = await this.readGoalFile(input.workspaceId);
const conflict =
this.conflictForExpectedGoalId(current, input.expectedGoalId) ??
this.conflictForStillCurrent(current, input) ??
this.conflictForReplacementGuard(current, input.replacementGuard);
Comment thread
ThomasK33 marked this conversation as resolved.
if (conflict) {
return Err(conflict);
Expand Down
27 changes: 27 additions & 0 deletions src/node/services/workspaceService.goalContinuation.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1019,6 +1019,33 @@ describe("automatic goal turns whose selected agent is unavailable (#5402)", ()
}
);

// #5452: a synchronous staleness probe (the root dispatcher's candidate identity) is evaluated
// in the same tick as the chat error, so no candidate replacement can slip in between.
test("a synchronous staleness probe stays adjacent to the refusal's chat error", async () => {
const t = await setup("researcher");
await t.deleteAgent();
const session = t.service.getOrCreateSession(workspaceId);
let replacedAfterProbe = false;
const replacedAtEmission: boolean[] = [];
spyOn(session, "emitChatEvent").mockImplementation(() => {
replacedAtEmission.push(replacedAfterProbe);
});

const refusal = await t.service.refuseUnavailableGoalTurnAgent(
workspaceId,
{ agentId: "researcher" },
() => {
queueMicrotask(() => {
replacedAfterProbe = true;
});
return true;
}
);

expect(refusal).not.toBeNull();
expect(replacedAtEmission).toEqual([false]);
});

test("a hidden saved selection (explore) still continues", async () => {
const t = await setup("explore");
const goals = t.goalService();
Expand Down
7 changes: 5 additions & 2 deletions src/node/services/workspaceService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19910,10 +19910,13 @@ export class WorkspaceService
async refuseUnavailableGoalTurnAgent(
workspaceId: string,
options: Pick<SendMessageOptions, "agentId" | "disableWorkspaceAgents">,
isCurrent: () => boolean = () => true
isCurrent: () => boolean | Promise<boolean> = () => true
): Promise<string | null> {
const refusal = await this.aiService.getAutomaticGoalTurnAgentRefusal(workspaceId, options);
if (refusal == null || !isCurrent()) return refusal;
if (refusal == null) return null;
// A synchronous probe stays adjacent to the emission (no await, so no microtask yield).
const current = isCurrent();
if (!(typeof current === "boolean" ? current : await current)) return refusal;
this.sessions.get(workspaceId)?.emitChatEvent(
createStreamErrorMessage({
messageId: createAssistantMessageId(),
Expand Down
Loading