perf: optimize app startup and background lifecycles
This commit is contained in:
@@ -16,6 +16,7 @@ import {
|
||||
updateImageWorkspaceGenerationQuote,
|
||||
} from '@/lib/image-workspace';
|
||||
import { useAuthStore } from '@/stores/auth';
|
||||
import { setDesktopBackgroundLease } from '@/lib/host-api';
|
||||
import {
|
||||
IMAGE_WORKSPACE_UNAVAILABLE_CODE,
|
||||
type DesignAssistantDeltaEvent,
|
||||
@@ -109,6 +110,18 @@ let projectDeletionGeneration = 0;
|
||||
let generationQuoteUpdateSequence = 0;
|
||||
let reconcileTaskStateOnNextConnect = false;
|
||||
|
||||
function updateBackgroundLease(input: {
|
||||
id: string;
|
||||
kind: string;
|
||||
active: boolean;
|
||||
}): void {
|
||||
try {
|
||||
void Promise.resolve(setDesktopBackgroundLease(input)).catch(() => undefined);
|
||||
} catch {
|
||||
// Renderer-only tests and embedded previews may not expose lifecycle IPC.
|
||||
}
|
||||
}
|
||||
|
||||
function taskRevisionKey(workspaceId: string, taskId: string): string {
|
||||
return `${workspaceId}:${taskId}`;
|
||||
}
|
||||
@@ -376,7 +389,70 @@ function replaceGenerationQuote(
|
||||
}
|
||||
|
||||
export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) => {
|
||||
const pendingAssistantDeltas: DesignAssistantDeltaEvent[] = [];
|
||||
let assistantDeltaFlushTimer: ReturnType<typeof setTimeout> | null = null;
|
||||
|
||||
const applyAssistantDelta = (
|
||||
state: ImageWorkspaceState,
|
||||
delta: DesignAssistantDeltaEvent,
|
||||
): ImageWorkspaceState => {
|
||||
if (state.activeWorkspaceId !== delta.workspaceId
|
||||
|| state.activeConversationId !== delta.conversationId
|
||||
|| (state.conversation?.turnRevision ?? 0) >= delta.turnRevision) {
|
||||
return state;
|
||||
}
|
||||
const pending = state.pendingTurn;
|
||||
if (!pending) {
|
||||
if (delta.chunkIndex !== 0) return state;
|
||||
return {
|
||||
...state,
|
||||
pendingTurn: {
|
||||
workspaceId: delta.workspaceId,
|
||||
conversationId: delta.conversationId,
|
||||
clientTurnId: delta.clientTurnId,
|
||||
turnRevision: delta.turnRevision,
|
||||
userText: null,
|
||||
assistantText: delta.delta,
|
||||
lastChunkIndex: delta.chunkIndex,
|
||||
createdAt: new Date().toISOString(),
|
||||
},
|
||||
};
|
||||
}
|
||||
if (pending.workspaceId !== delta.workspaceId
|
||||
|| pending.conversationId !== delta.conversationId
|
||||
|| pending.clientTurnId !== delta.clientTurnId
|
||||
|| pending.turnRevision !== delta.turnRevision
|
||||
|| delta.chunkIndex !== pending.lastChunkIndex + 1) {
|
||||
return state;
|
||||
}
|
||||
return {
|
||||
...state,
|
||||
pendingTurn: {
|
||||
...pending,
|
||||
assistantText: `${pending.assistantText}${delta.delta}`,
|
||||
lastChunkIndex: delta.chunkIndex,
|
||||
},
|
||||
};
|
||||
};
|
||||
|
||||
const flushAssistantDeltas = () => {
|
||||
assistantDeltaFlushTimer = null;
|
||||
if (pendingAssistantDeltas.length === 0) return;
|
||||
const queued = pendingAssistantDeltas.splice(0, pendingAssistantDeltas.length);
|
||||
set((state) => queued.reduce(applyAssistantDelta, state));
|
||||
};
|
||||
|
||||
const queueAssistantDelta = (delta: DesignAssistantDeltaEvent) => {
|
||||
pendingAssistantDeltas.push(delta);
|
||||
if (assistantDeltaFlushTimer === null) {
|
||||
assistantDeltaFlushTimer = setTimeout(flushAssistantDeltas, 50);
|
||||
}
|
||||
};
|
||||
|
||||
const stopTaskStream = () => {
|
||||
pendingAssistantDeltas.length = 0;
|
||||
if (assistantDeltaFlushTimer !== null) clearTimeout(assistantDeltaFlushTimer);
|
||||
assistantDeltaFlushTimer = null;
|
||||
closeTaskEventSource();
|
||||
set({ taskStreamState: 'idle' });
|
||||
};
|
||||
@@ -414,43 +490,7 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
if (!delta
|
||||
|| delta.workspaceId !== activeTaskEventWorkspaceId
|
||||
|| delta.conversationId !== activeTaskEventConversationId) return;
|
||||
set((state) => {
|
||||
if (state.activeWorkspaceId !== delta.workspaceId
|
||||
|| state.activeConversationId !== delta.conversationId
|
||||
|| (state.conversation?.turnRevision ?? 0) >= delta.turnRevision) {
|
||||
return state;
|
||||
}
|
||||
const pending = state.pendingTurn;
|
||||
if (!pending) {
|
||||
if (delta.chunkIndex !== 0) return state;
|
||||
return {
|
||||
pendingTurn: {
|
||||
workspaceId: delta.workspaceId,
|
||||
conversationId: delta.conversationId,
|
||||
clientTurnId: delta.clientTurnId,
|
||||
turnRevision: delta.turnRevision,
|
||||
userText: null,
|
||||
assistantText: delta.delta,
|
||||
lastChunkIndex: delta.chunkIndex,
|
||||
createdAt: new Date().toISOString(),
|
||||
},
|
||||
};
|
||||
}
|
||||
if (pending.workspaceId !== delta.workspaceId
|
||||
|| pending.conversationId !== delta.conversationId
|
||||
|| pending.clientTurnId !== delta.clientTurnId
|
||||
|| pending.turnRevision !== delta.turnRevision
|
||||
|| delta.chunkIndex !== pending.lastChunkIndex + 1) {
|
||||
return state;
|
||||
}
|
||||
return {
|
||||
pendingTurn: {
|
||||
...pending,
|
||||
assistantText: `${pending.assistantText}${delta.delta}`,
|
||||
lastChunkIndex: delta.chunkIndex,
|
||||
},
|
||||
};
|
||||
});
|
||||
queueAssistantDelta(delta);
|
||||
});
|
||||
source.addEventListener('design.conversation.snapshot', (event) => {
|
||||
const snapshot = parseConversationSnapshotEvent(event);
|
||||
@@ -579,7 +619,7 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
const message = messageOf(error);
|
||||
if (authenticationRequired(error)) {
|
||||
projectDeletionGeneration += 1;
|
||||
closeTaskEventSource();
|
||||
stopTaskStream();
|
||||
useAuthStore.getState().invalidateSession();
|
||||
set({
|
||||
status: 'auth-required',
|
||||
@@ -735,7 +775,7 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
)
|
||||
? currentId
|
||||
: bootstrap.workspaces[0]?.workspaceId ?? null;
|
||||
closeTaskEventSource();
|
||||
stopTaskStream();
|
||||
set({
|
||||
status: 'ready',
|
||||
bootstrap,
|
||||
@@ -753,7 +793,7 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
} catch (error) {
|
||||
const message = handleRequestError(error);
|
||||
if (!authenticationRequired(error)) {
|
||||
closeTaskEventSource();
|
||||
stopTaskStream();
|
||||
set({
|
||||
status: unavailable(error) ? 'unavailable' : 'error',
|
||||
bootstrap: null,
|
||||
@@ -877,7 +917,7 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
workspaceLoadGeneration += 1;
|
||||
conversationSelectionGeneration += 1;
|
||||
reconcileTaskStateOnNextConnect = false;
|
||||
closeTaskEventSource();
|
||||
stopTaskStream();
|
||||
const nextWorkspaceId = mostRecentlyUpdatedWorkspace(remainingWorkspaces)?.workspaceId
|
||||
?? null;
|
||||
set({
|
||||
@@ -1027,7 +1067,7 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
const shouldReconcile = reconcileTaskStateOnNextConnect;
|
||||
reconcileTaskStateOnNextConnect = false;
|
||||
if (shouldReconcile && get().taskStreamState !== 'connected') {
|
||||
closeTaskEventSource();
|
||||
stopTaskStream();
|
||||
}
|
||||
if (workspaceId && conversationId) startTaskStream(workspaceId, conversationId);
|
||||
if (!workspaceId || !conversationId || !shouldReconcile) return;
|
||||
@@ -1049,6 +1089,8 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
const userText = message.trim();
|
||||
const clientTurnId = createImageWorkspaceTurnId();
|
||||
if (!workspace || !conversation) throw new Error('请先选择设计会话');
|
||||
const leaseId = `canvas:${workspace.workspaceId}:${clientTurnId}`;
|
||||
updateBackgroundLease({ id: leaseId, kind: 'canvas-generation', active: true });
|
||||
try {
|
||||
set({
|
||||
pendingTurn: createPendingTurn(conversation, clientTurnId, userText),
|
||||
@@ -1087,6 +1129,8 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
: state.pendingTurn,
|
||||
}));
|
||||
return await recoverRevisionConflict(error);
|
||||
} finally {
|
||||
updateBackgroundLease({ id: leaseId, kind: 'canvas-generation', active: false });
|
||||
}
|
||||
},
|
||||
|
||||
@@ -1147,6 +1191,8 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
const conversation = get().conversation;
|
||||
if (!workspace || !conversation) throw new Error('请先选择设计会话');
|
||||
const clientTurnId = createImageWorkspaceTurnId();
|
||||
const leaseId = `canvas:${workspace.workspaceId}:${clientTurnId}`;
|
||||
updateBackgroundLease({ id: leaseId, kind: 'canvas-generation', active: true });
|
||||
const requestWorkspaceLoadGeneration = workspaceLoadGeneration;
|
||||
const requestConversationSelectionGeneration = conversationSelectionGeneration;
|
||||
const isActiveConfirmationWorkspace = (): boolean => (
|
||||
@@ -1247,6 +1293,8 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
: state.pendingTurn,
|
||||
}));
|
||||
return await recoverRevisionConflict(error);
|
||||
} finally {
|
||||
updateBackgroundLease({ id: leaseId, kind: 'canvas-generation', active: false });
|
||||
}
|
||||
},
|
||||
|
||||
@@ -1257,7 +1305,7 @@ export const useImageWorkspaceStore = create<ImageWorkspaceState>((set, get) =>
|
||||
projectDeletionGeneration += 1;
|
||||
generationQuoteUpdateSequence += 1;
|
||||
reconcileTaskStateOnNextConnect = false;
|
||||
closeTaskEventSource();
|
||||
stopTaskStream();
|
||||
set({
|
||||
status: 'idle',
|
||||
bootstrap: null,
|
||||
|
||||
@@ -1,5 +1,10 @@
|
||||
import { create } from 'zustand';
|
||||
import { createHostEventSource, ensureHostApiToken, hostApiFetch } from '@/lib/host-api';
|
||||
import {
|
||||
createHostEventSource,
|
||||
ensureHostApiToken,
|
||||
hostApiFetch,
|
||||
setDesktopBackgroundLease,
|
||||
} from '@/lib/host-api';
|
||||
import { dispatchWorksSquareTokenUsageStale } from '@/lib/works-square-usage-events';
|
||||
import { queueAgentSessionSync } from '@/lib/agent-session-sync';
|
||||
import {
|
||||
@@ -68,6 +73,21 @@ import type { ProjectType } from '../../shared/project-config';
|
||||
|
||||
export type { OpencodeErrorKind } from '../../shared/opencode-error-kind';
|
||||
|
||||
function updateBackgroundLease(input: {
|
||||
id: string;
|
||||
kind: string;
|
||||
active: boolean;
|
||||
}): void {
|
||||
// Renderer-only embeds and focused tests may provide a reduced Host API
|
||||
// surface. Lifecycle reporting is best effort and must never reject a
|
||||
// prompt when that optional channel is unavailable.
|
||||
try {
|
||||
void Promise.resolve(setDesktopBackgroundLease(input)).catch(() => undefined);
|
||||
} catch {
|
||||
// Main remains the source of truth when the lifecycle channel exists.
|
||||
}
|
||||
}
|
||||
|
||||
interface OpencodeActionResponse {
|
||||
success: boolean;
|
||||
status?: OpencodeStatus;
|
||||
@@ -196,6 +216,8 @@ interface OpencodeState {
|
||||
error: string | null;
|
||||
errorKind: OpencodeErrorKind | null;
|
||||
refreshStatus: () => Promise<void>;
|
||||
/** Close renderer-owned streams after Main has put the desktop into sleep. */
|
||||
pauseBackgroundStreams: () => void;
|
||||
start: () => Promise<void>;
|
||||
stop: () => Promise<void>;
|
||||
restart: () => Promise<void>;
|
||||
@@ -266,6 +288,16 @@ type SessionRunListenerRegistration = {
|
||||
source: EventSource;
|
||||
};
|
||||
const activeSessionRunListenerRegistrations = new Map<string, SessionRunListenerRegistration[]>();
|
||||
// Session history is durable in the OpenCode runtime. The renderer cache is
|
||||
// deliberately bounded so visiting many conversations does not retain every
|
||||
// full transcript for the lifetime of the window. Eviction only removes the
|
||||
// renderer copy; selecting an evicted session still hydrates it from the
|
||||
// runtime via `loadSessionMessages`.
|
||||
const MAX_RENDERER_SESSION_TRANSCRIPTS = 8;
|
||||
const sessionMessageCacheAccess = new Map<string, number>();
|
||||
let sessionMessageCacheAccessTick = 0;
|
||||
const activeSessionHydrations = new Map<string, Promise<void>>();
|
||||
const activeSessionMessageLoads = new Map<string, Promise<RawMessage[]>>();
|
||||
|
||||
function removeEventListenerIfSupported(
|
||||
source: EventSource,
|
||||
@@ -388,6 +420,10 @@ function isSessionRunCurrent(sessionId: string, token: number): boolean {
|
||||
|
||||
function closeActiveSessionEventSource(): void {
|
||||
cancelAllPromptStartWatchdogs();
|
||||
// Any in-flight snapshot belongs to the old runtime/project lifecycle. Do
|
||||
// not let it suppress a fresh read after a stop, restart, or project swap.
|
||||
activeSessionMessageLoads.clear();
|
||||
activeSessionHydrations.clear();
|
||||
activeSessionRunTokens.clear();
|
||||
for (const sessionId of pendingSessionRunStreamBatches.keys()) {
|
||||
cancelSessionRunStreamBatch(sessionId);
|
||||
@@ -540,6 +576,10 @@ export function resetOpencodeStoreEphemeralStateForTests(): void {
|
||||
commandLoadSequence = 0;
|
||||
runtimeStatusSequence = 0;
|
||||
sessionDiffLoadSequence = 0;
|
||||
sessionMessageCacheAccess.clear();
|
||||
sessionMessageCacheAccessTick = 0;
|
||||
activeSessionHydrations.clear();
|
||||
activeSessionMessageLoads.clear();
|
||||
}
|
||||
|
||||
function clearDrainingAbortedRunIfForPreviousRun(sessionId: string, runToken: number): void {
|
||||
@@ -1499,6 +1539,81 @@ function withoutKey<T>(record: Record<string, T>, key: string): Record<string, T
|
||||
return next;
|
||||
}
|
||||
|
||||
function getSessionCacheAccessKey(state: Pick<OpencodeState, 'activeProject'>, sessionId: string): string {
|
||||
return `${getProjectId(state.activeProject) ?? 'no-project'}:${sessionId}`;
|
||||
}
|
||||
|
||||
function touchSessionMessageCache(state: Pick<OpencodeState, 'activeProject'>, sessionId: string): void {
|
||||
sessionMessageCacheAccess.set(
|
||||
getSessionCacheAccessKey(state, sessionId),
|
||||
++sessionMessageCacheAccessTick,
|
||||
);
|
||||
}
|
||||
|
||||
function forgetSessionMessageCache(
|
||||
state: Pick<OpencodeState, 'activeProject'>,
|
||||
sessionId: string,
|
||||
): void {
|
||||
sessionMessageCacheAccess.delete(getSessionCacheAccessKey(state, sessionId));
|
||||
}
|
||||
|
||||
function getProtectedSessionMessageIds(state: OpencodeState): Set<string> {
|
||||
const protectedIds = new Set<string>();
|
||||
if (state.selectedSessionId) protectedIds.add(state.selectedSessionId);
|
||||
|
||||
for (const [sessionId, sending] of Object.entries(state.sendingSessionIds)) {
|
||||
if (sending) protectedIds.add(sessionId);
|
||||
}
|
||||
for (const [sessionId, runState] of Object.entries(state.sessionRunStates)) {
|
||||
if (runState.phase !== 'idle') protectedIds.add(sessionId);
|
||||
}
|
||||
for (const request of state.pendingQuestions) {
|
||||
if (request.sessionID) protectedIds.add(request.sessionID);
|
||||
}
|
||||
for (const request of state.pendingPermissions) {
|
||||
if (request.sessionID) protectedIds.add(request.sessionID);
|
||||
}
|
||||
for (const sessionId of Object.keys(state.streamingMessagesBySessionId)) {
|
||||
protectedIds.add(sessionId);
|
||||
}
|
||||
for (const sessionId of Object.keys(state.streamingToolsBySessionId)) {
|
||||
protectedIds.add(sessionId);
|
||||
}
|
||||
return protectedIds;
|
||||
}
|
||||
|
||||
function pruneSessionMessageCache(
|
||||
state: OpencodeState,
|
||||
messagesBySessionId: Record<string, RawMessage[]>,
|
||||
): Record<string, RawMessage[]> {
|
||||
const sessionIds = Object.keys(messagesBySessionId);
|
||||
if (sessionIds.length <= MAX_RENDERER_SESSION_TRANSCRIPTS) return messagesBySessionId;
|
||||
|
||||
const protectedIds = new Set(
|
||||
[...getProtectedSessionMessageIds(state)].filter((sessionId) => (
|
||||
Object.prototype.hasOwnProperty.call(messagesBySessionId, sessionId)
|
||||
)),
|
||||
);
|
||||
const orderedIds = [...sessionIds].sort((left, right) => (
|
||||
(sessionMessageCacheAccess.get(getSessionCacheAccessKey(state, right)) ?? 0)
|
||||
- (sessionMessageCacheAccess.get(getSessionCacheAccessKey(state, left)) ?? 0)
|
||||
));
|
||||
const keepIds = new Set<string>(protectedIds);
|
||||
for (const sessionId of orderedIds) {
|
||||
if (keepIds.size >= MAX_RENDERER_SESSION_TRANSCRIPTS) break;
|
||||
keepIds.add(sessionId);
|
||||
}
|
||||
|
||||
const next = { ...messagesBySessionId };
|
||||
for (const sessionId of sessionIds) {
|
||||
if (!keepIds.has(sessionId)) {
|
||||
delete next[sessionId];
|
||||
forgetSessionMessageCache(state, sessionId);
|
||||
}
|
||||
}
|
||||
return next;
|
||||
}
|
||||
|
||||
function getSessionMessagesForState(state: OpencodeState, sessionId: string): RawMessage[] {
|
||||
if (state.selectedSessionId === sessionId) return state.sessionMessages;
|
||||
return state.sessionMessagesBySessionId[sessionId] ?? [];
|
||||
@@ -1509,11 +1624,13 @@ function getSessionMessagePatch(
|
||||
sessionId: string,
|
||||
messages: RawMessage[],
|
||||
): Partial<OpencodeState> {
|
||||
touchSessionMessageCache(state, sessionId);
|
||||
const sessionMessagesBySessionId = pruneSessionMessageCache(state, {
|
||||
...state.sessionMessagesBySessionId,
|
||||
[sessionId]: messages,
|
||||
});
|
||||
return {
|
||||
sessionMessagesBySessionId: {
|
||||
...state.sessionMessagesBySessionId,
|
||||
[sessionId]: messages,
|
||||
},
|
||||
sessionMessagesBySessionId,
|
||||
...(state.selectedSessionId === sessionId ? { sessionMessages: messages } : {}),
|
||||
};
|
||||
}
|
||||
@@ -1779,6 +1896,7 @@ function deriveRevertedSessionIds(sessions: readonly OpencodeSession[]): Record<
|
||||
function resetProjectScopedState(sessions: OpencodeSession[] = []) {
|
||||
closeActiveSessionEventSource();
|
||||
sessionDiffLoadSequence += 1;
|
||||
sessionMessageCacheAccess.clear();
|
||||
return {
|
||||
sessions,
|
||||
selectedSessionId: null,
|
||||
@@ -2180,6 +2298,30 @@ async function hydrateProjectSession(
|
||||
): Promise<void> {
|
||||
const projectEventKey = getProjectEventKey(get());
|
||||
if (!projectEventKey) return;
|
||||
const hydrationKey = `${projectEventKey}:${sessionId}`;
|
||||
const activeHydration = activeSessionHydrations.get(hydrationKey);
|
||||
if (activeHydration) {
|
||||
await activeHydration;
|
||||
return;
|
||||
}
|
||||
|
||||
const hydration = hydrateProjectSessionSnapshot(set, get, sessionId, projectEventKey);
|
||||
activeSessionHydrations.set(hydrationKey, hydration);
|
||||
try {
|
||||
await hydration;
|
||||
} finally {
|
||||
if (activeSessionHydrations.get(hydrationKey) === hydration) {
|
||||
activeSessionHydrations.delete(hydrationKey);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function hydrateProjectSessionSnapshot(
|
||||
set: OpencodeSet,
|
||||
get: OpencodeGet,
|
||||
sessionId: string,
|
||||
projectEventKey: string,
|
||||
): Promise<void> {
|
||||
try {
|
||||
const response = await hostApiFetch<{ messages: unknown[] }>(
|
||||
`/api/opencode/sessions/${encodeURIComponent(sessionId)}/messages`,
|
||||
@@ -2210,11 +2352,37 @@ function getSessionTranscriptPatchFromRawMessages(
|
||||
return getHydratedSessionTranscriptPatch(state, sessionId, rawMessages);
|
||||
}
|
||||
|
||||
function getProjectSessionHydrationIds(state: OpencodeState): string[] {
|
||||
const sessionIds = new Set<string>();
|
||||
if (state.selectedSessionId) sessionIds.add(state.selectedSessionId);
|
||||
for (const [sessionId, sending] of Object.entries(state.sendingSessionIds)) {
|
||||
if (sending) sessionIds.add(sessionId);
|
||||
}
|
||||
for (const [sessionId, status] of Object.entries(state.sessionStatuses)) {
|
||||
if (status.type === 'busy' || status.type === 'retry') sessionIds.add(sessionId);
|
||||
}
|
||||
for (const [sessionId, runState] of Object.entries(state.sessionRunStates)) {
|
||||
if (runState.phase !== 'idle') {
|
||||
sessionIds.add(sessionId);
|
||||
}
|
||||
}
|
||||
for (const request of state.pendingQuestions) {
|
||||
if (request.sessionID) sessionIds.add(request.sessionID);
|
||||
}
|
||||
for (const request of state.pendingPermissions) {
|
||||
if (request.sessionID) sessionIds.add(request.sessionID);
|
||||
}
|
||||
return [...sessionIds].filter(Boolean);
|
||||
}
|
||||
|
||||
async function hydrateProjectSessions(set: OpencodeSet, get: OpencodeGet): Promise<void> {
|
||||
const sessions = get().sessions
|
||||
.map((session) => getSessionId(session))
|
||||
.filter((sessionId): sessionId is string => Boolean(sessionId));
|
||||
await Promise.allSettled(sessions.map((sessionId) => hydrateProjectSession(set, get, sessionId)));
|
||||
// Reconnecting the project stream used to fetch every transcript in
|
||||
// parallel. That made a project with hundreds of sessions turn a transient
|
||||
// network flap into a large burst of HTTP, parsing, and React work. The
|
||||
// selected and active sessions remain live; older sessions hydrate when the
|
||||
// user selects them, preserving complete history without the startup burst.
|
||||
const sessionIds = getProjectSessionHydrationIds(get());
|
||||
await Promise.allSettled(sessionIds.map((sessionId) => hydrateProjectSession(set, get, sessionId)));
|
||||
}
|
||||
|
||||
function applyProjectEventToStore(
|
||||
@@ -2578,6 +2746,12 @@ async function runSessionSubmission(
|
||||
const baseline = createSessionRunBaseline(getSessionMessagesForState(get(), sessionId));
|
||||
const runGeneration = get().runtimeGeneration;
|
||||
const runToken = beginSessionRun(sessionId);
|
||||
const backgroundLeaseId = `opencode:${sessionId}:${runToken}`;
|
||||
updateBackgroundLease({
|
||||
id: backgroundLeaseId,
|
||||
kind: submission.kind,
|
||||
active: true,
|
||||
});
|
||||
const promptId = submission.promptId;
|
||||
closeSessionEventSource(sessionId);
|
||||
set((state) => {
|
||||
@@ -3033,6 +3207,8 @@ async function runSessionSubmission(
|
||||
|
||||
const now = Date.now();
|
||||
const runListenerOwnsStream = hasSessionRunEventListener(sessionId);
|
||||
const liveProjectStreamHealthy = runListenerOwnsStream
|
||||
&& get().projectEventStreamState === 'connected';
|
||||
const observedRunState = getSessionRunStateForSession(get(), sessionId);
|
||||
const eventStatus = runListenerOwnsStream
|
||||
? observedRunState.runId === runToken
|
||||
@@ -3041,8 +3217,8 @@ async function runSessionSubmission(
|
||||
? IDLE_SESSION_STATUS
|
||||
: get().sessionStatuses[sessionId]
|
||||
: undefined;
|
||||
const shouldPollStatus = now >= nextStatusFallbackPollAt
|
||||
|| eventStatus?.type === 'idle';
|
||||
const shouldPollStatus = eventStatus?.type === 'idle'
|
||||
|| (!liveProjectStreamHealthy && now >= nextStatusFallbackPollAt);
|
||||
let statuses: OpencodeSessionStatusMap = {};
|
||||
if (shouldPollStatus) {
|
||||
const statusResponse = await awaitDuringPromptStart(
|
||||
@@ -3055,8 +3231,8 @@ async function runSessionSubmission(
|
||||
const observedStatus = shouldPollStatus
|
||||
? httpStatus ?? eventStatus
|
||||
: eventStatus ?? httpStatus ?? status;
|
||||
const shouldRefreshMessages = Date.now() >= nextMessageFallbackPollAt
|
||||
|| observedStatus?.type === 'idle';
|
||||
const shouldRefreshMessages = observedStatus?.type === 'idle'
|
||||
|| (!liveProjectStreamHealthy && Date.now() >= nextMessageFallbackPollAt);
|
||||
let refreshedMessages = messages;
|
||||
if (shouldRefreshMessages) {
|
||||
refreshedMessages = await awaitDuringPromptStart(fetchSessionMessages(sessionId));
|
||||
@@ -3285,9 +3461,87 @@ async function runSessionSubmission(
|
||||
throw error;
|
||||
} finally {
|
||||
cancelPromptStartWatchdog(sessionId, runToken);
|
||||
updateBackgroundLease({
|
||||
id: backgroundLeaseId,
|
||||
kind: submission.kind,
|
||||
active: false,
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
function getSessionMessageLoadKey(state: OpencodeState, sessionId: string): string {
|
||||
return `${getProjectEventKey(state) ?? getProjectId(state.activeProject) ?? 'no-project'}:${sessionId}`;
|
||||
}
|
||||
|
||||
function loadSessionMessagesDeduped(
|
||||
set: OpencodeSet,
|
||||
get: OpencodeGet,
|
||||
sessionId: string,
|
||||
): Promise<RawMessage[]> {
|
||||
const key = getSessionMessageLoadKey(get(), sessionId);
|
||||
const existing = activeSessionMessageLoads.get(key);
|
||||
if (existing) return existing;
|
||||
|
||||
const load = (async () => {
|
||||
set({ loading: true, ...errorState(null) });
|
||||
try {
|
||||
const response = await hostApiFetch<{ messages: unknown[] }>(
|
||||
`/api/opencode/sessions/${encodeURIComponent(sessionId)}/messages`,
|
||||
);
|
||||
const rawMessages = Array.isArray(response.messages) ? response.messages : [];
|
||||
const currentMessages = getSessionMessagesForState(get(), sessionId);
|
||||
const messages = mergeLocalUserAttachments(
|
||||
currentMessages,
|
||||
normalizeOpencodeSessionMessages(rawMessages),
|
||||
);
|
||||
if (isSessionRunning(get(), sessionId)) {
|
||||
set((state) => ({
|
||||
...getHydratedSessionTranscriptPatch(state, sessionId, rawMessages),
|
||||
...getSessionMessagePatch(state, sessionId, mergePendingOptimisticUserMessages(
|
||||
getSessionMessagesForState(state, sessionId),
|
||||
messages,
|
||||
)),
|
||||
loading: false,
|
||||
...errorState(null),
|
||||
}));
|
||||
return messages;
|
||||
}
|
||||
set((state) => ({
|
||||
...getHydratedSessionTranscriptPatch(state, sessionId, rawMessages),
|
||||
...getSessionMessagePatch({ ...state, selectedSessionId: sessionId }, sessionId, messages),
|
||||
selectedSessionId: sessionId,
|
||||
streamingMessage: state.streamingMessagesBySessionId[sessionId] ?? null,
|
||||
streamingTools: state.streamingToolsBySessionId[sessionId] ?? [],
|
||||
sendingSessionId: state.sendingSessionIds[sessionId] ? sessionId : null,
|
||||
loading: false,
|
||||
...errorState(null),
|
||||
}));
|
||||
return messages;
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
if (isSessionRunning(get(), sessionId)) {
|
||||
set(errorState(message));
|
||||
} else {
|
||||
set({
|
||||
loading: false,
|
||||
...errorState(message),
|
||||
});
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
})();
|
||||
activeSessionMessageLoads.set(key, load);
|
||||
void load.then(
|
||||
() => {
|
||||
if (activeSessionMessageLoads.get(key) === load) activeSessionMessageLoads.delete(key);
|
||||
},
|
||||
() => {
|
||||
if (activeSessionMessageLoads.get(key) === load) activeSessionMessageLoads.delete(key);
|
||||
},
|
||||
);
|
||||
return load;
|
||||
}
|
||||
|
||||
export const useOpencodeStore = create<OpencodeState>((set, get) => ({
|
||||
status: initialStatus,
|
||||
health: null,
|
||||
@@ -3356,6 +3610,11 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
|
||||
}));
|
||||
},
|
||||
|
||||
pauseBackgroundStreams() {
|
||||
closeActiveSessionEventSource();
|
||||
set({ projectEventStreamState: 'closed' });
|
||||
},
|
||||
|
||||
async start() {
|
||||
await runRuntimeAction('/api/opencode/start', set, get().status);
|
||||
set({ health: null, healthCheckedAt: null });
|
||||
@@ -3718,6 +3977,7 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
|
||||
const selectedSessionId = state.selectedSessionId;
|
||||
const deletingSelected = selectedSessionId === sessionId;
|
||||
closeSessionEventSource(sessionId);
|
||||
forgetSessionMessageCache(state, sessionId);
|
||||
set({
|
||||
sessionsByProjectId: cacheSessionsForProject(state.sessionsByProjectId, requestProjectId, nextSessions),
|
||||
...(requestStillActive
|
||||
@@ -3899,6 +4159,7 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
|
||||
|
||||
async selectSession(sessionId, options) {
|
||||
const state = get();
|
||||
touchSessionMessageCache(state, sessionId);
|
||||
sessionDiffLoadSequence += 1;
|
||||
const projectId = getProjectId(state.activeProject);
|
||||
if (projectId) {
|
||||
@@ -3955,52 +4216,7 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
|
||||
},
|
||||
|
||||
async loadSessionMessages(sessionId) {
|
||||
set({ loading: true, ...errorState(null) });
|
||||
try {
|
||||
const response = await hostApiFetch<{ messages: unknown[] }>(
|
||||
`/api/opencode/sessions/${encodeURIComponent(sessionId)}/messages`,
|
||||
);
|
||||
const rawMessages = Array.isArray(response.messages) ? response.messages : [];
|
||||
const currentMessages = getSessionMessagesForState(get(), sessionId);
|
||||
const messages = mergeLocalUserAttachments(
|
||||
currentMessages,
|
||||
normalizeOpencodeSessionMessages(rawMessages),
|
||||
);
|
||||
if (isSessionRunning(get(), sessionId)) {
|
||||
set((state) => ({
|
||||
...getHydratedSessionTranscriptPatch(state, sessionId, rawMessages),
|
||||
...getSessionMessagePatch(state, sessionId, mergePendingOptimisticUserMessages(
|
||||
getSessionMessagesForState(state, sessionId),
|
||||
messages,
|
||||
)),
|
||||
loading: false,
|
||||
...errorState(null),
|
||||
}));
|
||||
return messages;
|
||||
}
|
||||
set((state) => ({
|
||||
...getHydratedSessionTranscriptPatch(state, sessionId, rawMessages),
|
||||
...getSessionMessagePatch({ ...state, selectedSessionId: sessionId }, sessionId, messages),
|
||||
selectedSessionId: sessionId,
|
||||
streamingMessage: state.streamingMessagesBySessionId[sessionId] ?? null,
|
||||
streamingTools: state.streamingToolsBySessionId[sessionId] ?? [],
|
||||
sendingSessionId: state.sendingSessionIds[sessionId] ? sessionId : null,
|
||||
loading: false,
|
||||
...errorState(null),
|
||||
}));
|
||||
return messages;
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : String(error);
|
||||
if (isSessionRunning(get(), sessionId)) {
|
||||
set(errorState(message));
|
||||
} else {
|
||||
set({
|
||||
loading: false,
|
||||
...errorState(message),
|
||||
});
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
return await loadSessionMessagesDeduped(set, get, sessionId);
|
||||
},
|
||||
|
||||
async loadQuestions() {
|
||||
|
||||
Reference in New Issue
Block a user