Skip to content

Commit 34f581a

Browse files
committed
fix(clients): preserve provider-bound drafts
Refs #222
1 parent 8b9e7c5 commit 34f581a

20 files changed

Lines changed: 1296 additions & 118 deletions

apps/mobile/src/features/threads/NewTaskDraftScreen.tsx

Lines changed: 9 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -809,14 +809,16 @@ export function NewTaskDraftScreen(props: {
809809
return;
810810
}
811811
const draft = getComposerDraftSnapshot(draftKey);
812-
// Snapshot read keeps just-typed selector state; the availability gate
813-
// still applies so a stored selection on a disabled provider falls back
814-
// to the flow's resolved model.
812+
// Snapshot read keeps just-typed selector state. Ambient stale defaults
813+
// may fall back, but a human/recovered exact provider choice must remain
814+
// blocked rather than silently switch accounts.
815815
const modelSelection =
816-
resolveSelectableModelSelection(
817-
selectedEnvironmentServerConfig,
818-
draft.modelSelection ?? null,
819-
) ?? flow.selectedModel;
816+
draft.providerSelectionExplicit === true && draft.modelSelection !== undefined
817+
? resolveSelectableModelSelection(selectedEnvironmentServerConfig, draft.modelSelection)
818+
: (resolveSelectableModelSelection(
819+
selectedEnvironmentServerConfig,
820+
draft.modelSelection ?? null,
821+
) ?? flow.selectedModel);
820822
const workspaceMode = draft.workspaceSelection?.mode ?? flow.workspaceMode;
821823
const selectedBranchName = draft.workspaceSelection?.branch ?? flow.selectedBranchName;
822824
const selectedWorktreePath =

apps/mobile/src/features/threads/new-task-flow-provider.tsx

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -439,10 +439,14 @@ export function NewTaskFlowProvider(props: React.PropsWithChildren) {
439439
stickySelection: storedStickyModelSelection,
440440
},
441441
);
442-
const draftModelSelection = resolveSelectableModelSelection(
442+
const selectableDraftModelSelection = resolveSelectableModelSelection(
443443
selectedEnvironmentServerConfig,
444444
storedDraftModelSelection,
445445
);
446+
const draftModelSelection =
447+
selectedProjectDraft.providerSelectionExplicit === true && storedDraftModelSelection !== null
448+
? storedDraftModelSelection
449+
: selectableDraftModelSelection;
446450
const projectDefaultModelSelection = resolveDefaultableModelSelection(
447451
selectedEnvironmentServerConfig,
448452
storedProjectDefaultModelSelection,
@@ -530,6 +534,7 @@ export function NewTaskFlowProvider(props: React.PropsWithChildren) {
530534
const modelSelection = options ? { ...option.selection, options } : option.selection;
531535
updateComposerDraftSettings(selectedProjectDraftKey, {
532536
modelSelection,
537+
providerSelectionExplicit: true,
533538
runtimeMode: resolveModelSelectionRuntimeMode(
534539
selectedEnvironmentServerConfig,
535540
modelSelection,
@@ -559,6 +564,7 @@ export function NewTaskFlowProvider(props: React.PropsWithChildren) {
559564
};
560565
updateComposerDraftSettings(selectedProjectDraftKey, {
561566
modelSelection: nextSelection,
567+
providerSelectionExplicit: true,
562568
});
563569
setStickyComposerModelSelection(nextSelection);
564570
},
@@ -886,6 +892,7 @@ export function NewTaskFlowProvider(props: React.PropsWithChildren) {
886892
replaceComposerDraftAttachments(draftKey, message.attachments);
887893
updateComposerDraftSettings(draftKey, {
888894
modelSelection: message.modelSelection,
895+
providerSelectionExplicit: message.modelSelection !== undefined,
889896
runtimeMode: message.runtimeMode,
890897
interactionMode: message.interactionMode,
891898
workspaceSelection: {

apps/mobile/src/state/thread-outbox-manager.ts

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -201,6 +201,21 @@ export function createThreadOutboxManager(options: ThreadOutboxManagerOptions) {
201201
return true;
202202
});
203203

204+
/**
205+
* Compare-and-set rewrite for a caller holding an exact queue snapshot.
206+
* The revision fence also catches a replacement published while this
207+
* mutation waits for earlier durable writes.
208+
*/
209+
const updateIfCurrent = (
210+
expected: QueuedThreadMessage,
211+
replacement: QueuedThreadMessage,
212+
): Promise<boolean> => {
213+
if (!currentMessages().some((candidate) => candidate === expected)) {
214+
return Promise.resolve(false);
215+
}
216+
return update(replacement, revisions.get(expected.messageId) ?? 0);
217+
};
218+
204219
// `expectedRevision` makes the removal a compare-and-set too: an edit
205220
// accepted after the caller decided to remove (restore-to-composer reads
206221
// the payload it is about to delete) keeps the newer message queued.
@@ -267,6 +282,15 @@ export function createThreadOutboxManager(options: ThreadOutboxManagerOptions) {
267282
return removed;
268283
});
269284

285+
/** Remove only the exact queue snapshot the caller inspected. */
286+
const removeIfCurrent = (expected: QueuedThreadMessage): Promise<boolean> => {
287+
if (!currentMessages().some((candidate) => candidate === expected)) {
288+
return Promise.resolve(false);
289+
}
290+
const expectedRevision = revisions.get(expected.messageId) ?? 0;
291+
return remove(expected, expectedRevision).then((removed) => removed !== null);
292+
};
293+
270294
const clearEnvironment = (
271295
environmentId: EnvironmentId,
272296
): Promise<ReadonlyArray<QueuedThreadMessage>> => {
@@ -386,7 +410,9 @@ export function createThreadOutboxManager(options: ThreadOutboxManagerOptions) {
386410
/** Current write revision for a queued message; input to update's CAS. */
387411
revisionOf: (messageId: MessageId): number => revisions.get(messageId) ?? 0,
388412
update,
413+
updateIfCurrent,
389414
remove,
415+
removeIfCurrent,
390416
clearEnvironment,
391417
};
392418
}

apps/mobile/src/state/thread-outbox-model.ts

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -279,6 +279,36 @@ export function preserveQueuedThreadDeliveryHold(
279279
: undefined;
280280
}
281281

282+
/**
283+
* A held existing-thread send may be explicitly retargeted only to a selected
284+
* provider that is currently admissible and proves the thread binding's exact
285+
* continuation identity. The returned selection is never rewritten.
286+
*/
287+
export function resolveHeldSendSelectedProvider(input: {
288+
readonly boundInstanceId: ModelSelectionType["instanceId"] | undefined;
289+
readonly selectedModelSelection: ModelSelectionType | null | undefined;
290+
readonly providers: ReadonlyArray<ServerProvider> | null | undefined;
291+
}): ModelSelectionType | null {
292+
const selection = input.selectedModelSelection;
293+
if (input.boundInstanceId === undefined || selection == null) return null;
294+
const transition = resolveProviderContinuationTransition({
295+
providers: input.providers ?? [],
296+
currentInstanceId: input.boundInstanceId,
297+
targetInstanceId: selection.instanceId,
298+
});
299+
if (!transition.compatible) return null;
300+
const provider = input.providers?.find(
301+
(candidate) => candidate.instanceId === selection.instanceId,
302+
);
303+
return getProviderAdmissionAvailability({
304+
provider,
305+
instanceId: String(selection.instanceId),
306+
providerSnapshotKnown: input.providers !== null && input.providers !== undefined,
307+
}).status === "available"
308+
? selection
309+
: null;
310+
}
311+
282312
export function retryQueuedThreadMessage(
283313
message: QueuedThreadMessage,
284314
input: {
Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,110 @@
1+
import { describe, expect, it, vi } from "vite-plus/test";
2+
import {
3+
CommandId,
4+
EnvironmentId,
5+
MessageId,
6+
ProviderInstanceId,
7+
ThreadId,
8+
} from "@t3tools/contracts";
9+
10+
import { recoverPendingSendToComposer } from "./thread-outbox-recovery";
11+
import type { QueuedThreadMessage } from "./thread-outbox-model";
12+
13+
function heldMessage(): QueuedThreadMessage {
14+
return {
15+
environmentId: EnvironmentId.make("environment-1"),
16+
threadId: ThreadId.make("thread-1"),
17+
messageId: MessageId.make("message-held"),
18+
commandId: CommandId.make("command-held"),
19+
text: "exact held text",
20+
attachments: [
21+
{
22+
id: "held-image",
23+
previewUri: "file:///held.png",
24+
type: "image",
25+
name: "held.png",
26+
mimeType: "image/png",
27+
sizeBytes: 12,
28+
dataUrl: "data:image/png;base64,AQ==",
29+
},
30+
],
31+
modelSelection: {
32+
instanceId: ProviderInstanceId.make("codex_personal"),
33+
model: "gpt-5.4",
34+
options: [{ id: "reasoningEffort", value: "xhigh" }],
35+
},
36+
runtimeMode: "approval-required",
37+
interactionMode: "plan",
38+
deliveryHold: {
39+
kind: "provider-binding-mismatch",
40+
reason: "Choose a destination.",
41+
},
42+
createdAt: "2026-09-01T00:00:00.000Z",
43+
};
44+
}
45+
46+
describe("pending send composer recovery", () => {
47+
it("durably flushes the exact payload before CAS-removing the queue item", async () => {
48+
const message = heldMessage();
49+
const order: string[] = [];
50+
const restore = vi.fn((draftKey: string, snapshot: unknown) => {
51+
order.push(`restore:${draftKey}`);
52+
expect(snapshot).toEqual({
53+
text: message.text,
54+
attachments: message.attachments,
55+
modelSelection: message.modelSelection,
56+
runtimeMode: message.runtimeMode,
57+
interactionMode: message.interactionMode,
58+
});
59+
});
60+
const flushDraft = vi.fn(async () => {
61+
order.push("flush");
62+
});
63+
const removeIfCurrent = vi.fn(async (candidate: QueuedThreadMessage) => {
64+
order.push("remove");
65+
expect(candidate).toBe(message);
66+
return true;
67+
});
68+
69+
await expect(
70+
recoverPendingSendToComposer(
71+
{ message, draftKey: "environment-1:thread-1" },
72+
{ restore, flushDraft, removeIfCurrent },
73+
),
74+
).resolves.toBe("removed");
75+
expect(order).toEqual(["restore:environment-1:thread-1", "flush", "remove"]);
76+
});
77+
78+
it("keeps the queue item when the durable draft flush fails", async () => {
79+
const message = heldMessage();
80+
const flushError = new Error("draft disk full");
81+
const removeIfCurrent = vi.fn(async () => true);
82+
83+
await expect(
84+
recoverPendingSendToComposer(
85+
{ message, draftKey: "environment-1:thread-1" },
86+
{
87+
restore: () => {},
88+
flushDraft: async () => {
89+
throw flushError;
90+
},
91+
removeIfCurrent,
92+
},
93+
),
94+
).rejects.toBe(flushError);
95+
expect(removeIfCurrent).not.toHaveBeenCalled();
96+
});
97+
98+
it("reports a concurrent queue CAS loss without deleting the newer item", async () => {
99+
const message = heldMessage();
100+
const removeIfCurrent = vi.fn(async () => false);
101+
102+
await expect(
103+
recoverPendingSendToComposer(
104+
{ message, draftKey: "environment-1:thread-1" },
105+
{ restore: () => {}, flushDraft: async () => {}, removeIfCurrent },
106+
),
107+
).resolves.toBe("queue-changed");
108+
expect(removeIfCurrent).toHaveBeenCalledWith(message);
109+
});
110+
});
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
import type { QueuedThreadMessage } from "./thread-outbox-model";
2+
import {
3+
flushComposerDrafts,
4+
restorePendingSendComposerDraft,
5+
type ComposerDraftWorkspaceSelection,
6+
type PendingSendComposerSnapshot,
7+
} from "./use-composer-drafts";
8+
import { removeThreadOutboxMessageIfCurrent } from "./thread-outbox-removal";
9+
10+
export type PendingSendRecoveryResult = "removed" | "queue-changed";
11+
12+
export function pendingSendComposerSnapshot(
13+
message: QueuedThreadMessage,
14+
workspaceSelection?: ComposerDraftWorkspaceSelection,
15+
): PendingSendComposerSnapshot {
16+
return {
17+
text: message.text,
18+
attachments: message.attachments,
19+
...(message.modelSelection === undefined ? {} : { modelSelection: message.modelSelection }),
20+
...(message.runtimeMode === undefined ? {} : { runtimeMode: message.runtimeMode }),
21+
...(message.interactionMode === undefined ? {} : { interactionMode: message.interactionMode }),
22+
...(workspaceSelection === undefined ? {} : { workspaceSelection }),
23+
};
24+
}
25+
26+
/**
27+
* Recovery protocol for a held send:
28+
* 1. merge its exact payload and selection into the chosen composer;
29+
* 2. force that composer snapshot to durable storage;
30+
* 3. CAS-remove only the queue object the user opened.
31+
*
32+
* A crash or flush failure before step 3 leaves the original held item intact.
33+
* A concurrent retry makes step 3 return `queue-changed`, so the newer queued
34+
* item survives while the restored composer copy remains available to edit.
35+
*/
36+
export async function recoverPendingSendToComposer(
37+
input: {
38+
readonly message: QueuedThreadMessage;
39+
readonly draftKey: string;
40+
readonly workspaceSelection?: ComposerDraftWorkspaceSelection;
41+
},
42+
dependencies: {
43+
readonly restore: (draftKey: string, snapshot: PendingSendComposerSnapshot) => void;
44+
readonly flushDraft: () => Promise<void>;
45+
readonly removeIfCurrent: (message: QueuedThreadMessage) => Promise<boolean>;
46+
} = {
47+
restore: restorePendingSendComposerDraft,
48+
flushDraft: flushComposerDrafts,
49+
removeIfCurrent: removeThreadOutboxMessageIfCurrent,
50+
},
51+
): Promise<PendingSendRecoveryResult> {
52+
dependencies.restore(
53+
input.draftKey,
54+
pendingSendComposerSnapshot(input.message, input.workspaceSelection),
55+
);
56+
await dependencies.flushDraft();
57+
return (await dependencies.removeIfCurrent(input.message)) ? "removed" : "queue-changed";
58+
}

apps/mobile/src/state/thread-outbox-removal.ts

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,21 @@ export async function removeThreadOutboxMessage(
8181
return true;
8282
}
8383

84+
/**
85+
* Removes only the exact queue snapshot the caller inspected, then releases
86+
* the files owned by that snapshot.
87+
*/
88+
export async function removeThreadOutboxMessageIfCurrent(
89+
message: QueuedThreadMessage,
90+
): Promise<boolean> {
91+
const removed = await threadOutboxManager.removeIfCurrent(message);
92+
if (!removed) {
93+
return false;
94+
}
95+
await cleanUpRemovedMessages([message]);
96+
return true;
97+
}
98+
8499
/** Removes every queued message of an environment and releases their files. */
85100
export async function clearThreadOutboxEnvironment(environmentId: EnvironmentId): Promise<void> {
86101
// clearEnvironment loads and merges persisted messages itself and reports

0 commit comments

Comments
 (0)