server-session-v2-reducer.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509
  1. import type { OpenCodeEvent, SessionMessageInfo, SessionPendingMessage } from "@opencode-ai/client/promise"
  2. type Assistant = Extract<SessionMessageInfo, { type: "assistant" }>
  3. type Compaction = Extract<SessionMessageInfo, { type: "compaction" }>
  4. type Shell = Extract<SessionMessageInfo, { type: "shell" }>
  5. export type V2SessionReduction = {
  6. sessionID: string
  7. messages: SessionMessageInfo[]
  8. touched: string[]
  9. missing?: string
  10. }
  11. export function createV2SessionReducer() {
  12. const pending = new Map<string, SessionPendingMessage>()
  13. const reduce = (source: readonly SessionMessageInfo[], event: OpenCodeEvent): V2SessionReduction | undefined => {
  14. if (!("data" in event) || !("sessionID" in event.data) || typeof event.data.sessionID !== "string") return
  15. const sessionID = event.data.sessionID
  16. const result = (messages: SessionMessageInfo[], touched: string[] = []): V2SessionReduction => ({
  17. sessionID,
  18. messages,
  19. touched,
  20. })
  21. const append = (message: SessionMessageInfo) =>
  22. result(source.some((item) => item.id === message.id) ? [...source] : [...source, message], [message.id])
  23. switch (event.type) {
  24. case "session.input.admitted":
  25. pending.set(key(sessionID, event.data.inputID), event.data.input)
  26. return result([...source])
  27. case "session.input.promoted": {
  28. const input = pending.get(key(sessionID, event.data.inputID))
  29. pending.delete(key(sessionID, event.data.inputID))
  30. if (!input) return { ...result([...source]), missing: event.data.inputID }
  31. if (input.type === "user")
  32. return append({
  33. id: event.data.inputID,
  34. type: "user",
  35. metadata: input.data.metadata,
  36. text: input.data.text,
  37. files: input.data.files,
  38. agents: input.data.agents,
  39. time: { created: event.created },
  40. })
  41. return append({
  42. id: event.data.inputID,
  43. type: "synthetic",
  44. metadata: input.data.metadata,
  45. text: input.data.text,
  46. description: input.data.description,
  47. time: { created: event.created },
  48. })
  49. }
  50. case "session.agent.selected":
  51. return append({
  52. id: messageID(event.id),
  53. type: "agent-switched",
  54. metadata: event.metadata,
  55. agent: event.data.agent,
  56. time: { created: event.created },
  57. })
  58. case "session.model.selected":
  59. return append({
  60. id: messageID(event.id),
  61. type: "model-switched",
  62. metadata: event.metadata,
  63. model: event.data.model,
  64. previous: source.findLast(
  65. (item): item is Extract<SessionMessageInfo, { type: "model-switched" | "assistant" }> =>
  66. item.type === "model-switched" || item.type === "assistant",
  67. )?.model,
  68. time: { created: event.created },
  69. })
  70. case "session.synthetic":
  71. return append({
  72. id: messageID(event.id),
  73. type: "synthetic",
  74. metadata: event.data.metadata,
  75. text: event.data.text,
  76. description: event.data.description,
  77. time: { created: event.created },
  78. })
  79. case "session.skill.activated":
  80. return append({
  81. id: messageID(event.id),
  82. type: "skill",
  83. metadata: event.metadata,
  84. skill: event.data.id,
  85. name: event.data.name,
  86. text: event.data.text,
  87. time: { created: event.created },
  88. })
  89. case "session.shell.started":
  90. return append({
  91. id: messageID(event.id),
  92. type: "shell",
  93. metadata: event.metadata,
  94. shellID: event.data.shell.id,
  95. command: event.data.shell.command,
  96. status: event.data.shell.status,
  97. exit: event.data.shell.exit,
  98. time: { created: event.created },
  99. })
  100. case "session.shell.ended":
  101. return updateMessage<Shell>(
  102. source,
  103. (item): item is Shell => item.type === "shell" && item.shellID === event.data.shell.id,
  104. (item) => ({
  105. ...item,
  106. status: event.data.shell.status,
  107. exit: event.data.shell.exit,
  108. output: event.data.output,
  109. time: { ...item.time, completed: event.created },
  110. }),
  111. sessionID,
  112. )
  113. case "session.step.started": {
  114. const current = source.findLast((item): item is Assistant => item.type === "assistant" && !item.time.completed)
  115. const completed =
  116. current && current.id !== event.data.assistantMessageID
  117. ? update(source, current.id, (item) =>
  118. item.type === "assistant"
  119. ? { ...item, retry: undefined, time: { ...item.time, completed: event.created } }
  120. : item,
  121. )
  122. : [...source]
  123. const existing = completed.find((item) => item.id === event.data.assistantMessageID)
  124. if (existing?.type === "assistant")
  125. return result(
  126. update(completed, existing.id, (item) =>
  127. item.type === "assistant"
  128. ? {
  129. ...item,
  130. agent: event.data.agent,
  131. model: event.data.model,
  132. retry: undefined,
  133. error: undefined,
  134. finish: undefined,
  135. snapshot: event.data.snapshot ? { ...item.snapshot, start: event.data.snapshot } : item.snapshot,
  136. time: { ...item.time, completed: undefined },
  137. }
  138. : item,
  139. ),
  140. current && current.id !== existing.id ? [current.id, existing.id] : [existing.id],
  141. )
  142. return result(
  143. [
  144. ...completed,
  145. {
  146. id: event.data.assistantMessageID,
  147. type: "assistant",
  148. metadata: event.metadata,
  149. agent: event.data.agent,
  150. model: event.data.model,
  151. content: [],
  152. snapshot: event.data.snapshot ? { start: event.data.snapshot } : undefined,
  153. time: { created: event.created },
  154. },
  155. ],
  156. current ? [current.id, event.data.assistantMessageID] : [event.data.assistantMessageID],
  157. )
  158. }
  159. case "session.step.ended":
  160. return updateAssistant(source, event.data.assistantMessageID, sessionID, (item) => ({
  161. ...item,
  162. finish: event.data.finish,
  163. cost: event.data.cost,
  164. tokens: event.data.tokens,
  165. snapshot:
  166. event.data.snapshot || event.data.files
  167. ? { ...item.snapshot, end: event.data.snapshot, files: event.data.files }
  168. : item.snapshot,
  169. time: { ...item.time, completed: event.created },
  170. }))
  171. case "session.step.failed":
  172. return updateAssistant(source, event.data.assistantMessageID, sessionID, (item) => ({
  173. ...item,
  174. finish: "error",
  175. error: event.data.error,
  176. retry: undefined,
  177. cost: event.data.cost ?? item.cost,
  178. tokens: event.data.tokens ?? item.tokens,
  179. snapshot:
  180. event.data.snapshot || event.data.files
  181. ? { ...item.snapshot, end: event.data.snapshot, files: event.data.files }
  182. : item.snapshot,
  183. time: { ...item.time, completed: event.created },
  184. }))
  185. case "session.text.started":
  186. return updateAssistant(source, event.data.assistantMessageID, sessionID, (item) => ({
  187. ...item,
  188. content: insertOrdinal(item.content, "text", event.data.ordinal, { type: "text", text: "" }),
  189. }))
  190. case "session.text.delta":
  191. return updateContent(source, event.data.assistantMessageID, sessionID, "text", event.data.ordinal, (item) => ({
  192. ...item,
  193. text: item.text + event.data.delta,
  194. }))
  195. case "session.text.ended":
  196. return updateContent(source, event.data.assistantMessageID, sessionID, "text", event.data.ordinal, (item) => ({
  197. ...item,
  198. text: event.data.text,
  199. }))
  200. case "session.reasoning.started":
  201. return updateAssistant(source, event.data.assistantMessageID, sessionID, (item) => ({
  202. ...item,
  203. content: insertOrdinal(item.content, "reasoning", event.data.ordinal, {
  204. type: "reasoning",
  205. text: "",
  206. state: event.data.state,
  207. time: { created: event.created },
  208. }),
  209. }))
  210. case "session.reasoning.delta":
  211. return updateContent(
  212. source,
  213. event.data.assistantMessageID,
  214. sessionID,
  215. "reasoning",
  216. event.data.ordinal,
  217. (item) => ({
  218. ...item,
  219. text: item.text + event.data.delta,
  220. }),
  221. )
  222. case "session.reasoning.ended":
  223. return updateContent(
  224. source,
  225. event.data.assistantMessageID,
  226. sessionID,
  227. "reasoning",
  228. event.data.ordinal,
  229. (item) => ({
  230. ...item,
  231. text: event.data.text,
  232. state: event.data.state ?? item.state,
  233. time: { created: item.time?.created ?? event.created, completed: event.created },
  234. }),
  235. )
  236. case "session.tool.input.started":
  237. return updateAssistant(source, event.data.assistantMessageID, sessionID, (item) => ({
  238. ...item,
  239. content: item.content.some((content) => content.type === "tool" && content.id === event.data.callID)
  240. ? item.content
  241. : [
  242. ...item.content,
  243. {
  244. type: "tool",
  245. id: event.data.callID,
  246. name: event.data.name,
  247. state: { status: "streaming", input: "" },
  248. time: { created: event.created },
  249. },
  250. ],
  251. }))
  252. case "session.tool.input.delta":
  253. return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) =>
  254. tool.state.status === "streaming"
  255. ? { ...tool, state: { ...tool.state, input: tool.state.input + event.data.delta } }
  256. : tool,
  257. )
  258. case "session.tool.input.ended":
  259. return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) =>
  260. tool.state.status === "streaming" ? { ...tool, state: { ...tool.state, input: event.data.text } } : tool,
  261. )
  262. case "session.tool.called":
  263. return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) => ({
  264. ...tool,
  265. executed: event.data.executed,
  266. providerState: event.data.state,
  267. // structured: {}, content: []
  268. state: { status: "running", input: event.data.input, metadata: {} },
  269. time: { ...tool.time, ran: event.created },
  270. }))
  271. case "session.tool.progress":
  272. return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) =>
  273. tool.state.status === "running"
  274. ? {
  275. ...tool,
  276. // state: { ...tool.state, structured: event.data.structured, content: event.data.content },
  277. state: { ...tool.state, metadata: event.data.metadata },
  278. }
  279. : tool,
  280. )
  281. case "session.tool.success":
  282. return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) => {
  283. if (tool.state.status !== "running") return tool
  284. return {
  285. ...tool,
  286. executed: event.data.executed || tool.executed === true,
  287. providerResultState: event.data.resultState,
  288. state: {
  289. status: "completed",
  290. input: tool.state.input,
  291. // structured: event.data.structured,
  292. metadata: event.data.metadata,
  293. content: event.data.content,
  294. // result: event.data.result,
  295. },
  296. time: { ...tool.time, completed: event.created },
  297. }
  298. })
  299. case "session.tool.failed":
  300. return updateTool(source, event.data.assistantMessageID, event.data.callID, sessionID, (tool) => {
  301. if (tool.state.status !== "streaming" && tool.state.status !== "running") return tool
  302. return {
  303. ...tool,
  304. executed: event.data.executed || tool.executed === true,
  305. providerResultState: event.data.resultState,
  306. state: {
  307. status: "error",
  308. input: typeof tool.state.input === "string" ? {} : tool.state.input,
  309. // structured: tool.state.status === "running" ? tool.state.structured : {},
  310. metadata: event.data.metadata ?? (tool.state.status === "running" ? tool.state.metadata : {}),
  311. content: event.data.content,
  312. error: event.data.error,
  313. // result: event.data.result,
  314. },
  315. time: { ...tool.time, completed: event.created },
  316. }
  317. })
  318. case "session.retry.scheduled":
  319. return updateAssistant(source, event.data.assistantMessageID, sessionID, (item) => ({
  320. ...item,
  321. retry: { attempt: event.data.attempt, at: event.data.at, error: event.data.error },
  322. }))
  323. case "session.execution.succeeded":
  324. case "session.execution.failed":
  325. case "session.execution.interrupted": {
  326. const current = source.findLast((item): item is Assistant => item.type === "assistant" && !item.time.completed)
  327. if (!current?.retry) return result([...source])
  328. return updateAssistant(source, current.id, sessionID, (item) => ({ ...item, retry: undefined }))
  329. }
  330. case "session.compaction.started":
  331. return append({
  332. id: event.data.inputID ?? messageID(event.id),
  333. type: "compaction",
  334. status: "running",
  335. metadata: event.metadata,
  336. reason: event.data.reason,
  337. summary: "",
  338. recent: event.data.recent,
  339. time: { created: event.created },
  340. })
  341. case "session.compaction.delta":
  342. return updateMessage<Extract<Compaction, { status: "running" }>>(
  343. source,
  344. (item): item is Extract<Compaction, { status: "running" }> =>
  345. item.type === "compaction" && item.status === "running",
  346. (item) => ({
  347. ...item,
  348. summary: item.summary + event.data.text,
  349. }),
  350. sessionID,
  351. )
  352. case "session.compaction.ended": {
  353. const current = source.findLast(
  354. (item): item is Extract<Compaction, { status: "running" }> =>
  355. item.type === "compaction" && item.status === "running",
  356. )
  357. if (!current)
  358. return append({
  359. id: messageID(event.id),
  360. type: "compaction",
  361. status: "completed",
  362. metadata: event.metadata,
  363. reason: event.data.reason,
  364. summary: event.data.text,
  365. recent: event.data.recent,
  366. time: { created: event.created },
  367. })
  368. return result(
  369. update(source, current.id, () => ({
  370. ...current,
  371. status: "completed",
  372. reason: event.data.reason,
  373. summary: event.data.text,
  374. recent: event.data.recent,
  375. })),
  376. [current.id],
  377. )
  378. }
  379. case "session.compaction.failed": {
  380. const current = source.findLast(
  381. (item): item is Extract<Compaction, { status: "running" }> =>
  382. item.type === "compaction" && item.status === "running",
  383. )
  384. const failed: Extract<Compaction, { status: "failed" }> = {
  385. id: current?.id ?? event.data.inputID ?? messageID(event.id),
  386. type: "compaction",
  387. status: "failed",
  388. metadata: current?.metadata ?? event.metadata,
  389. reason: event.data.reason,
  390. error: event.data.error,
  391. time: current?.time ?? { created: event.created },
  392. }
  393. if (!current) return append(failed)
  394. return result(
  395. update(source, current.id, () => failed),
  396. [failed.id],
  397. )
  398. }
  399. default:
  400. return
  401. }
  402. }
  403. return {
  404. reduce,
  405. clear(sessionID: string) {
  406. for (const id of pending.keys()) {
  407. if (id.startsWith(`${sessionID}:`)) pending.delete(id)
  408. }
  409. },
  410. }
  411. }
  412. function key(sessionID: string, inputID: string) {
  413. return `${sessionID}:${inputID}`
  414. }
  415. function messageID(eventID: string) {
  416. return eventID.replace(/^evt_/, "msg_")
  417. }
  418. function update(
  419. source: readonly SessionMessageInfo[],
  420. id: string,
  421. apply: (item: SessionMessageInfo) => SessionMessageInfo,
  422. ) {
  423. return source.map((item) => (item.id === id ? apply(item) : item))
  424. }
  425. function updateMessage<T extends SessionMessageInfo>(
  426. source: readonly SessionMessageInfo[],
  427. matches: (item: SessionMessageInfo) => item is T,
  428. apply: (item: T) => T,
  429. sessionID: string,
  430. ): V2SessionReduction {
  431. const current = source.findLast(matches)
  432. if (!current) return { sessionID, messages: [...source], touched: [] }
  433. return {
  434. sessionID,
  435. messages: update(source, current.id, (item) => (matches(item) ? apply(item) : item)),
  436. touched: [current.id],
  437. }
  438. }
  439. function updateAssistant(
  440. source: readonly SessionMessageInfo[],
  441. id: string,
  442. sessionID: string,
  443. apply: (item: Assistant) => Assistant,
  444. ): V2SessionReduction {
  445. return {
  446. sessionID,
  447. messages: update(source, id, (item) => (item.type === "assistant" ? apply(item) : item)),
  448. touched: source.some((item) => item.id === id && item.type === "assistant") ? [id] : [],
  449. }
  450. }
  451. function updateContent<T extends "text" | "reasoning">(
  452. source: readonly SessionMessageInfo[],
  453. messageID: string,
  454. sessionID: string,
  455. type: T,
  456. ordinal: number,
  457. apply: (
  458. item: Extract<Assistant["content"][number], { type: T }>,
  459. ) => Extract<Assistant["content"][number], { type: T }>,
  460. ) {
  461. return updateAssistant(source, messageID, sessionID, (assistant) => {
  462. let index = -1
  463. return {
  464. ...assistant,
  465. content: assistant.content.map((item) => {
  466. if (item.type !== type || ++index !== ordinal) return item
  467. return apply(item as Extract<Assistant["content"][number], { type: T }>)
  468. }),
  469. }
  470. })
  471. }
  472. function updateTool(
  473. source: readonly SessionMessageInfo[],
  474. messageID: string,
  475. callID: string,
  476. sessionID: string,
  477. apply: (
  478. item: Extract<Assistant["content"][number], { type: "tool" }>,
  479. ) => Extract<Assistant["content"][number], { type: "tool" }>,
  480. ) {
  481. return updateAssistant(source, messageID, sessionID, (assistant) => ({
  482. ...assistant,
  483. content: assistant.content.map((item) => (item.type === "tool" && item.id === callID ? apply(item) : item)),
  484. }))
  485. }
  486. function insertOrdinal<T extends Assistant["content"][number]["type"]>(
  487. source: Assistant["content"],
  488. type: T,
  489. ordinal: number,
  490. item: Extract<Assistant["content"][number], { type: T }>,
  491. ) {
  492. const matches = source.filter((content) => content.type === type)
  493. if (matches[ordinal]) return source
  494. return [...source, item]
  495. }