fix(acp): show durable notice after gateway disconnect (#108827)

* fix(acp): record disconnect interruptions

* test(acp): tighten disconnect settlement coverage

Co-authored-by: NianJiuZst <180004567+NianJiuZst@users.noreply.github.com>

* fix(acp): preserve replay order under backpressure

Co-authored-by: NianJiuZst <180004567+NianJiuZst@users.noreply.github.com>

* test(acp): await accepted prompt ledger record

Co-authored-by: NianJiuZst <180004567+NianJiuZst@users.noreply.github.com>

* fix(acp): isolate ledger mutation queues

Co-authored-by: NianJiuZst <180004567+NianJiuZst@users.noreply.github.com>

---------

Co-authored-by: NianJiuZst <180004567+NianJiuZst@users.noreply.github.com>
Co-authored-by: Peter Steinberger <steipete@gmail.com>
This commit is contained in:
NianJiu
2026-07-17 01:01:10 +08:00
committed by GitHub
parent ba65a56ffd
commit 1eda44ee36
5 changed files with 447 additions and 70 deletions
+9
View File
@@ -115,6 +115,15 @@ describe("ACP translator event ledger replay", () => {
const promptPromise = firstAgent.prompt(createPromptRequest(created.sessionId, "Question"));
await waitForChatSend(firstRequestMock);
await vi.waitFor(async () => {
const replay = await eventLedger.readReplay({
sessionId: created.sessionId,
sessionKey: firstSession.sessionKey,
});
expect(
replay.events.some((event) => event.update.sessionUpdate === "user_message_chunk"),
).toBe(true);
});
const runId = firstSessionStore.getSession(created.sessionId)?.activeRunId;
if (!runId) {
throw new Error("Expected active ACP run");
+105 -1
View File
@@ -34,7 +34,7 @@ const update: SessionUpdate = {
availableCommands: [],
};
describe("AcpTranslatorSessionUpdates shutdown", () => {
describe("AcpTranslatorSessionUpdates", () => {
it("blocks ledger reads and writes after shutdown starts", async () => {
const ledger = createLedger();
const sessionUpdate = vi.fn(async () => {});
@@ -101,4 +101,108 @@ describe("AcpTranslatorSessionUpdates shutdown", () => {
expect(ledger.recordUpdate).not.toHaveBeenCalled();
});
it("preserves ledger order without waiting for ACP delivery", async () => {
let releaseFirstDelivery!: () => void;
const firstDelivery = new Promise<void>((resolve) => {
releaseFirstDelivery = resolve;
});
let deliveryCount = 0;
const ledger = createLedger();
const updates = createUpdates({
ledger,
sessionUpdate: vi.fn(async () => {
deliveryCount += 1;
if (deliveryCount === 1) {
await firstDelivery;
}
}),
});
const firstUpdate: SessionUpdate = {
sessionUpdate: "agent_message_chunk",
content: { type: "text", text: "first" },
};
const interruptionUpdate: SessionUpdate = {
sessionUpdate: "agent_message_chunk",
content: { type: "text", text: "interrupted" },
};
const firstEmit = updates.emit({
sessionId: "session-1",
sessionKey: "agent:main:session-1",
record: true,
update: firstUpdate,
});
await vi.waitFor(() => {
expect(ledger.recordUpdate).toHaveBeenCalledTimes(1);
});
await updates.emit({
sessionId: "session-1",
sessionKey: "agent:main:session-1",
record: true,
waitForDelivery: false,
update: interruptionUpdate,
});
expect(vi.mocked(ledger.recordUpdate).mock.calls.map(([params]) => params.update)).toEqual([
firstUpdate,
interruptionUpdate,
]);
await expect(Promise.race([firstEmit, Promise.resolve("pending")])).resolves.toBe("pending");
releaseFirstDelivery();
await firstEmit;
});
it("does not let a stalled ledger session block another session", async () => {
let releaseFirstWrite!: () => void;
const firstWrite = new Promise<void>((resolve) => {
releaseFirstWrite = resolve;
});
const ledger = createLedger();
vi.mocked(ledger.recordUpdate).mockImplementation(async ({ sessionId }) => {
if (sessionId === "session-1") {
await firstWrite;
}
});
const updates = createUpdates({ ledger });
const stalledEmit = updates.emit({
sessionId: "session-1",
sessionKey: "agent:main:session-1",
record: true,
waitForDelivery: false,
update,
});
await vi.waitFor(() => {
expect(ledger.recordUpdate).toHaveBeenCalledTimes(1);
});
const queuedSameSessionEmit = updates.emit({
sessionId: "session-1",
sessionKey: "agent:main:session-1",
record: true,
waitForDelivery: false,
update,
});
await updates.emit({
sessionId: "session-2",
sessionKey: "agent:main:session-2",
record: true,
waitForDelivery: false,
update,
});
expect(vi.mocked(ledger.recordUpdate).mock.calls.map(([params]) => params.sessionId)).toEqual([
"session-1",
"session-2",
]);
releaseFirstWrite();
await Promise.all([stalledEmit, queuedSameSessionEmit]);
expect(vi.mocked(ledger.recordUpdate).mock.calls.map(([params]) => params.sessionId)).toEqual([
"session-1",
"session-2",
"session-1",
]);
});
});
+77 -42
View File
@@ -33,6 +33,9 @@ function resolveLedgerSessionId(session: { sessionId: string; ledgerSessionId?:
/** Helper that keeps ACP client updates and replay ledger writes in sync. */
export class AcpTranslatorSessionUpdates {
private stopped = false;
// Queue each ledger session at emission time so a detached disconnect notice
// cannot overtake its older update or block unrelated session settlement.
private ledgerMutationTails = new Map<string, Promise<void>>();
constructor(private options: AcpTranslatorSessionUpdatesOptions) {}
@@ -108,22 +111,24 @@ export class AcpTranslatorSessionUpdates {
runId: string,
prompt: PromptRequest["prompt"],
): Promise<void> {
if (this.stopped) {
return;
}
try {
await this.options.eventLedger.recordUserPrompt({
sessionId: resolveLedgerSessionId(session),
sessionKey: session.sessionKey,
runId,
prompt,
});
} catch (err) {
this.options.log(
`event ledger prompt record failed for ${session.sessionId}: ${String(err)}`,
);
await this.markLedgerIncomplete(session);
}
await this.enqueueLedgerMutation(resolveLedgerSessionId(session), async () => {
if (this.stopped) {
return;
}
try {
await this.options.eventLedger.recordUserPrompt({
sessionId: resolveLedgerSessionId(session),
sessionKey: session.sessionKey,
runId,
prompt,
});
} catch (err) {
this.options.log(
`event ledger prompt record failed for ${session.sessionId}: ${String(err)}`,
);
await this.markLedgerIncomplete(session);
}
});
}
async emit(params: {
@@ -133,23 +138,33 @@ export class AcpTranslatorSessionUpdates {
runId?: string;
update: SessionUpdate;
record?: boolean;
waitForDelivery?: boolean;
}): Promise<void> {
if (this.stopped) {
return;
}
await this.options.connection.sessionUpdate({
const delivery = this.options.connection.sessionUpdate({
sessionId: params.sessionId,
update: params.update,
});
if (params.record && params.sessionKey) {
await this.recordLedgerUpdate({
sessionId: params.sessionId,
sessionKey: params.sessionKey,
...(params.ledgerSessionId ? { ledgerSessionId: params.ledgerSessionId } : {}),
...(params.runId ? { runId: params.runId } : {}),
update: params.update,
const recording =
params.record && params.sessionKey
? this.recordLedgerUpdate({
sessionId: params.sessionId,
sessionKey: params.sessionKey,
...(params.ledgerSessionId ? { ledgerSessionId: params.ledgerSessionId } : {}),
...(params.runId ? { runId: params.runId } : {}),
update: params.update,
})
: undefined;
if (params.waitForDelivery === false) {
void delivery.catch((err: unknown) => {
this.options.log(`session update delivery failed for ${params.sessionId}: ${String(err)}`);
});
} else {
await delivery;
}
await recording;
}
async sendAvailableCommands(
@@ -175,24 +190,44 @@ export class AcpTranslatorSessionUpdates {
runId?: string;
update: SessionUpdate;
}): Promise<void> {
if (this.stopped) {
return;
}
try {
await this.options.eventLedger.recordUpdate({
sessionId: params.ledgerSessionId ?? params.sessionId,
sessionKey: params.sessionKey,
...(params.runId ? { runId: params.runId } : {}),
update: params.update,
});
} catch (err) {
this.options.log(`event ledger update record failed for ${params.sessionId}: ${String(err)}`);
await this.markLedgerIncomplete({
sessionId: params.sessionId,
sessionKey: params.sessionKey,
...(params.ledgerSessionId ? { ledgerSessionId: params.ledgerSessionId } : {}),
});
}
await this.enqueueLedgerMutation(params.ledgerSessionId ?? params.sessionId, async () => {
if (this.stopped) {
return;
}
try {
await this.options.eventLedger.recordUpdate({
sessionId: params.ledgerSessionId ?? params.sessionId,
sessionKey: params.sessionKey,
...(params.runId ? { runId: params.runId } : {}),
update: params.update,
});
} catch (err) {
this.options.log(
`event ledger update record failed for ${params.sessionId}: ${String(err)}`,
);
await this.markLedgerIncomplete({
sessionId: params.sessionId,
sessionKey: params.sessionKey,
...(params.ledgerSessionId ? { ledgerSessionId: params.ledgerSessionId } : {}),
});
}
});
}
private enqueueLedgerMutation(
ledgerSessionId: string,
mutation: () => Promise<void>,
): Promise<void> {
const previous = this.ledgerMutationTails.get(ledgerSessionId) ?? Promise.resolve();
const pending = previous.then(mutation, mutation);
const tail = pending.catch(() => {});
this.ledgerMutationTails.set(ledgerSessionId, tail);
void tail.then(() => {
if (this.ledgerMutationTails.get(ledgerSessionId) === tail) {
this.ledgerMutationTails.delete(ledgerSessionId);
}
});
return pending;
}
private async markLedgerIncomplete(session: AcpTranslatorSessionRef): Promise<void> {
+204 -11
View File
@@ -3,6 +3,7 @@ import type { PromptRequest } from "@agentclientprotocol/sdk";
import { createInMemorySessionStore } from "@openclaw/acp-core/session";
import { describe, expect, it, vi } from "vitest";
import type { GatewayClient } from "../gateway/client.js";
import { createInMemoryAcpEventLedger } from "./event-ledger.js";
import { AcpGatewayAgent } from "./translator.js";
import {
createChatEvent,
@@ -38,6 +39,64 @@ function requireFirstRequestIdempotencyKey(requestMock: {
return idempotencyKey;
}
async function createDisconnectNoticeHarness(params: { sendAccepted: boolean }) {
const sessionId = "session-1";
const sessionKey = "agent:main:main";
const sessionStore = createInMemorySessionStore();
sessionStore.createSession({ sessionId, sessionKey, cwd: "/tmp" });
const eventLedger = createInMemoryAcpEventLedger();
await eventLedger.startSession({
sessionId,
sessionKey,
cwd: "/tmp",
complete: true,
});
const connection = createAcpConnection();
const sessionUpdate = vi.fn(
async (_params: Parameters<typeof connection.sessionUpdate>[0]) => {},
);
connection.sessionUpdate = sessionUpdate as typeof connection.sessionUpdate;
const request = vi.fn(async (method: string) => {
if (method === "chat.send") {
if (!params.sendAccepted) {
throw new Error("gateway closed (1006): connection lost");
}
return {};
}
if (method === "agent.wait") {
throw new Error("gateway closed (1006): connection lost");
}
return {};
}) as GatewayClient["request"];
const agent = new AcpGatewayAgent(connection, createAcpGateway(request), {
eventLedger,
sessionStore,
});
const promptPromise = promptAgent(agent, sessionId);
void promptPromise.catch(() => {});
await vi.waitFor(() => {
expect(request).toHaveBeenCalledWith("chat.send", expect.objectContaining({ sessionKey }), {
timeoutMs: null,
});
});
if (params.sendAccepted) {
await vi.waitFor(async () => {
const replay = await eventLedger.readReplay({ sessionId, sessionKey });
expect(
replay.events.some((event) => event.update.sessionUpdate === "user_message_chunk"),
).toBe(true);
});
}
return {
agent,
eventLedger,
promptPromise,
sessionId,
sessionKey,
sessionUpdate,
};
}
describe("acp translator stop reason mapping", () => {
it("error state resolves as end_turn, not refusal", async () => {
const { agent, promptPromise, runId } = await createPendingPromptHarness();
@@ -169,32 +228,83 @@ describe("acp translator stop reason mapping", () => {
it("rejects in-flight prompts when the gateway does not reconnect before the grace window", async () => {
vi.useFakeTimers();
try {
const { agent, promptPromise } = await createPendingPromptHarness();
void promptPromise.catch(() => {});
const { agent, eventLedger, promptPromise, sessionId, sessionKey, sessionUpdate } =
await createDisconnectNoticeHarness({ sendAccepted: true });
agent.handleGatewayDisconnect("1006: connection lost");
await vi.advanceTimersByTimeAsync(5_000);
await expect(promptPromise).rejects.toThrow("Gateway disconnected: 1006: connection lost");
const expectedText =
"[OpenClaw interruption] The Gateway disconnected after accepting this message, so its final outcome is unknown. Check the session before retrying.";
expect(sessionUpdate).toHaveBeenCalledWith({
sessionId,
update: {
sessionUpdate: "agent_message_chunk",
content: { type: "text", text: expectedText },
},
});
const replay = await eventLedger.readReplay({ sessionId, sessionKey });
const recordedNotices = replay.events.filter(
(event) =>
event.update.sessionUpdate === "agent_message_chunk" &&
event.update.content.type === "text" &&
event.update.content.text === expectedText,
);
expect(recordedNotices).toHaveLength(1);
} finally {
vi.useRealTimers();
}
});
it("makes the disconnect notice durable without waiting for ACP delivery", async () => {
vi.useFakeTimers();
let releaseDelivery: (() => void) | undefined;
try {
const { agent, eventLedger, promptPromise, sessionUpdate } =
await createDisconnectNoticeHarness({ sendAccepted: true });
const settleSpy = observeSettlement(promptPromise);
const originalRecordUpdate = eventLedger.recordUpdate;
let releaseRecord: (() => void) | undefined;
const recordBlocked = new Promise<void>((resolve) => {
releaseRecord = resolve;
});
let noticeRecordStarted = false;
eventLedger.recordUpdate = async (params) => {
if (params.update.sessionUpdate === "agent_message_chunk") {
noticeRecordStarted = true;
await recordBlocked;
}
await originalRecordUpdate(params);
};
const deliveryBlocked = new Promise<void>((resolve) => {
releaseDelivery = resolve;
});
sessionUpdate.mockImplementation(async () => await deliveryBlocked);
agent.handleGatewayDisconnect("1006: connection lost");
await vi.advanceTimersByTimeAsync(5_000);
await vi.waitFor(() => {
expect(noticeRecordStarted).toBe(true);
expect(sessionUpdate).toHaveBeenCalledTimes(1);
});
expect(settleSpy).not.toHaveBeenCalled();
releaseRecord?.();
await expect(promptPromise).rejects.toThrow("Gateway disconnected: 1006: connection lost");
} finally {
releaseDelivery?.();
vi.useRealTimers();
}
});
it("keeps pre-ack send disconnects inside the reconnect grace window", async () => {
vi.useFakeTimers();
try {
const request = vi.fn(async (method: string) => {
if (method === "chat.send") {
throw new Error("gateway closed (1006): connection lost");
}
return {};
}) as GatewayClient["request"];
const { agent, sessionId } = createSessionAgentHarness(request);
const promptPromise = promptAgent(agent, sessionId);
const { agent, eventLedger, promptPromise, sessionId, sessionKey, sessionUpdate } =
await createDisconnectNoticeHarness({ sendAccepted: false });
const settleSpy = observeSettlement(promptPromise);
await Promise.resolve();
expect(settleSpy).not.toHaveBeenCalled();
agent.handleGatewayDisconnect("1006: connection lost");
@@ -203,6 +313,24 @@ describe("acp translator stop reason mapping", () => {
await vi.advanceTimersByTimeAsync(1);
await expect(promptPromise).rejects.toThrow("Gateway disconnected: 1006: connection lost");
const expectedText =
"[OpenClaw interruption] The Gateway disconnected before OpenClaw could confirm whether this message was accepted, so its final outcome is unknown. Check the session before retrying.";
expect(sessionUpdate).toHaveBeenCalledWith({
sessionId,
update: {
sessionUpdate: "agent_message_chunk",
content: { type: "text", text: expectedText },
},
});
const replay = await eventLedger.readReplay({ sessionId, sessionKey });
expect(
replay.events.filter(
(event) =>
event.update.sessionUpdate === "agent_message_chunk" &&
event.update.content.type === "text" &&
event.update.content.text === expectedText,
),
).toHaveLength(1);
} finally {
vi.useRealTimers();
}
@@ -620,6 +748,71 @@ describe("acp translator stop reason mapping", () => {
await expect(promptPromise).resolves.toEqual({ stopReason: "end_turn" });
});
it("keeps a replacement prompt while a stale send failure races disconnect rejection", async () => {
vi.useFakeTimers();
try {
let rejectFirstSend: ((error: Error) => void) | undefined;
const firstSend = new Promise<never>((_, reject) => {
rejectFirstSend = reject;
});
const runIds: string[] = [];
const request = vi.fn(
(method: string, params?: Record<string, unknown>): Promise<unknown> => {
if (method === "chat.send") {
const runId = params?.idempotencyKey;
if (typeof runId === "string") {
runIds.push(runId);
}
return runIds.length === 1 ? firstSend : Promise.resolve({});
}
if (method === "agent.wait") {
return Promise.reject(new Error("gateway closed (1006): connection lost"));
}
return Promise.resolve({});
},
) as GatewayClient["request"];
const sessionStore = createInMemorySessionStore();
const sessionId = "session-1";
const sessionKey = "agent:main:main";
sessionStore.createSession({ sessionId, sessionKey, cwd: "/tmp" });
const agent = new AcpGatewayAgent(createAcpConnection(), createAcpGateway(request), {
sessionStore,
});
const firstPrompt = promptAgent(agent, sessionId, "first");
void firstPrompt.catch(() => {});
await vi.waitFor(() => {
expect(runIds).toHaveLength(1);
});
agent.handleGatewayDisconnect("1006: connection lost");
await vi.advanceTimersByTimeAsync(5_000);
await expect(firstPrompt).rejects.toThrow("Gateway disconnected: 1006: connection lost");
const secondPrompt = promptAgent(agent, sessionId, "second");
await vi.waitFor(() => {
expect(runIds).toHaveLength(2);
});
rejectFirstSend?.(new Error("gateway closed (1006): connection lost"));
await Promise.resolve();
await expect(Promise.race([secondPrompt, Promise.resolve("pending")])).resolves.toBe(
"pending",
);
await agent.handleGatewayEvent(
createChatEvent({
runId: runIds[1],
sessionKey,
seq: 1,
state: "final",
}),
);
await expect(secondPrompt).resolves.toEqual({ stopReason: "end_turn" });
} finally {
vi.useRealTimers();
}
});
it("does not let a stale disconnect deadline reject a newer prompt on the same session", async () => {
vi.useFakeTimers();
try {
+52 -16
View File
@@ -721,7 +721,7 @@ export class AcpGatewayAgent implements Agent {
if (status === "error") {
const current = pending();
if (current) {
this.rejectPendingPrompt(
await this.rejectPendingPrompt(
current,
new Error("Chat failed before the run started; try again."),
);
@@ -775,21 +775,21 @@ export class AcpGatewayAgent implements Agent {
}
};
void sendWithProvenanceFallback().catch((err: unknown) => {
void sendWithProvenanceFallback().catch(async (err: unknown) => {
const promptKey = this.pendingPromptKey(params.sessionId, runId);
if (
isGatewayCloseError(err) &&
(this.getPendingPrompt(params.sessionId, runId) || this.settlingPromptKeys.has(promptKey))
) {
if (this.settlingPromptKeys.has(promptKey)) {
return;
}
this.clearApprovalRelaysForPrompt(params.sessionId, runId, { denyActive: true });
this.pendingPrompts.delete(params.sessionId);
this.sessionStore.clearActiveRun(params.sessionId);
if (this.pendingPrompts.size === 0) {
this.clearDisconnectTimer();
if (isGatewayCloseError(err) && this.getPendingPrompt(params.sessionId, runId)) {
return;
}
reject(err instanceof Error ? err : new Error(String(err)));
const error = err instanceof Error ? err : new Error(String(err));
const current = this.getPendingPrompt(params.sessionId, runId);
if (current) {
await this.rejectPendingPrompt(current, error);
return;
}
reject(error);
});
});
}
@@ -1299,11 +1299,19 @@ export class AcpGatewayAgent implements Agent {
this.disconnectTimer.unref?.();
}
private rejectPendingPrompt(pending: PendingPrompt, error: Error): void {
private async rejectPendingPrompt(
pending: PendingPrompt,
error: Error,
options: { recordDisconnectNotice?: boolean } = {},
): Promise<void> {
const currentPending = this.getPendingPrompt(pending.sessionId, pending.idempotencyKey);
if (currentPending !== pending) {
return;
}
const promptKey = this.pendingPromptKey(pending.sessionId, pending.idempotencyKey);
// Claim before emitting so late Gateway events cannot settle this prompt twice.
this.settlingPromptKeys.add(promptKey);
this.clearApprovalRelaysForPrompt(pending.sessionId, pending.idempotencyKey, {
denyActive: true,
});
@@ -1312,7 +1320,33 @@ export class AcpGatewayAgent implements Agent {
if (this.pendingPrompts.size === 0) {
this.clearDisconnectTimer();
}
pending.reject(error);
try {
if (options.recordDisconnectNotice) {
const text = pending.sendAccepted
? "[OpenClaw interruption] The Gateway disconnected after accepting this message, so its final outcome is unknown. Check the session before retrying."
: "[OpenClaw interruption] The Gateway disconnected before OpenClaw could confirm whether this message was accepted, so its final outcome is unknown. Check the session before retrying.";
// Make replay durable before rejecting, but do not let ACP client backpressure
// extend the disconnect deadline indefinitely.
await this.sessionUpdates.emit({
sessionId: pending.sessionId,
sessionKey: pending.sessionKey,
...(pending.ledgerSessionId ? { ledgerSessionId: pending.ledgerSessionId } : {}),
runId: pending.idempotencyKey,
record: true,
waitForDelivery: false,
update: {
sessionUpdate: "agent_message_chunk",
content: { type: "text", text },
},
});
}
} catch (noticeError) {
this.log(`disconnect notice failed for ${pending.idempotencyKey}: ${String(noticeError)}`);
} finally {
pending.reject(error);
this.settlingPromptKeys.delete(promptKey);
}
}
private clearPendingDisconnectState(
@@ -1394,9 +1428,10 @@ export class AcpGatewayAgent implements Agent {
this.log(`agent.wait reconcile failed for ${pending.idempotencyKey}: ${String(err)}`);
if (deadlineExpired) {
if (this.shouldRejectPendingAtDisconnectDeadline(pending, disconnectContext)) {
this.rejectPendingPrompt(
await this.rejectPendingPrompt(
pending,
new Error(`Gateway disconnected: ${disconnectContext.reason}`),
{ recordDisconnectNotice: true },
);
return false;
}
@@ -1424,9 +1459,10 @@ export class AcpGatewayAgent implements Agent {
if (!currentDisconnectContext) {
return false;
}
this.rejectPendingPrompt(
await this.rejectPendingPrompt(
currentPending,
new Error(`Gateway disconnected: ${currentDisconnectContext.reason}`),
{ recordDisconnectNotice: true },
);
return false;
}