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
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
### Fixed
- Keep session-initiated lifecycle boundaries separate from queued SDK turns so a delayed predecessor completion cannot publish the successor's receipt with the wrong final answer. Preserve attached invocations and existing deadline, abort, and retired-owner isolation.
- Preserve a single public lifecycle boundary across accepted retries that suppress their predecessor's end, including tokenless interactive retries, so reopening the session retains terminal delivery and idle checkpoints.
2 changes: 2 additions & 0 deletions packages/coding-agent/src/extensibility/extensions/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -764,6 +764,8 @@ export interface AgentStartEvent extends SharedAgentStartEvent {
sdkRunToken?: string;
/** Complete set of SDK queue owners consumed together by this run. */
sdkRunTokens?: string[];
/** Internal public lifecycle boundary; accepted retries with a suppressed predecessor end retain it. */
lifecycleScope?: AttemptScope;
}

/** Fired when an agent run fails before emitting agent_end. The error is the
Expand Down
251 changes: 251 additions & 0 deletions packages/coding-agent/src/sdk/host/session-runtime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { chmod, mkdir, mkdtemp, readdir, readFile, rm, writeFile } from "node:fs
import * as os from "node:os";
import * as path from "node:path";
import { Agent, markNonDispatchedToolEvent, ThinkingLevel } from "@gajae-code/agent-core";
import { createAttemptMinter } from "@gajae-code/agent-core/attempt-scope";
import type { UserMessage } from "@gajae-code/ai/core";
import { createMockModel } from "@gajae-code/ai/providers/mock";
import { logger } from "@gajae-code/utils";
Expand Down Expand Up @@ -5139,6 +5140,54 @@ test("SDK turn.steer preserves its expected run token and propagates a stale-run
});

describe("post-acceptance invocation terminalization", () => {
test("preserves external prompt correlation across a retryable provider error", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "gjc-retryable-provider-correlation-"));
let harness: InvocationHarness | undefined;
try {
harness = await invocationHarness("retryable-provider-correlation", cwd, {
sendUserMessage: async (_content, options) => {
await options?.onPreflightAcceptCommit?.();
await neverSettlingPromise();
},
});
const accepted = await harness.control("turn.prompt", {
text: "recover after a provider overload",
clientRef: "retryable-provider-correlation",
});
expect(accepted.ok).toBe(true);
const correlation = {
commandId: accepted.result?.commandId,
turnId: accepted.result?.turnId,
};
const retryScope = {};

for (let attempt = 0; attempt < 3; attempt++) {
await harness.emit("agent_start", { lifecycleScope: retryScope });
}
await harness.emit("agent_end", {
messages: [{ role: "assistant", stopReason: "stop", content: [{ type: "text", text: "recovered" }] }],
});

const terminal = await settledStatus(harness, "turn.result", { kind: "prompt", ...correlation });
expect(terminal).toMatchObject({
status: "terminal_ok",
commandId: correlation.commandId,
turnId: correlation.turnId,
outcome: { kind: "stopped", reason: "end_turn" },
});
expect(harness.broadcasts.filter(frame => frame.kind === "agent_end").at(-1)).toMatchObject({
payload: {
commandId: correlation.commandId,
turnId: correlation.turnId,
outcome: { kind: "stopped", reason: "end_turn" },
},
});
} finally {
await harness?.stop();
await rm(cwd, { recursive: true, force: true });
}
});

test.each([
{ status: 402, code: "provider_http_402" },
{ status: 429, code: "provider_http_429" },
Expand Down Expand Up @@ -5786,6 +5835,85 @@ describe("post-acceptance invocation terminalization", () => {
}
});

test("tokenless retry boundaries release replacement terminal publication and replay", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "sdk-tokenless-retry-reopen-"));
const sessionId = "tokenless-retry-reopen";
const harness = await invocationHarness(sessionId, cwd, {});
const scopes = createAttemptMinter();
const activity = spyOn(SessionSdkSessionRuntime.prototype, "reportActivity");
try {
// Accepted ordinary retries suppress predecessor ends. Three attempts
// therefore share one public lifecycle boundary, not three unmatched ends.
const lifecycleScope = scopes.mint("main");
for (let attempt = 0; attempt < 3; attempt++) {
await harness.emit("agent_start", { lifecycleScope });
}
await harness.emit("agent_end");
for (let turn = 0; turn < 2; turn++) {
await harness.emit("agent_start", { lifecycleScope: scopes.mint("main") });
await harness.emit("agent_end");
}
await harness.switchSession("different-session");
await harness.switchSession(sessionId);
const priorEnds = harness.broadcasts.filter(frame => frame.kind === "agent_end").length;
await harness.emit("agent_start", { lifecycleScope: scopes.mint("main") });
await harness.emit("agent_end");
expect(harness.broadcasts.filter(frame => frame.kind === "agent_end")).toHaveLength(priorEnds + 1);
expect(activity.mock.calls.at(-1)?.[0]).toBe("idle");
expect(await harness.requestOnSession(sessionId, { type: "event_replay", sinceSeq: 0 })).toMatchObject({
type: "event_replay_result",
events: expect.arrayContaining([expect.objectContaining({ kind: "agent_end" })]),
});
} finally {
activity.mockRestore();
await harness.stop();
await rm(cwd, { recursive: true, force: true });
}
});

test("tokenless retry preserves attached owners alongside newly promoted SDK work", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "sdk-tokenless-retry-owners-"));
const promotions: Array<(promotion: { startsOwnRun: boolean }) => void> = [];
const harness = await invocationHarness("tokenless-retry-owners", cwd, {
sendUserMessage: async (_content, options) => {
await options?.onPreflightAcceptCommit?.();
if (options?.onQueuedPromoted) promotions.push(options.onQueuedPromoted);
options?.onPreflightAccepted?.();
},
});
try {
const lifecycleScope = createAttemptMinter().mint("main");
await harness.emit("agent_start", { lifecycleScope });
const attached = await harness.control("turn.follow_up", { text: "attached work" });
const queued = await harness.control("turn.follow_up", { text: "promoted work" });
expect(attached.ok).toBe(true);
expect(queued.ok).toBe(true);
expect(promotions).toHaveLength(2);
promotions[0]!({ startsOwnRun: false });
promotions[1]!({ startsOwnRun: true });
await harness.emit("agent_start", { lifecycleScope });
await harness.emit("agent_end", { messages: [{ role: "assistant", content: "shared retry final" }] });
for (const accepted of [attached, queued]) {
expect(
await harness.query("turn.result", {
kind: "prompt",
commandId: accepted.result?.commandId,
turnId: accepted.result?.turnId,
}),
).toMatchObject({
result: {
status: "terminal_ok",
receiptState: "present",
content: { text: "shared retry final" },
},
});
}
} finally {
await harness.stop();
await rm(cwd, { recursive: true, force: true });
}
});

test("a prompt killed by a provider stream interrupt reports a terminal failed status", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "gjc-terminalize-prompt-"));
try {
Expand Down Expand Up @@ -5851,6 +5979,129 @@ describe("post-acceptance invocation terminalization", () => {
await rm(cwd, { recursive: true, force: true });
}
});
test("unowned predecessor end cannot publish a queued successor receipt", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "sdk-unowned-predecessor-"));
const promotions: Array<(promotion: { startsOwnRun: boolean }) => void> = [];
const harness = await invocationHarness("unowned-predecessor", cwd, {
sendUserMessage: async (_content, options) => {
await options?.onPreflightAcceptCommit?.();
if (options?.onQueuedPromoted) promotions.push(options.onQueuedPromoted);
options?.onPreflightAccepted?.();
},
});
try {
await harness.emit("agent_start");
const accepted = await harness.control("turn.follow_up", { text: "queued verification" });
expect(accepted.ok).toBe(true);
expect(promotions).toHaveLength(1);
promotions[0]!({ startsOwnRun: true });
await harness.emit("agent_start");
const selector = {
kind: "prompt",
commandId: accepted.result?.commandId,
turnId: accepted.result?.turnId,
};
await harness.emit("agent_end", { messages: [{ role: "assistant", content: "predecessor final" }] });
expect(await harness.query("turn.result", selector)).toMatchObject({
result: { status: "in_flight", receiptState: "absent" },
});
await harness.emit("agent_end", { messages: [{ role: "assistant", content: "verification final" }] });
expect(await harness.query("turn.result", selector)).toMatchObject({
result: {
status: "terminal_ok",
turnId: accepted.result?.turnId,
receiptState: "present",
content: { text: "verification final" },
},
});
} finally {
await harness.stop();
await rm(cwd, { recursive: true, force: true });
}
});

test("unowned run stays paired after predecessor SDK references retire", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "sdk-unowned-retirement-"));
const promotions: Array<(promotion: { startsOwnRun: boolean }) => void> = [];
const harness = await invocationHarness("unowned-retirement", cwd, {
sendUserMessage: async (_content, options) => {
await options?.onPreflightAcceptCommit?.();
if (options?.onQueuedPromoted) promotions.push(options.onQueuedPromoted);
options?.onPreflightAccepted?.();
},
});
try {
const first = await harness.control("turn.follow_up", { text: "first" });
expect(promotions).toHaveLength(1);
promotions[0]!({ startsOwnRun: true });
await harness.emit("agent_start");
await harness.emit("agent_start");
await harness.emit("agent_end", { messages: [{ role: "assistant", content: "first final" }] });
expect(
await harness.query("turn.result", {
kind: "prompt",
commandId: first.result?.commandId,
turnId: first.result?.turnId,
}),
).toMatchObject({ result: { status: "terminal_ok", content: { text: "first final" } } });
const next = await harness.control("turn.follow_up", { text: "next" });
expect(promotions).toHaveLength(2);
promotions[1]!({ startsOwnRun: true });
await harness.emit("agent_start");
const selector = { kind: "prompt", commandId: next.result?.commandId, turnId: next.result?.turnId };
await harness.emit("agent_end", { messages: [{ role: "assistant", content: "unowned final" }] });
expect(await harness.query("turn.result", selector)).toMatchObject({
result: { status: "in_flight", receiptState: "absent" },
});
await harness.emit("agent_end", { messages: [{ role: "assistant", content: "next final" }] });
expect(await harness.query("turn.result", selector)).toMatchObject({
result: { status: "terminal_ok", content: { text: "next final" } },
});
} finally {
await harness.stop();
await rm(cwd, { recursive: true, force: true });
}
});

test("unowned run retains attached invocation across predecessor completion", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "sdk-unowned-attached-"));
const promotions: Array<(promotion: { startsOwnRun: boolean }) => void> = [];
const harness = await invocationHarness("unowned-attached", cwd, {
sendUserMessage: async (_content, options) => {
await options?.onPreflightAcceptCommit?.();
if (options?.onQueuedPromoted) promotions.push(options.onQueuedPromoted);
options?.onPreflightAccepted?.();
},
});
try {
await harness.control("turn.follow_up", { text: "predecessor" });
expect(promotions).toHaveLength(1);
promotions[0]!({ startsOwnRun: true });
await harness.emit("agent_start");
await harness.emit("agent_start");
const attached = await harness.control("turn.follow_up", { text: "attach to current unowned run" });
expect(attached.ok).toBe(true);
expect(promotions).toHaveLength(2);
promotions[1]!({ startsOwnRun: false });
const selector = {
kind: "prompt",
commandId: attached.result?.commandId,
turnId: attached.result?.turnId,
};
await harness.emit("agent_end", { messages: [{ role: "assistant", content: "predecessor final" }] });
expect(await harness.query("turn.result", selector)).toMatchObject({
result: { status: expect.stringMatching(/accepted|in_flight/), receiptState: "absent" },
});
await harness.emit("agent_end", { messages: [{ role: "assistant", content: "attached final" }] });
expect(await harness.query("turn.result", selector)).toMatchObject({
result: { status: "terminal_ok", content: { text: "attached final" } },
});
} finally {
await harness.stop();
await rm(cwd, { recursive: true, force: true });
}
});

test("immediate prompt after abort ack is not terminalized by the aborted turn's delayed agent_end", async () => {
const cwd = await mkdtemp(path.join(os.tmpdir(), "gjc-abort-immediate-prompt-"));
try {
Expand Down
Loading
Loading