session.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427
  1. import { Agent } from "@/agent/agent"
  2. import { Bus } from "@/bus"
  3. import { Command } from "@/command"
  4. import { Permission } from "@/permission"
  5. import { PermissionID } from "@/permission/schema"
  6. import { SessionShare } from "@/share/session"
  7. import { Session } from "@/session/session"
  8. import { SessionCompaction } from "@/session/compaction"
  9. import { MessageV2 } from "@/session/message-v2"
  10. import { SessionPrompt } from "@/session/prompt"
  11. import { SessionRevert } from "@/session/revert"
  12. import { SessionRunState } from "@/session/run-state"
  13. import { SessionStatus } from "@/session/status"
  14. import { SessionSummary } from "@/session/summary"
  15. import { Todo } from "@/session/todo"
  16. import { MessageID, PartID, SessionID } from "@/session/schema"
  17. import { NamedError } from "@opencode-ai/core/util/error"
  18. import { Cause, Effect, Option, Schema, Scope } from "effect"
  19. import * as Stream from "effect/Stream"
  20. import { HttpServerRequest, HttpServerResponse } from "effect/unstable/http"
  21. import { HttpApiBuilder, HttpApiError, HttpApiSchema } from "effect/unstable/httpapi"
  22. import { InstanceHttpApi } from "../api"
  23. import {
  24. CommandPayload,
  25. DiffQuery,
  26. ForkPayload,
  27. InitPayload,
  28. ListQuery,
  29. MessagesQuery,
  30. PermissionResponsePayload,
  31. PromptPayload,
  32. RevertPayload,
  33. ShellPayload,
  34. SummarizePayload,
  35. UpdatePayload,
  36. } from "../groups/session"
  37. import * as SessionError from "./session-errors"
  38. export const sessionHandlers = HttpApiBuilder.group(InstanceHttpApi, "session", (handlers) =>
  39. Effect.gen(function* () {
  40. const session = yield* Session.Service
  41. const shareSvc = yield* SessionShare.Service
  42. const promptSvc = yield* SessionPrompt.Service
  43. const revertSvc = yield* SessionRevert.Service
  44. const compactSvc = yield* SessionCompaction.Service
  45. const runState = yield* SessionRunState.Service
  46. const agentSvc = yield* Agent.Service
  47. const permissionSvc = yield* Permission.Service
  48. const statusSvc = yield* SessionStatus.Service
  49. const todoSvc = yield* Todo.Service
  50. const summary = yield* SessionSummary.Service
  51. const bus = yield* Bus.Service
  52. const scope = yield* Scope.Scope
  53. const mapBusy = <A, E, R>(effect: Effect.Effect<A, E, R>): Effect.Effect<A, E | HttpApiError.BadRequest, R> =>
  54. effect.pipe(
  55. Effect.catchCause((cause): Effect.Effect<never, E | HttpApiError.BadRequest> => {
  56. if (Cause.squash(cause) instanceof Session.BusyError) return Effect.fail(new HttpApiError.BadRequest({}))
  57. return Effect.failCause(cause)
  58. }),
  59. )
  60. const list = Effect.fn("SessionHttpApi.list")(function* (ctx: { query: typeof ListQuery.Type }) {
  61. return yield* session.list({
  62. directory: ctx.query.scope === "project" ? undefined : ctx.query.directory,
  63. scope: ctx.query.scope,
  64. path: ctx.query.path,
  65. roots: ctx.query.roots,
  66. start: ctx.query.start,
  67. search: ctx.query.search,
  68. limit: ctx.query.limit,
  69. })
  70. })
  71. const status = Effect.fn("SessionHttpApi.status")(function* () {
  72. return Object.fromEntries(yield* statusSvc.list())
  73. })
  74. const requireSession = Effect.fn("SessionHttpApi.requireSession")(function* (sessionID: SessionID) {
  75. return yield* SessionError.mapStorageNotFound(session.get(sessionID))
  76. })
  77. const get = Effect.fn("SessionHttpApi.get")(function* (ctx: { params: { sessionID: SessionID } }) {
  78. return yield* requireSession(ctx.params.sessionID)
  79. })
  80. const children = Effect.fn("SessionHttpApi.children")(function* (ctx: { params: { sessionID: SessionID } }) {
  81. yield* requireSession(ctx.params.sessionID)
  82. return yield* session.children(ctx.params.sessionID)
  83. })
  84. const todo = Effect.fn("SessionHttpApi.todo")(function* (ctx: { params: { sessionID: SessionID } }) {
  85. yield* requireSession(ctx.params.sessionID)
  86. return yield* todoSvc.get(ctx.params.sessionID)
  87. })
  88. const diff = Effect.fn("SessionHttpApi.diff")(function* (ctx: {
  89. params: { sessionID: SessionID }
  90. query: typeof DiffQuery.Type
  91. }) {
  92. return yield* summary.diff({ sessionID: ctx.params.sessionID, messageID: ctx.query.messageID })
  93. })
  94. const messages = Effect.fn("SessionHttpApi.messages")(function* (ctx: {
  95. params: { sessionID: SessionID }
  96. query: typeof MessagesQuery.Type
  97. }) {
  98. if (ctx.query.before && ctx.query.limit === undefined) return yield* new HttpApiError.BadRequest({})
  99. if (ctx.query.before) {
  100. const before = ctx.query.before
  101. yield* Effect.try({
  102. try: () => MessageV2.cursor.decode(before),
  103. catch: () => new HttpApiError.BadRequest({}),
  104. })
  105. }
  106. yield* requireSession(ctx.params.sessionID)
  107. if (ctx.query.limit === undefined || ctx.query.limit === 0) {
  108. return yield* SessionError.mapStorageNotFound(session.messages({ sessionID: ctx.params.sessionID }))
  109. }
  110. const page = yield* SessionError.mapStorageNotFound(
  111. MessageV2.page({
  112. sessionID: ctx.params.sessionID,
  113. limit: ctx.query.limit,
  114. before: ctx.query.before,
  115. }),
  116. )
  117. if (!page.cursor) return page.items
  118. const request = yield* HttpServerRequest.HttpServerRequest
  119. // toURL() honors the Host + x-forwarded-proto headers, so the Link
  120. // header echoes the real origin instead of a hard-coded localhost.
  121. const url = Option.getOrElse(HttpServerRequest.toURL(request), () => new URL(request.url, "http://localhost"))
  122. url.searchParams.set("limit", ctx.query.limit.toString())
  123. url.searchParams.set("before", page.cursor)
  124. return HttpServerResponse.jsonUnsafe(page.items, {
  125. headers: {
  126. "Access-Control-Expose-Headers": "Link, X-Next-Cursor",
  127. Link: `<${url.toString()}>; rel="next"`,
  128. "X-Next-Cursor": page.cursor,
  129. },
  130. })
  131. })
  132. const message = Effect.fn("SessionHttpApi.message")(function* (ctx: {
  133. params: { sessionID: SessionID; messageID: MessageID }
  134. }) {
  135. return yield* SessionError.mapStorageNotFound(
  136. MessageV2.get({ sessionID: ctx.params.sessionID, messageID: ctx.params.messageID }),
  137. )
  138. })
  139. const create = Effect.fn("SessionHttpApi.create")(function* (ctx: { payload?: Session.CreateInput }) {
  140. return yield* shareSvc.create(ctx.payload)
  141. })
  142. const createRaw = Effect.fn("SessionHttpApi.createRaw")(function* (ctx: {
  143. request: HttpServerRequest.HttpServerRequest
  144. }) {
  145. const body = yield* Effect.orDie(ctx.request.text)
  146. if (body.trim().length === 0) return yield* create({})
  147. const json = yield* Effect.try({
  148. try: () => JSON.parse(body) as unknown,
  149. catch: () => new HttpApiError.BadRequest({}),
  150. })
  151. const payload = yield* Schema.decodeUnknownEffect(Session.CreateInput)(json).pipe(
  152. Effect.mapError(() => new HttpApiError.BadRequest({})),
  153. )
  154. return yield* create({ payload })
  155. })
  156. const remove = Effect.fn("SessionHttpApi.remove")(function* (ctx: { params: { sessionID: SessionID } }) {
  157. yield* SessionError.mapStorageNotFound(session.remove(ctx.params.sessionID))
  158. return true
  159. })
  160. const update = Effect.fn("SessionHttpApi.update")(function* (ctx: {
  161. params: { sessionID: SessionID }
  162. payload: typeof UpdatePayload.Type
  163. }) {
  164. const current = yield* requireSession(ctx.params.sessionID)
  165. if (ctx.payload.title !== undefined) {
  166. yield* session.setTitle({ sessionID: ctx.params.sessionID, title: ctx.payload.title })
  167. }
  168. if (ctx.payload.permission !== undefined) {
  169. yield* session.setPermission({
  170. sessionID: ctx.params.sessionID,
  171. permission: Permission.merge(current.permission ?? [], ctx.payload.permission),
  172. })
  173. }
  174. if (ctx.payload.time?.archived !== undefined) {
  175. yield* session.setArchived({ sessionID: ctx.params.sessionID, time: ctx.payload.time.archived })
  176. }
  177. return yield* requireSession(ctx.params.sessionID)
  178. })
  179. const fork = Effect.fn("SessionHttpApi.fork")(function* (ctx: {
  180. params: { sessionID: SessionID }
  181. payload?: typeof ForkPayload.Type
  182. }) {
  183. return yield* SessionError.mapStorageNotFound(
  184. session.fork({ sessionID: ctx.params.sessionID, messageID: ctx.payload?.messageID }),
  185. )
  186. })
  187. const forkRaw = Effect.fn("SessionHttpApi.forkRaw")(function* (ctx: {
  188. params: { sessionID: SessionID }
  189. request: HttpServerRequest.HttpServerRequest
  190. }) {
  191. const body = yield* Effect.orDie(ctx.request.text)
  192. if (body.trim().length === 0) return yield* fork({ params: ctx.params })
  193. const json = yield* Effect.try({
  194. try: () => JSON.parse(body) as unknown,
  195. catch: () => new HttpApiError.BadRequest({}),
  196. })
  197. const payload = yield* Schema.decodeUnknownEffect(ForkPayload)(json).pipe(
  198. Effect.mapError(() => new HttpApiError.BadRequest({})),
  199. )
  200. return yield* fork({ params: ctx.params, payload })
  201. })
  202. const abort = Effect.fn("SessionHttpApi.abort")(function* (ctx: { params: { sessionID: SessionID } }) {
  203. yield* promptSvc.cancel(ctx.params.sessionID)
  204. return true
  205. })
  206. const init = Effect.fn("SessionHttpApi.init")(function* (ctx: {
  207. params: { sessionID: SessionID }
  208. payload: typeof InitPayload.Type
  209. }) {
  210. yield* requireSession(ctx.params.sessionID)
  211. yield* promptSvc
  212. .command({
  213. sessionID: ctx.params.sessionID,
  214. messageID: ctx.payload.messageID,
  215. model: `${ctx.payload.providerID}/${ctx.payload.modelID}`,
  216. command: Command.Default.INIT,
  217. arguments: "",
  218. })
  219. .pipe(Effect.mapError(() => new HttpApiError.BadRequest({})))
  220. return true
  221. })
  222. // share/unshare errors aren't all client-induced — storage and network
  223. // failures from SessionShare are real possibilities. Map to a typed 500
  224. // (matches the legacy route behavior which routed any failure through
  225. // ErrorMiddleware → NamedError.Unknown 500) instead of blanket-mapping
  226. // every failure to a 400 BadRequest.
  227. const share = Effect.fn("SessionHttpApi.share")(function* (ctx: { params: { sessionID: SessionID } }) {
  228. yield* requireSession(ctx.params.sessionID)
  229. yield* shareSvc.share(ctx.params.sessionID).pipe(Effect.mapError(() => new HttpApiError.InternalServerError({})))
  230. return yield* requireSession(ctx.params.sessionID)
  231. })
  232. const unshare = Effect.fn("SessionHttpApi.unshare")(function* (ctx: { params: { sessionID: SessionID } }) {
  233. yield* requireSession(ctx.params.sessionID)
  234. yield* shareSvc
  235. .unshare(ctx.params.sessionID)
  236. .pipe(Effect.mapError(() => new HttpApiError.InternalServerError({})))
  237. return yield* requireSession(ctx.params.sessionID)
  238. })
  239. const summarize = Effect.fn("SessionHttpApi.summarize")(function* (ctx: {
  240. params: { sessionID: SessionID }
  241. payload: typeof SummarizePayload.Type
  242. }) {
  243. yield* revertSvc.cleanup(yield* requireSession(ctx.params.sessionID))
  244. const messages = yield* SessionError.mapStorageNotFound(session.messages({ sessionID: ctx.params.sessionID }))
  245. const defaultAgent = yield* agentSvc.defaultAgent()
  246. const currentAgent = messages.findLast((message) => message.info.role === "user")?.info.agent ?? defaultAgent
  247. yield* compactSvc.create({
  248. sessionID: ctx.params.sessionID,
  249. agent: currentAgent,
  250. model: {
  251. providerID: ctx.payload.providerID,
  252. modelID: ctx.payload.modelID,
  253. },
  254. auto: ctx.payload.auto ?? false,
  255. })
  256. yield* promptSvc.loop({ sessionID: ctx.params.sessionID })
  257. return true
  258. })
  259. const prompt = Effect.fn("SessionHttpApi.prompt")(function* (ctx: {
  260. params: { sessionID: SessionID }
  261. payload: typeof PromptPayload.Type
  262. }) {
  263. yield* requireSession(ctx.params.sessionID)
  264. const message = yield* promptSvc
  265. .prompt({
  266. ...ctx.payload,
  267. sessionID: ctx.params.sessionID,
  268. })
  269. .pipe(Effect.mapError(() => new HttpApiError.BadRequest({})))
  270. return HttpServerResponse.stream(Stream.make(JSON.stringify(message)).pipe(Stream.encodeText), {
  271. contentType: "application/json",
  272. })
  273. })
  274. const promptAsync = Effect.fn("SessionHttpApi.promptAsync")(function* (ctx: {
  275. params: { sessionID: SessionID }
  276. payload: typeof PromptPayload.Type
  277. }) {
  278. yield* requireSession(ctx.params.sessionID)
  279. yield* promptSvc.prompt({ ...ctx.payload, sessionID: ctx.params.sessionID }).pipe(
  280. Effect.catchCause((cause) =>
  281. Effect.gen(function* () {
  282. yield* Effect.logError("prompt_async failed").pipe(
  283. Effect.annotateLogs({ sessionID: ctx.params.sessionID, cause }),
  284. )
  285. yield* bus.publish(Session.Event.Error, {
  286. sessionID: ctx.params.sessionID,
  287. error: new NamedError.Unknown({ message: Cause.pretty(cause) }).toObject(),
  288. })
  289. }),
  290. ),
  291. Effect.forkIn(scope, { startImmediately: true }),
  292. )
  293. return HttpApiSchema.NoContent.make()
  294. })
  295. const command = Effect.fn("SessionHttpApi.command")(function* (ctx: {
  296. params: { sessionID: SessionID }
  297. payload: typeof CommandPayload.Type
  298. }) {
  299. yield* requireSession(ctx.params.sessionID)
  300. return yield* promptSvc
  301. .command({ ...ctx.payload, sessionID: ctx.params.sessionID })
  302. .pipe(Effect.mapError(() => new HttpApiError.BadRequest({})))
  303. })
  304. const shell = Effect.fn("SessionHttpApi.shell")(function* (ctx: {
  305. params: { sessionID: SessionID }
  306. payload: typeof ShellPayload.Type
  307. }) {
  308. yield* requireSession(ctx.params.sessionID)
  309. return yield* mapBusy(promptSvc.shell({ ...ctx.payload, sessionID: ctx.params.sessionID }))
  310. })
  311. const revert = Effect.fn("SessionHttpApi.revert")(function* (ctx: {
  312. params: { sessionID: SessionID }
  313. payload: typeof RevertPayload.Type
  314. }) {
  315. yield* requireSession(ctx.params.sessionID)
  316. return yield* mapBusy(revertSvc.revert({ sessionID: ctx.params.sessionID, ...ctx.payload }))
  317. })
  318. const unrevert = Effect.fn("SessionHttpApi.unrevert")(function* (ctx: { params: { sessionID: SessionID } }) {
  319. yield* requireSession(ctx.params.sessionID)
  320. return yield* mapBusy(revertSvc.unrevert({ sessionID: ctx.params.sessionID }))
  321. })
  322. const permissionRespond = Effect.fn("SessionHttpApi.permissionRespond")(function* (ctx: {
  323. params: { sessionID: SessionID; permissionID: PermissionID }
  324. payload: typeof PermissionResponsePayload.Type
  325. }) {
  326. yield* requireSession(ctx.params.sessionID)
  327. yield* permissionSvc.reply({ requestID: ctx.params.permissionID, reply: ctx.payload.response })
  328. return true
  329. })
  330. const deleteMessage = Effect.fn("SessionHttpApi.deleteMessage")(function* (ctx: {
  331. params: { sessionID: SessionID; messageID: MessageID }
  332. }) {
  333. yield* requireSession(ctx.params.sessionID)
  334. yield* mapBusy(runState.assertNotBusy(ctx.params.sessionID))
  335. yield* session.removeMessage(ctx.params)
  336. return true
  337. })
  338. const deletePart = Effect.fn("SessionHttpApi.deletePart")(function* (ctx: {
  339. params: { sessionID: SessionID; messageID: MessageID; partID: PartID }
  340. }) {
  341. yield* requireSession(ctx.params.sessionID)
  342. yield* session.removePart(ctx.params)
  343. return true
  344. })
  345. const updatePart = Effect.fn("SessionHttpApi.updatePart")(function* (ctx: {
  346. params: { sessionID: SessionID; messageID: MessageID; partID: PartID }
  347. payload: typeof MessageV2.Part.Type
  348. }) {
  349. yield* requireSession(ctx.params.sessionID)
  350. const payload = ctx.payload as MessageV2.Part
  351. if (
  352. payload.id !== ctx.params.partID ||
  353. payload.messageID !== ctx.params.messageID ||
  354. payload.sessionID !== ctx.params.sessionID
  355. ) {
  356. return yield* new HttpApiError.BadRequest({})
  357. }
  358. return yield* session.updatePart(payload)
  359. })
  360. return handlers
  361. .handle("list", list)
  362. .handle("status", status)
  363. .handle("get", get)
  364. .handle("children", children)
  365. .handle("todo", todo)
  366. .handle("diff", diff)
  367. .handle("messages", messages)
  368. .handle("message", message)
  369. .handleRaw("create", createRaw)
  370. .handle("remove", remove)
  371. .handle("update", update)
  372. .handleRaw("fork", forkRaw)
  373. .handle("abort", abort)
  374. .handle("init", init)
  375. .handle("share", share)
  376. .handle("unshare", unshare)
  377. .handle("summarize", summarize)
  378. .handle("prompt", prompt)
  379. .handle("promptAsync", promptAsync)
  380. .handle("command", command)
  381. .handle("shell", shell)
  382. .handle("revert", revert)
  383. .handle("unrevert", unrevert)
  384. .handle("permissionRespond", permissionRespond)
  385. .handle("deleteMessage", deleteMessage)
  386. .handle("deletePart", deletePart)
  387. .handle("updatePart", updatePart)
  388. }),
  389. )