fix(ui): retire submitted attachments from newer drafts (#162898)

Keep pending attachment claims with the submit guard until durable admission. Retire only unchanged files in the captured current composer scope while preserving newer drafts and overlapping submissions.

Related: #137427. Broader pre-admission handoff work remains separate.
This commit is contained in:
Peter Steinberger 2026-10-01 11:47:16 -07:00 • committed by GitHub
parent 06b95be505
commit b820a807e9
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
6 changed files with 539 additions and 169 deletions

View file

@ -260,6 +260,12 @@ group-targeted New Session remains immediate.
Panes share outbox recovery for the same conversation. Activity in another
conversation does not restart that recovery; reconnecting checks every saved outbox.
Files remain in the composer while a submitted message waits for attachment
storage. If you begin another draft during that wait, saving the first message
removes its unchanged files unless another pending submission still owns them.
Your newer text and added or replaced files stay in the composer. A storage
failure keeps those files available for retry.
On wide desktop panes, a compact rail of horizontal marks sits in the transcript's left gutter. Hover for a short message preview, or click a mark to jump to that message. Tab focuses the rail; arrow keys move between marks, Enter or Space jumps, Home and End select the endpoints, and Escape dismisses the preview. At rest, all marks are identical 8 × 2px strokes at 12px spacing. They stay faint; marks for messages currently visible in the transcript light up together as you scroll. Hovering a mark grows it to 32px and lights only that mark in text color, with progressively shorter strokes across three neighbors on either side. The wave and highlighted mark stay in place while the pointer moves onto the preview card. Leaving the rail and card, or pressing Escape, clears the wave. The other marks keep their resting colors. Outside that hover range, widths stay fixed. An empty message preview shows “Preview unavailable.” Each visible user message has a mark, and assistant messages from the same run share one mark. Tool calls, results, and progress alone do not create marks. An assistant mark jumps to the first currently displayed response in its run and previews the latest displayed response. Its identity stays stable as streaming output becomes persisted history, and its current-position highlight follows later response content in the same run. The rail covers loaded history; messages without run identity retain their transcript grouping. Long rails scroll internally within 45% of the viewport height, with fades only at ends that hide more messages. Scrolling the transcript keeps the current mark visible; you can also scroll the rail to explore other messages. The rail stays hidden on mobile, in narrow or short panes, and when your saved message width leaves too little gutter space. A jump briefly tints the target message with a soft background, fading over 1.2 seconds without a border or ring. Reduced motion disables mark transitions and shows the target tint statically for one second.
Session dashboards follow the selected conversation's agent, including when multiple agents each use a `global` session. Split panes keep their owners separate; panes showing the same agent and conversation share dashboard updates.

View file

@ -10,6 +10,8 @@ import {
createChatFlowE2eSuite,
expectRequestCountStable,
installMockGateway,
requireRecord,
requireString,
} from "./chat-flow.test-support.ts";
import {
holdOutboxPreviewReads,
@ -637,7 +639,22 @@ suite.define(() => {
}),
);
expect(await composerFor(page).inputValue()).toBe("Mock Gateway: newer input must survive");
expect(await paneFor(page).locator(".chat-attachment-thumb").count()).toBe(1);
expect(await paneFor(page).locator(".chat-attachment-thumb").count()).toBe(0);
const firstRunId = requireString(
requireRecord(sent.params).idempotencyKey,
"first send run id",
);
await gateway.resolveDeferred("chat.send");
await gateway.emitChatFinal({
runId: firstRunId,
text: "Mock Gateway: first message completed.",
});
await paneFor(page).getByRole("button", { name: "Send message", exact: true }).click();
const next = await gateway.waitForRequest("chat.send", { after: 1 });
const nextParams = requireRecord(next.params);
expect(nextParams.message).toBe("Mock Gateway: newer input must survive");
expect(nextParams).not.toHaveProperty("attachments");
expect(requireString(nextParams.idempotencyKey, "second send run id")).not.toBe(firstRunId);
});
});
it("keeps another credential owner isolated and retains a corrupt bundle without sending partial content", async () => {

View file

@ -0,0 +1,386 @@
// @vitest-environment jsdom
import { expectDefined } from "@openclaw/normalization-core";
import { describe, expect, it, onTestFinished, vi } from "vitest";
import { createDeferred } from "../../../../test/helpers/promise.js";
import * as outboxPayloadStore from "../../lib/chat/outbox-payload-store.runtime.ts";
import {
readStoredOutboxStore,
storageTargetForGateway,
subscribeStoredChatOutboxChanges,
} from "../../lib/chat/outbox-store.ts";
import { getChatAttachmentDataUrl } from "./attachment-payload-store.ts";
import {
createDeliveryAttachmentBatch,
createStagedAttachment,
} from "./chat-delivery-attachments.test-support.ts";
import { makeChatHost } from "./chat-host.test-support.ts";
import { handleSendChat } from "./chat-send-submit.ts";
import { listStoredChatOutboxes } from "./composer-persistence.ts";
import { useChatSendBrowserFixture } from "./outbox-browser.test-support.ts";
import { prepareOutboxPayload } from "./outbox-payloads.ts";
useChatSendBrowserFixture();
describe("chat attachment admission", () => {
it.each(["connection", "recovery owner", "selected agent"])(
"ignores an attachment admission from a replaced %s",
async (changedOwner) => {
const { attachments, dataUrls } = createDeliveryAttachmentBatch();
const host = makeChatHost({
requestHandlers: {},
connected: false,
chatMessage: "Stale submission",
chatAttachments: attachments,
});
if (changedOwner === "selected agent") {
host.sessionKey = "global";
host.assistantAgentId = "main";
host.agentsList = { defaultId: "main", mainKey: "main", scope: "global" };
}
const entered = [createDeferred(), createDeferred()];
const releases = [createDeferred(), createDeferred()];
let writes = 0;
const writePayload = outboxPayloadStore.writeOutboxPayload;
vi.spyOn(outboxPayloadStore, "writeOutboxPayload").mockImplementation(async (...args) => {
const index = writes++;
entered[index]?.resolve();
await releases[index]?.promise;
return writePayload(...args);
});
const stale = handleSendChat(host);
let current: ReturnType<typeof handleSendChat> | undefined;
try {
await Promise.race([
entered[0]!.promise,
stale.then(() => {
throw new Error("Stale submission ended before payload write");
}),
]);
if (changedOwner === "connection") {
host.connectionEpoch = (host.connectionEpoch ?? 0) + 1;
} else if (changedOwner === "selected agent") {
host.assistantAgentId = "other";
} else {
vi.spyOn(expectDefined(host.client, "client"), "recoveryScope", "get").mockReturnValue(
"new-synthetic-principal",
);
}
host.chatMessage = "Current submission";
current = handleSendChat(host);
await Promise.race([
entered[1]!.promise,
current.then(() => {
throw new Error("Current submission ended before payload write");
}),
]);
host.chatMessage = "Newer draft";
releases[1]!.resolve();
await current;
expect(host.chatMessage).toBe("Newer draft");
expect(host.chatAttachments).toEqual([]);
} finally {
releases.forEach((release) => release.resolve());
await Promise.all([stale, current]);
}
const queued = listStoredChatOutboxes(host)[0]?.queue ?? [];
expect(queued.map((item) => item.text)).toEqual(["Current submission"]);
const hydrated = await prepareOutboxPayload(host, expectDefined(queued[0], "current input"));
expect(
hydrated.status === "ready"
? hydrated.update.attachments?.map(getChatAttachmentDataUrl)
: [],
).toEqual(dataUrls);
expect(host.request).not.toHaveBeenCalled();
},
);
it("clears the matching composer after overlapping attachment submissions settle", async () => {
const { attachments } = createDeliveryAttachmentBatch();
const host = makeChatHost({
requestHandlers: {},
connected: false,
chatMessage: "First overlapping prompt",
chatAttachments: attachments,
});
const entered = [createDeferred(), createDeferred()];
const release = createDeferred();
let writes = 0;
const writePayload = outboxPayloadStore.writeOutboxPayload;
vi.spyOn(outboxPayloadStore, "writeOutboxPayload").mockImplementation(async (...args) => {
const writeIndex = writes++;
entered[writeIndex]?.resolve();
await release.promise;
return writePayload(...args);
});
const first = handleSendChat(host);
let second: ReturnType<typeof handleSendChat> | undefined;
try {
await Promise.race([
entered[0]!.promise,
first.then(() => {
throw new Error("First submission ended before payload write");
}),
]);
host.chatMessage = "Second overlapping prompt";
second = handleSendChat(host);
await Promise.race([
entered[1]!.promise,
second.then(() => {
throw new Error("Second submission ended before payload write");
}),
]);
} finally {
release.resolve();
await Promise.all([first, second]);
}
expect(listStoredChatOutboxes(host)[0]?.queue.map((item) => item.text)).toEqual([
"First overlapping prompt",
"Second overlapping prompt",
]);
expect(host.chatMessage).toBe("");
expect(host.chatAttachments).toEqual([]);
expect(host.request).not.toHaveBeenCalled();
});
it.each(["committed", "payload failure", "metadata failure", "unrelated send"])(
"preserves a newer draft while settling submitted attachment ownership (%s)",
async (outcome) => {
const { attachments, dataUrls } = createDeliveryAttachmentBatch();
const added = createStagedAttachment("newer-draft-file");
const replacement = { ...attachments[1]!, fileName: "replacement.pdf" };
const newerReply = { messageId: "newer-quote", text: "Newer quote" };
const newerMentions = [{ profileId: "alex", start: 0, end: 5 }];
const olderStarted = createDeferred();
const olderAck = createDeferred<{ runId: string; status: string }>();
const host = makeChatHost({
requestHandlers: {
"chat.send": () => {
olderStarted.resolve();
return olderAck.promise;
},
},
connected: false,
chatMessage: "First prompt",
chatAttachments: attachments,
});
if (outcome === "unrelated send") {
host.connected = true;
host.chatMessage = "Older text-only send";
host.chatAttachments = [];
const older = handleSendChat(host);
onTestFinished(async () => {
olderAck.resolve({ runId: "older-run", status: "started" });
await older;
});
await Promise.race([
olderStarted.promise,
older.then(() => {
throw new Error("Older send ended before transport");
}),
]);
host.chatMessage = "First prompt";
host.chatAttachments = attachments;
}
const started = createDeferred();
const release = createDeferred();
const writePayload = outboxPayloadStore.writeOutboxPayload;
vi.spyOn(outboxPayloadStore, "writeOutboxPayload").mockImplementationOnce(async (...args) => {
started.resolve();
await release.promise;
return outcome === "payload failure"
? { status: "failed", reason: "unavailable" }
: writePayload(...args);
});
const sending = handleSendChat(host);
try {
await Promise.race([
started.promise,
sending.then(() => {
throw new Error("Submission ended before payload write");
}),
]);
host.chatMessage = "@Alex newer input";
host.chatMentions = newerMentions;
host.chatReplyTarget = newerReply;
host.chatAttachments = [attachments[0]!, replacement, added];
expect(
listStoredChatOutboxes(host).flatMap((outbox) => outbox.queue.map((item) => item.text)),
).toEqual(outcome === "unrelated send" ? ["Older text-only send"] : []);
expect(host.chatAttachments.map(getChatAttachmentDataUrl)).toEqual([
...dataUrls,
getChatAttachmentDataUrl(added),
]);
if (outcome === "metadata failure") {
const write = sessionStorage.setItem.bind(sessionStorage);
const target = storageTargetForGateway(host.settings?.gatewayUrl);
vi.spyOn(sessionStorage, "setItem").mockImplementation((key, value) => {
if (key === target.key) {
throw new DOMException("quota exceeded", "QuotaExceededError");
}
write(key, value);
});
}
} finally {
release.resolve();
await sending;
}
expect(host.chatMessage).toBe("@Alex newer input");
expect(host.chatMentions).toEqual(newerMentions);
expect(host.chatReplyTarget).toEqual(newerReply);
expect(host.request).toHaveBeenCalledTimes(outcome === "unrelated send" ? 1 : 0);
if (outcome === "payload failure" || outcome === "metadata failure") {
expect(host.chatAttachments).toEqual([attachments[0], replacement, added]);
expect(listStoredChatOutboxes(host)).toEqual([]);
return;
}
expect(host.chatAttachments).toEqual([replacement, added]);
const first = expectDefined(
listStoredChatOutboxes(host)[0]?.queue.find((item) => item.text === "First prompt"),
"first submission",
);
const hydrated = await prepareOutboxPayload(host, first);
expect(
hydrated.status === "ready"
? hydrated.update.attachments?.map(getChatAttachmentDataUrl)
: [],
).toEqual(dataUrls);
await handleSendChat(host);
const queue =
listStoredChatOutboxes(host)[0]?.queue.filter(
(item) => item.text !== "Older text-only send",
) ?? [];
expect(queue.map((item) => item.id)).toEqual([first.id, expect.any(String)]);
expect(queue[1]?.id).not.toBe(first.id);
expect(queue[1]?.text).toContain("@Alex newer input");
expect(queue[1]?.attachments?.map((attachment) => attachment.id)).toEqual([
replacement.id,
added.id,
]);
},
);
it.each(["defaults", "route", "recovery owner", "reply"])(
"keeps the creation-time destination and input while payload admission awaits changed %s",
async (change) => {
const { attachments, dataUrls } = createDeliveryAttachmentBatch();
const replyTarget = {
messageId: "original-quote",
sourceMessageId: "original-entry",
text: "Original quote",
};
const newerReply = { messageId: "newer-quote", text: "Newer quote" };
const host = makeChatHost({
requestHandlers: {},
connected: false,
sessionKey: "main",
agentsList: { defaultId: "main", mainKey: "main", scope: "per-sender" },
chatMessage: "original destination",
chatAttachments: attachments,
chatReplyTarget: replyTarget,
});
const started = createDeferred();
const release = createDeferred();
const writePayload = outboxPayloadStore.writeOutboxPayload;
vi.spyOn(outboxPayloadStore, "writeOutboxPayload").mockImplementationOnce(async (...args) => {
started.resolve();
await release.promise;
return writePayload(...args);
});
const sending = handleSendChat(host);
try {
await Promise.race([
started.promise,
sending.then(() => {
throw new Error("Submission ended before payload write");
}),
]);
if (change === "defaults") {
host.agentsList = { defaultId: "main", mainKey: "current", scope: "per-sender" };
} else if (change === "route") {
host.sessionKey = "agent:main:elsewhere";
} else if (change === "recovery owner") {
vi.spyOn(
expectDefined(host.client, "payload client"),
"recoveryScope",
"get",
).mockReturnValue("different-owner");
} else {
host.chatReplyTarget = newerReply;
}
if (change !== "reply") {
host.chatMessage = "newer input";
}
} finally {
release.resolve();
await sending;
}
const expectedDraft = change === "reply" ? "original destination" : "newer input";
expect(host.chatMessage).toBe(expectedDraft);
expect(host.chatReplyTarget).toEqual(change === "reply" ? newerReply : replyTarget);
expect(host.request).not.toHaveBeenCalled();
if (change !== "defaults" && change !== "reply") {
expect(listStoredChatOutboxes(host)).toEqual([]);
expect(host.chatAttachments.map(getChatAttachmentDataUrl)).toEqual(dataUrls);
return;
}
const stored = expectDefined(listStoredChatOutboxes(host)[0], "captured outbox");
expect(stored).toMatchObject({ sessionKey: "agent:main:main", agentId: "main" });
expect(stored.queue[0]).toMatchObject({
sessionKey: "agent:main:main",
sendAttempts: 0,
replyToId: "original-entry",
});
const hydrated = await prepareOutboxPayload(
host,
expectDefined(stored.queue[0], "stored input"),
);
expect(
hydrated.status === "ready"
? hydrated.update.attachments?.map(getChatAttachmentDataUrl)
: [],
).toEqual(dataUrls);
expect(host.chatMessage).toBe(expectedDraft);
expect(host.chatAttachments).toEqual(change === "defaults" ? attachments : []);
expect(host.request).not.toHaveBeenCalled();
},
);
it("keeps a verified Blob admission when its notification changes recovery owner", async () => {
const { attachments, dataUrls } = createDeliveryAttachmentBatch();
const host = makeChatHost({
requestHandlers: {},
connected: false,
chatMessage: "committed input",
chatAttachments: attachments,
});
const client = expectDefined(host.client, "recovery client");
const target = storageTargetForGateway(host.settings?.gatewayUrl);
const recovery = vi.spyOn(client, "recoveryScope", "get");
const originalRecovery = client.recoveryScope;
const cleanup = vi.spyOn(outboxPayloadStore, "removeOutboxPayloads");
const stop = subscribeStoredChatOutboxChanges(() => {
recovery.mockReturnValue("new-synthetic-principal");
host.chatMessage = "newer input";
});
try {
await handleSendChat(host);
} finally {
stop();
}
const raw = readStoredOutboxStore(sessionStorage, target);
const queued = expectDefined(
Object.values(raw.sessions).flatMap((session) => session.queue ?? [])[0],
"verified committed input",
);
expect(queued).toMatchObject({ text: "committed input", sendAttempts: 0 });
expect(listStoredChatOutboxes(host)).toEqual([]);
expect(cleanup).not.toHaveBeenCalled();
recovery.mockReturnValue(originalRecovery);
const hydrated = await prepareOutboxPayload(host, queued);
expect(
hydrated.status === "ready" ? hydrated.update.attachments?.map(getChatAttachmentDataUrl) : [],
).toEqual(dataUrls);
expect(host.chatMessage).toBe("newer input");
expect(host.request).not.toHaveBeenCalled();
});
});

View file

@ -62,7 +62,11 @@ import {
} from "./chat-send-support.ts";
import { recordChatSendTiming } from "./chat-send-timing.ts";
import { getPendingChatPickerPatch } from "./chat-settings-patches.ts";
import { withChatSubmitGuard, withChatSubmitHandoff } from "./chat-submit-guard.ts";
import {
withChatSubmitGuard,
withChatSubmitHandoff,
type ChatSubmitGuard,
} from "./chat-submit-guard.ts";
import { attachmentBatchRejection } from "./components/chat-attachment-admission.ts";
import { recordNonTranscriptInputHistory } from "./input-history.ts";
import {
@ -136,6 +140,15 @@ export async function handleSendChat(
const attachmentsToSend = snapshotChatAttachments(
messageOverride == null ? host.chatAttachments : (opts?.attachmentsOverride ?? []),
);
const submitGuardOptions = {
attachments: messageOverride == null ? attachmentsToSend : [],
scope: resolveUiConversationIdentity(host, submittedSessionKey),
isCurrent: () =>
submittedOwnerIsCurrent() &&
host.client === submittedClient &&
host.connectionEpoch === submittedEpoch &&
host.sessionKey === submittedSessionKey,
};
const clearComposer = (retainAttachments: "none" | "annotations" | "all" = "none") =>
messageOverride == null
? clearSubmittedComposerState(
@ -246,7 +259,7 @@ export async function handleSendChat(
host.chatRunError = null;
const question = extractCompanionCommandQuestion(userMessage);
const submitKey = chatSubmitKey(host, "local", message, []);
await withChatSubmitGuard(host, submitKey, async () => {
await withChatSubmitGuard(host, submitKey, submitGuardOptions, async () => {
if (messageOverride == null) {
recordNonTranscriptInputHistory(host, userMessage);
clearComposer("all");
@ -266,30 +279,35 @@ export async function handleSendChat(
dispatchClientPresentation
) {
const submitKey = chatSubmitKey(host, "local", message, []);
const presentationResult = await withChatSubmitGuard(host, submitKey, async () => {
if (host.sessionKey !== submittedSessionKey) {
return "not-handled" as const;
}
let handled = false;
try {
handled = await dispatchClientPresentation(clientPresentation.action);
} catch {
// Presentation failures retain the established remote command path.
}
if (!handled) {
return "not-handled" as const;
}
// The awaited action may outlive its submitted session; never mutate a newly selected one.
if (host.sessionKey !== submittedSessionKey) {
const presentationResult = await withChatSubmitGuard(
host,
submitKey,
submitGuardOptions,
async () => {
if (host.sessionKey !== submittedSessionKey) {
return "not-handled" as const;
}
let handled = false;
try {
handled = await dispatchClientPresentation(clientPresentation.action);
} catch {
// Presentation failures retain the established remote command path.
}
if (!handled) {
return "not-handled" as const;
}
// The awaited action may outlive its submitted session; never mutate a newly selected one.
if (host.sessionKey !== submittedSessionKey) {
return "handled" as const;
}
host.chatRunError = null;
if (messageOverride == null) {
clearComposer();
recordNonTranscriptInputHistory(host, message);
}
return "handled" as const;
}
host.chatRunError = null;
if (messageOverride == null) {
clearComposer();
recordNonTranscriptInputHistory(host, message);
}
return "handled" as const;
});
},
);
// An in-flight identical submit is already deciding whether to handle or fall through.
if (presentationResult !== "not-handled") {
return undefined;
@ -311,7 +329,7 @@ export async function handleSendChat(
(isChatBusy(host) || isInitialChatHistoryUnavailable(host))
) {
const submitKey = chatSubmitKey(host, "detached", message, attachmentsToSend);
await withChatSubmitGuard(host, submitKey, async () => {
await withChatSubmitGuard(host, submitKey, submitGuardOptions, async () => {
if (!(await waitForSubmittedRoute(host, submittedSessionKey))) {
return;
}
@ -339,7 +357,7 @@ export async function handleSendChat(
return undefined;
}
const submitKey = chatSubmitKey(host, "local", message, attachmentsToSend);
await withChatSubmitGuard(host, submitKey, async () => {
await withChatSubmitGuard(host, submitKey, submitGuardOptions, async () => {
const admission = captureChatOutboxAdmission(host, host.sessionKey);
if (messageOverride == null) {
recordNonTranscriptInputHistory(host, userMessage);
@ -429,7 +447,7 @@ export async function handleSendChat(
};
if (waitsForPicker) {
const submitKey = chatSubmitKey(host, "local", message, attachmentsToSend);
await withChatSubmitGuard(host, submitKey, dispatchLocalCommand);
await withChatSubmitGuard(host, submitKey, submitGuardOptions, dispatchLocalCommand);
} else {
await dispatchLocalCommand();
}
@ -480,7 +498,7 @@ export async function handleSendChat(
effectiveMentions,
);
let accepted = false;
const submitMessage = async () => {
const submitMessage = async (guard: ChatSubmitGuard) => {
if (host.chatLoading && (intent || rawParsedCommand || isInlineEditSubmission)) {
// Commands and row edits retain their draft until history resolves.
if (!(await loadChatHistory(host))) {
@ -493,10 +511,7 @@ export async function handleSendChat(
}
const submittedAgentId = scopedAgentIdForSession(host, submittedSessionKey);
const submissionOwnerIsCurrent = () =>
submittedOwnerIsCurrent() &&
host.client === submittedClient &&
host.connectionEpoch === submittedEpoch &&
host.sessionKey === submittedSessionKey &&
submitGuardOptions.isCurrent() &&
visibleSessionMatches(host, submittedSessionKey, submittedAgentId);
if (!visibleSessionMatches(host, submittedSessionKey, submittedAgentId)) {
setChatError(host, t("mcpServers.sessionUnavailable"));
@ -608,6 +623,15 @@ export async function handleSendChat(
: undefined,
);
const admittedDurably = admissionResult === "admitted";
if (admittedDurably) {
guard.releaseAttachments();
if (messageOverride == null && !rawParsedCommand && !intent && submissionOwnerIsCurrent()) {
// Pending admissions retain their files until their own composer snapshot can settle.
host.chatAttachments = host.chatAttachments.filter(
(attachment) => !guard.canRetireAttachment(attachment),
);
}
}
if (resumedEdit) {
retireEditedQueuedMessageSource(host, admittedDurably, queued.attachments, resumedEdit);
}
@ -666,6 +690,14 @@ export async function handleSendChat(
recordChatSendTiming(host, pending, "queued-busy", submittedAtMs);
}
};
await withChatSubmitGuard(host, submitKey, submitMessage, submissionAction);
await withChatSubmitGuard(
host,
submitKey,
{
...submitGuardOptions,
action: submissionAction,
},
submitMessage,
);
return accepted;
}

View file

@ -15,12 +15,7 @@ import {
} from "../../lib/chat/commands.ts";
import { extractText } from "../../lib/chat/message-extract.ts";
import * as outboxPayloadStore from "../../lib/chat/outbox-payload-store.runtime.ts";
import {
captureChatOutboxAdmission,
readStoredOutboxStore,
storageTargetForGateway,
subscribeStoredChatOutboxChanges,
} from "../../lib/chat/outbox-store.ts";
import { captureChatOutboxAdmission } from "../../lib/chat/outbox-store.ts";
import {
createGatewayHarness,
createTestSessionCapability,
@ -3106,130 +3101,6 @@ describe("handleSendChat", () => {
);
});
it.each(["defaults", "route", "recovery owner", "reply"])(
"keeps the creation-time destination and input while payload admission awaits changed %s",
async (change) => {
const { attachments, dataUrls } = createDeliveryAttachmentBatch();
const replyTarget = {
messageId: "original-quote",
sourceMessageId: "original-entry",
text: "Original quote",
};
const newerReply = { messageId: "newer-quote", text: "Newer quote" };
const host = makeChatHost({
requestHandlers: {},
connected: false,
sessionKey: "main",
agentsList: { defaultId: "main", mainKey: "main", scope: "per-sender" },
chatMessage: "original destination",
chatAttachments: attachments,
chatReplyTarget: replyTarget,
});
const started = createDeferred();
const release = createDeferred();
const writePayload = outboxPayloadStore.writeOutboxPayload;
vi.spyOn(outboxPayloadStore, "writeOutboxPayload").mockImplementationOnce(async (...args) => {
started.resolve();
await release.promise;
return writePayload(...args);
});
const sending = handleSendChat(host);
try {
await Promise.race([
started.promise,
sending.then(() => {
throw new Error("Submission ended before payload write");
}),
]);
if (change === "defaults") {
host.agentsList = { defaultId: "main", mainKey: "current", scope: "per-sender" };
} else if (change === "route") {
host.sessionKey = "agent:main:elsewhere";
} else if (change === "recovery owner") {
vi.spyOn(
expectDefined(host.client, "payload client"),
"recoveryScope",
"get",
).mockReturnValue("different-owner");
} else {
host.chatReplyTarget = newerReply;
}
if (change !== "reply") {
host.chatMessage = "newer input";
}
} finally {
release.resolve();
await sending;
}
const expectedDraft = change === "reply" ? "original destination" : "newer input";
expect(host.chatMessage).toBe(expectedDraft);
expect(host.chatReplyTarget).toEqual(change === "reply" ? newerReply : replyTarget);
expect(host.request).not.toHaveBeenCalled();
if (change !== "defaults" && change !== "reply") {
expect(listStoredChatOutboxes(host)).toEqual([]);
expect(host.chatAttachments.map(getChatAttachmentDataUrl)).toEqual(dataUrls);
return;
}
const stored = expectDefined(listStoredChatOutboxes(host)[0], "captured outbox");
expect(stored).toMatchObject({ sessionKey: "agent:main:main", agentId: "main" });
expect(stored.queue[0]).toMatchObject({
sessionKey: "agent:main:main",
sendAttempts: 0,
replyToId: "original-entry",
});
const hydrated = await prepareOutboxPayload(
host,
expectDefined(stored.queue[0], "stored input"),
);
expect(
hydrated.status === "ready"
? hydrated.update.attachments?.map(getChatAttachmentDataUrl)
: [],
).toEqual(dataUrls);
expect(host.chatMessage).toBe(expectedDraft);
expect(host.request).not.toHaveBeenCalled();
},
);
it("keeps a verified Blob admission when its notification changes recovery owner", async () => {
const { attachments, dataUrls } = createDeliveryAttachmentBatch();
const host = makeChatHost({
requestHandlers: {},
connected: false,
chatMessage: "committed input",
chatAttachments: attachments,
});
const client = expectDefined(host.client, "recovery client");
const target = storageTargetForGateway(host.settings?.gatewayUrl);
const recovery = vi.spyOn(client, "recoveryScope", "get");
const originalRecovery = client.recoveryScope;
const cleanup = vi.spyOn(outboxPayloadStore, "removeOutboxPayloads");
const stop = subscribeStoredChatOutboxChanges(() => {
recovery.mockReturnValue("new-synthetic-principal");
host.chatMessage = "newer input";
});
try {
await handleSendChat(host);
} finally {
stop();
}
const raw = readStoredOutboxStore(sessionStorage, target);
const queued = expectDefined(
Object.values(raw.sessions).flatMap((session) => session.queue ?? [])[0],
"verified committed input",
);
expect(queued).toMatchObject({ text: "committed input", sendAttempts: 0 });
expect(listStoredChatOutboxes(host)).toEqual([]);
expect(cleanup).not.toHaveBeenCalled();
recovery.mockReturnValue(originalRecovery);
const hydrated = await prepareOutboxPayload(host, queued);
expect(
hydrated.status === "ready" ? hydrated.update.attachments?.map(getChatAttachmentDataUrl) : [],
).toEqual(dataUrls);
expect(host.chatMessage).toBe("newer input");
expect(host.request).not.toHaveBeenCalled();
});
it.each([
{ caller: "wrong queue", reason: "missing" },
{ caller: "wrong source tab", reason: "missing" },

View file

@ -1,15 +1,31 @@
import type { ChatQueueItem } from "../../lib/chat/chat-types.ts";
import type { ChatAttachment, ChatQueueItem } from "../../lib/chat/chat-types.ts";
import { sameQueuedDeliveryVersion } from "../../lib/chat/outbox-store-codec.ts";
import {
storedChatOutboxScopeKey,
type StoredChatOutboxScope,
} from "../../lib/chat/outbox-store-scope.ts";
import { visibleSessionMatches } from "../../lib/sessions/index.ts";
import { resolveUiConversationIdentity } from "../../lib/sessions/session-key.ts";
import { generateUUID } from "../../lib/uuid.ts";
import { isInitialChatHistoryUnavailable } from "./chat-history-state.ts";
import type { QueuedChatSendResult } from "./chat-outbox-drain.ts";
import { chatOutboxOwner } from "./chat-outbox-owner.ts";
import { readQueuedMessageById } from "./chat-queue.ts";
import type { ChatHost } from "./chat-send-contract.ts";
import { chatAttachmentDraftSignature } from "./durable-composer-persistence.ts";
import { hasDirectSessionRun, isChatBusy } from "./run-lifecycle.ts";
const submissionActionIds = new WeakMap<Event, string>();
type AttachmentAdmission = {
signatures: ReadonlySet<string>;
isCurrent(): boolean;
};
const pendingAttachmentAdmissions = new WeakMap<ChatHost, Set<AttachmentAdmission>>();
export type ChatSubmitGuard = {
releaseAttachments(): void;
canRetireAttachment(attachment: ChatAttachment): boolean;
};
function yieldChatSubmitToInput(): Promise<void> {
return new Promise<void>((resolve) => {
@ -96,10 +112,16 @@ export async function withChatSubmitHandoff(
export async function withChatSubmitGuard<T>(
host: ChatHost,
key: string,
run: () => Promise<T>,
action?: Event,
options: {
action?: Event;
attachments?: readonly ChatAttachment[];
scope: StoredChatOutboxScope;
isCurrent(): boolean;
},
run: (guard: ChatSubmitGuard) => Promise<T>,
): Promise<T | undefined> {
let guardKey = key;
const { action, scope } = options;
if (action) {
const actionId = submissionActionIds.get(action) ?? generateUUID();
submissionActionIds.set(action, actionId);
@ -110,9 +132,45 @@ export async function withChatSubmitGuard<T>(
return undefined;
}
guards.add(guardKey);
const scopeKey = storedChatOutboxScopeKey(scope);
const attachments: AttachmentAdmission = {
signatures: new Set(
options.attachments?.map((attachment) => chatAttachmentDraftSignature("", [attachment])),
),
isCurrent: () =>
options.isCurrent() &&
storedChatOutboxScopeKey(resolveUiConversationIdentity(host, host.sessionKey)) === scopeKey,
};
let pending = pendingAttachmentAdmissions.get(host);
if (attachments.signatures.size) {
pending ??= new Set();
pending.add(attachments);
pendingAttachmentAdmissions.set(host, pending);
}
const releaseAttachments = () => {
if (!pending?.delete(attachments)) {
return;
}
if (!pending.size) {
pendingAttachmentAdmissions.delete(host);
}
};
try {
return await run();
return await run({
releaseAttachments,
canRetireAttachment: (attachment) => {
const signature = chatAttachmentDraftSignature("", [attachment]);
return (
attachments.isCurrent() &&
attachments.signatures.has(signature) &&
![...(pendingAttachmentAdmissions.get(host) ?? [])].some(
(claim) => claim.isCurrent() && claim.signatures.has(signature),
)
);
},
});
} finally {
releaseAttachments();
guards.delete(guardKey);
}
}