|
|
@@ -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,
|
|
|
}
|
|
|
}
|
|
|
|