fix(coding-agent): require a fresh yield after async-result deliveries in the quiescence barrier
The barrier in driveSessionToYield was unreachable for terminal yields:
the yield tool's shouldTerminate fired requestAbort("terminate"), so
abortSignal was always aborted before the barrier's loop condition ran,
and a run with pending owner jobs completed immediately with whatever
the pre-async yield said (Codex review on #6119).
- Split "stop the free-running turn after yield" from "terminate the
run": a terminal yield with pending owner async work now parks the
run with a recoverable session abort (budget-stop precedent) via
requestYieldTurnStop; only a quiescent yield terminates.
- An async-result follow-up injected after a recorded yield un-latches
it (transcript-ordered, in the run monitor) and re-runs the reminder
ladder, so the run only completes on a yield that postdates every
delivered result — including results injected during the notice turn.
- A run that never refreshes a superseded yield fails (exit 1) with an
explicit reason; the stale payload ships only as failed-run salvage
through the existing failed-after-yield finalize path.
- Rewrote subagent-async-pending.md: the "your current yield stands"
option contradicted the enforced contract.
- Regression tests: parked yield -> injected result -> fresh yield wins;
refusal -> stale payload fails; no-async fast path unchanged.
This commit is contained in:
@@ -10,7 +10,7 @@
|
||||
|
||||
### Changed
|
||||
|
||||
- Subagents now inherit `async.enabled` and `bash.autoBackground.enabled` from the parent instead of having both force-disabled. Subagent runs complete only after their own background jobs settle (results are folded into the run as async-result follow-ups, with a one-time notice offering `hub` wait/cancel), and teardown cancels and awaits surviving jobs before isolation worktree capture and cleanup.
|
||||
- Subagents now inherit `async.enabled` and `bash.autoBackground.enabled` from the parent instead of having both force-disabled. Subagent runs complete only after their own background jobs settle and the agent submits a `yield` that postdates every delivered result: a terminal yield with jobs still pending parks the run (recoverable turn stop) instead of completing it, async results are folded in as follow-up turns (with a one-time notice offering `hub` wait/cancel), a result delivered after a yield supersedes that yield and re-runs the yield reminder ladder, and a run that never refreshes a superseded yield fails with the stale payload preserved as salvage. Teardown cancels and awaits surviving jobs before isolation worktree capture and cleanup.
|
||||
|
||||
### Fixed
|
||||
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
Your yield was recorded, but {{count}} background job{{#if multiple}}s{{/if}} you own {{#if multiple}}are{{else}}is{{/if}} still running: {{jobs}}.
|
||||
|
||||
This run completes only after these jobs settle; their results arrive as follow-up messages. Decide now:
|
||||
This run completes only after these jobs settle AND you submit a fresh `yield` that accounts for their results. Job results arrive as follow-up messages; a result that arrives after your yield supersedes it — your current yield will NOT be accepted as the final report. Decide now:
|
||||
- Need the results? Wait for them (`hub` op:"wait"), then submit a fresh `yield` that incorporates them.
|
||||
- Job no longer needed? Cancel it (`hub` op:"cancel", ids:[...]) and re-yield.
|
||||
- Otherwise say nothing further; your current yield stands once the jobs finish.
|
||||
- Otherwise stand by; when each result arrives, submit a fresh `yield` (repeat your report unchanged if the result does not affect it).
|
||||
|
||||
@@ -13,6 +13,13 @@ import type { AsyncJob } from "../async";
|
||||
import asyncResultTemplate from "../prompts/tools/async-result.md" with { type: "text" };
|
||||
import type { CustomMessage } from "./messages";
|
||||
|
||||
/**
|
||||
* `customType` of the injected async-result follow-up message. The task
|
||||
* executor's run monitor matches on it to invalidate a previously recorded
|
||||
* yield: a result injected after the yield supersedes that yield's payload.
|
||||
*/
|
||||
export const ASYNC_RESULT_MESSAGE_TYPE = "async-result";
|
||||
|
||||
/** Result payloads longer than this spill to an artifact with an inline preview. */
|
||||
export const ASYNC_INLINE_RESULT_MAX_CHARS = 12_000;
|
||||
export const ASYNC_PREVIEW_MAX_CHARS = 4_000;
|
||||
@@ -54,7 +61,7 @@ export function buildAsyncResultBatchMessage(entries: AsyncResultEntry[]): Custo
|
||||
};
|
||||
return {
|
||||
role: "custom",
|
||||
customType: "async-result",
|
||||
customType: ASYNC_RESULT_MESSAGE_TYPE,
|
||||
content: prompt.render(asyncResultTemplate, {
|
||||
multiple: jobs.length > 1,
|
||||
jobs,
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
*/
|
||||
|
||||
import path from "node:path";
|
||||
import type { AgentEvent, AgentIdentity, AgentTelemetryConfig } from "@oh-my-pi/pi-agent-core";
|
||||
import type { AgentEvent, AgentIdentity, AgentMessage, AgentTelemetryConfig } from "@oh-my-pi/pi-agent-core";
|
||||
import { recordHandoff, resolveTelemetry } from "@oh-my-pi/pi-agent-core";
|
||||
import type { Api, Model, ServiceTierByFamily, Usage } from "@oh-my-pi/pi-ai";
|
||||
import { logger, popLoopPhase, prompt, pushLoopPhase, untilAborted } from "@oh-my-pi/pi-utils";
|
||||
@@ -42,6 +42,7 @@ import { AgentRegistry } from "../registry/agent-registry";
|
||||
import { type CreateAgentSessionOptions, createAgentSession, discoverAuthStorage } from "../sdk";
|
||||
import type { AgentSession, AgentSessionEvent, Prewalk } from "../session/agent-session";
|
||||
import type { ArtifactManager } from "../session/artifacts";
|
||||
import { ASYNC_RESULT_MESSAGE_TYPE } from "../session/async-job-delivery";
|
||||
import type { AuthStorage } from "../session/auth-storage";
|
||||
import { SKILL_PROMPT_MESSAGE_TYPE, USER_INTERRUPT_LABEL } from "../session/messages";
|
||||
import { SessionManager } from "../session/session-manager";
|
||||
@@ -877,6 +878,20 @@ interface SubagentRunMonitor {
|
||||
budgetStopRequested(): boolean;
|
||||
/** Resolves when the budget-stop session abort has settled (immediately when no stop fired). */
|
||||
waitForBudgetStop(): Promise<void>;
|
||||
/**
|
||||
* True when a recorded yield was invalidated by a later async-result
|
||||
* injection and no fresh yield has landed since: the yield payload
|
||||
* predates background job outcomes the model was shown.
|
||||
*/
|
||||
yieldInvalidatedByAsync(): boolean;
|
||||
/**
|
||||
* True once a terminal yield with pending owner async work stopped the
|
||||
* free-running turn (recoverable, like a budget stop) instead of
|
||||
* terminating the run. Cleared when {@link waitForYieldTurnStop} settles.
|
||||
*/
|
||||
yieldTurnStopRequested(): boolean;
|
||||
/** Resolves when the yield turn-stop session abort has settled (immediately when none fired). */
|
||||
waitForYieldTurnStop(): Promise<void>;
|
||||
/** The abort kind for this run, when an abort was requested. */
|
||||
abortKind(): AbortReason | undefined;
|
||||
/** True when the abort carries a precise external reason (signal / wall-clock / budget). */
|
||||
@@ -903,6 +918,15 @@ interface SubagentRunMonitor {
|
||||
finish(): void;
|
||||
}
|
||||
|
||||
/**
|
||||
* True when `message` is the session-injected async-result follow-up
|
||||
* ({@link ASYNC_RESULT_MESSAGE_TYPE}): the transcript-ordered signal that a
|
||||
* background job outcome landed after whatever the model said before it.
|
||||
*/
|
||||
function isAsyncResultInjection(message: AgentMessage | undefined): boolean {
|
||||
return message?.role === "custom" && message.customType === ASYNC_RESULT_MESSAGE_TYPE;
|
||||
}
|
||||
|
||||
function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor {
|
||||
const {
|
||||
index,
|
||||
@@ -954,6 +978,9 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor {
|
||||
let activeSession: AgentSession | null = null;
|
||||
let yieldCalled = false;
|
||||
let yieldCallPending = false;
|
||||
let yieldInvalidatedByAsync = false;
|
||||
let yieldTurnStopRequested = false;
|
||||
let yieldTurnStopPromise: Promise<void> | null = null;
|
||||
|
||||
// Accumulate usage incrementally from message_end events (no memory for streaming events)
|
||||
const accumulatedUsage: Usage = {
|
||||
@@ -1026,6 +1053,30 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor {
|
||||
: Promise.resolve();
|
||||
};
|
||||
|
||||
// Yield turn-stop: a terminal yield recorded while owner async work is
|
||||
// still pending is a scheduling pause, not run completion. Stop the
|
||||
// free-running turn exactly like a budget stop (session abort, monitor
|
||||
// signal untouched) so driveSessionToYield's quiescence barrier can settle
|
||||
// the jobs, fold their results in, and demand a fresh yield. Terminating
|
||||
// here instead would abort the run signal and make the barrier
|
||||
// unreachable, completing the run with a payload that predates the job
|
||||
// outcomes.
|
||||
const requestYieldTurnStop = () => {
|
||||
if (yieldTurnStopRequested || abortSent || resolved) return;
|
||||
yieldTurnStopRequested = true;
|
||||
const session = activeSession;
|
||||
yieldTurnStopPromise = session
|
||||
? session.abort().catch(error => {
|
||||
logger.debug("Subagent yield turn-stop abort failed", {
|
||||
error: error instanceof Error ? error.message : String(error),
|
||||
});
|
||||
})
|
||||
: Promise.resolve();
|
||||
};
|
||||
|
||||
/** Owner async work that can still re-wake the run (quiescence barrier predicate). */
|
||||
const sessionHasPendingAsyncWork = (): boolean => activeSession?.hasPendingAsyncWork?.() ?? false;
|
||||
|
||||
// Handle abort signal
|
||||
if (signal) {
|
||||
signal.addEventListener(
|
||||
@@ -1229,6 +1280,7 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor {
|
||||
if (toolName === "yield") {
|
||||
yieldCalled = true;
|
||||
yieldCallPending = false;
|
||||
yieldInvalidatedByAsync = false;
|
||||
}
|
||||
};
|
||||
|
||||
@@ -1242,6 +1294,16 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor {
|
||||
if (event.message?.role === "assistant") {
|
||||
resetRecentOutput();
|
||||
}
|
||||
// An async-result follow-up injected after a recorded yield
|
||||
// supersedes that yield: its payload predates the job outcome the
|
||||
// model is now being shown. Un-latch so the quiescence barrier's
|
||||
// reminder ladder demands a fresh yield. Guarded on the run signal:
|
||||
// once the run is completing, late injections must not destabilize
|
||||
// the settled classification.
|
||||
if (yieldCalled && !abortSignal.aborted && isAsyncResultInjection(event.message)) {
|
||||
yieldCalled = false;
|
||||
yieldInvalidatedByAsync = true;
|
||||
}
|
||||
break;
|
||||
|
||||
case "tool_execution_start": {
|
||||
@@ -1325,7 +1387,14 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor {
|
||||
isError: event.isError,
|
||||
})
|
||||
) {
|
||||
requestAbort("terminate");
|
||||
if (event.toolName === "yield" && sessionHasPendingAsyncWork()) {
|
||||
// Terminal yield with owner jobs still pending: park the
|
||||
// run behind the quiescence barrier instead of completing
|
||||
// it (see requestYieldTurnStop).
|
||||
requestYieldTurnStop();
|
||||
} else {
|
||||
requestAbort("terminate");
|
||||
}
|
||||
}
|
||||
}
|
||||
flushProgress = true;
|
||||
@@ -1567,6 +1636,25 @@ function createSubagentRunMonitor(args: RunMonitorArgs): SubagentRunMonitor {
|
||||
abortReason === "signal" || runtimeLimitExceeded || budgetLimitExceeded || budgetStopRequested,
|
||||
budgetStopRequested: () => budgetStopRequested,
|
||||
waitForBudgetStop: () => budgetStopAbortPromise ?? Promise.resolve(),
|
||||
yieldInvalidatedByAsync: () => yieldInvalidatedByAsync,
|
||||
yieldTurnStopRequested: () => yieldTurnStopRequested,
|
||||
waitForYieldTurnStop: async () => {
|
||||
const pending = yieldTurnStopPromise;
|
||||
if (!pending) {
|
||||
yieldTurnStopRequested = false;
|
||||
return;
|
||||
}
|
||||
try {
|
||||
await pending;
|
||||
} finally {
|
||||
// Clear only after the abort settled so the idempotence gate in
|
||||
// requestYieldTurnStop stays closed while it is in flight.
|
||||
if (yieldTurnStopPromise === pending) {
|
||||
yieldTurnStopPromise = null;
|
||||
yieldTurnStopRequested = false;
|
||||
}
|
||||
}
|
||||
},
|
||||
// A soft stop that never escalated still identifies as a budget abort so
|
||||
// the lifecycle can park the agent as resumable instead of killing it.
|
||||
abortKind: () => abortReason ?? (budgetStopRequested ? "budget" : undefined),
|
||||
@@ -1664,83 +1752,106 @@ async function driveSessionToYield(
|
||||
await awaitAbortable(session.prompt(task, { attribution: "agent" }));
|
||||
await awaitAbortable(session.waitForIdle());
|
||||
} catch (err) {
|
||||
// A budget stop cancels the free-running turn by aborting the
|
||||
// session, which can surface here as a rejected prompt. Swallow it
|
||||
// and drive the forced final yield below; real caller/timeout
|
||||
// aborts (monitor signal) and genuine failures keep the old path.
|
||||
if (!monitor.budgetStopRequested() || abortSignal.aborted) throw err;
|
||||
// A budget stop or a yield turn-stop (terminal yield parked behind
|
||||
// the async quiescence barrier) cancels the free-running turn by
|
||||
// aborting the session, which can surface here as a rejected
|
||||
// prompt. Swallow it and drive the barrier/forced final yield
|
||||
// below; real caller/timeout aborts (monitor signal) and genuine
|
||||
// failures keep the old path.
|
||||
const recoverableStop = monitor.budgetStopRequested() || monitor.yieldTurnStopRequested();
|
||||
if (!recoverableStop || abortSignal.aborted) throw err;
|
||||
}
|
||||
|
||||
const reminderToolChoice = buildNamedToolChoice("yield", session.model);
|
||||
|
||||
let retryCount = 0;
|
||||
while (!monitor.yieldCalled() && retryCount < MAX_YIELD_RETRIES && !abortSignal.aborted) {
|
||||
// A budget stop collapses the reminder ladder to a single forced
|
||||
// final yield: wait for the stop's session abort to settle, then
|
||||
// prompt once with the wrap-up reminder + named tool choice.
|
||||
const budgetStop = monitor.budgetStopRequested();
|
||||
if (budgetStop) {
|
||||
retryCount = MAX_YIELD_RETRIES - 1;
|
||||
await monitor.waitForBudgetStop();
|
||||
if (monitor.yieldCalled() || abortSignal.aborted) break;
|
||||
}
|
||||
// Skip reminders when the model returned a terminal error (e.g.
|
||||
// rate-limit cap hit, auth failure). Re-prompting would just
|
||||
// hit the same wall, multiplying the failure noise without
|
||||
// any chance of producing a yield.
|
||||
const lastBeforeReminder = session.getLastAssistantMessage();
|
||||
if (lastBeforeReminder?.stopReason === "error") break;
|
||||
try {
|
||||
retryCount++;
|
||||
const reminder = prompt.render(submitReminderTemplate, {
|
||||
retryCount,
|
||||
maxRetries: MAX_YIELD_RETRIES,
|
||||
budgetStop,
|
||||
});
|
||||
|
||||
const isFinalRetry = retryCount >= MAX_YIELD_RETRIES;
|
||||
await awaitAbortable(
|
||||
session.prompt(reminder, {
|
||||
attribution: "agent",
|
||||
synthetic: true,
|
||||
...(isFinalRetry && reminderToolChoice ? { toolChoice: reminderToolChoice } : {}),
|
||||
}),
|
||||
);
|
||||
await awaitAbortable(session.waitForIdle());
|
||||
} catch (err) {
|
||||
if (abortSignal.aborted || err instanceof ToolAbortError) {
|
||||
// Benign control-flow exit — user cancel (^C) or compaction aborting
|
||||
// pending operations both surface here as ToolAbortError. The outer
|
||||
// catch and finally already mark the run aborted; logging at ERROR
|
||||
// would spam operator dashboards with non-failures.
|
||||
logger.debug("Subagent prompt aborted");
|
||||
} else {
|
||||
logger.error("Subagent prompt failed", {
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
const runYieldLadder = async (): Promise<void> => {
|
||||
let retryCount = 0;
|
||||
while (!monitor.yieldCalled() && retryCount < MAX_YIELD_RETRIES && !abortSignal.aborted) {
|
||||
// A budget stop collapses the reminder ladder to a single forced
|
||||
// final yield: wait for the stop's session abort to settle, then
|
||||
// prompt once with the wrap-up reminder + named tool choice.
|
||||
const budgetStop = monitor.budgetStopRequested();
|
||||
if (budgetStop) {
|
||||
retryCount = MAX_YIELD_RETRIES - 1;
|
||||
await monitor.waitForBudgetStop();
|
||||
if (monitor.yieldCalled() || abortSignal.aborted) break;
|
||||
}
|
||||
// Skip reminders when the model returned a terminal error (e.g.
|
||||
// rate-limit cap hit, auth failure). Re-prompting would just
|
||||
// hit the same wall, multiplying the failure noise without
|
||||
// any chance of producing a yield.
|
||||
const lastBeforeReminder = session.getLastAssistantMessage();
|
||||
if (lastBeforeReminder?.stopReason === "error") break;
|
||||
try {
|
||||
retryCount++;
|
||||
const reminder = prompt.render(submitReminderTemplate, {
|
||||
retryCount,
|
||||
maxRetries: MAX_YIELD_RETRIES,
|
||||
budgetStop,
|
||||
});
|
||||
|
||||
const isFinalRetry = retryCount >= MAX_YIELD_RETRIES;
|
||||
await awaitAbortable(
|
||||
session.prompt(reminder, {
|
||||
attribution: "agent",
|
||||
synthetic: true,
|
||||
...(isFinalRetry && reminderToolChoice ? { toolChoice: reminderToolChoice } : {}),
|
||||
}),
|
||||
);
|
||||
await awaitAbortable(session.waitForIdle());
|
||||
} catch (err) {
|
||||
if (abortSignal.aborted || err instanceof ToolAbortError) {
|
||||
// Benign control-flow exit — user cancel (^C) or compaction aborting
|
||||
// pending operations both surface here as ToolAbortError. The outer
|
||||
// catch and finally already mark the run aborted; logging at ERROR
|
||||
// would spam operator dashboards with non-failures.
|
||||
logger.debug("Subagent prompt aborted");
|
||||
} else {
|
||||
logger.error("Subagent prompt failed", {
|
||||
error: err instanceof Error ? err.message : String(err),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
// Quiescence barrier (structured concurrency): a final yield with owner
|
||||
// background jobs still running or undelivered is a scheduling pause,
|
||||
// not run completion. Wait for the jobs to settle, let their
|
||||
// async-result follow-up turns run to idle, and only then classify the
|
||||
// terminal state — the isolation runner captures and destroys the
|
||||
// worktree right after this run resolves, so no owner job that could
|
||||
// still re-wake the session may outlive it. Suppressed (acknowledged /
|
||||
// hub-watched) jobs never re-wake the run and are reaped at teardown.
|
||||
// Yield ladder + quiescence barrier (structured concurrency), one
|
||||
// loop: each iteration first demands a yield — initially, and again
|
||||
// whenever an async-result delivery un-latched the previous one
|
||||
// (including during the notice turn) — then either completes on
|
||||
// quiescence or settles one generation of owner async work.
|
||||
//
|
||||
// A final yield with owner background jobs still running or
|
||||
// undelivered is a scheduling pause, not run completion — the monitor
|
||||
// parks such a yield with a recoverable turn-stop instead of
|
||||
// terminating the run. Jobs are settled and their results folded into
|
||||
// the run as async-result follow-up turns; each delivered result
|
||||
// supersedes the yield it postdates, so the reminder ladder re-runs
|
||||
// to demand a fresh yield that accounts for it. Only a yield with no
|
||||
// pending owner work left is terminal — the isolation runner captures
|
||||
// and destroys the worktree right after this run resolves, so no
|
||||
// owner job that could still re-wake the session may outlive it.
|
||||
// Suppressed (acknowledged / hub-watched) jobs never re-wake the run
|
||||
// and are reaped at teardown.
|
||||
//
|
||||
// Before blocking on running jobs, tell the model ONCE what it is
|
||||
// waiting on so it can `hub` wait/cancel instead of sitting silent
|
||||
// until the jobs (or the runtime limit) expire. Guarded to a single
|
||||
// notice per run — a model that yields again with jobs still running
|
||||
// is then waited out passively. Runs that never yielded (ladder
|
||||
// exhausted / terminal model error) skip the barrier entirely — more
|
||||
// until the jobs (or the runtime limit) expire. Runs that never yield
|
||||
// (ladder exhausted / terminal model error) skip the barrier — more
|
||||
// injected turns just multiply the failure noise; the teardown reap
|
||||
// still cancels and awaits their jobs before worktree capture.
|
||||
let asyncPendingNoticeSent = false;
|
||||
while (monitor.yieldCalled() && !abortSignal.aborted && session.hasPendingAsyncWork()) {
|
||||
while (!abortSignal.aborted) {
|
||||
if (!monitor.yieldCalled()) {
|
||||
await runYieldLadder();
|
||||
// Ladder exhausted / terminal model error: classified below
|
||||
// (missing yield, or stale yield when one was invalidated).
|
||||
if (!monitor.yieldCalled()) break;
|
||||
}
|
||||
// Let the parked yield's turn-stop session abort settle before
|
||||
// prompting again (mirrors waitForBudgetStop).
|
||||
await awaitAbortable(monitor.waitForYieldTurnStop());
|
||||
if (!session.hasPendingAsyncWork()) break;
|
||||
if (!asyncPendingNoticeSent) {
|
||||
asyncPendingNoticeSent = true;
|
||||
const running = session.getAsyncJobSnapshot()?.running ?? [];
|
||||
@@ -1763,11 +1874,13 @@ async function driveSessionToYield(
|
||||
});
|
||||
}
|
||||
// Re-evaluate: the notice turn may have cancelled, watched, or
|
||||
// absorbed the jobs.
|
||||
// absorbed the jobs — or already re-yielded.
|
||||
continue;
|
||||
}
|
||||
}
|
||||
await awaitAbortable(session.settleAsyncWork());
|
||||
// Results delivered during the settle invalidated the recorded
|
||||
// yield: the next iteration's ladder demands a fresh one.
|
||||
}
|
||||
|
||||
if (monitor.yieldCalled()) {
|
||||
@@ -1806,6 +1919,18 @@ async function driveSessionToYield(
|
||||
abortReasonText ??= monitor.resolveAbortReasonText();
|
||||
exitCode = 1;
|
||||
}
|
||||
|
||||
// A recorded yield that async-result deliveries superseded and the
|
||||
// model never refreshed is stale: fail the run instead of letting the
|
||||
// parent act on a payload that predates the background job outcomes
|
||||
// the model was shown. The stale payload still ships through
|
||||
// finalizeSubprocessOutput's failed-after-yield path (exit 1 + stderr,
|
||||
// output preserved as salvage).
|
||||
if (monitor.yieldInvalidatedByAsync() && !abortSignal.aborted) {
|
||||
exitCode = 1;
|
||||
error ??=
|
||||
"Background job results arrived after the subagent's last yield; it did not submit a refreshed yield covering them.";
|
||||
}
|
||||
} catch (err) {
|
||||
if (abortSignal.aborted && monitor.yieldCalled() && !monitor.runtimeLimitExceeded()) {
|
||||
exitCode = 0;
|
||||
|
||||
@@ -333,7 +333,6 @@ describe("AsyncJobManager", () => {
|
||||
});
|
||||
|
||||
test("scoped delivery drain times out while a matching delivery callback is in flight", async () => {
|
||||
let mainJobId = "";
|
||||
let targetJobId = "";
|
||||
let releaseMainDelivery = (): void => {};
|
||||
let notifyMainDeliveryStarted = (): void => {};
|
||||
@@ -363,7 +362,7 @@ describe("AsyncJobManager", () => {
|
||||
completions.push(jobId);
|
||||
});
|
||||
|
||||
mainJobId = manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" });
|
||||
manager.register("task", "main job", async () => "main result", { ownerId: "0-Main" });
|
||||
targetJobId = manager.register("task", "subagent job", async () => "subagent result", {
|
||||
ownerId: "3-AuthLoader",
|
||||
});
|
||||
|
||||
@@ -0,0 +1,253 @@
|
||||
/**
|
||||
* Quiescence barrier fresh-yield contract (PR #6119 review): a terminal
|
||||
* `yield` recorded while owner background jobs are still pending parks the
|
||||
* run instead of terminating it, and an async-result delivered after that
|
||||
* yield supersedes it — the run only completes on a yield that postdates
|
||||
* every delivered result. A model that never refreshes its yield must fail
|
||||
* the run rather than surface the stale payload as a clean success.
|
||||
*/
|
||||
import { afterEach, describe, expect, it, vi } from "bun:test";
|
||||
import type { AssistantMessage } from "@oh-my-pi/pi-ai";
|
||||
import type { LoadExtensionsResult } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/types";
|
||||
import type { CreateAgentSessionResult } from "@oh-my-pi/pi-coding-agent/sdk";
|
||||
import * as sdkModule from "@oh-my-pi/pi-coding-agent/sdk";
|
||||
import type { AgentSession, AgentSessionEvent } from "@oh-my-pi/pi-coding-agent/session/agent-session";
|
||||
import { runSubprocess } from "@oh-my-pi/pi-coding-agent/task/executor";
|
||||
import type { AgentDefinition } from "@oh-my-pi/pi-coding-agent/task/types";
|
||||
import { EventBus } from "@oh-my-pi/pi-coding-agent/utils/event-bus";
|
||||
|
||||
const baseAgent: AgentDefinition = { name: "task", description: "test", systemPrompt: "test", source: "bundled" };
|
||||
|
||||
function assistantStopMessage(text: string): AssistantMessage {
|
||||
return {
|
||||
role: "assistant",
|
||||
content: [{ type: "text", text }],
|
||||
api: "openai-responses",
|
||||
provider: "openai",
|
||||
model: "mock",
|
||||
usage: {
|
||||
input: 0,
|
||||
output: 0,
|
||||
cacheRead: 0,
|
||||
cacheWrite: 0,
|
||||
totalTokens: 0,
|
||||
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
||||
},
|
||||
stopReason: "stop",
|
||||
timestamp: Date.now(),
|
||||
};
|
||||
}
|
||||
|
||||
interface AsyncQuiescenceHarness {
|
||||
session: AgentSession;
|
||||
prompts: string[];
|
||||
abortCalls: () => number;
|
||||
settleCalls: () => number;
|
||||
emitTerminalYield: (data: unknown) => void;
|
||||
finishJob: () => void;
|
||||
}
|
||||
|
||||
/**
|
||||
* Mock session with the owner-async surface the barrier drives:
|
||||
* `hasPendingAsyncWork` / `getAsyncJobSnapshot` / `settleAsyncWork`. The job
|
||||
* "finishes" during the first settle, which injects the async-result
|
||||
* follow-up (custom message_start) and a plain assistant reaction WITHOUT a
|
||||
* fresh yield — exactly the review's stale-yield scenario.
|
||||
*/
|
||||
function createAsyncSession(
|
||||
onPrompt: (params: { text: string; promptIndex: number; harness: AsyncQuiescenceHarness }) => void,
|
||||
): AsyncQuiescenceHarness {
|
||||
const listeners: Array<(event: AgentSessionEvent) => void> = [];
|
||||
const state = { messages: [] as AssistantMessage[] };
|
||||
const prompts: string[] = [];
|
||||
let abortCount = 0;
|
||||
let settleCount = 0;
|
||||
let pendingAsync = true;
|
||||
let runningJobs: Array<{ id: string; label?: string }> = [{ id: "job-1", label: "background build" }];
|
||||
let toolCallSeq = 0;
|
||||
|
||||
const emit = (event: AgentSessionEvent) => {
|
||||
for (const listener of [...listeners]) listener(event);
|
||||
};
|
||||
|
||||
const emitTerminalYield = (data: unknown) => {
|
||||
toolCallSeq += 1;
|
||||
emit({
|
||||
type: "tool_execution_end",
|
||||
toolCallId: `yield-${toolCallSeq}`,
|
||||
toolName: "yield",
|
||||
result: {
|
||||
content: [{ type: "text", text: "Result submitted." }],
|
||||
details: { status: "success", data },
|
||||
},
|
||||
} as AgentSessionEvent);
|
||||
};
|
||||
|
||||
const finishJob = () => {
|
||||
pendingAsync = false;
|
||||
runningJobs = [];
|
||||
// Owner job completed: the session injects the async-result follow-up
|
||||
// turn. The model reacts with text only — no fresh yield.
|
||||
emit({
|
||||
type: "message_start",
|
||||
message: {
|
||||
role: "custom",
|
||||
customType: "async-result",
|
||||
content: "<system-notice>Background job job-1 has completed.\nexit 1: build FAILED</system-notice>",
|
||||
display: true,
|
||||
attribution: "agent",
|
||||
timestamp: Date.now(),
|
||||
},
|
||||
} as AgentSessionEvent);
|
||||
const reaction = assistantStopMessage("The background build failed after I yielded.");
|
||||
state.messages.push(reaction);
|
||||
emit({ type: "message_end", message: reaction } as AgentSessionEvent);
|
||||
};
|
||||
|
||||
const harness: AsyncQuiescenceHarness = {
|
||||
session: undefined as unknown as AgentSession,
|
||||
prompts,
|
||||
abortCalls: () => abortCount,
|
||||
settleCalls: () => settleCount,
|
||||
emitTerminalYield,
|
||||
finishJob,
|
||||
};
|
||||
|
||||
const session = {
|
||||
state,
|
||||
agent: { state: { systemPrompt: ["test"] } },
|
||||
model: undefined,
|
||||
extensionRunner: undefined,
|
||||
sessionManager: { appendSessionInit: () => {} },
|
||||
getActiveToolNames: () => ["read", "yield"],
|
||||
getEnabledToolNames: () => ["read", "yield"],
|
||||
setActiveToolsByName: async (_toolNames: string[]) => {},
|
||||
subscribe: (listener: (event: AgentSessionEvent) => void) => {
|
||||
listeners.push(listener);
|
||||
return () => {
|
||||
const index = listeners.indexOf(listener);
|
||||
if (index >= 0) listeners.splice(index, 1);
|
||||
};
|
||||
},
|
||||
prompt: async (text: string) => {
|
||||
prompts.push(text);
|
||||
onPrompt({ text, promptIndex: prompts.length, harness });
|
||||
},
|
||||
waitForIdle: async () => {},
|
||||
getLastAssistantMessage: () => state.messages[state.messages.length - 1],
|
||||
hasPendingAsyncWork: () => pendingAsync,
|
||||
getAsyncJobSnapshot: () => ({ running: runningJobs, recent: [] }),
|
||||
settleAsyncWork: async () => {
|
||||
settleCount += 1;
|
||||
harness.finishJob();
|
||||
},
|
||||
abort: async () => {
|
||||
abortCount += 1;
|
||||
},
|
||||
dispose: async () => {},
|
||||
};
|
||||
harness.session = session as unknown as AgentSession;
|
||||
return harness;
|
||||
}
|
||||
|
||||
function mockCreateAgentSession(session: AgentSession) {
|
||||
return vi.spyOn(sdkModule, "createAgentSession").mockResolvedValue({
|
||||
session,
|
||||
extensionsResult: {} as unknown as LoadExtensionsResult,
|
||||
setToolUIContext: () => {},
|
||||
eventBus: new EventBus(),
|
||||
} as CreateAgentSessionResult);
|
||||
}
|
||||
|
||||
describe("runSubprocess async quiescence fresh-yield contract", () => {
|
||||
afterEach(() => {
|
||||
vi.restoreAllMocks();
|
||||
});
|
||||
|
||||
it("parks a pending yield, injects the result, and completes on the fresh yield", async () => {
|
||||
const harness = createAsyncSession(({ promptIndex, harness: h }) => {
|
||||
if (promptIndex === 1) {
|
||||
// Terminal yield while the background job is still running.
|
||||
h.emitTerminalYield({ report: "STALE: build passing (job still running)" });
|
||||
return;
|
||||
}
|
||||
if (promptIndex === 2) {
|
||||
// Async-pending notice: the model stands by. The job then
|
||||
// finishes during the barrier's settle.
|
||||
return;
|
||||
}
|
||||
// Reminder ladder after the async-result invalidated the yield:
|
||||
// submit the fresh yield that accounts for the job outcome.
|
||||
h.emitTerminalYield({ report: "FRESH: build failed, see job-1" });
|
||||
});
|
||||
mockCreateAgentSession(harness.session);
|
||||
|
||||
const result = await runSubprocess({
|
||||
cwd: "/tmp",
|
||||
agent: baseAgent,
|
||||
task: "do the work",
|
||||
index: 0,
|
||||
id: "quiescence-fresh-yield",
|
||||
});
|
||||
|
||||
// Run did not terminate on the parked yield: the barrier noticed, the
|
||||
// job settled, and the ladder demanded exactly one more prompt.
|
||||
expect(harness.prompts).toHaveLength(3);
|
||||
expect(harness.prompts[1]).toContain("yield was recorded");
|
||||
expect(harness.settleCalls()).toBe(1);
|
||||
// The parked yield stopped the turn without killing the run.
|
||||
expect(harness.abortCalls()).toBeGreaterThanOrEqual(1);
|
||||
// The fresh yield — not the stale one — is the result of record.
|
||||
expect(result.exitCode).toBe(0);
|
||||
expect(result.output).toContain("FRESH: build failed");
|
||||
expect(result.output).not.toContain("STALE");
|
||||
});
|
||||
|
||||
it("fails the run when the model never refreshes the superseded yield", async () => {
|
||||
const harness = createAsyncSession(({ promptIndex, harness: h }) => {
|
||||
if (promptIndex === 1) {
|
||||
h.emitTerminalYield({ report: "STALE: build passing (job still running)" });
|
||||
}
|
||||
// Notice and every reminder: the model never yields again.
|
||||
});
|
||||
mockCreateAgentSession(harness.session);
|
||||
|
||||
const result = await runSubprocess({
|
||||
cwd: "/tmp",
|
||||
agent: baseAgent,
|
||||
task: "do the work",
|
||||
index: 0,
|
||||
id: "quiescence-stale-refusal",
|
||||
});
|
||||
|
||||
// task + notice + full reminder ladder (3).
|
||||
expect(harness.prompts).toHaveLength(5);
|
||||
// Stale payload must not read as success; it ships only as failed-run
|
||||
// salvage with an explicit reason.
|
||||
expect(result.exitCode).toBe(1);
|
||||
expect(result.error).toContain("refreshed yield");
|
||||
expect(result.output).toContain("STALE: build passing");
|
||||
});
|
||||
|
||||
it("terminates immediately on yield when no owner async work is pending", async () => {
|
||||
const harness = createAsyncSession(({ promptIndex, harness: h }) => {
|
||||
if (promptIndex === 1) {
|
||||
h.finishJob();
|
||||
h.emitTerminalYield({ report: "done" });
|
||||
}
|
||||
});
|
||||
mockCreateAgentSession(harness.session);
|
||||
|
||||
const result = await runSubprocess({
|
||||
cwd: "/tmp",
|
||||
agent: baseAgent,
|
||||
task: "do the work",
|
||||
index: 0,
|
||||
id: "quiescence-no-async",
|
||||
});
|
||||
|
||||
expect(harness.prompts).toHaveLength(1);
|
||||
expect(result.exitCode).toBe(0);
|
||||
expect(result.output).toContain("done");
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user