server-session.ts 43 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105
  1. import { Binary } from "@opencode-ai/core/util/binary"
  2. import { retry } from "@opencode-ai/core/util/retry"
  3. import type {
  4. Message,
  5. OpencodeClient,
  6. Part,
  7. PermissionRequest,
  8. QuestionRequest,
  9. Session,
  10. SessionStatus,
  11. SnapshotFileDiff,
  12. Todo,
  13. } from "@opencode-ai/sdk/v2/client"
  14. import { batch } from "solid-js"
  15. import { createStore, produce, reconcile } from "solid-js/store"
  16. import { diffs as cleanDiffs, message as cleanMessage } from "@/utils/diffs"
  17. import { rootSession } from "@/utils/session-route"
  18. import { dropSessionCaches, pickSessionCacheEvictions, SESSION_CACHE_LIMIT } from "./global-sync/session-cache"
  19. const cmp = (a: string, b: string) => (a < b ? -1 : a > b ? 1 : 0)
  20. const cmpMessage = (a: Message, b: Message) => a.time.created - b.time.created || cmp(a.id, b.id)
  21. const SKIP_PARTS = new Set(["patch", "step-start", "step-finish"])
  22. const initialMessagePageSize = 2
  23. const historyMessagePageSize = 200
  24. const sessionInfoLimit = 2_048
  25. const emptyIDs: ReadonlySet<string> = new Set()
  26. type OptimisticItem = {
  27. message: Message
  28. parts: Part[]
  29. confirmedParts?: Part[]
  30. confirmedMessage?: boolean
  31. }
  32. type MessagePage = {
  33. session: Message[]
  34. part: { id: string; part: Part[] }[]
  35. cursor?: string
  36. complete: boolean
  37. }
  38. // Most markers describe the current HTTP attempt; deltaParts persists non-durable stream state across retries.
  39. type MessageLoadState = {
  40. touchedMessages: Set<string>
  41. removedMessages: Set<string>
  42. retainedMessages: Set<string>
  43. touchedParts: Map<string, Set<string>>
  44. deltaParts: Map<string, Set<string>>
  45. carriedDeltaParts: Map<string, Set<string>>
  46. removedParts: Map<string, Set<string>>
  47. optimisticParts: Map<string, Set<string>>
  48. orphanParents: Set<string>
  49. clearedMessageParts: Set<string>
  50. }
  51. function mergeOptimisticPage(page: MessagePage, items: OptimisticItem[]) {
  52. if (items.length === 0) return { ...page, observed: [] as { messageID: string; parts: Part[] }[] }
  53. const session = [...page.session]
  54. const part = new Map(page.part.map((item) => [item.id, item.part]))
  55. const observed: { messageID: string; parts: Part[] }[] = []
  56. for (const item of items) {
  57. const result = Binary.search(session, item.message.id, (message) => message.id)
  58. if (!result.found) session.splice(result.index, 0, item.message)
  59. const current = part.get(item.message.id)
  60. const confirmed = result.found
  61. ? item.parts.filter((part) => Binary.search(current ?? [], part.id, (value) => value.id).found)
  62. : []
  63. if (result.found) observed.push({ messageID: item.message.id, parts: confirmed })
  64. part.set(
  65. item.message.id,
  66. merge(
  67. result.found ? (current ?? []) : merge(item.confirmedParts ?? [], current ?? []),
  68. item.parts.filter((part) => !confirmed.includes(part)),
  69. ),
  70. )
  71. }
  72. return {
  73. ...page,
  74. session,
  75. part: [...part.entries()].sort((a, b) => cmp(a[0], b[0])).map(([id, parts]) => ({ id, part: parts })),
  76. observed,
  77. }
  78. }
  79. function runInflight(map: Map<string, Promise<void>>, key: string, task: () => Promise<void>) {
  80. const pending = map.get(key)
  81. if (pending) return pending
  82. const promise = task().finally(() => {
  83. if (map.get(key) === promise) map.delete(key)
  84. })
  85. map.set(key, promise)
  86. return promise
  87. }
  88. function merge<T extends { id: string }>(a: readonly T[], b: readonly T[]) {
  89. const items = new Map(a.map((item) => [item.id, item] as const))
  90. for (const item of b) items.set(item.id, item)
  91. return [...items.values()].sort((x, y) => cmp(x.id, y.id))
  92. }
  93. function reconcileFetched<T extends { id: string }>(
  94. fetched: T[],
  95. current: readonly T[],
  96. options: {
  97. touched?: ReadonlySet<string>
  98. retained?: ReadonlySet<string>
  99. preserveUnfetched?: boolean | ((item: T) => boolean)
  100. } = {},
  101. ) {
  102. const result = new Map(fetched.map((item) => [item.id, item]))
  103. const live = new Map(current.map((item) => [item.id, item]))
  104. if (options.preserveUnfetched) {
  105. for (const item of current) {
  106. if (!result.has(item.id) && (options.preserveUnfetched === true || options.preserveUnfetched(item)))
  107. result.set(item.id, item)
  108. }
  109. }
  110. for (const id of options.retained ?? emptyIDs) {
  111. if (result.has(id)) continue
  112. const item = live.get(id)
  113. if (item) result.set(id, item)
  114. }
  115. // Events observed while the request is pending are the freshest client state for those identities.
  116. for (const id of options.touched ?? emptyIDs) {
  117. const item = live.get(id)
  118. if (item) result.set(id, item)
  119. if (!item) result.delete(id)
  120. }
  121. return [...result.values()].sort((a, b) => cmp(a.id, b.id))
  122. }
  123. export function createServerSession(client: OpencodeClient, options?: { retry?: typeof retry }) {
  124. const [data, setData] = createStore({
  125. info: {} as Record<string, Session | undefined>,
  126. session_status: {} as Record<string, SessionStatus>,
  127. session_diff: {} as Record<string, SnapshotFileDiff[]>,
  128. todo: {} as Record<string, Todo[]>,
  129. permission: {} as Record<string, PermissionRequest[]>,
  130. question: {} as Record<string, QuestionRequest[]>,
  131. message: {} as Record<string, Message[]>,
  132. part: {} as Record<string, Part[]>,
  133. part_text_accum_delta: {} as Record<string, string>,
  134. session_working(id: string) {
  135. return (this.session_status[id]?.type ?? "idle") !== "idle"
  136. },
  137. })
  138. const requests = new Map<string, Promise<Session>>()
  139. const inflight = new Map<string, Promise<void>>()
  140. const inflightDiff = new Map<string, Promise<void>>()
  141. const inflightTodo = new Map<string, Promise<void>>()
  142. const optimistic = new Map<string, Map<string, OptimisticItem>>()
  143. const messageLoads = new Map<string, MessageLoadState>()
  144. const pendingParts = new Map<string, Map<string, Set<string>>>()
  145. const orphanParts = new Map<string, Set<string>>()
  146. const removedMessages = new Map<string, Set<string>>()
  147. const deltaBases = new Map<string, { base: string; sessionID: string }>()
  148. const deleteMessageParts = (
  149. cache: { part: Record<string, Part[] | undefined>; part_text_accum_delta: Record<string, string | undefined> },
  150. messageID: string,
  151. ) => {
  152. for (const part of cache.part[messageID] ?? []) {
  153. delete cache.part_text_accum_delta[part.id]
  154. deltaBases.delete(part.id)
  155. }
  156. delete cache.part[messageID]
  157. }
  158. const seen = new Set<string>()
  159. const infoSeen = new Set<string>()
  160. const pinned = new Map<string, number>()
  161. const generations = new Map<string, object>()
  162. const generation = (sessionID: string) => {
  163. const current = generations.get(sessionID)
  164. if (current) return current
  165. const created = {}
  166. generations.set(sessionID, created)
  167. return created
  168. }
  169. const [meta, setMeta] = createStore({
  170. limit: {} as Record<string, number | undefined>,
  171. cursor: {} as Record<string, string | undefined>,
  172. complete: {} as Record<string, boolean | undefined>,
  173. loading: {} as Record<string, boolean | undefined>,
  174. at: {} as Record<string, number | undefined>,
  175. })
  176. const remember = (session: Session) => {
  177. setData("info", session.id, reconcile(session))
  178. infoSeen.delete(session.id)
  179. infoSeen.add(session.id)
  180. if (infoSeen.size > sessionInfoLimit) {
  181. const preserve = new Set([
  182. ...pinned.keys(),
  183. ...requests.keys(),
  184. ...inflight.keys(),
  185. ...inflightDiff.keys(),
  186. ...inflightTodo.keys(),
  187. ...messageLoads.keys(),
  188. ...optimistic.keys(),
  189. ...Object.entries(data.permission)
  190. .filter(([, items]) => items.length > 0)
  191. .map(([sessionID]) => sessionID),
  192. ...Object.entries(data.question)
  193. .filter(([, items]) => items.length > 0)
  194. .map(([sessionID]) => sessionID),
  195. ...Object.entries(data.session_status)
  196. .filter(([, status]) => status.type !== "idle")
  197. .map(([sessionID]) => sessionID),
  198. ])
  199. for (const sessionID of preserve) {
  200. let current = data.info[sessionID]
  201. while (current) {
  202. preserve.add(current.id)
  203. current = current.parentID ? data.info[current.parentID] : undefined
  204. }
  205. }
  206. const stale: string[] = []
  207. for (const sessionID of infoSeen) {
  208. if (infoSeen.size - stale.length <= sessionInfoLimit) break
  209. if (!preserve.has(sessionID)) stale.push(sessionID)
  210. }
  211. stale.forEach((sessionID) => infoSeen.delete(sessionID))
  212. stale.forEach((sessionID) => generations.delete(sessionID))
  213. setData(
  214. "info",
  215. produce((draft) => stale.forEach((sessionID) => delete draft[sessionID])),
  216. )
  217. }
  218. return session
  219. }
  220. const resolve = (sessionID: string, options?: { force?: boolean }) => {
  221. const cached = data.info[sessionID]
  222. if (cached && !options?.force) return Promise.resolve(cached)
  223. const pending = requests.get(sessionID)
  224. if (pending) return pending
  225. const active = generation(sessionID)
  226. const request = client.session.get({ sessionID }).then((result) => {
  227. if (!result.data) throw new Error(`Session not found: ${sessionID}`)
  228. if (generations.get(sessionID) !== active) return result.data
  229. return remember(result.data)
  230. })
  231. requests.set(sessionID, request)
  232. const cleanup = () => {
  233. if (requests.get(sessionID) === request) requests.delete(sessionID)
  234. if (
  235. generations.get(sessionID) === active &&
  236. !data.info[sessionID] &&
  237. !requests.has(sessionID) &&
  238. !messageLoads.has(sessionID) &&
  239. !inflight.has(sessionID) &&
  240. !inflightDiff.has(sessionID) &&
  241. !inflightTodo.has(sessionID)
  242. )
  243. generations.delete(sessionID)
  244. }
  245. void request.then(cleanup, cleanup)
  246. return request
  247. }
  248. const peekLineage = (sessionID: string) => {
  249. const session = data.info[sessionID]
  250. if (!session) return
  251. const seen = new Set([session.id])
  252. let root = session
  253. while (root.parentID) {
  254. if (seen.has(root.parentID)) throw new Error(`Session parent cycle: ${root.parentID}`)
  255. seen.add(root.parentID)
  256. const parent = data.info[root.parentID]
  257. if (!parent) return
  258. root = parent
  259. }
  260. return { session, root }
  261. }
  262. const clearOptimistic = (sessionID: string, messageID?: string) => {
  263. if (!messageID) {
  264. optimistic.delete(sessionID)
  265. return
  266. }
  267. const items = optimistic.get(sessionID)
  268. if (!items) return
  269. items.delete(messageID)
  270. if (items.size === 0) optimistic.delete(sessionID)
  271. }
  272. const clearOptimisticPart = (sessionID: string, messageID: string, partID: string) => {
  273. const items = optimistic.get(sessionID)
  274. const item = items?.get(messageID)
  275. if (!items || !item) return
  276. const parts = item.parts.filter((part) => part.id !== partID)
  277. const confirmedParts = item.confirmedParts?.filter((part) => part.id !== partID)
  278. if (parts.length === 0) {
  279. clearOptimistic(sessionID, messageID)
  280. return
  281. }
  282. items.set(messageID, { ...item, parts, confirmedParts, confirmedMessage: true })
  283. }
  284. const confirmOptimisticPart = (sessionID: string, messageID: string, part: Part) => {
  285. const items = optimistic.get(sessionID)
  286. const item = items?.get(messageID)
  287. if (!items || !item) return
  288. const parts = item.parts.filter((value) => value.id !== part.id)
  289. if (parts.length === 0) {
  290. clearOptimistic(sessionID, messageID)
  291. return
  292. }
  293. items.set(messageID, {
  294. ...item,
  295. parts,
  296. confirmedParts: merge(item.confirmedParts ?? [], [part]),
  297. confirmedMessage: true,
  298. })
  299. }
  300. const confirmOptimistic = (sessionID: string, messageID: string, confirmedParts: Part[]) => {
  301. const items = optimistic.get(sessionID)
  302. const item = items?.get(messageID)
  303. if (!items || !item) return
  304. const confirmed = new Set(confirmedParts.map((part) => part.id))
  305. const parts = item.parts.filter((part) => !confirmed.has(part.id))
  306. if (parts.length === 0) {
  307. clearOptimistic(sessionID, messageID)
  308. return
  309. }
  310. items.set(messageID, {
  311. ...item,
  312. parts,
  313. confirmedParts: merge(item.confirmedParts ?? [], confirmedParts),
  314. confirmedMessage: true,
  315. })
  316. }
  317. const trackPartChange = (sessionID: string, messageID: string, partID: string) => {
  318. const load = messageLoads.get(sessionID)
  319. if (!load) return
  320. // A part event keeps an existing parent when the fetched page omits it without overriding fetched metadata.
  321. const messages = data.message[sessionID]
  322. if (messages && Binary.search(messages, messageID, (message) => message.id).found)
  323. load.retainedMessages.add(messageID)
  324. const parts = load.touchedParts.get(messageID)
  325. if (parts) {
  326. parts.add(partID)
  327. return
  328. }
  329. load.touchedParts.set(messageID, new Set([partID]))
  330. }
  331. const resetMessageLoad = (sessionID: string, load: MessageLoadState) => {
  332. load.touchedMessages.clear()
  333. load.retainedMessages.clear()
  334. load.touchedParts.clear()
  335. load.carriedDeltaParts.clear()
  336. load.clearedMessageParts.clear()
  337. for (const messageID of load.removedMessages) {
  338. load.touchedMessages.add(messageID)
  339. load.clearedMessageParts.add(messageID)
  340. }
  341. for (const [messageID, parts] of load.deltaParts) {
  342. load.touchedParts.set(messageID, new Set(parts))
  343. load.carriedDeltaParts.set(messageID, new Set(parts))
  344. const messages = data.message[sessionID]
  345. if (messages && Binary.search(messages, messageID, (message) => message.id).found)
  346. load.retainedMessages.add(messageID)
  347. }
  348. for (const [messageID, parts] of load.removedParts) {
  349. const touched = load.touchedParts.get(messageID) ?? new Set<string>()
  350. parts.forEach((partID) => touched.add(partID))
  351. load.touchedParts.set(messageID, touched)
  352. const messages = data.message[sessionID]
  353. if (messages && Binary.search(messages, messageID, (message) => message.id).found)
  354. load.retainedMessages.add(messageID)
  355. }
  356. for (const [messageID, parts] of load.optimisticParts) {
  357. load.removedMessages.delete(messageID)
  358. load.clearedMessageParts.add(messageID)
  359. load.touchedMessages.add(messageID)
  360. const touched = load.touchedParts.get(messageID) ?? new Set<string>()
  361. parts.forEach((partID) => touched.add(partID))
  362. load.touchedParts.set(messageID, touched)
  363. }
  364. }
  365. const evict = (sessionIDs: string[]) => {
  366. if (sessionIDs.length === 0) return
  367. const evicted = new Set(sessionIDs)
  368. for (const [partID, item] of deltaBases) {
  369. if (evicted.has(item.sessionID)) deltaBases.delete(partID)
  370. }
  371. sessionIDs.forEach((sessionID) => {
  372. generations.delete(sessionID)
  373. clearOptimistic(sessionID)
  374. requests.delete(sessionID)
  375. inflight.delete(sessionID)
  376. inflightDiff.delete(sessionID)
  377. inflightTodo.delete(sessionID)
  378. messageLoads.delete(sessionID)
  379. pendingParts.delete(sessionID)
  380. orphanParts.delete(sessionID)
  381. removedMessages.delete(sessionID)
  382. })
  383. setData(
  384. produce((draft) => {
  385. dropSessionCaches(draft, sessionIDs)
  386. }),
  387. )
  388. setMeta(
  389. produce((draft) => {
  390. for (const sessionID of sessionIDs) {
  391. delete draft.limit[sessionID]
  392. delete draft.cursor[sessionID]
  393. delete draft.complete[sessionID]
  394. delete draft.loading[sessionID]
  395. delete draft.at[sessionID]
  396. }
  397. }),
  398. )
  399. }
  400. const protectedSessions = () =>
  401. new Set([
  402. ...pinned.keys(),
  403. ...requests.keys(),
  404. ...inflight.keys(),
  405. ...inflightDiff.keys(),
  406. ...inflightTodo.keys(),
  407. ...messageLoads.keys(),
  408. ...optimistic.keys(),
  409. ...Object.entries(data.permission)
  410. .filter(([, items]) => items.length > 0)
  411. .map(([sessionID]) => sessionID),
  412. ...Object.entries(data.question)
  413. .filter(([, items]) => items.length > 0)
  414. .map(([sessionID]) => sessionID),
  415. ...Object.entries(data.session_status)
  416. .filter(([, status]) => status.type !== "idle")
  417. .map(([sessionID]) => sessionID),
  418. ])
  419. const touch = (sessionID: string) =>
  420. evict(
  421. pickSessionCacheEvictions({ seen, keep: sessionID, limit: SESSION_CACHE_LIMIT, preserve: protectedSessions() }),
  422. )
  423. const fetchMessages = async (sessionID: string, limit: number, before?: string, onAttempt?: () => void) => {
  424. const response = await (options?.retry ?? retry)(() => {
  425. onAttempt?.()
  426. return client.session.messages({ sessionID, limit, before })
  427. })
  428. const items = (response.data ?? []).filter((item) => !!item?.info?.id)
  429. return {
  430. session: items.map((item) => cleanMessage(item.info)).sort((a, b) => cmp(a.id, b.id)),
  431. part: items.map((item) => ({
  432. id: item.info.id,
  433. part: item.parts.filter((part) => !!part?.id).sort((a, b) => cmp(a.id, b.id)),
  434. })),
  435. cursor: response.response.headers.get("x-next-cursor") ?? undefined,
  436. complete: !response.response.headers.get("x-next-cursor"),
  437. }
  438. }
  439. const replaceMessages = (sessionID: string, messages: Message[]) => {
  440. const messageIDs = new Set(messages.map((message) => message.id))
  441. const dropped = (data.message[sessionID] ?? []).filter((message) => !messageIDs.has(message.id))
  442. setData("message", sessionID, reconcile(messages, { key: "id" }))
  443. setData(
  444. produce((draft) => {
  445. for (const message of dropped) deleteMessageParts(draft, message.id)
  446. }),
  447. )
  448. return messageIDs
  449. }
  450. const replaceParts = (
  451. sessionID: string,
  452. items: MessagePage["part"],
  453. messageIDs: Set<string>,
  454. load?: MessageLoadState,
  455. ) => {
  456. for (const item of items) {
  457. if (!messageIDs.has(item.id)) continue
  458. const fetched = load?.clearedMessageParts.has(item.id)
  459. ? []
  460. : item.part.filter((part) => !SKIP_PARTS.has(part.type))
  461. const fetchedIDs = new Set(fetched.map((part) => part.id))
  462. const pending = pendingParts.get(sessionID)?.get(item.id)
  463. const touched = new Set([...(load?.touchedParts.get(item.id) ?? []), ...(pending ?? [])])
  464. for (const part of fetched) {
  465. const accumulated = data.part_text_accum_delta[part.id]
  466. const base = deltaBases.get(part.id)?.base
  467. const preserveDelta =
  468. base !== undefined &&
  469. accumulated !== undefined &&
  470. "text" in part &&
  471. typeof part.text === "string" &&
  472. part.text.startsWith(base) &&
  473. accumulated.startsWith(part.text) &&
  474. accumulated !== part.text
  475. if (preserveDelta) touched.add(part.id)
  476. if (load?.carriedDeltaParts.get(item.id)?.has(part.id) && !preserveDelta) touched.delete(part.id)
  477. }
  478. for (const partID of load?.carriedDeltaParts.get(item.id) ?? []) {
  479. if (!fetchedIDs.has(partID)) touched.delete(partID)
  480. }
  481. const parts = reconcileFetched(fetched, data.part[item.id] ?? [], { touched })
  482. if (!parts.length) {
  483. orphanParts.get(sessionID)?.delete(item.id)
  484. setData(produce((draft) => deleteMessageParts(draft, item.id)))
  485. continue
  486. }
  487. const partIDs = new Set(parts.map((part) => part.id))
  488. setData(
  489. "part_text_accum_delta",
  490. produce((draft) => {
  491. for (const part of data.part[item.id] ?? []) {
  492. if (!partIDs.has(part.id) || !touched.has(part.id)) {
  493. delete draft[part.id]
  494. deltaBases.delete(part.id)
  495. }
  496. }
  497. }),
  498. )
  499. setData("part", item.id, reconcile(parts, { key: "id" }))
  500. orphanParts.get(sessionID)?.delete(item.id)
  501. }
  502. }
  503. const applyMessagePage = (
  504. sessionID: string,
  505. page: MessagePage,
  506. load: MessageLoadState | undefined,
  507. preserveUnfetched: boolean | ((message: Message) => boolean),
  508. cleanupOrphans: boolean,
  509. ) => {
  510. const merged = mergeOptimisticPage(page, [...(optimistic.get(sessionID)?.values() ?? [])])
  511. merged.observed.forEach((item) => {
  512. if (!load?.clearedMessageParts.has(item.messageID)) confirmOptimistic(sessionID, item.messageID, item.parts)
  513. })
  514. const touchedMessages = new Set([...(load?.touchedMessages ?? []), ...(removedMessages.get(sessionID) ?? [])])
  515. const messages = reconcileFetched(merged.session, data.message[sessionID] ?? [], {
  516. touched: touchedMessages,
  517. retained: load?.retainedMessages,
  518. preserveUnfetched,
  519. })
  520. batch(() => {
  521. const messageIDs = replaceMessages(sessionID, messages)
  522. replaceParts(sessionID, merged.part, messageIDs, load)
  523. const orphans = orphanParts.get(sessionID)
  524. if (cleanupOrphans && page.complete && orphans) {
  525. for (const messageID of orphans) {
  526. if (!messageIDs.has(messageID)) setData(produce((draft) => deleteMessageParts(draft, messageID)))
  527. }
  528. orphanParts.delete(sessionID)
  529. }
  530. setMeta("limit", sessionID, messages.length)
  531. setMeta("cursor", sessionID, merged.cursor)
  532. setMeta("complete", sessionID, merged.complete)
  533. setMeta("at", sessionID, Date.now())
  534. })
  535. }
  536. const loadMessages = async (sessionID: string, limit: number, before?: string, mode?: "replace" | "prepend") => {
  537. if (meta.loading[sessionID]) return
  538. const active = generation(sessionID)
  539. const load: MessageLoadState = {
  540. touchedMessages: new Set(),
  541. removedMessages: new Set(),
  542. retainedMessages: new Set(),
  543. touchedParts: new Map(),
  544. deltaParts: new Map(),
  545. carriedDeltaParts: new Map(),
  546. removedParts: new Map(),
  547. optimisticParts: new Map(),
  548. orphanParents: new Set(),
  549. clearedMessageParts: new Set(),
  550. }
  551. messageLoads.set(sessionID, load)
  552. setMeta("loading", sessionID, true)
  553. let applied = false
  554. await fetchMessages(sessionID, limit, before, () => resetMessageLoad(sessionID, load))
  555. .then((page) => {
  556. if (generations.get(sessionID) !== active) return
  557. const first = page.session.reduce<Message | undefined>(
  558. (oldest, message) => (!oldest || cmpMessage(message, oldest) < 0 ? message : oldest),
  559. undefined,
  560. )
  561. const preserveUnfetched =
  562. mode === "prepend" || (!page.complete && (!first || ((message: Message) => cmpMessage(message, first) < 0)))
  563. applyMessagePage(
  564. sessionID,
  565. page,
  566. messageLoads.get(sessionID) === load ? load : undefined,
  567. preserveUnfetched,
  568. mode !== "prepend",
  569. )
  570. applied = true
  571. })
  572. .finally(() => {
  573. if (!applied && generations.get(sessionID) === active && messageLoads.get(sessionID) === load) {
  574. for (const messageID of load.orphanParents) {
  575. if (!orphanParts.get(sessionID)?.has(messageID)) continue
  576. setData(produce((draft) => deleteMessageParts(draft, messageID)))
  577. orphanParts.get(sessionID)?.delete(messageID)
  578. }
  579. if (orphanParts.get(sessionID)?.size === 0) orphanParts.delete(sessionID)
  580. }
  581. if (messageLoads.get(sessionID) === load) messageLoads.delete(sessionID)
  582. if (generations.get(sessionID) === active) setMeta("loading", sessionID, false)
  583. })
  584. }
  585. const sync = (sessionID: string, options?: { force?: boolean; messageLimit?: number }) => {
  586. touch(sessionID)
  587. return runInflight(inflight, sessionID, async () => {
  588. const cached = data.message[sessionID] !== undefined && meta.limit[sessionID] !== undefined
  589. if (cached && data.info[sessionID] && !options?.force) return
  590. await Promise.all([
  591. resolve(sessionID, options),
  592. cached && !options?.force
  593. ? Promise.resolve()
  594. : loadMessages(sessionID, options?.messageLimit ?? meta.limit[sessionID] ?? initialMessagePageSize),
  595. ])
  596. })
  597. }
  598. const prefetch = async (sessionID: string, limit: number) => {
  599. touch(sessionID)
  600. await inflight.get(sessionID)
  601. if (
  602. Date.now() - (meta.at[sessionID] ?? 0) <= 15_000 &&
  603. (meta.complete[sessionID] || (data.message[sessionID]?.length ?? 0) >= limit)
  604. )
  605. return
  606. await runInflight(inflight, sessionID, () => loadMessages(sessionID, limit))
  607. }
  608. const eventSessionID = (event: { type: string; properties?: unknown }) => {
  609. const properties = event.properties
  610. if (!properties || typeof properties !== "object") return
  611. if ("sessionID" in properties && typeof properties.sessionID === "string") return properties.sessionID
  612. if (
  613. "info" in properties &&
  614. properties.info &&
  615. typeof properties.info === "object" &&
  616. "sessionID" in properties.info &&
  617. typeof properties.info.sessionID === "string"
  618. )
  619. return properties.info.sessionID
  620. if (
  621. "part" in properties &&
  622. properties.part &&
  623. typeof properties.part === "object" &&
  624. "sessionID" in properties.part &&
  625. typeof properties.part.sessionID === "string"
  626. )
  627. return properties.part.sessionID
  628. }
  629. const apply = (event: { type: string; properties?: unknown }) => {
  630. const eventID = eventSessionID(event)
  631. if (eventID) {
  632. touch(eventID)
  633. if (
  634. !data.info[eventID] &&
  635. event.type !== "session.created" &&
  636. event.type !== "session.updated" &&
  637. event.type !== "session.deleted"
  638. )
  639. void resolve(eventID).catch(() => {})
  640. }
  641. switch (event.type) {
  642. case "session.created":
  643. remember((event.properties as { info: Session }).info)
  644. return
  645. case "session.updated": {
  646. const info = (event.properties as { info: Session }).info
  647. remember(info)
  648. if (info.time.archived) evict([info.id])
  649. return
  650. }
  651. case "session.deleted": {
  652. const sessionID = (event.properties as { info: Session }).info.id
  653. infoSeen.delete(sessionID)
  654. setData(
  655. "info",
  656. produce((draft) => void delete draft[sessionID]),
  657. )
  658. evict([sessionID])
  659. return
  660. }
  661. case "session.diff": {
  662. const props = event.properties as { sessionID: string; diff: SnapshotFileDiff[] }
  663. setData("session_diff", props.sessionID, reconcile(cleanDiffs(props.diff), { key: "file" }))
  664. return
  665. }
  666. case "todo.updated": {
  667. const props = event.properties as { sessionID: string; todos: Todo[] }
  668. setData("todo", props.sessionID, reconcile(props.todos, { key: "id" }))
  669. return
  670. }
  671. case "session.status": {
  672. const props = event.properties as { sessionID: string; status: SessionStatus }
  673. setData("session_status", props.sessionID, reconcile(props.status))
  674. return
  675. }
  676. case "message.updated": {
  677. const info = cleanMessage((event.properties as { info: Message }).info)
  678. const load = messageLoads.get(info.sessionID)
  679. load?.touchedMessages.add(info.id)
  680. load?.removedMessages.delete(info.id)
  681. const items = optimistic.get(info.sessionID)
  682. const item = items?.get(info.id)
  683. if (items && item) {
  684. if (item.parts.length === 0) clearOptimistic(info.sessionID, info.id)
  685. if (item.parts.length > 0) items.set(info.id, { ...item, confirmedMessage: true })
  686. }
  687. const orphans = orphanParts.get(info.sessionID)
  688. orphans?.delete(info.id)
  689. if (orphans?.size === 0) orphanParts.delete(info.sessionID)
  690. const removedMessagesForSession = removedMessages.get(info.sessionID)
  691. removedMessagesForSession?.delete(info.id)
  692. if (removedMessagesForSession?.size === 0) removedMessages.delete(info.sessionID)
  693. const messages = data.message[info.sessionID]
  694. if (!messages) {
  695. setData("message", info.sessionID, [info])
  696. return
  697. }
  698. const result = Binary.search(messages, info.id, (message) => message.id)
  699. if (result.found) setData("message", info.sessionID, result.index, reconcile(info))
  700. if (!result.found)
  701. setData("message", info.sessionID, (value = []) => {
  702. const next = value.slice()
  703. next.splice(result.index, 0, info)
  704. return next
  705. })
  706. return
  707. }
  708. case "message.removed": {
  709. const props = event.properties as { sessionID: string; messageID: string }
  710. const load = messageLoads.get(props.sessionID)
  711. load?.touchedMessages.add(props.messageID)
  712. load?.removedMessages.add(props.messageID)
  713. load?.clearedMessageParts.add(props.messageID)
  714. load?.deltaParts.delete(props.messageID)
  715. load?.carriedDeltaParts.delete(props.messageID)
  716. load?.removedParts.delete(props.messageID)
  717. load?.optimisticParts.delete(props.messageID)
  718. pendingParts.get(props.sessionID)?.delete(props.messageID)
  719. if (pendingParts.get(props.sessionID)?.size === 0) pendingParts.delete(props.sessionID)
  720. const removedMessagesForSession = removedMessages.get(props.sessionID) ?? new Set<string>()
  721. removedMessagesForSession.add(props.messageID)
  722. removedMessages.set(props.sessionID, removedMessagesForSession)
  723. clearOptimistic(props.sessionID, props.messageID)
  724. setData(
  725. produce((draft) => {
  726. const messages = draft.message[props.sessionID]
  727. if (messages) {
  728. const result = Binary.search(messages, props.messageID, (message) => message.id)
  729. if (result.found) messages.splice(result.index, 1)
  730. }
  731. deleteMessageParts(draft, props.messageID)
  732. }),
  733. )
  734. return
  735. }
  736. case "message.part.updated": {
  737. const part = (event.properties as { part: Part }).part
  738. if (SKIP_PARTS.has(part.type)) return
  739. const messages = data.message[part.sessionID]
  740. const load = messageLoads.get(part.sessionID)
  741. const missing = !messages || !Binary.search(messages, part.messageID, (message) => message.id).found
  742. // Outside a page load, accepting a part without its ordered parent event would create an unbounded orphan.
  743. if (
  744. missing &&
  745. (!load ||
  746. load.clearedMessageParts.has(part.messageID) ||
  747. removedMessages.get(part.sessionID)?.has(part.messageID))
  748. )
  749. return
  750. if (missing) {
  751. const orphans = orphanParts.get(part.sessionID) ?? new Set<string>()
  752. orphans.add(part.messageID)
  753. orphanParts.set(part.sessionID, orphans)
  754. load?.orphanParents.add(part.messageID)
  755. }
  756. const deltas = load?.deltaParts.get(part.messageID)
  757. deltas?.delete(part.id)
  758. if (deltas?.size === 0) load?.deltaParts.delete(part.messageID)
  759. const carried = load?.carriedDeltaParts.get(part.messageID)
  760. carried?.delete(part.id)
  761. if (carried?.size === 0) load?.carriedDeltaParts.delete(part.messageID)
  762. const removed = load?.removedParts.get(part.messageID)
  763. removed?.delete(part.id)
  764. if (removed?.size === 0) load?.removedParts.delete(part.messageID)
  765. const pending = pendingParts.get(part.sessionID)?.get(part.messageID)
  766. pending?.delete(part.id)
  767. if (pending?.size === 0) pendingParts.get(part.sessionID)?.delete(part.messageID)
  768. if (pendingParts.get(part.sessionID)?.size === 0) pendingParts.delete(part.sessionID)
  769. const optimistic = load?.optimisticParts.get(part.messageID)
  770. optimistic?.delete(part.id)
  771. if (optimistic?.size === 0) load?.optimisticParts.delete(part.messageID)
  772. deltaBases.delete(part.id)
  773. trackPartChange(part.sessionID, part.messageID, part.id)
  774. confirmOptimisticPart(part.sessionID, part.messageID, part)
  775. setData(
  776. "part_text_accum_delta",
  777. produce((draft) => void delete draft[part.id]),
  778. )
  779. const parts = data.part[part.messageID]
  780. if (!parts) {
  781. setData("part", part.messageID, [part])
  782. return
  783. }
  784. const result = Binary.search(parts, part.id, (item) => item.id)
  785. if (result.found) setData("part", part.messageID, result.index, reconcile(part))
  786. if (!result.found)
  787. setData("part", part.messageID, (value = []) => {
  788. const next = value.slice()
  789. next.splice(result.index, 0, part)
  790. return next
  791. })
  792. return
  793. }
  794. case "message.part.removed": {
  795. const props = event.properties as { sessionID: string; messageID: string; partID: string }
  796. // Part removal is event-only on the server, so its tombstone lasts until a later update or eviction.
  797. const pending = pendingParts.get(props.sessionID) ?? new Map<string, Set<string>>()
  798. const parts = pending.get(props.messageID) ?? new Set<string>()
  799. parts.add(props.partID)
  800. pending.set(props.messageID, parts)
  801. pendingParts.set(props.sessionID, pending)
  802. const deltas = messageLoads.get(props.sessionID)?.deltaParts.get(props.messageID)
  803. deltas?.delete(props.partID)
  804. if (deltas?.size === 0) messageLoads.get(props.sessionID)?.deltaParts.delete(props.messageID)
  805. const load = messageLoads.get(props.sessionID)
  806. const carried = load?.carriedDeltaParts.get(props.messageID)
  807. carried?.delete(props.partID)
  808. if (carried?.size === 0) load?.carriedDeltaParts.delete(props.messageID)
  809. if (load) {
  810. const parts = load.removedParts.get(props.messageID) ?? new Set<string>()
  811. parts.add(props.partID)
  812. load.removedParts.set(props.messageID, parts)
  813. const optimistic = load.optimisticParts.get(props.messageID)
  814. optimistic?.delete(props.partID)
  815. if (optimistic?.size === 0) load.optimisticParts.delete(props.messageID)
  816. }
  817. trackPartChange(props.sessionID, props.messageID, props.partID)
  818. clearOptimisticPart(props.sessionID, props.messageID, props.partID)
  819. setData(
  820. produce((draft) => {
  821. delete draft.part_text_accum_delta[props.partID]
  822. deltaBases.delete(props.partID)
  823. const parts = draft.part[props.messageID]
  824. if (!parts) return
  825. const result = Binary.search(parts, props.partID, (part) => part.id)
  826. if (result.found) parts.splice(result.index, 1)
  827. if (parts.length === 0) delete draft.part[props.messageID]
  828. }),
  829. )
  830. return
  831. }
  832. case "message.part.delta": {
  833. const props = event.properties as {
  834. sessionID: string
  835. messageID: string
  836. partID: string
  837. field: string
  838. delta: string
  839. }
  840. const parts = data.part[props.messageID]
  841. if (!parts) return
  842. const result = Binary.search(parts, props.partID, (part) => part.id)
  843. if (!result.found) return
  844. trackPartChange(props.sessionID, props.messageID, props.partID)
  845. const load = messageLoads.get(props.sessionID)
  846. if (load) {
  847. const parts = load.deltaParts.get(props.messageID) ?? new Set<string>()
  848. parts.add(props.partID)
  849. load.deltaParts.set(props.messageID, parts)
  850. const carried = load.carriedDeltaParts.get(props.messageID)
  851. carried?.delete(props.partID)
  852. if (carried?.size === 0) load.carriedDeltaParts.delete(props.messageID)
  853. }
  854. const field = props.field as keyof (typeof parts)[number]
  855. const current = parts[result.index]?.[field]
  856. if (!deltaBases.has(props.partID) && typeof current === "string")
  857. deltaBases.set(props.partID, { base: current, sessionID: props.sessionID })
  858. setData(
  859. "part_text_accum_delta",
  860. props.partID,
  861. (value) => (value ?? (typeof current === "string" ? current : "")) + props.delta,
  862. )
  863. setData(
  864. "part",
  865. props.messageID,
  866. produce((draft) => {
  867. if (!draft) return
  868. const part = draft[result.index]
  869. const field = props.field as keyof typeof part
  870. ;(part[field] as string) = ((part[field] as string | undefined) ?? "") + props.delta
  871. }),
  872. )
  873. return
  874. }
  875. case "permission.asked": {
  876. const permission = event.properties as PermissionRequest
  877. const permissions = data.permission[permission.sessionID]
  878. if (!permissions) {
  879. setData("permission", permission.sessionID, [permission])
  880. return
  881. }
  882. const result = Binary.search(permissions, permission.id, (item) => item.id)
  883. if (result.found) setData("permission", permission.sessionID, result.index, reconcile(permission))
  884. if (!result.found)
  885. setData(
  886. "permission",
  887. permission.sessionID,
  888. produce((draft) => void draft.splice(result.index, 0, permission)),
  889. )
  890. return
  891. }
  892. case "permission.replied": {
  893. const props = event.properties as { sessionID: string; requestID: string }
  894. setData(
  895. "permission",
  896. props.sessionID,
  897. produce((draft) => {
  898. if (!draft) return
  899. const result = Binary.search(draft, props.requestID, (item) => item.id)
  900. if (result.found) draft.splice(result.index, 1)
  901. }),
  902. )
  903. return
  904. }
  905. case "question.asked": {
  906. const question = event.properties as QuestionRequest
  907. const questions = data.question[question.sessionID]
  908. if (!questions) {
  909. setData("question", question.sessionID, [question])
  910. return
  911. }
  912. const result = Binary.search(questions, question.id, (item) => item.id)
  913. if (result.found) setData("question", question.sessionID, result.index, reconcile(question))
  914. if (!result.found)
  915. setData(
  916. "question",
  917. question.sessionID,
  918. produce((draft) => void draft.splice(result.index, 0, question)),
  919. )
  920. return
  921. }
  922. case "question.replied":
  923. case "question.rejected": {
  924. const props = event.properties as { sessionID: string; requestID: string }
  925. setData(
  926. "question",
  927. props.sessionID,
  928. produce((draft) => {
  929. if (!draft) return
  930. const result = Binary.search(draft, props.requestID, (item) => item.id)
  931. if (result.found) draft.splice(result.index, 1)
  932. }),
  933. )
  934. }
  935. }
  936. }
  937. return {
  938. data,
  939. set: setData,
  940. get: (sessionID: string) => data.info[sessionID],
  941. peek: (sessionID: string) => data.info[sessionID],
  942. remember,
  943. resolve,
  944. lineage: {
  945. peek: peekLineage,
  946. async resolve(sessionID: string) {
  947. const session = await resolve(sessionID)
  948. return { session, root: await rootSession(session, resolve) }
  949. },
  950. },
  951. sync,
  952. prefetch,
  953. shouldPrefetch(sessionID: string, limit: number) {
  954. if (data.message[sessionID] === undefined) return true
  955. if (Date.now() - (meta.at[sessionID] ?? 0) > 15_000) return true
  956. if (meta.complete[sessionID]) return false
  957. return (meta.limit[sessionID] ?? 0) <= limit
  958. },
  959. fresh(sessionID: string, ttl: number) {
  960. return Date.now() - (meta.at[sessionID] ?? 0) <= ttl
  961. },
  962. optimistic: {
  963. add(input: { sessionID: string; message: Message; parts: Part[] }) {
  964. const parts = input.parts
  965. .filter((part) => !!part?.id && !SKIP_PARTS.has(part.type))
  966. .sort((a, b) => cmp(a.id, b.id))
  967. const load = messageLoads.get(input.sessionID)
  968. if (load?.clearedMessageParts.has(input.message.id)) {
  969. const touched = load.touchedParts.get(input.message.id) ?? new Set<string>()
  970. parts.forEach((part) => touched.add(part.id))
  971. load.touchedParts.set(input.message.id, touched)
  972. }
  973. if (load) {
  974. load.removedMessages.delete(input.message.id)
  975. load.optimisticParts.set(input.message.id, new Set(parts.map((part) => part.id)))
  976. }
  977. const items = optimistic.get(input.sessionID)
  978. const removedMessagesForSession = removedMessages.get(input.sessionID)
  979. removedMessagesForSession?.delete(input.message.id)
  980. if (removedMessagesForSession?.size === 0) removedMessages.delete(input.sessionID)
  981. if (items) items.set(input.message.id, { ...input, parts, confirmedParts: [] })
  982. if (!items)
  983. optimistic.set(input.sessionID, new Map([[input.message.id, { ...input, parts, confirmedParts: [] }]]))
  984. setData("message", input.sessionID, (messages = []) => merge(messages, [input.message]))
  985. setData(
  986. "part_text_accum_delta",
  987. produce((draft) => {
  988. for (const part of [...(data.part[input.message.id] ?? []), ...parts]) {
  989. delete draft[part.id]
  990. deltaBases.delete(part.id)
  991. }
  992. }),
  993. )
  994. setData("part", input.message.id, parts)
  995. },
  996. remove(input: { sessionID: string; messageID: string }) {
  997. const item = optimistic.get(input.sessionID)?.get(input.messageID)
  998. if (!item) return
  999. messageLoads.get(input.sessionID)?.optimisticParts.delete(input.messageID)
  1000. clearOptimistic(input.sessionID, input.messageID)
  1001. if (item.confirmedMessage) {
  1002. const partIDs = new Set(item.parts.map((part) => part.id))
  1003. setData(
  1004. produce((draft) => {
  1005. for (const part of item.parts) {
  1006. delete draft.part_text_accum_delta[part.id]
  1007. deltaBases.delete(part.id)
  1008. }
  1009. const parts = draft.part[input.messageID]
  1010. if (!parts) return
  1011. draft.part[input.messageID] = parts.filter((part) => !partIDs.has(part.id))
  1012. if (draft.part[input.messageID]?.length === 0) delete draft.part[input.messageID]
  1013. }),
  1014. )
  1015. return
  1016. }
  1017. setData("message", input.sessionID, (messages) => messages?.filter((message) => message.id !== input.messageID))
  1018. setData(produce((draft) => deleteMessageParts(draft, input.messageID)))
  1019. },
  1020. },
  1021. diff(sessionID: string, options?: { force?: boolean }) {
  1022. touch(sessionID)
  1023. if (data.session_diff[sessionID] !== undefined && !options?.force) return Promise.resolve()
  1024. return runInflight(inflightDiff, sessionID, () => {
  1025. const active = generation(sessionID)
  1026. return retry(() => client.session.diff({ sessionID })).then((result) => {
  1027. if (generations.get(sessionID) !== active) return
  1028. setData("session_diff", sessionID, reconcile(cleanDiffs(result.data), { key: "file" }))
  1029. })
  1030. })
  1031. },
  1032. todo(sessionID: string, options?: { force?: boolean }) {
  1033. touch(sessionID)
  1034. if (data.todo[sessionID] !== undefined && !options?.force) return Promise.resolve()
  1035. return runInflight(inflightTodo, sessionID, () => {
  1036. const active = generation(sessionID)
  1037. return retry(() => client.session.todo({ sessionID })).then((result) => {
  1038. if (generations.get(sessionID) !== active) return
  1039. setData("todo", sessionID, reconcile(result.data ?? [], { key: "id" }))
  1040. })
  1041. })
  1042. },
  1043. history: {
  1044. more: (sessionID: string) =>
  1045. data.message[sessionID] !== undefined &&
  1046. meta.limit[sessionID] !== undefined &&
  1047. !meta.complete[sessionID] &&
  1048. !!meta.cursor[sessionID],
  1049. loading: (sessionID: string) => meta.loading[sessionID] ?? false,
  1050. async loadMore(sessionID: string, count = historyMessagePageSize) {
  1051. touch(sessionID)
  1052. if (meta.loading[sessionID] || meta.complete[sessionID] || !meta.cursor[sessionID]) return
  1053. await loadMessages(sessionID, count, meta.cursor[sessionID], "prepend")
  1054. },
  1055. },
  1056. evict(sessionID: string) {
  1057. if (protectedSessions().has(sessionID)) return
  1058. seen.delete(sessionID)
  1059. evict([sessionID])
  1060. },
  1061. pin(sessionID: string) {
  1062. pinned.set(sessionID, (pinned.get(sessionID) ?? 0) + 1)
  1063. touch(sessionID)
  1064. },
  1065. unpin(sessionID: string) {
  1066. const count = pinned.get(sessionID)
  1067. if (!count || count === 1) pinned.delete(sessionID)
  1068. if (count && count > 1) pinned.set(sessionID, count - 1)
  1069. },
  1070. apply,
  1071. }
  1072. }
  1073. export type ServerSession = ReturnType<typeof createServerSession>