实现 Codex 风格上下文压缩时间线交互

This commit is contained in:
2026-08-15 08:43:33 +08:00
parent 953b0f491e
commit fd9b5b46a9
13 changed files with 2044 additions and 73 deletions

View File

@@ -13,7 +13,10 @@ import {
removeStreamingPartFromMessage,
} from '@/lib/opencode-message';
import {
completeCompaction,
createPendingCompaction,
hydrateOpenCodeSession,
removeCompaction,
reduceOpenCodeEvent,
type OpenCodeEventEnvelope,
} from '@/lib/opencode-session-state';
@@ -162,7 +165,6 @@ interface OpencodeState {
sessionStatuses: OpencodeSessionStatusMap;
sessionTranscriptBySessionId: Record<string, OpencodeSessionTranscriptState>;
projectEventStreamState: 'closed' | 'connecting' | 'connected' | 'reconnecting';
compactingSessionIds: Record<string, boolean>;
sessionMessages: RawMessage[];
sessionMessagesBySessionId: Record<string, RawMessage[]>;
streamingMessage: RawMessage | null;
@@ -587,8 +589,12 @@ function getRuntimeStatusPatch(
status: OpencodeStatus,
): Pick<
OpencodeState,
'status' | 'runtimeGeneration' | 'commandsLoading' | 'commandsError'
> {
| 'status'
| 'runtimeGeneration'
| 'commandsLoading'
| 'commandsError'
| 'sessionTranscriptBySessionId'
> {
const identityChanged = runtimeIdentity(state.status) !== runtimeIdentity(status);
return {
status,
@@ -599,6 +605,9 @@ function getRuntimeStatusPatch(
),
commandsLoading: identityChanged ? false : state.commandsLoading,
commandsError: identityChanged ? null : state.commandsError,
sessionTranscriptBySessionId: identityChanged
? removeRunningCompactionsFromTranscripts(state.sessionTranscriptBySessionId)
: state.sessionTranscriptBySessionId,
};
}
@@ -1712,7 +1721,6 @@ function resetProjectScopedState(sessions: OpencodeSession[] = []) {
sessionStatuses: {},
sessionTranscriptBySessionId: {},
projectEventStreamState: 'closed' as const,
compactingSessionIds: {},
sessionMessages: [],
sessionMessagesBySessionId: {},
streamingMessage: null,
@@ -1877,6 +1885,85 @@ function getSessionTranscriptPatch(
};
}
function getOrCreateSessionTranscript(
state: OpencodeState,
sessionId: string,
): OpencodeSessionTranscriptState {
return state.sessionTranscriptBySessionId[sessionId]
?? hydrateOpenCodeSession(
sessionId,
[],
state.sessionStatuses[sessionId]?.type ?? 'idle',
);
}
function getCompactionEventIdentity(payload: Record<string, unknown>): string | undefined {
if (typeof payload.eventID === 'string' && payload.eventID.trim()) return payload.eventID;
if (typeof payload.eventId === 'string' && payload.eventId.trim()) return payload.eventId;
if (typeof payload.sequence === 'string' && payload.sequence.trim()) return payload.sequence;
return typeof payload.sequence === 'number' && Number.isFinite(payload.sequence)
? String(payload.sequence)
: undefined;
}
function withCompactionRunIdentity(
payload: Record<string, unknown>,
runToken: number,
generation: number,
source: 'manual' | 'automatic',
): Record<string, unknown> {
const part = payload.part;
if (!part || typeof part !== 'object' || Array.isArray(part)) return payload;
const partRecord = part as Record<string, unknown>;
if (partRecord.type !== 'compaction') return payload;
return {
...payload,
part: {
...partRecord,
runID: String(runToken),
generation,
...(source === 'manual' ? { auto: false } : {}),
},
};
}
function removeRunningCompactionsForRun(
transcript: OpencodeSessionTranscriptState,
runToken: number,
generation?: number,
): OpencodeSessionTranscriptState {
const runID = String(runToken);
let next = transcript;
for (const id of transcript.compactionOrder) {
const event = next.compactionsById[id];
if (
event?.status === 'running'
&& event.runID === runID
&& (generation === undefined || event.generation === generation)
) {
next = removeCompaction(next, { id });
}
}
return next;
}
function removeRunningCompactionsFromTranscripts(
transcripts: Record<string, OpencodeSessionTranscriptState>,
): Record<string, OpencodeSessionTranscriptState> {
let changed = false;
const next: Record<string, OpencodeSessionTranscriptState> = {};
for (const [sessionId, transcript] of Object.entries(transcripts)) {
let cleaned = transcript;
for (const id of transcript.compactionOrder) {
if (cleaned.compactionsById[id]?.status !== 'running') continue;
cleaned = removeCompaction(cleaned, { id });
}
changed ||= cleaned !== transcript;
next[sessionId] = cleaned;
}
return changed ? next : transcripts;
}
function getSessionStreamingEventPatch(
state: OpencodeState,
sessionId: string,
@@ -1907,19 +1994,9 @@ function getSessionStreamingEventPatch(
const streamingTools = nextToolStatus
? mergeStreamingToolStatuses(currentTools, nextToolStatus)
: currentTools;
const part = payload.part;
const compacting = Boolean(
part
&& typeof part === 'object'
&& !Array.isArray(part)
&& (part as Record<string, unknown>).type === 'compaction',
);
return {
...transcriptPatch,
...getSessionStreamingPatch(state, sessionId, streamingMessage, streamingTools),
compactingSessionIds: compacting
? { ...state.compactingSessionIds, [sessionId]: true }
: state.compactingSessionIds,
};
}
@@ -2009,7 +2086,12 @@ function getHydratedSessionTranscriptPatch(
return {
sessionTranscriptBySessionId: {
...state.sessionTranscriptBySessionId,
[sessionId]: hydrateOpenCodeSession(sessionId, rawMessages, status),
[sessionId]: hydrateOpenCodeSession(
sessionId,
rawMessages,
status,
state.sessionTranscriptBySessionId[sessionId],
),
},
};
}
@@ -2105,19 +2187,18 @@ function applyProjectEventToStore(
sessionStatuses: { ...state.sessionStatuses, [sessionId]: status },
};
}
if (type === 'session.idle' || type === 'session.compacted') {
if (type === 'session.idle') {
return {
...transcriptPatch,
sessionStatuses: { ...state.sessionStatuses, [sessionId]: IDLE_SESSION_STATUS },
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
...(active ? {} : clearSessionStreamingPatch(state, sessionId)),
};
}
if (type === 'session.compacted') return transcriptPatch;
if (type === 'session.error') {
return {
...transcriptPatch,
sessionStatuses: { ...state.sessionStatuses, [sessionId]: IDLE_SESSION_STATUS },
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
...(active ? {} : clearSessionStreamingPatch(state, sessionId)),
};
}
@@ -2328,7 +2409,6 @@ function finishSessionRunSuccessfully(
...statuses,
[sessionId]: IDLE_SESSION_STATUS,
},
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
sendingSessionIds,
sendingSessionId: getSendingSessionIdForSelection(nextState),
loading: false,
@@ -2411,6 +2491,7 @@ async function runSessionSubmission(
? setMessageDeliveryStatus(submission.optimisticUserMessage, 'sending')
: undefined;
const baseline = createSessionRunBaseline(getSessionMessagesForState(get(), sessionId));
const runGeneration = get().runtimeGeneration;
const runToken = beginSessionRun(sessionId);
const promptId = submission.promptId;
closeSessionEventSource(sessionId);
@@ -2435,6 +2516,15 @@ async function runSessionSubmission(
runId: runToken,
promptId,
});
const transcript = submission.kind === 'compact'
? createPendingCompaction(getOrCreateSessionTranscript(state, sessionId), {
id: promptId,
source: 'manual',
runID: String(runToken),
generation: runGeneration,
startedAt: Date.now(),
})
: undefined;
return {
...getSessionMessagePatch(stateForMessagePatch, sessionId, sessionMessages),
...clearSessionStreamingPatch(state, sessionId),
@@ -2447,6 +2537,14 @@ async function runSessionSubmission(
...state.sessionStatuses,
[sessionId]: { type: 'busy' },
},
...(transcript
? {
sessionTranscriptBySessionId: {
...state.sessionTranscriptBySessionId,
[sessionId]: transcript,
},
}
: {}),
};
});
@@ -2497,19 +2595,35 @@ async function runSessionSubmission(
...state.sessionStatuses,
[sessionId]: IDLE_SESSION_STATUS,
},
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
}));
});
registerSessionRunEventListener(sessionId, source, 'session.compacted', (event) => {
const payload = JSON.parse((event as MessageEvent<string>).data) as Record<string, unknown>;
if (!isSessionRunCurrent(sessionId, runToken)) return;
if (get().runtimeGeneration !== runGeneration) return;
if (getOpencodeEventSessionId(payload) !== sessionId) return;
flushSessionRunStreamBatch(set, sessionId, runToken);
set((state) => ({
...getSessionTranscriptPatch(state, 'session.compacted', payload),
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
}));
set((state) => {
const transcript = completeCompaction(
getOrCreateSessionTranscript(state, sessionId),
{ runID: String(runToken), generation: runGeneration },
{
nativeEventID: getCompactionEventIdentity(payload),
runID: String(runToken),
generation: runGeneration,
completedAt: Date.now(),
},
);
return transcript === state.sessionTranscriptBySessionId[sessionId]
? {}
: {
sessionTranscriptBySessionId: {
...state.sessionTranscriptBySessionId,
[sessionId]: transcript,
},
};
});
});
registerSessionRunEventListener(sessionId, source, 'session.error', (event) => {
@@ -2556,7 +2670,14 @@ async function runSessionSubmission(
...state.sessionStatuses,
[sessionId]: IDLE_SESSION_STATUS,
},
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
sessionTranscriptBySessionId: {
...state.sessionTranscriptBySessionId,
[sessionId]: removeRunningCompactionsForRun(
getOrCreateSessionTranscript(state, sessionId),
runToken,
runGeneration,
),
},
sendingSessionIds,
sendingSessionId: getSendingSessionIdForSelection(nextState),
loading: false,
@@ -2590,7 +2711,14 @@ async function runSessionSubmission(
...state.sessionStatuses,
[sessionId]: IDLE_SESSION_STATUS,
},
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
sessionTranscriptBySessionId: {
...state.sessionTranscriptBySessionId,
[sessionId]: removeRunningCompactionsForRun(
getOrCreateSessionTranscript(state, sessionId),
runToken,
runGeneration,
),
},
sendingSessionIds,
sendingSessionId: getSendingSessionIdForSelection(nextState),
loading: false,
@@ -2657,8 +2785,16 @@ async function runSessionSubmission(
});
registerSessionRunEventListener(sessionId, source, 'message.part.updated', (event) => {
const payload = JSON.parse((event as MessageEvent<string>).data) as Record<string, unknown>;
const rawPayload = JSON.parse((event as MessageEvent<string>).data) as Record<string, unknown>;
if (!isSessionRunCurrent(sessionId, runToken)) return;
const payload = get().runtimeGeneration === runGeneration
? withCompactionRunIdentity(
rawPayload,
runToken,
runGeneration,
submission.kind === 'compact' ? 'manual' : 'automatic',
)
: rawPayload;
if (getOpencodeEventSessionId(payload) !== sessionId) return;
queueSessionRunStreamEvent(set, sessionId, runToken, 'message.part.updated', payload);
});
@@ -2744,8 +2880,13 @@ async function runSessionSubmission(
const now = Date.now();
const runListenerOwnsStream = hasSessionRunEventListener(sessionId);
const observedRunState = getSessionRunStateForSession(get(), sessionId);
const eventStatus = runListenerOwnsStream
? get().sessionStatuses[sessionId]
? observedRunState.runId === runToken
&& observedRunState.phase === 'idle'
&& observedRunState.terminalReason === 'completed'
? IDLE_SESSION_STATUS
: get().sessionStatuses[sessionId]
: undefined;
const shouldPollStatus = now >= nextStatusFallbackPollAt
|| eventStatus?.type === 'idle';
@@ -2797,9 +2938,13 @@ async function runSessionSubmission(
// A missing HTTP status is common while the runtime is transitioning.
// Prefer a newly observed assistant response over the optimistic/event
// busy marker, otherwise a completed run can remain busy forever.
const polledStatus = shouldPollStatus
? httpStatus ?? inferredStatus
: eventStatus ?? status ?? inferredStatus;
const polledStatus = submission.kind === 'compact'
? shouldPollStatus
? httpStatus ?? eventStatus ?? status
: eventStatus ?? status
: shouldPollStatus
? httpStatus ?? inferredStatus
: eventStatus ?? status ?? inferredStatus;
status = assistantError || assistantAborted ? IDLE_SESSION_STATUS : polledStatus;
const currentMessages = getSessionMessagesForState(get(), sessionId);
messages = assistantAborted || status.type === 'idle'
@@ -2848,6 +2993,14 @@ async function runSessionSubmission(
...statuses,
[sessionId]: IDLE_SESSION_STATUS,
},
sessionTranscriptBySessionId: {
...state.sessionTranscriptBySessionId,
[sessionId]: removeRunningCompactionsForRun(
getOrCreateSessionTranscript(state, sessionId),
runToken,
runGeneration,
),
},
sendingSessionIds,
sendingSessionId: getSendingSessionIdForSelection(nextState),
loading: false,
@@ -2945,7 +3098,14 @@ async function runSessionSubmission(
...state.sessionStatuses,
[sessionId]: IDLE_SESSION_STATUS,
},
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
sessionTranscriptBySessionId: {
...state.sessionTranscriptBySessionId,
[sessionId]: removeRunningCompactionsForRun(
getOrCreateSessionTranscript(state, sessionId),
runToken,
runGeneration,
),
},
sendingSessionIds,
sendingSessionId: getSendingSessionIdForSelection(nextState),
loading: false,
@@ -2987,7 +3147,6 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
sessionRunStates: {},
pendingQuestions: [],
pendingPermissions: [],
compactingSessionIds: {},
sessionTodos: [],
sessionTodosBySessionId: {},
sessionTodosLoading: false,
@@ -3397,7 +3556,6 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
sessionStatuses: Object.fromEntries(
Object.entries(state.sessionStatuses).filter(([key]) => key !== sessionId),
),
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
selectedSessionId: deletingSelected ? null : selectedSessionId,
sessionMessages: deletingSelected ? [] : state.sessionMessages,
sessionMessagesBySessionId: withoutKey(state.sessionMessagesBySessionId, sessionId),
@@ -3437,6 +3595,7 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
const abortedRunId = getSessionRunStateForSession(get(), sessionId).runId
?? activeSessionRunTokens.get(sessionId)
?? null;
const abortedRunGeneration = get().runtimeGeneration;
invalidateSessionRun(sessionId);
closeSessionEventSource(sessionId);
set((state) => {
@@ -3456,7 +3615,18 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
...state.sessionStatuses,
[sessionId]: IDLE_SESSION_STATUS,
},
compactingSessionIds: withoutKey(state.compactingSessionIds, sessionId),
...(abortedRunId === null
? {}
: {
sessionTranscriptBySessionId: {
...state.sessionTranscriptBySessionId,
[sessionId]: removeRunningCompactionsForRun(
getOrCreateSessionTranscript(state, sessionId),
abortedRunId,
abortedRunGeneration,
),
},
}),
sendingSessionIds,
sendingSessionId: getSendingSessionIdForSelection(nextState),
loading: true,
@@ -4075,7 +4245,7 @@ export const useOpencodeStore = create<OpencodeState>((set, get) => ({
try {
return await runSessionSubmission(set, get, sessionId, {
kind: 'compact',
completion: 'post-success',
completion: 'observed',
promptId: `compact-${sessionId}-${Date.now()}`,
post: () => hostApiFetch<OpencodeSessionMessageActionResponse>(
`/api/opencode/sessions/${encodeURIComponent(sessionId)}/summarize`,