File size: 7,844 Bytes
5f40163 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 | import { create } from "zustand";
export type PendingUserMessageStatus = "sending" | "error";
/**
* How long a pending message is allowed to stay in "sending" state before we
* give up and flip it to "error" with a retry link. This guards against the
* "server crashed / websocket dropped after our send resolved, echo never
* arrives" scenario where the message would otherwise hang forever.
*
* Exported so tests can override it via vi.fakeTimers without hard-coding the
* value.
*/
export const PENDING_MESSAGE_TIMEOUT_MS = 150_000;
export interface PendingUserMessage {
id: string;
/**
* The conversation this pending message belongs to. The chat UI filters the
* global queue by the active conversation id so messages enqueued in one
* conversation never leak into another when the user switches.
*/
conversationId: string;
/** User-visible bubble text (what the user typed; no file annotations). */
text: string;
/**
* The exact string sent to the server (may include the appended
* "Files uploaded: …" prompt when attachments are present). Used as the
* primary key when matching against the echoed `UserMessageEvent`.
*/
content: string;
status: PendingUserMessageStatus;
imageUrls: string[];
fileUrls: string[];
timestamp: string;
errorMessage?: string;
}
interface OptimisticUserMessageState {
pendingMessages: PendingUserMessage[];
}
export interface EnqueuePendingMessagePayload {
conversationId: string;
/** User-visible text for the bubble. */
text: string;
/**
* The exact string sent to the server. Defaults to `text` for call sites
* that don't transform the content (e.g. git-control-bar, task-card).
*/
content?: string;
imageUrls?: string[];
fileUrls?: string[];
timestamp?: string;
}
interface OptimisticUserMessageActions {
/**
* Append a new user message to the queue with status "sending".
* Returns the locally-generated id for later updates. Schedules a
* `PENDING_MESSAGE_TIMEOUT_MS` watchdog that flips the entry to "error" if
* it's still in "sending" state when the timer fires.
*/
enqueuePendingMessage: (payload: EnqueuePendingMessagePayload) => string;
/** Mark a pending message as failed (the API rejected it). */
markPendingMessageError: (id: string, errorMessage?: string) => void;
/** Mark a pending message as sending again (used when retrying). */
markPendingMessageSending: (id: string) => void;
/** Drop a pending message from the queue (e.g., after success/cancellation). */
removePendingMessage: (id: string) => void;
/**
* Remove the pending message that matches the given echoed `content` in
* the given conversation. Matching is done by exact content equality on
* messages still in "sending" state; if no match exists we fall back to
* removing the oldest "sending" entry in that conversation so that an echo
* with a slightly munged body (e.g. trailing-whitespace stripped by the
* server) still clears its bubble. Scoping by `conversationId` ensures a
* stale ack for one conversation never pops a pending entry belonging to
* another.
*/
consumeMatchingPendingMessage: (
conversationId: string,
content: string,
) => PendingUserMessage | null;
/** Wipe all queued messages (e.g., when changing conversations). */
clearPendingMessages: () => void;
/**
* Move pending entries from a provisional task URL (`task-{uuid}`) to the
* real conversation id once cloud provisioning finishes.
*/
reassignPendingMessages: (
fromConversationId: string,
toConversationId: string,
) => void;
}
type OptimisticUserMessageStore = OptimisticUserMessageState &
OptimisticUserMessageActions;
const initialState: OptimisticUserMessageState = {
pendingMessages: [],
};
// Use a timestamp + random suffix instead of a module-level counter so ids
// stay unique across test resets and don't accumulate state between runs.
// `crypto.randomUUID` would be ideal but isn't available in older test
// environments, so a base36 random suffix is a safe lowest-common-denominator.
const generatePendingId = (): string =>
`pending-${Date.now()}-${Math.random().toString(36).slice(2, 10)}`;
export const useOptimisticUserMessageStore = create<OptimisticUserMessageStore>(
(set, get) => ({
...initialState,
enqueuePendingMessage: (payload) => {
const id = generatePendingId();
const message: PendingUserMessage = {
id,
conversationId: payload.conversationId,
text: payload.text,
content: payload.content ?? payload.text,
status: "sending",
imageUrls: payload.imageUrls ?? [],
fileUrls: payload.fileUrls ?? [],
timestamp: payload.timestamp ?? new Date().toISOString(),
};
set((state) => ({
pendingMessages: [...state.pendingMessages, message],
}));
// Watchdog: if the server echo never lands (WS dropped, server crashed,
// network partition), flip this entry to "error" so the user gets a
// retry link instead of a permanently-pinned "Sending…" bubble.
setTimeout(() => {
const current = get().pendingMessages.find((m) => m.id === id);
if (current?.status === "sending") {
get().markPendingMessageError(id, "Send timed out");
}
}, PENDING_MESSAGE_TIMEOUT_MS);
return id;
},
markPendingMessageError: (id, errorMessage) =>
set((state) => ({
pendingMessages: state.pendingMessages.map((message) =>
message.id === id
? { ...message, status: "error", errorMessage }
: message,
),
})),
markPendingMessageSending: (id) =>
set((state) => ({
pendingMessages: state.pendingMessages.map((message) =>
message.id === id
? { ...message, status: "sending", errorMessage: undefined }
: message,
),
})),
removePendingMessage: (id) =>
set((state) => ({
pendingMessages: state.pendingMessages.filter(
(message) => message.id !== id,
),
})),
consumeMatchingPendingMessage: (conversationId, content) => {
// Single atomic `set` so the find + filter can't observe an interleaved
// mutation from another action. We prefer an exact content match (this
// is what makes out-of-order echoes safe: an echo of "world" will pop
// the "world" bubble, not the older "hello" one). If no exact match
// exists — e.g. the server slightly munged the body — fall back to the
// oldest "sending" entry in this conversation so the user doesn't end
// up with a permanently-stuck bubble in the happy-path single-message
// case.
let consumed: PendingUserMessage | null = null;
set((state) => {
const sending = state.pendingMessages
.map((m, i) => ({ m, i }))
.filter(
({ m }) =>
m.status === "sending" && m.conversationId === conversationId,
);
if (sending.length === 0) return state;
const exact = sending.find(({ m }) => m.content === content);
const target = exact ?? sending[0];
consumed = target.m;
return {
pendingMessages: [
...state.pendingMessages.slice(0, target.i),
...state.pendingMessages.slice(target.i + 1),
],
};
});
return consumed;
},
clearPendingMessages: () => set(() => ({ ...initialState })),
reassignPendingMessages: (fromConversationId, toConversationId) =>
set((state) => ({
pendingMessages: state.pendingMessages.map((message) =>
message.conversationId === fromConversationId
? { ...message, conversationId: toConversationId }
: message,
),
})),
}),
);
|