Files
openchamber/packages/ui/src/stores/messageQueueStore.ts
T
Bohdan Triapitsyn 237cae16b3 fix: stop the composer re-sending a queued message already in flight
A queued message is removed from the queue only after its send resolves,
so between dispatch and resolution it stays visible to every reader — and
a composer submit merges the whole queue into its own send. Over a relay
that window is seconds, long enough to deliver the same message twice.

The queue now tracks which entries are awaiting the server. Dispatchers
skip them, clearQueue retains them so the pending send can still remove
or restore its own entry, and the flag is not persisted because a restart
has no in-flight sends.
2026-08-03 23:14:14 +03:00

334 lines
14 KiB
TypeScript

import { create } from 'zustand';
import { devtools, persist } from 'zustand/middleware';
import { createDeferredSafeJSONStorage } from './utils/safeStorage';
import type { AttachedFile } from './types/sessionTypes';
import { updateDesktopSettings } from '@/lib/persistence';
import { getRuntimeKey } from '@/lib/runtime-switch';
import { normalizePath } from '@/lib/pathNormalization';
export type FollowUpBehavior = 'steer' | 'queue';
export const DEFAULT_FOLLOW_UP_BEHAVIOR: FollowUpBehavior = 'queue';
export const isFollowUpBehavior = (value: unknown): value is FollowUpBehavior => (
value === 'steer' || value === 'queue'
);
export const normalizeFollowUpBehavior = (
value: unknown,
legacyQueueModeEnabled?: boolean | null,
): FollowUpBehavior => {
// "immediate" was removed: on a busy session it was wire-identical to
// "steer" (OpenCode only supports delivery "steer" | "queue", defaulting
// to "steer"), so collapse any persisted/legacy "immediate" onto "steer".
if (value === 'immediate') {
return 'steer';
}
if (isFollowUpBehavior(value)) {
return value;
}
if (legacyQueueModeEnabled === false) {
return 'steer';
}
if (legacyQueueModeEnabled === true) {
return 'queue';
}
return DEFAULT_FOLLOW_UP_BEHAVIOR;
};
export interface QueuedMessage {
id: string;
content: string;
attachments?: AttachedFile[];
createdAt: number;
/** Send config captured at queue time — used as-is when auto-sending */
sendConfig?: {
providerID: string;
modelID: string;
agent?: string;
variant?: string;
};
}
export type MessageQueueTarget = {
runtimeKey: string;
directory: string;
sessionId: string;
};
const MAX_QUEUE_TARGETS = 50;
const MAX_MESSAGES_PER_QUEUE = 20;
export const createMessageQueueTarget = (
sessionId: string,
directory: string | null | undefined,
runtimeKey: string = getRuntimeKey(),
): MessageQueueTarget | null => {
const normalizedDirectory = normalizePath(directory);
if (!runtimeKey || !normalizedDirectory || !sessionId) return null;
return { runtimeKey, directory: normalizedDirectory, sessionId };
};
export const getMessageQueueKey = (target: MessageQueueTarget): string =>
`${target.runtimeKey}\n${target.directory}\n${target.sessionId}`;
export const parseMessageQueueKey = (key: string): MessageQueueTarget | null => {
const [runtimeKey, directory, ...sessionParts] = key.split('\n');
return createMessageQueueTarget(sessionParts.join('\n'), directory, runtimeKey);
};
interface MessageQueueState {
queuedMessages: Record<string, QueuedMessage[]>; // runtime + directory + session → queue
quarantinedLegacyMessages: Record<string, QueuedMessage[]>;
followUpBehavior: FollowUpBehavior;
/**
* Queued messages whose send is currently awaiting the server, per target.
*
* A queued item is removed only after its send resolves, so between
* dispatch and resolution it is still visible to every other reader — and
* a composer submit merges the whole queue into its own send. Over a relay
* that window is seconds, long enough for the same message to be delivered
* twice. Dispatchers must skip entries listed here.
*
* Never persisted: a restart has no in-flight sends, and a stale flag would
* strand a queued message permanently.
*/
sendingIds: Record<string, string[]>;
}
interface MessageQueueActions {
addToQueue: (target: MessageQueueTarget, message: Omit<QueuedMessage, 'id' | 'createdAt'>) => void;
removeFromQueue: (target: MessageQueueTarget, messageId: string) => void;
reorderQueue: (target: MessageQueueTarget, fromId: string, toId: string) => void;
popToInput: (target: MessageQueueTarget, messageId: string) => QueuedMessage | null;
clearQueue: (target: MessageQueueTarget) => void;
clearAllQueues: () => void;
markSending: (target: MessageQueueTarget, messageId: string) => void;
clearSending: (target: MessageQueueTarget, messageId: string) => void;
getSendableQueue: (target: MessageQueueTarget) => QueuedMessage[];
setFollowUpBehavior: (behavior: FollowUpBehavior) => void;
getQueueForTarget: (target: MessageQueueTarget) => QueuedMessage[];
}
type MessageQueueStore = MessageQueueState & MessageQueueActions;
type PersistedMessageQueueState = {
queuedMessages?: Record<string, QueuedMessage[]>;
quarantinedLegacyMessages?: Record<string, QueuedMessage[]>;
followUpBehavior?: FollowUpBehavior;
queueModeEnabled?: boolean;
};
export const migrateMessageQueueState = (persistedState: unknown, version: number): Partial<MessageQueueStore> => {
const state = (persistedState ?? {}) as PersistedMessageQueueState;
const legacyQueues = version < 2 ? (state.queuedMessages ?? {}) : {};
return {
queuedMessages: version < 2 ? {} : (state.queuedMessages ?? {}),
quarantinedLegacyMessages: {
...(state.quarantinedLegacyMessages ?? {}),
...legacyQueues,
},
followUpBehavior: normalizeFollowUpBehavior(state.followUpBehavior, state.queueModeEnabled ?? null),
};
};
export const useMessageQueueStore = create<MessageQueueStore>()(
devtools(
persist(
(set, get) => ({
queuedMessages: {},
quarantinedLegacyMessages: {},
followUpBehavior: DEFAULT_FOLLOW_UP_BEHAVIOR,
sendingIds: {},
addToQueue: (target, message) => {
const key = getMessageQueueKey(target);
const id = `queued-${Date.now()}-${Math.random().toString(36).substring(2, 9)}`;
const queuedMessage: QueuedMessage = {
id,
content: message.content,
attachments: message.attachments,
createdAt: Date.now(),
sendConfig: message.sendConfig,
};
set((state) => {
const currentQueue = state.queuedMessages[key] ?? [];
const queuedMessages = {
...state.queuedMessages,
[key]: [...currentQueue, queuedMessage].slice(-MAX_MESSAGES_PER_QUEUE),
};
const keys = Object.keys(queuedMessages);
if (keys.length > MAX_QUEUE_TARGETS) {
keys.sort((left, right) => (
(queuedMessages[left]?.[0]?.createdAt ?? 0) - (queuedMessages[right]?.[0]?.createdAt ?? 0)
));
for (const staleKey of keys.slice(0, keys.length - MAX_QUEUE_TARGETS)) delete queuedMessages[staleKey];
}
return {
queuedMessages,
};
});
},
removeFromQueue: (target, messageId) => {
const key = getMessageQueueKey(target);
set((state) => {
const currentQueue = state.queuedMessages[key] ?? [];
const newQueue = currentQueue.filter((m) => m.id !== messageId);
if (newQueue.length === 0) {
const { [key]: _removed, ...rest } = state.queuedMessages;
void _removed;
return { queuedMessages: rest };
}
return {
queuedMessages: {
...state.queuedMessages,
[key]: newQueue,
},
};
});
},
reorderQueue: (target, fromId, toId) => {
if (fromId === toId) return;
const key = getMessageQueueKey(target);
set((state) => {
const currentQueue = state.queuedMessages[key];
if (!currentQueue) return state;
const fromIndex = currentQueue.findIndex((m) => m.id === fromId);
const toIndex = currentQueue.findIndex((m) => m.id === toId);
if (fromIndex === -1 || toIndex === -1) return state;
const newQueue = currentQueue.slice();
const [moved] = newQueue.splice(fromIndex, 1);
newQueue.splice(toIndex, 0, moved);
return {
queuedMessages: {
...state.queuedMessages,
[key]: newQueue,
},
};
});
},
popToInput: (target, messageId) => {
const key = getMessageQueueKey(target);
const state = get();
const currentQueue = state.queuedMessages[key] ?? [];
const message = currentQueue.find((m) => m.id === messageId);
if (!message) {
return null;
}
// Remove from queue
set((prevState) => {
const queue = prevState.queuedMessages[key] ?? [];
const newQueue = queue.filter((m) => m.id !== messageId);
if (newQueue.length === 0) {
const { [key]: _removed, ...rest } = prevState.queuedMessages;
void _removed;
return { queuedMessages: rest };
}
return {
queuedMessages: {
...prevState.queuedMessages,
[key]: newQueue,
},
};
});
return message;
},
clearQueue: (target) => {
const key = getMessageQueueKey(target);
set((state) => {
// Clearing drops what is still queued, never a message
// already handed to the server: that send will resolve
// and must find its entry to remove or restore.
const sending = state.sendingIds[key] ?? [];
const retained = (state.queuedMessages[key] ?? []).filter((m) => sending.includes(m.id));
if (retained.length > 0) {
return { queuedMessages: { ...state.queuedMessages, [key]: retained } };
}
const { [key]: _removed, ...rest } = state.queuedMessages;
void _removed;
return { queuedMessages: rest };
});
},
clearAllQueues: () => {
set({ queuedMessages: {}, sendingIds: {} });
},
markSending: (target, messageId) => {
const key = getMessageQueueKey(target);
set((state) => {
const current = state.sendingIds[key] ?? [];
if (current.includes(messageId)) return state;
return { sendingIds: { ...state.sendingIds, [key]: [...current, messageId] } };
});
},
clearSending: (target, messageId) => {
const key = getMessageQueueKey(target);
set((state) => {
const current = state.sendingIds[key];
if (!current || !current.includes(messageId)) return state;
const next = current.filter((id) => id !== messageId);
if (next.length === 0) {
const { [key]: _removed, ...rest } = state.sendingIds;
void _removed;
return { sendingIds: rest };
}
return { sendingIds: { ...state.sendingIds, [key]: next } };
});
},
getSendableQueue: (target) => {
const key = getMessageQueueKey(target);
const state = get();
const queue = state.queuedMessages[key] ?? [];
const sending = state.sendingIds[key];
if (!sending || sending.length === 0) return queue;
return queue.filter((message) => !sending.includes(message.id));
},
setFollowUpBehavior: (behavior) => {
set({ followUpBehavior: behavior });
void updateDesktopSettings({ followUpBehavior: behavior });
},
getQueueForTarget: (target) => {
return get().queuedMessages[getMessageQueueKey(target)] ?? [];
},
}),
{
name: 'message-queue-store',
version: 2,
storage: createDeferredSafeJSONStorage(),
partialize: (state) => ({
queuedMessages: state.queuedMessages,
quarantinedLegacyMessages: state.quarantinedLegacyMessages,
followUpBehavior: state.followUpBehavior,
}),
migrate: migrateMessageQueueState,
}
),
{
name: 'message-queue-store',
}
)
);