|
|
@@ -18,6 +18,13 @@ import { compareMessages, messageKey, normalizeSessionMessages } from "@/utils/s
|
|
|
import { dropSessionCaches, pickSessionCacheEvictions, SESSION_CACHE_LIMIT } from "./global-sync/session-cache"
|
|
|
import { createV2SessionReducer, type V2SessionReduction } from "./server-session-v2-reducer"
|
|
|
import type { ServerApi } from "@/utils/server"
|
|
|
+import {
|
|
|
+ createCommentMetadata,
|
|
|
+ formatCommentNote,
|
|
|
+ parseCommentNote,
|
|
|
+ readCommentMetadata,
|
|
|
+ type PromptComment,
|
|
|
+} from "@/utils/comment-note"
|
|
|
|
|
|
type MessageApi = ServerApi["message"]
|
|
|
|
|
|
@@ -28,31 +35,6 @@ const historyMessagePageSize = 200
|
|
|
const sessionInfoLimit = 2_048
|
|
|
const emptyIDs: ReadonlySet<string> = new Set()
|
|
|
|
|
|
-function projectMessageSource(message: Message): SessionMessageInfo[] {
|
|
|
- if (message.role === "user") {
|
|
|
- return [
|
|
|
- { id: `${message.id}:agent`, type: "agent-switched", agent: message.agent, time: message.time },
|
|
|
- {
|
|
|
- id: `${message.id}:model`,
|
|
|
- type: "model-switched",
|
|
|
- model: { id: message.model.modelID, providerID: message.model.providerID, variant: message.model.variant },
|
|
|
- time: message.time,
|
|
|
- },
|
|
|
- { id: message.id, type: "user", text: "", time: message.time },
|
|
|
- ]
|
|
|
- }
|
|
|
- return [
|
|
|
- {
|
|
|
- id: message.id,
|
|
|
- type: "assistant",
|
|
|
- agent: message.agent ?? message.mode,
|
|
|
- model: { id: message.modelID, providerID: message.providerID, variant: message.variant },
|
|
|
- content: [],
|
|
|
- time: message.time,
|
|
|
- },
|
|
|
- ]
|
|
|
-}
|
|
|
-
|
|
|
function needsOlderTurnRoot(source: readonly SessionMessageInfo[]) {
|
|
|
const boundary = source.find(
|
|
|
(message) =>
|
|
|
@@ -64,13 +46,6 @@ function needsOlderTurnRoot(source: readonly SessionMessageInfo[]) {
|
|
|
return boundary?.type === "assistant"
|
|
|
}
|
|
|
|
|
|
-type OptimisticItem = {
|
|
|
- message: Message
|
|
|
- parts: Part[]
|
|
|
- confirmedParts?: Part[]
|
|
|
- confirmedMessage?: boolean
|
|
|
-}
|
|
|
-
|
|
|
type MessagePage = {
|
|
|
session: Message[]
|
|
|
part: { id: string; part: Part[] }[]
|
|
|
@@ -81,6 +56,18 @@ type MessagePage = {
|
|
|
complete: boolean
|
|
|
}
|
|
|
|
|
|
+export type PromptEcho = {
|
|
|
+ sessionID: string
|
|
|
+ messageID: string
|
|
|
+ text: string
|
|
|
+ displayText: string
|
|
|
+ agent: string
|
|
|
+ model: { providerID: string; modelID: string; variant?: string }
|
|
|
+ files?: { uri: string; mime: string; name?: string; mention?: { start: number; end: number; text: string } }[]
|
|
|
+ agents?: { name: string; mention?: { start: number; end: number; text: string } }[]
|
|
|
+ comments: PromptComment[]
|
|
|
+}
|
|
|
+
|
|
|
// Most markers describe the current HTTP attempt; deltaParts persists non-durable stream state across retries.
|
|
|
type MessageLoadState = {
|
|
|
touchedMessages: Set<string>
|
|
|
@@ -90,7 +77,6 @@ type MessageLoadState = {
|
|
|
deltaParts: Map<string, Set<string>>
|
|
|
carriedDeltaParts: Map<string, Set<string>>
|
|
|
removedParts: Map<string, Set<string>>
|
|
|
- optimisticParts: Map<string, Set<string>>
|
|
|
orphanParents: Set<string>
|
|
|
clearedMessageParts: Set<string>
|
|
|
touchedSource: Set<string>
|
|
|
@@ -101,34 +87,6 @@ type MessageLoadBaseline = Pick<
|
|
|
"touchedMessages" | "retainedMessages" | "touchedParts" | "clearedMessageParts"
|
|
|
>
|
|
|
|
|
|
-function mergeOptimisticPage(page: MessagePage, items: OptimisticItem[]) {
|
|
|
- if (items.length === 0) return { ...page, observed: [] as { messageID: string; parts: Part[] }[] }
|
|
|
- const session = [...page.session]
|
|
|
- const part = new Map(page.part.map((item) => [item.id, item.part]))
|
|
|
- const observed: { messageID: string; parts: Part[] }[] = []
|
|
|
- for (const item of items) {
|
|
|
- const result = Binary.search(session, messageKey(item.message), messageKey)
|
|
|
- const found = result.found
|
|
|
- if (!found) session.splice(result.index, 0, item.message)
|
|
|
- const current = part.get(item.message.id)
|
|
|
- const confirmed = found ? item.parts.filter((part) => current?.some((value) => value.id === part.id)) : []
|
|
|
- if (found) observed.push({ messageID: item.message.id, parts: confirmed })
|
|
|
- part.set(
|
|
|
- item.message.id,
|
|
|
- merge(
|
|
|
- found ? (current ?? []) : merge(item.confirmedParts ?? [], current ?? []),
|
|
|
- item.parts.filter((part) => !confirmed.includes(part)),
|
|
|
- ),
|
|
|
- )
|
|
|
- }
|
|
|
- return {
|
|
|
- ...page,
|
|
|
- session,
|
|
|
- part: [...part.entries()].sort((a, b) => cmp(a[0], b[0])).map(([id, parts]) => ({ id, part: parts })),
|
|
|
- observed,
|
|
|
- }
|
|
|
-}
|
|
|
-
|
|
|
function runInflight(map: Map<string, Promise<void>>, key: string, task: () => Promise<void>) {
|
|
|
const pending = map.get(key)
|
|
|
if (pending) return pending
|
|
|
@@ -212,7 +170,6 @@ export function createServerSession(
|
|
|
const requests = new Map<string, Promise<SessionInfo>>()
|
|
|
const inflight = new Map<string, Promise<void>>()
|
|
|
const inflightTodo = new Map<string, Promise<void>>()
|
|
|
- const optimistic = new Map<string, Map<string, OptimisticItem>>()
|
|
|
const v2 = createV2SessionReducer()
|
|
|
const pendingRevision = new Map<string, number>()
|
|
|
const formRevision = new Map<string, number>()
|
|
|
@@ -223,7 +180,45 @@ export function createServerSession(
|
|
|
const pendingParts = new Map<string, Map<string, Set<string>>>()
|
|
|
const orphanParts = new Map<string, Set<string>>()
|
|
|
const removedMessages = new Map<string, Set<string>>()
|
|
|
+ const echoes = new Map<string, Map<string, "sending" | "admitted">>()
|
|
|
+ const messageSnapshots = new Map<string, Set<string>>()
|
|
|
+ const settledInputs = new Map<string, Set<string>>()
|
|
|
const deltaBases = new Map<string, { base: string; sessionID: string }>()
|
|
|
+ const markEcho = (sessionID: string, messageID: string) => {
|
|
|
+ const messages = echoes.get(sessionID) ?? new Map<string, "sending" | "admitted">()
|
|
|
+ messages.set(messageID, "sending")
|
|
|
+ echoes.set(sessionID, messages)
|
|
|
+ }
|
|
|
+ const confirmEcho = (sessionID: string, messageID: string) => {
|
|
|
+ const messages = echoes.get(sessionID)
|
|
|
+ if (!messages?.has(messageID)) return false
|
|
|
+ messages.set(messageID, "admitted")
|
|
|
+ return true
|
|
|
+ }
|
|
|
+ const releaseEcho = (sessionID: string, messageID: string) => {
|
|
|
+ const messages = echoes.get(sessionID)
|
|
|
+ const state = messages?.get(messageID)
|
|
|
+ if (!messages || !state) return
|
|
|
+ messages.delete(messageID)
|
|
|
+ if (messages.size === 0) echoes.delete(sessionID)
|
|
|
+ return state
|
|
|
+ }
|
|
|
+ const present = (messageID: string, parts: Part[]) => {
|
|
|
+ const local = data.part[messageID] ?? []
|
|
|
+ const comments = local.filter(
|
|
|
+ (part) =>
|
|
|
+ part.type === "text" &&
|
|
|
+ part.synthetic &&
|
|
|
+ (readCommentMetadata(part.metadata) !== undefined || parseCommentNote(part.text) !== undefined),
|
|
|
+ )
|
|
|
+ if (!comments.length) return parts
|
|
|
+ const text = local.find((part) => part.type === "text" && !part.synthetic)
|
|
|
+ const projected = parts.flatMap((part) => {
|
|
|
+ if (part.id !== `${messageID}:text:0` || part.type !== "text") return [part]
|
|
|
+ return text?.type === "text" && text.text ? [{ ...part, text: text.text }] : []
|
|
|
+ })
|
|
|
+ return merge(projected, comments)
|
|
|
+ }
|
|
|
const deleteMessageParts = (
|
|
|
cache: { part: Record<string, Part[] | undefined>; part_text_accum_delta: Record<string, string | undefined> },
|
|
|
messageID: string,
|
|
|
@@ -253,18 +248,6 @@ export function createServerSession(
|
|
|
at: {} as Record<string, number | undefined>,
|
|
|
})
|
|
|
|
|
|
- const indexProjectedMessage = (message: Message) => {
|
|
|
- const current = data.session_message[message.sessionID] ?? []
|
|
|
- if (current.some((item) => item.id === message.id)) return
|
|
|
- const projected = projectMessageSource(message)
|
|
|
- const projectedIDs = new Set(projected.map((item) => item.id))
|
|
|
- setData(
|
|
|
- "session_message",
|
|
|
- message.sessionID,
|
|
|
- reconcile([...current.filter((item) => !projectedIDs.has(item.id)), ...projected]),
|
|
|
- )
|
|
|
- }
|
|
|
-
|
|
|
const remember = (session: SessionInfo) => {
|
|
|
setData("info", session.id, reconcile(session))
|
|
|
infoSeen.delete(session.id)
|
|
|
@@ -276,7 +259,7 @@ export function createServerSession(
|
|
|
...inflight.keys(),
|
|
|
...inflightTodo.keys(),
|
|
|
...messageLoads.keys(),
|
|
|
- ...optimistic.keys(),
|
|
|
+ ...echoes.keys(),
|
|
|
...Object.entries(data.permission)
|
|
|
.filter(([, items]) => items.length > 0)
|
|
|
.map(([sessionID]) => sessionID),
|
|
|
@@ -352,65 +335,6 @@ export function createServerSession(
|
|
|
return { session, root }
|
|
|
}
|
|
|
|
|
|
- const clearOptimistic = (sessionID: string, messageID?: string) => {
|
|
|
- if (!messageID) {
|
|
|
- optimistic.delete(sessionID)
|
|
|
- return
|
|
|
- }
|
|
|
- const items = optimistic.get(sessionID)
|
|
|
- if (!items) return
|
|
|
- items.delete(messageID)
|
|
|
- if (items.size === 0) optimistic.delete(sessionID)
|
|
|
- }
|
|
|
-
|
|
|
- const clearOptimisticPart = (sessionID: string, messageID: string, partID: string) => {
|
|
|
- const items = optimistic.get(sessionID)
|
|
|
- const item = items?.get(messageID)
|
|
|
- if (!items || !item) return
|
|
|
- const parts = item.parts.filter((part) => part.id !== partID)
|
|
|
- const confirmedParts = item.confirmedParts?.filter((part) => part.id !== partID)
|
|
|
- if (parts.length === 0) {
|
|
|
- clearOptimistic(sessionID, messageID)
|
|
|
- return
|
|
|
- }
|
|
|
- items.set(messageID, { ...item, parts, confirmedParts, confirmedMessage: true })
|
|
|
- }
|
|
|
-
|
|
|
- const confirmOptimisticPart = (sessionID: string, messageID: string, part: Part) => {
|
|
|
- const items = optimistic.get(sessionID)
|
|
|
- const item = items?.get(messageID)
|
|
|
- if (!items || !item) return
|
|
|
- const parts = item.parts.filter((value) => value.id !== part.id)
|
|
|
- if (parts.length === 0) {
|
|
|
- clearOptimistic(sessionID, messageID)
|
|
|
- return
|
|
|
- }
|
|
|
- items.set(messageID, {
|
|
|
- ...item,
|
|
|
- parts,
|
|
|
- confirmedParts: merge(item.confirmedParts ?? [], [part]),
|
|
|
- confirmedMessage: true,
|
|
|
- })
|
|
|
- }
|
|
|
-
|
|
|
- const confirmOptimistic = (sessionID: string, messageID: string, confirmedParts: Part[]) => {
|
|
|
- const items = optimistic.get(sessionID)
|
|
|
- const item = items?.get(messageID)
|
|
|
- if (!items || !item) return
|
|
|
- const confirmed = new Set(confirmedParts.map((part) => part.id))
|
|
|
- const parts = item.parts.filter((part) => !confirmed.has(part.id))
|
|
|
- if (parts.length === 0) {
|
|
|
- clearOptimistic(sessionID, messageID)
|
|
|
- return
|
|
|
- }
|
|
|
- items.set(messageID, {
|
|
|
- ...item,
|
|
|
- parts,
|
|
|
- confirmedParts: merge(item.confirmedParts ?? [], confirmedParts),
|
|
|
- confirmedMessage: true,
|
|
|
- })
|
|
|
- }
|
|
|
-
|
|
|
const trackPartChange = (sessionID: string, messageID: string, partID: string) => {
|
|
|
const load = messageLoads.get(sessionID)
|
|
|
if (!load) return
|
|
|
@@ -448,14 +372,6 @@ export function createServerSession(
|
|
|
const messages = data.message[sessionID]
|
|
|
if (messages?.some((message) => message.id === messageID)) load.retainedMessages.add(messageID)
|
|
|
}
|
|
|
- for (const [messageID, parts] of load.optimisticParts) {
|
|
|
- load.removedMessages.delete(messageID)
|
|
|
- load.clearedMessageParts.add(messageID)
|
|
|
- load.touchedMessages.add(messageID)
|
|
|
- const touched = load.touchedParts.get(messageID) ?? new Set<string>()
|
|
|
- parts.forEach((partID) => touched.add(partID))
|
|
|
- load.touchedParts.set(messageID, touched)
|
|
|
- }
|
|
|
baseline?.touchedMessages.forEach((messageID) => load.touchedMessages.add(messageID))
|
|
|
baseline?.retainedMessages.forEach((messageID) => load.retainedMessages.add(messageID))
|
|
|
baseline?.clearedMessageParts.forEach((messageID) => load.clearedMessageParts.add(messageID))
|
|
|
@@ -486,7 +402,9 @@ export function createServerSession(
|
|
|
sessionIDs.forEach((sessionID) => {
|
|
|
messageHydrationRevision.set(sessionID, (messageHydrationRevision.get(sessionID) ?? 0) + 1)
|
|
|
generations.delete(sessionID)
|
|
|
- clearOptimistic(sessionID)
|
|
|
+ echoes.delete(sessionID)
|
|
|
+ messageSnapshots.delete(sessionID)
|
|
|
+ settledInputs.delete(sessionID)
|
|
|
requests.delete(sessionID)
|
|
|
inflight.delete(sessionID)
|
|
|
inflightTodo.delete(sessionID)
|
|
|
@@ -521,7 +439,7 @@ export function createServerSession(
|
|
|
...inflight.keys(),
|
|
|
...inflightTodo.keys(),
|
|
|
...messageLoads.keys(),
|
|
|
- ...optimistic.keys(),
|
|
|
+ ...echoes.keys(),
|
|
|
...Object.entries(data.permission)
|
|
|
.filter(([, items]) => items.length > 0)
|
|
|
.map(([sessionID]) => sessionID),
|
|
|
@@ -598,9 +516,10 @@ export function createServerSession(
|
|
|
) => {
|
|
|
for (const item of items) {
|
|
|
if (!messageIDs.has(item.id)) continue
|
|
|
- const fetched = load?.clearedMessageParts.has(item.id)
|
|
|
- ? []
|
|
|
- : item.part.filter((part) => !SKIP_PARTS.has(part.type))
|
|
|
+ const fetched = present(
|
|
|
+ item.id,
|
|
|
+ load?.clearedMessageParts.has(item.id) ? [] : item.part.filter((part) => !SKIP_PARTS.has(part.type)),
|
|
|
+ )
|
|
|
const fetchedIDs = new Set(fetched.map((part) => part.id))
|
|
|
const pending = pendingParts.get(sessionID)?.get(item.id)
|
|
|
const touched = new Set([...(load?.touchedParts.get(item.id) ?? []), ...(pending ?? [])])
|
|
|
@@ -651,25 +570,37 @@ export function createServerSession(
|
|
|
preserveUnfetched: boolean | ((message: Message) => boolean),
|
|
|
cleanupOrphans: boolean,
|
|
|
) => {
|
|
|
+ if (page.sourceMode === "latest")
|
|
|
+ messageSnapshots.set(sessionID, new Set((page.source ?? []).map((message) => message.id)))
|
|
|
+ page.source?.forEach((message) => releaseEcho(sessionID, message.id))
|
|
|
const source = page.source
|
|
|
? (() => {
|
|
|
const incoming = new Map(page.source.map((message) => [message.id, message]))
|
|
|
const existing = data.session_message[sessionID] ?? []
|
|
|
const boundary = Math.min(...page.source.map((message) => message.time.created))
|
|
|
+ const inbox = new Set(data.input[sessionID] ?? [])
|
|
|
const current = existing.filter(
|
|
|
(message) =>
|
|
|
!incoming.has(message.id) &&
|
|
|
+ !inbox.has(message.id) &&
|
|
|
(page.sourceMode === "older" ||
|
|
|
load?.touchedSource.has(message.id) ||
|
|
|
(!page.complete && message.time.created < boundary)),
|
|
|
)
|
|
|
+ // message.list never returns admitted-but-undelivered inbox entries; keep them after the
|
|
|
+ // fetched history until a delivered or cancelled event resolves them.
|
|
|
+ const admitted = existing.filter((message) => !incoming.has(message.id) && inbox.has(message.id))
|
|
|
+ const combined =
|
|
|
+ page.sourceMode === "older"
|
|
|
+ ? [...page.source, ...current, ...admitted]
|
|
|
+ : [...current, ...page.source, ...admitted]
|
|
|
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),
|
|
|
+ return combined.map((message) =>
|
|
|
+ load?.touchedSource.has(message.id) ? (live.get(message.id) ?? message) : message,
|
|
|
)
|
|
|
})()
|
|
|
: undefined
|
|
|
- const projected =
|
|
|
+ const merged =
|
|
|
page.projectSource && source
|
|
|
? (() => {
|
|
|
const normalized = normalizeSessionMessages(sessionID, source)
|
|
|
@@ -682,16 +613,15 @@ export function createServerSession(
|
|
|
}
|
|
|
})()
|
|
|
: 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)
|
|
|
- })
|
|
|
const touchedMessages = new Set([...(load?.touchedMessages ?? []), ...(removedMessages.get(sessionID) ?? [])])
|
|
|
const messages = reconcileFetched(merged.session, data.message[sessionID] ?? [], {
|
|
|
touched: touchedMessages,
|
|
|
retained: load?.retainedMessages,
|
|
|
removed: load?.removedMessages,
|
|
|
- preserveUnfetched,
|
|
|
+ preserveUnfetched: (message) =>
|
|
|
+ echoes.get(sessionID)?.has(message.id) === true ||
|
|
|
+ preserveUnfetched === true ||
|
|
|
+ (typeof preserveUnfetched === "function" && preserveUnfetched(message)),
|
|
|
compare: compareMessages,
|
|
|
})
|
|
|
batch(() => {
|
|
|
@@ -723,7 +653,6 @@ export function createServerSession(
|
|
|
deltaParts: new Map(),
|
|
|
carriedDeltaParts: new Map(),
|
|
|
removedParts: new Map(),
|
|
|
- optimisticParts: new Map(),
|
|
|
orphanParents: new Set(),
|
|
|
clearedMessageParts: new Set(),
|
|
|
touchedSource: new Set(),
|
|
|
@@ -744,11 +673,7 @@ export function createServerSession(
|
|
|
const users = new Set([
|
|
|
...page.session.filter((message) => message.role === "user").map((message) => message.id),
|
|
|
...(data.message[sessionID] ?? [])
|
|
|
- .filter((message) => {
|
|
|
- if (message.role !== "user") return false
|
|
|
- const item = optimistic.get(sessionID)?.get(message.id)
|
|
|
- return load.touchedMessages.has(message.id) && (!item || item.confirmedMessage === true)
|
|
|
- })
|
|
|
+ .filter((message) => message.role === "user" && load.touchedMessages.has(message.id))
|
|
|
.map((message) => message.id),
|
|
|
])
|
|
|
const parentIDs = [
|
|
|
@@ -893,7 +818,7 @@ export function createServerSession(
|
|
|
apply({ type: "message.updated", properties: { sessionID: reduction.sessionID, info: message } })
|
|
|
}
|
|
|
for (const messageID of touched) {
|
|
|
- const next = normalized.parts.get(messageID) ?? []
|
|
|
+ const next = present(messageID, 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 } })
|
|
|
@@ -926,6 +851,67 @@ export function createServerSession(
|
|
|
.catch(() => {})
|
|
|
}
|
|
|
|
|
|
+ const removeEcho = (sessionID: string, messageID: string) => {
|
|
|
+ if (!releaseEcho(sessionID, messageID)) return false
|
|
|
+ pendingRevision.set(sessionID, (pendingRevision.get(sessionID) ?? 0) + 1)
|
|
|
+ const load = messageLoads.get(sessionID)
|
|
|
+ load?.touchedMessages.add(messageID)
|
|
|
+ load?.removedMessages.add(messageID)
|
|
|
+ load?.clearedMessageParts.add(messageID)
|
|
|
+ batch(() => {
|
|
|
+ setData("pending", sessionID, (items) => items?.filter((item) => item.id !== messageID))
|
|
|
+ setData("input", sessionID, (items) => items?.filter((id) => id !== messageID))
|
|
|
+ setData("message", sessionID, (messages) => messages?.filter((message) => message.id !== messageID))
|
|
|
+ setData(produce((draft) => deleteMessageParts(draft, messageID)))
|
|
|
+ })
|
|
|
+ return true
|
|
|
+ }
|
|
|
+
|
|
|
+ const confirmInbox = (item: SessionInboxInfo) => {
|
|
|
+ if (!confirmEcho(item.sessionID, item.id)) return false
|
|
|
+ v2.confirm(item)
|
|
|
+ pendingRevision.set(item.sessionID, (pendingRevision.get(item.sessionID) ?? 0) + 1)
|
|
|
+ const current = data.pending[item.sessionID] ?? []
|
|
|
+ const index = current.findIndex((entry) => entry.id === item.id)
|
|
|
+ if (index < 0) setData("pending", item.sessionID, [...current, item])
|
|
|
+ if (index >= 0) setData("pending", item.sessionID, index, reconcile(item))
|
|
|
+ return true
|
|
|
+ }
|
|
|
+
|
|
|
+ const reconcileInbox = (sessionID: string) => {
|
|
|
+ const pending = new Set((data.pending[sessionID] ?? []).map((item) => item.id))
|
|
|
+ const fetched = messageSnapshots.get(sessionID) ?? new Set<string>()
|
|
|
+ const removed = [...(settledInputs.get(sessionID) ?? [])].filter(
|
|
|
+ (messageID) => !pending.has(messageID) && !fetched.has(messageID),
|
|
|
+ )
|
|
|
+ settledInputs.delete(sessionID)
|
|
|
+ if (removed.length) {
|
|
|
+ const ids = new Set(removed)
|
|
|
+ const source = data.session_message[sessionID] ?? []
|
|
|
+ projectV2({
|
|
|
+ sessionID,
|
|
|
+ messages: source.filter((message) => !ids.has(message.id)),
|
|
|
+ touched: [],
|
|
|
+ removed: source.filter((message) => ids.has(message.id)).map((message) => message.id),
|
|
|
+ })
|
|
|
+ }
|
|
|
+
|
|
|
+ const messages = echoes.get(sessionID)
|
|
|
+ if (!messages) return
|
|
|
+ const projected = new Set((data.session_message[sessionID] ?? []).map((message) => message.id))
|
|
|
+ for (const [messageID, state] of messages) {
|
|
|
+ if (projected.has(messageID)) {
|
|
|
+ releaseEcho(sessionID, messageID)
|
|
|
+ continue
|
|
|
+ }
|
|
|
+ if (pending.has(messageID)) {
|
|
|
+ confirmEcho(sessionID, messageID)
|
|
|
+ continue
|
|
|
+ }
|
|
|
+ if (state === "admitted") removeEcho(sessionID, messageID)
|
|
|
+ }
|
|
|
+ }
|
|
|
+
|
|
|
const applyV2 = (event: OpenCodeEvent) => {
|
|
|
if (event.type === "form.created") {
|
|
|
formRevision.set(event.data.form.sessionID, (formRevision.get(event.data.form.sessionID) ?? 0) + 1)
|
|
|
@@ -949,6 +935,9 @@ export function createServerSession(
|
|
|
}
|
|
|
if (!("data" in event) || !("sessionID" in event.data) || typeof event.data.sessionID !== "string") return
|
|
|
const sessionID = event.data.sessionID
|
|
|
+ if (event.type === "session.inbox.enqueued" || event.type === "session.inbox.delivered")
|
|
|
+ releaseEcho(sessionID, event.data.inboxID)
|
|
|
+ if (event.type === "session.inbox.cancelled") removeEcho(sessionID, event.data.inboxID)
|
|
|
if (
|
|
|
event.type === "session.inbox.enqueued" ||
|
|
|
event.type === "session.inbox.delivery.changed" ||
|
|
|
@@ -960,11 +949,10 @@ export function createServerSession(
|
|
|
pendingRevision.set(sessionID, (pendingRevision.get(sessionID) ?? 0) + 1)
|
|
|
if (event.type === "session.inbox.enqueued") {
|
|
|
const current = data.pending[sessionID] ?? []
|
|
|
- if (!current.some((item) => item.id === event.data.inboxID))
|
|
|
- setData("pending", sessionID, [
|
|
|
- ...current,
|
|
|
- { id: event.data.inboxID, sessionID, timeCreated: event.created, ...event.data.item },
|
|
|
- ])
|
|
|
+ const item = { id: event.data.inboxID, sessionID, timeCreated: event.created, ...event.data.item }
|
|
|
+ const index = current.findIndex((entry) => entry.id === event.data.inboxID)
|
|
|
+ if (index < 0) setData("pending", sessionID, [...current, item])
|
|
|
+ if (index >= 0) setData("pending", sessionID, index, reconcile(item))
|
|
|
if (event.data.item.type !== "compaction" && !data.input[sessionID]?.includes(event.data.inboxID))
|
|
|
setData("input", sessionID, [...(data.input[sessionID] ?? []), event.data.inboxID])
|
|
|
}
|
|
|
@@ -1101,16 +1089,9 @@ export function createServerSession(
|
|
|
}
|
|
|
case "message.updated": {
|
|
|
const info = (event.properties as { info: Message }).info
|
|
|
- indexProjectedMessage(info)
|
|
|
const load = messageLoads.get(info.sessionID)
|
|
|
load?.touchedMessages.add(info.id)
|
|
|
load?.removedMessages.delete(info.id)
|
|
|
- const items = optimistic.get(info.sessionID)
|
|
|
- const item = items?.get(info.id)
|
|
|
- if (items && item) {
|
|
|
- if (item.parts.length === 0) clearOptimistic(info.sessionID, info.id)
|
|
|
- if (item.parts.length > 0) items.set(info.id, { ...item, confirmedMessage: true })
|
|
|
- }
|
|
|
const orphans = orphanParts.get(info.sessionID)
|
|
|
orphans?.delete(info.id)
|
|
|
if (orphans?.size === 0) orphanParts.delete(info.sessionID)
|
|
|
@@ -1123,13 +1104,18 @@ export function createServerSession(
|
|
|
return
|
|
|
}
|
|
|
const result = Binary.search(messages, messageKey(info), messageKey)
|
|
|
- if (result.found) setData("message", info.sessionID, result.index, reconcile(info))
|
|
|
- if (!result.found)
|
|
|
- setData("message", info.sessionID, (value = []) => {
|
|
|
- const next = value.slice()
|
|
|
- next.splice(result.index, 0, info)
|
|
|
- return next
|
|
|
- })
|
|
|
+ if (result.found) {
|
|
|
+ setData("message", info.sessionID, result.index, reconcile(info))
|
|
|
+ return
|
|
|
+ }
|
|
|
+ // Delivery rewrites time.created, changing the sort key; reposition instead of duplicating.
|
|
|
+ setData("message", info.sessionID, (value = []) => {
|
|
|
+ const next = value.slice()
|
|
|
+ const moved = next.findIndex((message) => message.id === info.id)
|
|
|
+ if (moved >= 0) next.splice(moved, 1)
|
|
|
+ next.splice(moved >= 0 && moved < result.index ? result.index - 1 : result.index, 0, info)
|
|
|
+ return next
|
|
|
+ })
|
|
|
return
|
|
|
}
|
|
|
case "message.removed": {
|
|
|
@@ -1144,13 +1130,11 @@ export function createServerSession(
|
|
|
load?.deltaParts.delete(props.messageID)
|
|
|
load?.carriedDeltaParts.delete(props.messageID)
|
|
|
load?.removedParts.delete(props.messageID)
|
|
|
- load?.optimisticParts.delete(props.messageID)
|
|
|
pendingParts.get(props.sessionID)?.delete(props.messageID)
|
|
|
if (pendingParts.get(props.sessionID)?.size === 0) pendingParts.delete(props.sessionID)
|
|
|
const removedMessagesForSession = removedMessages.get(props.sessionID) ?? new Set<string>()
|
|
|
removedMessagesForSession.add(props.messageID)
|
|
|
removedMessages.set(props.sessionID, removedMessagesForSession)
|
|
|
- clearOptimistic(props.sessionID, props.messageID)
|
|
|
setData(
|
|
|
produce((draft) => {
|
|
|
const messages = draft.message[props.sessionID]
|
|
|
@@ -1196,12 +1180,8 @@ export function createServerSession(
|
|
|
pending?.delete(part.id)
|
|
|
if (pending?.size === 0) pendingParts.get(part.sessionID)?.delete(part.messageID)
|
|
|
if (pendingParts.get(part.sessionID)?.size === 0) pendingParts.delete(part.sessionID)
|
|
|
- const optimistic = load?.optimisticParts.get(part.messageID)
|
|
|
- optimistic?.delete(part.id)
|
|
|
- if (optimistic?.size === 0) load?.optimisticParts.delete(part.messageID)
|
|
|
deltaBases.delete(part.id)
|
|
|
trackPartChange(part.sessionID, part.messageID, part.id)
|
|
|
- confirmOptimisticPart(part.sessionID, part.messageID, part)
|
|
|
setData(
|
|
|
"part_text_accum_delta",
|
|
|
produce((draft) => void delete draft[part.id]),
|
|
|
@@ -1240,12 +1220,8 @@ export function createServerSession(
|
|
|
const parts = load.removedParts.get(props.messageID) ?? new Set<string>()
|
|
|
parts.add(props.partID)
|
|
|
load.removedParts.set(props.messageID, parts)
|
|
|
- const optimistic = load.optimisticParts.get(props.messageID)
|
|
|
- optimistic?.delete(props.partID)
|
|
|
- if (optimistic?.size === 0) load.optimisticParts.delete(props.messageID)
|
|
|
}
|
|
|
trackPartChange(props.sessionID, props.messageID, props.partID)
|
|
|
- clearOptimisticPart(props.sessionID, props.messageID, props.partID)
|
|
|
setData(
|
|
|
produce((draft) => {
|
|
|
delete draft.part_text_accum_delta[props.partID]
|
|
|
@@ -1354,25 +1330,30 @@ export function createServerSession(
|
|
|
while (true) {
|
|
|
const pendingAt = pendingRevision.get(sessionID) ?? 0
|
|
|
const formAt = formRevision.get(sessionID) ?? 0
|
|
|
+ const previous = new Set(data.input[sessionID] ?? [])
|
|
|
const result = await load()
|
|
|
const pendingStable = (pendingRevision.get(sessionID) ?? 0) === pendingAt
|
|
|
const formStable = (formRevision.get(sessionID) ?? 0) === formAt
|
|
|
if (pendingStable) {
|
|
|
+ const current = new Set(result.pending.filter((item) => item.type !== "compaction").map((item) => item.id))
|
|
|
+ const settled = settledInputs.get(sessionID) ?? new Set<string>()
|
|
|
+ previous.forEach((messageID) => {
|
|
|
+ if (!current.has(messageID)) settled.add(messageID)
|
|
|
+ })
|
|
|
+ if (settled.size) settledInputs.set(sessionID, settled)
|
|
|
+ result.pending.forEach(v2.confirm)
|
|
|
setData("pending", sessionID, reconcile(result.pending))
|
|
|
- setData(
|
|
|
- "input",
|
|
|
- sessionID,
|
|
|
- reconcile(result.pending.filter((item) => item.type !== "compaction").map((item) => item.id)),
|
|
|
- )
|
|
|
+ setData("input", sessionID, reconcile([...current]))
|
|
|
}
|
|
|
if (formStable) setData("form", sessionID, reconcile(result.forms))
|
|
|
if (pendingStable && formStable) return
|
|
|
}
|
|
|
},
|
|
|
refreshPinned(hydrateTransient: (sessionID: string) => Promise<void>) {
|
|
|
+ const sessions = [...pinned.keys()]
|
|
|
return Promise.all(
|
|
|
- [...pinned.keys()].flatMap((sessionID) => [sync(sessionID, { force: true }), hydrateTransient(sessionID)]),
|
|
|
- ).then(() => undefined)
|
|
|
+ sessions.flatMap((sessionID) => [sync(sessionID, { force: true }), hydrateTransient(sessionID)]),
|
|
|
+ ).then(() => sessions.forEach(reconcileInbox))
|
|
|
},
|
|
|
invalidate() {
|
|
|
invalidationRevision += 1
|
|
|
@@ -1390,68 +1371,72 @@ export function createServerSession(
|
|
|
fresh(sessionID: string, ttl: number) {
|
|
|
return Date.now() - (meta.at[sessionID] ?? 0) <= ttl
|
|
|
},
|
|
|
- optimistic: {
|
|
|
- add(input: { sessionID: string; message: Message; parts: Part[] }) {
|
|
|
- const parts = input.parts
|
|
|
- .filter((part) => !!part?.id && !SKIP_PARTS.has(part.type))
|
|
|
- .sort((a, b) => cmp(a.id, b.id))
|
|
|
- const load = messageLoads.get(input.sessionID)
|
|
|
- if (load?.clearedMessageParts.has(input.message.id)) {
|
|
|
- const touched = load.touchedParts.get(input.message.id) ?? new Set<string>()
|
|
|
- parts.forEach((part) => touched.add(part.id))
|
|
|
- load.touchedParts.set(input.message.id, touched)
|
|
|
- }
|
|
|
- if (load) {
|
|
|
- load.removedMessages.delete(input.message.id)
|
|
|
- load.optimisticParts.set(input.message.id, new Set(parts.map((part) => part.id)))
|
|
|
+ inbox: {
|
|
|
+ echo(input: PromptEcho) {
|
|
|
+ const created = Date.now()
|
|
|
+ const files = input.files?.map((file) => ({
|
|
|
+ data: "",
|
|
|
+ mime: file.mime,
|
|
|
+ source: { type: "uri" as const, uri: file.uri },
|
|
|
+ name: file.name,
|
|
|
+ mention: file.mention,
|
|
|
+ }))
|
|
|
+ const item: SessionInboxInfo = {
|
|
|
+ id: input.messageID,
|
|
|
+ sessionID: input.sessionID,
|
|
|
+ timeCreated: created,
|
|
|
+ type: "user",
|
|
|
+ delivery: "steer",
|
|
|
+ payload: { text: input.text, files, agents: input.agents },
|
|
|
}
|
|
|
- const items = optimistic.get(input.sessionID)
|
|
|
- const removedMessagesForSession = removedMessages.get(input.sessionID)
|
|
|
- removedMessagesForSession?.delete(input.message.id)
|
|
|
- if (removedMessagesForSession?.size === 0) removedMessages.delete(input.sessionID)
|
|
|
- if (items) items.set(input.message.id, { ...input, parts, confirmedParts: [] })
|
|
|
- if (!items)
|
|
|
- optimistic.set(input.sessionID, new Map([[input.message.id, { ...input, parts, confirmedParts: [] }]]))
|
|
|
- indexProjectedMessage(input.message)
|
|
|
- setData("message", input.sessionID, (messages = []) => merge(messages, [input.message]).sort(compareMessages))
|
|
|
- setData(
|
|
|
- "part_text_accum_delta",
|
|
|
- produce((draft) => {
|
|
|
- for (const part of [...(data.part[input.message.id] ?? []), ...parts]) {
|
|
|
- delete draft[part.id]
|
|
|
- deltaBases.delete(part.id)
|
|
|
- }
|
|
|
- }),
|
|
|
- )
|
|
|
- setData("part", input.message.id, parts)
|
|
|
+ const projected = normalizeSessionMessages(input.sessionID, [
|
|
|
+ { id: `${input.messageID}:agent`, type: "agent-switched", agent: input.agent, time: { created } },
|
|
|
+ {
|
|
|
+ id: `${input.messageID}:model`,
|
|
|
+ type: "model-switched",
|
|
|
+ model: {
|
|
|
+ id: input.model.modelID,
|
|
|
+ providerID: input.model.providerID,
|
|
|
+ variant: input.model.variant,
|
|
|
+ },
|
|
|
+ time: { created },
|
|
|
+ },
|
|
|
+ {
|
|
|
+ id: input.messageID,
|
|
|
+ type: "user",
|
|
|
+ text: input.displayText,
|
|
|
+ files,
|
|
|
+ agents: input.agents,
|
|
|
+ time: { created },
|
|
|
+ },
|
|
|
+ ])
|
|
|
+ const message = projected.messages[0]!
|
|
|
+ const comments: Part[] = input.comments.map((comment, index) => ({
|
|
|
+ id: `${input.messageID}:comment:${index}`,
|
|
|
+ sessionID: input.sessionID,
|
|
|
+ messageID: input.messageID,
|
|
|
+ type: "text",
|
|
|
+ text: formatCommentNote(comment),
|
|
|
+ synthetic: true,
|
|
|
+ metadata: createCommentMetadata(comment),
|
|
|
+ }))
|
|
|
+ const parts = merge(projected.parts.get(input.messageID) ?? [], comments).sort((a, b) => cmp(a.id, b.id))
|
|
|
+ removedMessages.get(input.sessionID)?.delete(input.messageID)
|
|
|
+ markEcho(input.sessionID, input.messageID)
|
|
|
+ pendingRevision.set(input.sessionID, (pendingRevision.get(input.sessionID) ?? 0) + 1)
|
|
|
+ batch(() => {
|
|
|
+ setData("pending", input.sessionID, (items = []) => [...items.filter((entry) => entry.id !== item.id), item])
|
|
|
+ if (!data.input[input.sessionID]?.includes(input.messageID))
|
|
|
+ setData("input", input.sessionID, [...(data.input[input.sessionID] ?? []), input.messageID])
|
|
|
+ setData("message", input.sessionID, (messages = []) => merge(messages, [message]).sort(compareMessages))
|
|
|
+ setData("part", input.messageID, parts)
|
|
|
+ })
|
|
|
},
|
|
|
- remove(input: { sessionID: string; messageID: string }) {
|
|
|
- const item = optimistic.get(input.sessionID)?.get(input.messageID)
|
|
|
- if (!item) return
|
|
|
- messageLoads.get(input.sessionID)?.optimisticParts.delete(input.messageID)
|
|
|
- clearOptimistic(input.sessionID, input.messageID)
|
|
|
- if (item.confirmedMessage) {
|
|
|
- const partIDs = new Set(item.parts.map((part) => part.id))
|
|
|
- setData(
|
|
|
- produce((draft) => {
|
|
|
- for (const part of item.parts) {
|
|
|
- delete draft.part_text_accum_delta[part.id]
|
|
|
- deltaBases.delete(part.id)
|
|
|
- }
|
|
|
- const parts = draft.part[input.messageID]
|
|
|
- if (!parts) return
|
|
|
- draft.part[input.messageID] = parts.filter((part) => !partIDs.has(part.id))
|
|
|
- if (draft.part[input.messageID]?.length === 0) delete draft.part[input.messageID]
|
|
|
- }),
|
|
|
- )
|
|
|
- return
|
|
|
- }
|
|
|
- const projectedIDs = new Set(projectMessageSource(item.message).map((message) => message.id))
|
|
|
- setData("session_message", input.sessionID, (messages) =>
|
|
|
- messages?.filter((message) => !projectedIDs.has(message.id)),
|
|
|
- )
|
|
|
- setData("message", input.sessionID, (messages) => messages?.filter((message) => message.id !== input.messageID))
|
|
|
- setData(produce((draft) => deleteMessageParts(draft, input.messageID)))
|
|
|
+ confirm: confirmInbox,
|
|
|
+ reconcile: reconcileInbox,
|
|
|
+ clearEcho(input: { sessionID: string; messageID: string }) {
|
|
|
+ if (echoes.get(input.sessionID)?.get(input.messageID) !== "sending") return false
|
|
|
+ return removeEcho(input.sessionID, input.messageID)
|
|
|
},
|
|
|
},
|
|
|
async todo(sessionID: string, request?: { force?: boolean }) {
|