feat(app): project current server state (#38459)

This commit is contained in:
Brendan Allan 2026-07-24 10:49:16 +08:00 committed by GitHub
commit 37c263e153
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
38 changed files with 2922 additions and 477 deletions

View file

@ -1,5 +1,6 @@
import { Binary } from "@opencode-ai/core/util/binary"
import { retry } from "@opencode-ai/core/util/retry"
import type { MessageApi, OpenCodeEvent, SessionApi, SessionMessageInfo } from "@opencode-ai/client/promise"
import type {
Message,
OpencodeClient,
@ -8,15 +9,18 @@ import type {
QuestionRequest,
Session,
SessionStatus,
SnapshotFileDiff,
Todo,
} from "@opencode-ai/sdk/v2/client"
import type { FileDiffInfo } from "@opencode-ai/client/promise"
import { batch } from "solid-js"
import { createStore, produce, reconcile } from "solid-js/store"
import { diffs as cleanDiffs, message as cleanMessage } from "@/utils/diffs"
import { message as cleanMessage } from "@/utils/diffs"
import { sessionNotFoundError } from "@/utils/server-errors"
import { rootSession } from "@/utils/session-route"
import { normalizeSessionInfo } from "@/utils/session"
import { normalizeSessionMessages } from "@/utils/session-message"
import { dropSessionCaches, pickSessionCacheEvictions, SESSION_CACHE_LIMIT } from "./global-sync/session-cache"
import { createV2SessionReducer, type V2SessionReduction } from "./server-session-v2-reducer"
const cmp = (a: string, b: string) => (a < b ? -1 : a > b ? 1 : 0)
const cmpMessage = (a: Message, b: Message) => a.time.created - b.time.created || cmp(a.id, b.id)
@ -36,10 +40,37 @@ type OptimisticItem = {
type MessagePage = {
session: Message[]
part: { id: string; part: Part[] }[]
source?: SessionMessageInfo[]
sourceMode?: "latest" | "older"
projectSource?: boolean
cursor?: string
complete: boolean
}
function legacyMessageSource(items: { info: Message; parts: Part[] }[]): SessionMessageInfo[] {
return items
.slice()
.sort((a, b) => cmp(a.info.id, b.info.id))
.map((item) => {
if (item.info.role === "user") {
return {
id: item.info.id,
type: "user" as const,
text: item.parts.flatMap((part) => (part.type === "text" ? [part.text] : [])).join("\n"),
time: item.info.time,
}
}
return {
id: item.info.id,
type: "assistant" as const,
agent: item.info.agent ?? item.info.mode,
model: { id: item.info.modelID, providerID: item.info.providerID, variant: item.info.variant },
content: [],
time: item.info.time,
}
})
}
// Most markers describe the current HTTP attempt; deltaParts persists non-durable stream state across retries.
type MessageLoadState = {
touchedMessages: Set<string>
@ -52,6 +83,7 @@ type MessageLoadState = {
optimisticParts: Map<string, Set<string>>
orphanParents: Set<string>
clearedMessageParts: Set<string>
touchedSource: Set<string>
}
type MessageLoadBaseline = Pick<
@ -137,15 +169,25 @@ function reconcileFetched<T extends { id: string }>(
return [...result.values()].sort((a, b) => cmp(a.id, b.id))
}
export function createServerSession(client: OpencodeClient, options?: { retry?: typeof retry }) {
type ServerSessionOptions = { retry?: typeof retry; protocol?: Promise<"v1" | "v2"> }
export function createServerSession(
client: OpencodeClient,
sessionApiOrOptions?: SessionApi | ServerSessionOptions,
messageApi?: MessageApi,
currentOptions?: ServerSessionOptions,
) {
const sessionApi = messageApi ? (sessionApiOrOptions as SessionApi) : undefined
const options = messageApi ? currentOptions : (sessionApiOrOptions as ServerSessionOptions | undefined)
const [data, setData] = createStore({
info: {} as Record<string, Session | undefined>,
session_status: {} as Record<string, SessionStatus>,
session_diff: {} as Record<string, SnapshotFileDiff[]>,
session_diff: {} as Record<string, FileDiffInfo[]>,
todo: {} as Record<string, Todo[]>,
permission: {} as Record<string, PermissionRequest[]>,
question: {} as Record<string, QuestionRequest[]>,
message: {} as Record<string, Message[]>,
session_message: {} as Record<string, SessionMessageInfo[]>,
part: {} as Record<string, Part[]>,
part_text_accum_delta: {} as Record<string, string>,
session_working(id: string) {
@ -154,9 +196,9 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
})
const requests = new Map<string, Promise<Session>>()
const inflight = new Map<string, Promise<void>>()
const inflightDiff = new Map<string, Promise<void>>()
const inflightTodo = new Map<string, Promise<void>>()
const optimistic = new Map<string, Map<string, OptimisticItem>>()
const v2 = createV2SessionReducer()
const messageLoads = new Map<string, MessageLoadState>()
const pendingParts = new Map<string, Map<string, Set<string>>>()
const orphanParts = new Map<string, Set<string>>()
@ -191,6 +233,16 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
at: {} as Record<string, number | undefined>,
})
const indexLegacyMessage = (message: Message) => {
const current = data.session_message[message.sessionID] ?? []
if (current.some((item) => item.id === message.id)) return
setData(
"session_message",
message.sessionID,
reconcile([...current, ...legacyMessageSource([{ info: message, parts: [] }])]),
)
}
const remember = (session: Session) => {
setData("info", session.id, reconcile(session))
infoSeen.delete(session.id)
@ -200,7 +252,6 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
...pinned.keys(),
...requests.keys(),
...inflight.keys(),
...inflightDiff.keys(),
...inflightTodo.keys(),
...messageLoads.keys(),
...optimistic.keys(),
@ -242,27 +293,31 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
const pending = requests.get(sessionID)
if (pending) return pending
const active = generation(sessionID)
const request = client.session.get({ sessionID }).then((result) => {
if (!result.data) throw sessionNotFoundError(sessionID)
if (generations.get(sessionID) !== active) return result.data
return remember(result.data)
const request = sessionApi
? sessionApi.get({ sessionID }).then(normalizeSessionInfo)
: client.session.get({ sessionID }).then((result) => {
if (!result.data) throw sessionNotFoundError(sessionID)
return result.data
})
const resolved = request.then((result) => {
if (generations.get(sessionID) !== active) return result
return remember(result)
})
requests.set(sessionID, request)
requests.set(sessionID, resolved)
const cleanup = () => {
if (requests.get(sessionID) === request) requests.delete(sessionID)
if (requests.get(sessionID) === resolved) requests.delete(sessionID)
if (
generations.get(sessionID) === active &&
!data.info[sessionID] &&
!requests.has(sessionID) &&
!messageLoads.has(sessionID) &&
!inflight.has(sessionID) &&
!inflightDiff.has(sessionID) &&
!inflightTodo.has(sessionID)
)
generations.delete(sessionID)
}
void request.then(cleanup, cleanup)
return request
void resolved.then(cleanup, cleanup)
return resolved
}
const peekLineage = (sessionID: string) => {
@ -419,9 +474,9 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
clearOptimistic(sessionID)
requests.delete(sessionID)
inflight.delete(sessionID)
inflightDiff.delete(sessionID)
inflightTodo.delete(sessionID)
messageLoads.delete(sessionID)
v2.clear(sessionID)
pendingParts.delete(sessionID)
orphanParts.delete(sessionID)
removedMessages.delete(sessionID)
@ -449,7 +504,6 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
...pinned.keys(),
...requests.keys(),
...inflight.keys(),
...inflightDiff.keys(),
...inflightTodo.keys(),
...messageLoads.keys(),
...optimistic.keys(),
@ -470,6 +524,25 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
)
const fetchMessages = async (sessionID: string, limit: number, before?: string, onAttempt?: () => void) => {
if (messageApi && (await options?.protocol) !== "v1") {
const response = await (options?.retry ?? retry)(() => {
onAttempt?.()
return messageApi.list(before ? { sessionID, limit, cursor: before } : { sessionID, limit, order: "desc" })
})
const source = [...response.data].reverse()
const normalized = normalizeSessionMessages(sessionID, source)
return {
session: normalized.messages.sort((a, b) => cmp(a.id, b.id)),
part: [...normalized.parts.entries()]
.map(([id, part]) => ({ id, part: part.sort((a, b) => cmp(a.id, b.id)) }))
.sort((a, b) => cmp(a.id, b.id)),
source,
sourceMode: before ? ("older" as const) : ("latest" as const),
projectSource: true,
cursor: response.cursor.next ?? undefined,
complete: response.data.length === 0,
}
}
const response = await (options?.retry ?? retry)(() => {
onAttempt?.()
return client.session.messages({ sessionID, limit, before })
@ -481,12 +554,24 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
id: item.info.id,
part: item.parts.filter((part) => !!part?.id).sort((a, b) => cmp(a.id, b.id)),
})),
source: legacyMessageSource(items),
sourceMode: before ? ("older" as const) : ("latest" as const),
cursor: response.response.headers.get("x-next-cursor") ?? undefined,
complete: !response.response.headers.get("x-next-cursor"),
}
}
const fetchMessage = async (sessionID: string, messageID: string, onAttempt?: () => void) => {
if (sessionApi && (await options?.protocol) !== "v1") {
const response = await (options?.retry ?? retry)(() => {
onAttempt?.()
return sessionApi.message({ sessionID, messageID })
})
const normalized = normalizeSessionMessages(sessionID, [response])
const message = normalized.messages[0]
if (!message) throw new Error(`Message not found: ${messageID}`)
return { message, parts: normalized.parts.get(messageID) ?? [] }
}
const response = await (options?.retry ?? retry)(() => {
onAttempt?.()
return client.session.message({ sessionID, messageID })
@ -571,7 +656,31 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
preserveUnfetched: boolean | ((message: Message) => boolean),
cleanupOrphans: boolean,
) => {
const merged = mergeOptimisticPage(page, [...(optimistic.get(sessionID)?.values() ?? [])])
const source = page.source
? (() => {
const incoming = new Map(page.source.map((message) => [message.id, message]))
const existing = data.session_message[sessionID] ?? []
const current = existing.filter((message) => !incoming.has(message.id))
const live = new Map(existing.map((message) => [message.id, message]))
return (page.sourceMode === "older" ? [...page.source, ...current] : [...current, ...page.source]).map(
(message) => (load?.touchedSource.has(message.id) ? (live.get(message.id) ?? message) : message),
)
})()
: undefined
const projected =
page.projectSource && source
? (() => {
const normalized = normalizeSessionMessages(sessionID, source)
return {
...page,
session: normalized.messages.sort((a, b) => cmp(a.id, b.id)),
part: [...normalized.parts.entries()]
.map(([id, part]) => ({ id, part: part.sort((a, b) => cmp(a.id, b.id)) }))
.sort((a, b) => cmp(a.id, b.id)),
}
})()
: page
const merged = mergeOptimisticPage(projected, [...(optimistic.get(sessionID)?.values() ?? [])])
merged.observed.forEach((item) => {
if (!load?.clearedMessageParts.has(item.messageID)) confirmOptimistic(sessionID, item.messageID, item.parts)
})
@ -583,6 +692,7 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
preserveUnfetched,
})
batch(() => {
if (source) setData("session_message", sessionID, reconcile(source))
const messageIDs = replaceMessages(sessionID, messages)
replaceParts(sessionID, merged.part, messageIDs, load)
const orphans = orphanParts.get(sessionID)
@ -613,6 +723,7 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
optimisticParts: new Map(),
orphanParents: new Set(),
clearedMessageParts: new Set(),
touchedSource: new Set(),
}
messageLoads.set(sessionID, load)
setMeta("loading", sessionID, true)
@ -747,6 +858,109 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
return properties.part.sessionID
}
const projectV2 = (reduction: V2SessionReduction) => {
reduction.touched.forEach((messageID) => messageLoads.get(reduction.sessionID)?.touchedSource.add(messageID))
setData("session_message", reduction.sessionID, reconcile(reduction.messages))
if (reduction.touched.length === 0) return
const touched = new Set(reduction.touched)
let parentID: string | undefined
for (const message of reduction.messages) {
if (message.type === "user" || (message.type === "synthetic" && message.description?.trim()))
parentID = message.id
if (message.type === "shell") {
if (touched.has(message.id)) touched.add(`${message.id}:assistant`)
parentID = undefined
}
if (message.type === "assistant" && touched.has(message.id) && parentID) touched.add(parentID)
if (message.type === "compaction" && touched.has(message.id) && parentID) touched.add(parentID)
}
const normalized = normalizeSessionMessages(reduction.sessionID, reduction.messages)
batch(() => {
for (const message of normalized.messages) {
if (!touched.has(message.id)) continue
apply({ type: "message.updated", properties: { sessionID: reduction.sessionID, info: message } })
}
for (const messageID of touched) {
const next = normalized.parts.get(messageID) ?? []
const nextIDs = new Set(next.map((part) => part.id))
for (const part of next) {
apply({ type: "message.part.updated", properties: { sessionID: reduction.sessionID, part } })
}
for (const part of data.part[messageID] ?? []) {
if (nextIDs.has(part.id)) continue
apply({
type: "message.part.removed",
properties: { sessionID: reduction.sessionID, messageID, partID: part.id },
})
}
}
})
}
const hydrateV2Message = (sessionID: string, messageID: string) => {
if (!sessionApi) return
void sessionApi
.message({ sessionID, messageID })
.then((message) => {
const current = data.session_message[sessionID] ?? []
const messages = [...current.filter((item) => item.id !== message.id), message].sort((a, b) => cmp(a.id, b.id))
projectV2({ sessionID, messages, touched: [message.id] })
})
.catch(() => {})
}
const applyV2 = (event: OpenCodeEvent) => {
if (!("data" in event) || !("sessionID" in event.data) || typeof event.data.sessionID !== "string") return
const sessionID = event.data.sessionID
const reduction = v2.reduce(data.session_message[sessionID] ?? [], event)
if (reduction) {
projectV2(reduction)
if (reduction.missing) hydrateV2Message(sessionID, reduction.missing)
}
const info = data.info[sessionID]
if (event.type === "session.renamed" && info)
remember({ ...info, title: event.data.title, time: { ...info.time, updated: event.created } })
if (event.type === "session.moved" && info)
remember({
...info,
projectID: event.data.projectID ?? info.projectID,
workspaceID: event.data.location.workspaceID,
directory: event.data.location.directory,
path: event.data.subpath,
time: { ...info.time, updated: event.created },
})
if (event.type === "session.usage.updated" && info)
remember({ ...info, cost: event.data.cost, tokens: event.data.tokens })
if (event.type === "session.archived") {
if (info) remember({ ...info, time: { ...info.time, archived: event.created, updated: event.created } })
evict([sessionID])
}
if (event.type === "session.execution.started") setData("session_status", sessionID, { type: "busy" })
if (
event.type === "session.execution.succeeded" ||
event.type === "session.execution.failed" ||
event.type === "session.execution.interrupted"
)
setData("session_status", sessionID, { type: "idle" })
if (event.type === "session.retry.scheduled")
setData("session_status", sessionID, {
type: "retry",
attempt: event.data.attempt,
message: event.data.error.message,
next: event.data.at,
})
if (event.type === "session.forked") void resolve(sessionID, { force: true }).catch(() => {})
if (
event.type === "session.revert.staged" ||
event.type === "session.revert.cleared" ||
event.type === "session.revert.committed"
)
void resolve(sessionID, { force: true }).catch(() => {})
}
const apply = (event: { type: string; properties?: unknown }) => {
const eventID = eventSessionID(event)
if (eventID) {
@ -770,7 +984,9 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
return
}
case "session.deleted": {
const sessionID = (event.properties as { info: Session }).info.id
const properties = event.properties as { sessionID?: string; info?: Session }
const sessionID = properties.info?.id ?? properties.sessionID
if (!sessionID) return
infoSeen.delete(sessionID)
setData(
"info",
@ -779,11 +995,6 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
evict([sessionID])
return
}
case "session.diff": {
const props = event.properties as { sessionID: string; diff: SnapshotFileDiff[] }
setData("session_diff", props.sessionID, reconcile(cleanDiffs(props.diff), { key: "file" }))
return
}
case "todo.updated": {
const props = event.properties as { sessionID: string; todos: Todo[] }
setData("todo", props.sessionID, reconcile(props.todos, { key: "id" }))
@ -796,6 +1007,7 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
}
case "message.updated": {
const info = cleanMessage((event.properties as { info: Message }).info)
indexLegacyMessage(info)
const load = messageLoads.get(info.sessionID)
load?.touchedMessages.add(info.id)
load?.removedMessages.delete(info.id)
@ -828,6 +1040,9 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
}
case "message.removed": {
const props = event.properties as { sessionID: string; messageID: string }
setData("session_message", props.sessionID, (messages) =>
messages?.filter((message) => message.id !== props.messageID),
)
const load = messageLoads.get(props.sessionID)
load?.touchedMessages.add(props.messageID)
load?.removedMessages.add(props.messageID)
@ -1140,23 +1355,16 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
setData(produce((draft) => deleteMessageParts(draft, input.messageID)))
},
},
diff(sessionID: string, options?: { force?: boolean }) {
async todo(sessionID: string, request?: { force?: boolean }) {
touch(sessionID)
if (data.session_diff[sessionID] !== undefined && !options?.force) return Promise.resolve()
return runInflight(inflightDiff, sessionID, () => {
const active = generation(sessionID)
return retry(() => client.session.diff({ sessionID })).then((result) => {
if (generations.get(sessionID) !== active) return
setData("session_diff", sessionID, reconcile(cleanDiffs(result.data), { key: "file" }))
})
})
},
todo(sessionID: string, options?: { force?: boolean }) {
touch(sessionID)
if (data.todo[sessionID] !== undefined && !options?.force) return Promise.resolve()
if (data.todo[sessionID] !== undefined && !request?.force) return
if ((await options?.protocol) === "v2") {
setData("todo", sessionID, [])
return
}
return runInflight(inflightTodo, sessionID, () => {
const active = generation(sessionID)
return retry(() => client.session.todo({ sessionID })).then((result) => {
return (options?.retry ?? retry)(() => client.session.todo({ sessionID })).then((result) => {
if (generations.get(sessionID) !== active) return
setData("todo", sessionID, reconcile(result.data ?? [], { key: "id" }))
})
@ -1190,6 +1398,7 @@ export function createServerSession(client: OpencodeClient, options?: { retry?:
if (count && count > 1) pinned.set(sessionID, count - 1)
},
apply,
applyV2,
}
}