Files
openchamber/packages/ui/src/hooks/useQueuedMessageAutoSend.ts
T

201 lines
6.3 KiB
TypeScript
Raw Normal View History

import React from 'react';
import { useMessageQueueStore, type QueuedMessage } from '@/stores/messageQueueStore';
import { useSessionUIStore } from '@/sync/session-ui-store';
import { useSelectionStore } from '@/sync/selection-store';
import { useConfigStore } from '@/stores/useConfigStore';
import { useContextStore } from '@/stores/contextStore';
import { parseAgentMentions } from '@/lib/messages/agentMentions';
import { getSyncSessionStatus } from '@/sync/sync-refs';
import { useDirectorySync } from '@/sync/sync-context';
type SessionStatusType = 'idle' | 'busy' | 'retry';
const RECENT_ABORT_WINDOW_MS = 2000;
const hasRecentAbort = (sessionId: string): boolean => {
const abortRecord = useSessionUIStore.getState().sessionAbortFlags.get(sessionId);
if (!abortRecord) {
return false;
}
return Date.now() - abortRecord.timestamp < RECENT_ABORT_WINDOW_MS;
};
export const buildQueuedAutoSendPayload = (queue: QueuedMessage[]) => {
const queued = queue[0];
if (!queued) {
return null;
}
const agents = useConfigStore.getState().getVisibleAgents();
const { sanitizedText, mention } = parseAgentMentions(queued.content, agents);
return {
queuedMessageId: queued.id,
primaryText: sanitizedText,
primaryAttachments: queued.attachments ?? [],
agentMentionName: mention?.name,
sendConfig: queued.sendConfig,
};
};
type QueuedAutoSendPayload = NonNullable<ReturnType<typeof buildQueuedAutoSendPayload>>;
type ResolvedQueuedSendConfig = {
providerID: string;
modelID: string;
agent?: string;
variant?: string;
};
export const sendQueuedAutoSendPayload = (
sessionId: string,
payload: QueuedAutoSendPayload,
resolved: ResolvedQueuedSendConfig,
) => {
return useSessionUIStore.getState().sendMessage(
payload.primaryText,
resolved.providerID,
resolved.modelID,
resolved.agent,
payload.primaryAttachments,
payload.agentMentionName,
undefined,
resolved.variant,
'normal',
{ sessionId },
);
};
const resolveSessionSendConfig = (sessionId: string) => {
const context = useContextStore.getState();
const config = useConfigStore.getState();
const selection = useSelectionStore.getState();
const selectedAgent =
context.getSessionAgentSelection(sessionId)
?? context.getCurrentAgent(sessionId)
?? config.currentAgentName
?? undefined;
const sessionModel = context.getSessionModelSelection(sessionId);
const agentModel = selectedAgent
? context.getAgentModelForSession(sessionId, selectedAgent)
: null;
const providerID =
agentModel?.providerId
?? sessionModel?.providerId
?? config.currentProviderId
?? selection.lastUsedProvider?.providerID;
const modelID =
agentModel?.modelId
?? sessionModel?.modelId
?? config.currentModelId
?? selection.lastUsedProvider?.modelID;
const variant =
selectedAgent && providerID && modelID
? (selection.getAgentModelVariantForSession(sessionId, selectedAgent, providerID, modelID)
?? context.getAgentModelVariantForSession(sessionId, selectedAgent, providerID, modelID))
: undefined;
return {
providerID,
modelID,
agent: selectedAgent,
variant,
};
};
export function useQueuedMessageAutoSend(enabledOrOptions?: boolean | { enabled?: boolean }) {
const enabled = typeof enabledOrOptions === 'boolean' ? enabledOrOptions : (enabledOrOptions?.enabled ?? true);
const queuedMessages = useMessageQueueStore((state) => state.queuedMessages);
const sessionStatusRecord = useDirectorySync((state) => state.session_status);
const inFlightSessionsRef = React.useRef<Set<string>>(new Set());
const previousStatusRef = React.useRef<Map<string, SessionStatusType>>(new Map());
React.useEffect(() => {
if (!enabled) {
return;
}
const dispatchSessionQueue = async (sessionId: string, queueSnapshot: QueuedMessage[]) => {
if (queueSnapshot.length === 0) {
return;
}
if (inFlightSessionsRef.current.has(sessionId)) {
return;
}
if (hasRecentAbort(sessionId)) {
return;
}
const currentStatus = getSyncSessionStatus(sessionId)?.type ?? 'idle';
if (currentStatus !== 'idle') {
return;
}
const payload = buildQueuedAutoSendPayload(queueSnapshot);
if (!payload) {
return;
}
if (!payload.primaryText && payload.primaryAttachments.length === 0) {
return;
}
// Use send config captured at queue time; fall back to current config
const captured = payload.sendConfig;
const resolved = captured?.providerID && captured?.modelID
? captured
: resolveSessionSendConfig(sessionId);
if (!resolved.providerID || !resolved.modelID) {
return;
}
inFlightSessionsRef.current.add(sessionId);
try {
await sendQueuedAutoSendPayload(sessionId, payload, {
providerID: resolved.providerID,
modelID: resolved.modelID,
agent: resolved.agent,
variant: resolved.variant,
});
const removeFromQueue = useMessageQueueStore.getState().removeFromQueue;
removeFromQueue(sessionId, payload.queuedMessageId);
} catch (error) {
console.warn('[queue] queued auto-send failed:', error);
} finally {
inFlightSessionsRef.current.delete(sessionId);
}
};
const statusRecord = sessionStatusRecord ?? {};
const nextStatusMap = new Map(previousStatusRef.current);
for (const [sessionId, status] of Object.entries(statusRecord)) {
if (status) {
nextStatusMap.set(sessionId, status.type as SessionStatusType);
}
}
const queueEntries = Object.entries(queuedMessages);
queueEntries.forEach(([sessionId, queue]) => {
const currentStatusType = (statusRecord[sessionId]?.type ?? 'idle') as SessionStatusType;
const previousStatusType = previousStatusRef.current.get(sessionId);
const becameIdle =
(previousStatusType === 'busy' || previousStatusType === 'retry')
&& currentStatusType === 'idle';
const firstSeenIdle = previousStatusType === undefined && currentStatusType === 'idle';
if (queue.length > 0 && (becameIdle || firstSeenIdle)) {
void dispatchSessionQueue(sessionId, queue);
}
nextStatusMap.set(sessionId, currentStatusType);
});
previousStatusRef.current = nextStatusMap;
}, [enabled, queuedMessages, sessionStatusRecord]);
}