session.ts 38 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118
  1. import { PermissionV1 } from "@opencode-ai/core/v1/permission"
  2. import { Slug } from "@opencode-ai/core/util/slug"
  3. import { SessionV1 } from "@opencode-ai/core/v1/session"
  4. import { serviceUse } from "@opencode-ai/core/effect/service-use"
  5. import path from "path"
  6. import { BackgroundJob } from "@/background/job"
  7. import { Decimal } from "decimal.js"
  8. import type { ProviderMetadata, Usage } from "@opencode-ai/llm"
  9. import { InstallationVersion } from "@opencode-ai/core/installation/version"
  10. import { Database } from "@opencode-ai/core/database/database"
  11. import { makeRuntime } from "@opencode-ai/core/effect/runtime"
  12. import { EventV2Bridge } from "@/event-v2-bridge"
  13. import { EventV2 } from "@opencode-ai/core/event"
  14. import { SessionV2 } from "@opencode-ai/core/session"
  15. import { SessionExecution } from "@opencode-ai/core/session/execution"
  16. import { NotFoundError } from "@/storage/storage"
  17. import { eq } from "drizzle-orm"
  18. import { and } from "drizzle-orm"
  19. import { gte } from "drizzle-orm"
  20. import { isNull } from "drizzle-orm"
  21. import { desc } from "drizzle-orm"
  22. import { like } from "drizzle-orm"
  23. import { sql } from "drizzle-orm"
  24. import { inArray } from "drizzle-orm"
  25. import { lt } from "drizzle-orm"
  26. import { or } from "drizzle-orm"
  27. import type { SQL } from "drizzle-orm"
  28. import { PartTable, SessionTable } from "@opencode-ai/core/session/sql"
  29. import { ProjectTable } from "@opencode-ai/core/project/sql"
  30. import { Log } from "@opencode-ai/core/util/log"
  31. import { MessageV2 } from "./message-v2"
  32. import type { InstanceContext } from "../project/instance-context"
  33. import { InstanceState } from "@/effect/instance-state"
  34. import { Snapshot } from "@/snapshot"
  35. import { ProjectV2 } from "@opencode-ai/core/project"
  36. import { WorkspaceV2 } from "@opencode-ai/core/workspace"
  37. import { SessionID, MessageID, PartID } from "./schema"
  38. import type { Provider } from "@/provider/provider"
  39. import { Permission } from "@/permission"
  40. import { Global } from "@opencode-ai/core/global"
  41. import { Effect, Layer, Option, Context, Schema, Types } from "effect"
  42. import { NonNegativeInt, optionalOmitUndefined } from "@opencode-ai/core/schema"
  43. import { RuntimeFlags } from "@/effect/runtime-flags"
  44. import { ProviderV2 } from "@opencode-ai/core/provider"
  45. import { ModelV2 } from "@opencode-ai/core/model"
  46. const log = Log.create({ service: "session" })
  47. const runtime = makeRuntime(Database.Service, Database.defaultLayer)
  48. const parentTitlePrefix = "New session - "
  49. const childTitlePrefix = "Child session - "
  50. export function isDefaultTitle(title: string) {
  51. return new RegExp(
  52. `^(${parentTitlePrefix}|${childTitlePrefix})\\d{4}-\\d{2}-\\d{2}T\\d{2}:\\d{2}:\\d{2}\\.\\d{3}Z$`,
  53. ).test(title)
  54. }
  55. type SessionRow = typeof SessionTable.$inferSelect
  56. export function fromRow(row: SessionRow): Info {
  57. const summary =
  58. row.summary_additions !== null || row.summary_deletions !== null || row.summary_files !== null
  59. ? {
  60. additions: row.summary_additions ?? 0,
  61. deletions: row.summary_deletions ?? 0,
  62. files: row.summary_files ?? 0,
  63. diffs: row.summary_diffs ?? undefined,
  64. }
  65. : undefined
  66. const share = row.share_url ? { url: row.share_url } : undefined
  67. const revert = row.revert ?? undefined
  68. return {
  69. id: row.id,
  70. slug: row.slug,
  71. projectID: row.project_id,
  72. workspaceID: row.workspace_id ?? undefined,
  73. directory: row.directory,
  74. path: row.path ?? undefined,
  75. parentID: row.parent_id ?? undefined,
  76. title: row.title,
  77. agent: row.agent ?? undefined,
  78. model: row.model
  79. ? {
  80. id: ModelV2.ID.make(row.model.id),
  81. providerID: ProviderV2.ID.make(row.model.providerID),
  82. variant: row.model.variant,
  83. }
  84. : undefined,
  85. version: row.version,
  86. summary,
  87. cost: row.cost,
  88. tokens: {
  89. input: row.tokens_input,
  90. output: row.tokens_output,
  91. reasoning: row.tokens_reasoning,
  92. cache: {
  93. read: row.tokens_cache_read,
  94. write: row.tokens_cache_write,
  95. },
  96. },
  97. share,
  98. metadata: row.metadata ?? undefined,
  99. revert,
  100. permission: row.permission ? [...row.permission] : undefined,
  101. time: {
  102. created: row.time_created,
  103. updated: row.time_updated,
  104. compacting: row.time_compacting ?? undefined,
  105. archived: row.time_archived ?? undefined,
  106. },
  107. }
  108. }
  109. export function toRow(info: Info) {
  110. return {
  111. id: info.id,
  112. project_id: info.projectID,
  113. workspace_id: info.workspaceID,
  114. parent_id: info.parentID,
  115. slug: info.slug,
  116. directory: info.directory,
  117. path: info.path,
  118. title: info.title,
  119. agent: info.agent,
  120. model: info.model,
  121. version: info.version,
  122. share_url: info.share?.url,
  123. summary_additions: info.summary?.additions,
  124. summary_deletions: info.summary?.deletions,
  125. summary_files: info.summary?.files,
  126. summary_diffs: info.summary?.diffs,
  127. metadata: info.metadata,
  128. cost: info.cost ?? 0,
  129. tokens_input: (info.tokens ?? EmptyTokens).input,
  130. tokens_output: (info.tokens ?? EmptyTokens).output,
  131. tokens_reasoning: (info.tokens ?? EmptyTokens).reasoning,
  132. tokens_cache_read: (info.tokens ?? EmptyTokens).cache.read,
  133. tokens_cache_write: (info.tokens ?? EmptyTokens).cache.write,
  134. revert: info.revert ?? null,
  135. permission: info.permission,
  136. time_created: info.time.created,
  137. time_updated: info.time.updated,
  138. time_compacting: info.time.compacting,
  139. time_archived: info.time.archived,
  140. }
  141. }
  142. function getForkedTitle(title: string): string {
  143. const match = title.match(/^(.+) \(fork #(\d+)\)$/)
  144. if (match) {
  145. const base = match[1]
  146. const num = parseInt(match[2], 10)
  147. return `${base} (fork #${num + 1})`
  148. }
  149. return `${title} (fork #1)`
  150. }
  151. function sessionPath(worktree: string, cwd: string) {
  152. return path.relative(path.resolve(worktree), cwd).replaceAll("\\", "/")
  153. }
  154. const Summary = Schema.Struct({
  155. additions: Schema.Finite,
  156. deletions: Schema.Finite,
  157. files: Schema.Finite,
  158. diffs: optionalOmitUndefined(Schema.Array(Snapshot.FileDiff)),
  159. })
  160. const Tokens = Schema.Struct({
  161. input: Schema.Finite,
  162. output: Schema.Finite,
  163. reasoning: Schema.Finite,
  164. cache: Schema.Struct({
  165. read: Schema.Finite,
  166. write: Schema.Finite,
  167. }),
  168. })
  169. const EmptyTokens = { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }
  170. const Share = Schema.Struct({
  171. url: Schema.String,
  172. })
  173. // Legacy HTTP accepted negative values here. Keep archive timestamps permissive
  174. // while excluding non-finite values that cannot round-trip through JSON.
  175. export const ArchivedTimestamp = Schema.Finite
  176. const Time = Schema.Struct({
  177. created: NonNegativeInt,
  178. updated: NonNegativeInt,
  179. compacting: optionalOmitUndefined(NonNegativeInt),
  180. archived: optionalOmitUndefined(ArchivedTimestamp),
  181. })
  182. const Revert = Schema.Struct({
  183. messageID: MessageID,
  184. partID: optionalOmitUndefined(PartID),
  185. snapshot: optionalOmitUndefined(Schema.String),
  186. diff: optionalOmitUndefined(Schema.String),
  187. })
  188. const Model = Schema.Struct({
  189. id: ModelV2.ID,
  190. providerID: ProviderV2.ID,
  191. variant: optionalOmitUndefined(Schema.String),
  192. })
  193. export const Metadata = Schema.Record(Schema.String, Schema.Any)
  194. export const Info = Schema.Struct({
  195. id: SessionID,
  196. slug: Schema.String,
  197. projectID: ProjectV2.ID,
  198. workspaceID: optionalOmitUndefined(WorkspaceV2.ID),
  199. directory: Schema.String,
  200. path: optionalOmitUndefined(Schema.String),
  201. parentID: optionalOmitUndefined(SessionID),
  202. summary: optionalOmitUndefined(Summary),
  203. cost: optionalOmitUndefined(Schema.Finite),
  204. tokens: optionalOmitUndefined(Tokens),
  205. share: optionalOmitUndefined(Share),
  206. title: Schema.String,
  207. agent: optionalOmitUndefined(Schema.String),
  208. model: optionalOmitUndefined(Model),
  209. version: Schema.String,
  210. metadata: optionalOmitUndefined(Metadata),
  211. time: Time,
  212. permission: optionalOmitUndefined(PermissionV1.Ruleset),
  213. revert: optionalOmitUndefined(Revert),
  214. }).annotate({ identifier: "Session" })
  215. export type Info = Types.DeepMutable<Schema.Schema.Type<typeof Info>>
  216. export const ProjectInfo = Schema.Struct({
  217. id: ProjectV2.ID,
  218. name: optionalOmitUndefined(Schema.String),
  219. worktree: Schema.String,
  220. }).annotate({ identifier: "ProjectSummary" })
  221. export type ProjectInfo = Types.DeepMutable<Schema.Schema.Type<typeof ProjectInfo>>
  222. export const GlobalInfo = Schema.Struct({
  223. ...Info.fields,
  224. project: Schema.NullOr(ProjectInfo),
  225. }).annotate({ identifier: "GlobalSession" })
  226. export type GlobalInfo = Types.DeepMutable<Schema.Schema.Type<typeof GlobalInfo>>
  227. export const CreateInput = Schema.optional(
  228. Schema.Struct({
  229. parentID: Schema.optional(SessionID),
  230. title: Schema.optional(Schema.String),
  231. agent: Schema.optional(Schema.String),
  232. model: Schema.optional(Model),
  233. metadata: Schema.optional(Metadata),
  234. permission: Schema.optional(PermissionV1.Ruleset),
  235. workspaceID: Schema.optional(WorkspaceV2.ID),
  236. }),
  237. )
  238. export type CreateInput = Types.DeepMutable<Schema.Schema.Type<typeof CreateInput>>
  239. export const ForkInput = Schema.Struct({
  240. sessionID: SessionID,
  241. messageID: Schema.optional(MessageID),
  242. })
  243. export const GetInput = SessionID
  244. export const ChildrenInput = SessionID
  245. export const RemoveInput = SessionID
  246. export const SetTitleInput = Schema.Struct({ sessionID: SessionID, title: Schema.String })
  247. export const SetArchivedInput = Schema.Struct({
  248. sessionID: SessionID,
  249. time: Schema.optional(ArchivedTimestamp),
  250. })
  251. export const SetMetadataInput = Schema.Struct({
  252. sessionID: SessionID,
  253. metadata: Metadata,
  254. })
  255. export const SetPermissionInput = Schema.Struct({
  256. sessionID: SessionID,
  257. permission: PermissionV1.Ruleset,
  258. })
  259. export const SetRevertInput = Schema.Struct({
  260. sessionID: SessionID,
  261. revert: Schema.optional(Revert),
  262. summary: Schema.optional(Summary),
  263. })
  264. export const MessagesInput = Schema.Struct({
  265. sessionID: SessionID,
  266. limit: Schema.optional(NonNegativeInt),
  267. })
  268. export type ListInput = {
  269. directory?: string
  270. scope?: "project"
  271. path?: string
  272. workspaceID?: WorkspaceV2.ID
  273. roots?: boolean
  274. start?: number
  275. search?: string
  276. limit?: number
  277. }
  278. export type GlobalListInput = {
  279. directory?: string
  280. roots?: boolean
  281. start?: number
  282. cursor?: number
  283. search?: string
  284. limit?: number
  285. archived?: boolean
  286. }
  287. const CreatedEventSchema = Schema.Struct({
  288. sessionID: SessionID,
  289. info: Info,
  290. })
  291. const UpdatedShare = Schema.Struct({
  292. url: Schema.optional(Schema.NullOr(Schema.String)),
  293. })
  294. const UpdatedTime = Schema.Struct({
  295. created: Schema.optional(Schema.NullOr(NonNegativeInt)),
  296. updated: Schema.optional(Schema.NullOr(NonNegativeInt)),
  297. compacting: Schema.optional(Schema.NullOr(NonNegativeInt)),
  298. archived: Schema.optional(Schema.NullOr(ArchivedTimestamp)),
  299. })
  300. const UpdatedInfo = Schema.Struct({
  301. id: Schema.optional(Schema.NullOr(SessionID)),
  302. slug: Schema.optional(Schema.NullOr(Schema.String)),
  303. projectID: Schema.optional(Schema.NullOr(ProjectV2.ID)),
  304. workspaceID: Schema.optional(Schema.NullOr(WorkspaceV2.ID)),
  305. directory: Schema.optional(Schema.NullOr(Schema.String)),
  306. path: Schema.optional(Schema.NullOr(Schema.String)),
  307. parentID: Schema.optional(Schema.NullOr(SessionID)),
  308. summary: Schema.optional(Schema.NullOr(Summary)),
  309. cost: Schema.optional(Schema.Finite),
  310. tokens: Schema.optional(Tokens),
  311. share: Schema.optional(UpdatedShare),
  312. title: Schema.optional(Schema.NullOr(Schema.String)),
  313. agent: Schema.optional(Schema.NullOr(Schema.String)),
  314. model: Schema.optional(Schema.NullOr(Model)),
  315. version: Schema.optional(Schema.NullOr(Schema.String)),
  316. metadata: Schema.optional(Schema.NullOr(Metadata)),
  317. time: Schema.optional(UpdatedTime),
  318. permission: Schema.optional(Schema.NullOr(PermissionV1.Ruleset)),
  319. revert: Schema.optional(Schema.NullOr(Revert)),
  320. })
  321. const UpdatedEventSchema = Schema.Struct({
  322. sessionID: SessionID,
  323. info: UpdatedInfo,
  324. })
  325. export const Event = {
  326. Created: SessionV1.Event.Created,
  327. Updated: SessionV1.Event.Updated,
  328. Deleted: SessionV1.Event.Deleted,
  329. Diff: EventV2.define({
  330. type: "session.diff",
  331. schema: {
  332. sessionID: SessionID,
  333. diff: Schema.Array(Snapshot.FileDiff),
  334. },
  335. }),
  336. Error: EventV2.define({
  337. type: "session.error",
  338. schema: {
  339. sessionID: Schema.optional(SessionID),
  340. // Reuses SessionV1.Assistant.fields.error (already Schema.optional) so
  341. // the derived schema keeps the same discriminated-union shape on the event stream.
  342. error: SessionV1.Assistant.fields.error,
  343. },
  344. }),
  345. }
  346. export function plan(input: { slug: string; time: { created: number } }, instance: InstanceContext) {
  347. const base = instance.project.vcs
  348. ? path.join(instance.worktree, ".opencode", "plans")
  349. : path.join(Global.Path.data, "plans")
  350. return path.join(base, [input.time.created, input.slug].join("-") + ".md")
  351. }
  352. export const getUsage = (input: { model: Provider.Model; usage: Usage; metadata?: ProviderMetadata }) => {
  353. const safe = (value: number) => {
  354. if (!Number.isFinite(value)) return 0
  355. return Math.max(0, value)
  356. }
  357. const inputTokens = safe(input.usage.inputTokens ?? 0)
  358. const outputTokens = safe(input.usage.outputTokens ?? 0)
  359. const reasoningTokens = safe(input.usage.reasoningTokens ?? 0)
  360. const cacheReadInputTokens = safe(input.usage.cacheReadInputTokens ?? 0)
  361. const cacheWriteInputTokens = safe(
  362. Number(
  363. input.usage.cacheWriteInputTokens ??
  364. input.metadata?.["anthropic"]?.["cacheCreationInputTokens"] ??
  365. // google-vertex-anthropic returns metadata under "vertex" key
  366. // (AnthropicMessagesLanguageModel custom provider key from 'vertex.anthropic.messages')
  367. input.metadata?.["vertex"]?.["cacheCreationInputTokens"] ??
  368. // @ts-expect-error
  369. input.metadata?.["bedrock"]?.["usage"]?.["cacheWriteInputTokens"] ??
  370. // @ts-expect-error
  371. input.metadata?.["venice"]?.["usage"]?.["cacheCreationInputTokens"] ??
  372. 0,
  373. ),
  374. )
  375. // AI SDK v6 normalized inputTokens to include cached tokens across all providers
  376. // (including Anthropic/Bedrock which previously excluded them). Always subtract cache
  377. // tokens to get the non-cached input count for separate cost calculation.
  378. const adjustedInputTokens = safe(inputTokens - cacheReadInputTokens - cacheWriteInputTokens)
  379. const total = input.usage.totalTokens
  380. const tokens = {
  381. total,
  382. input: adjustedInputTokens,
  383. output: safe(outputTokens - reasoningTokens),
  384. reasoning: reasoningTokens,
  385. cache: {
  386. write: cacheWriteInputTokens,
  387. read: cacheReadInputTokens,
  388. },
  389. }
  390. const contextTokens = inputTokens
  391. const costInfo =
  392. input.model.cost?.tiers
  393. ?.filter((item) => item.tier.type === "context" && contextTokens > item.tier.size)
  394. .sort((a, b) => b.tier.size - a.tier.size)[0] ??
  395. (input.model.cost?.experimentalOver200K && contextTokens > 200_000
  396. ? input.model.cost.experimentalOver200K
  397. : input.model.cost)
  398. const totalNanoAiu = input.metadata?.["copilot"]?.["totalNanoAiu"]
  399. return {
  400. cost:
  401. typeof totalNanoAiu === "number" && Number.isFinite(totalNanoAiu) && totalNanoAiu >= 0
  402. ? new Decimal(totalNanoAiu).div(100_000_000_000).toNumber()
  403. : safe(
  404. new Decimal(0)
  405. .add(new Decimal(tokens.input).mul(costInfo?.input ?? 0).div(1_000_000))
  406. .add(new Decimal(tokens.output).mul(costInfo?.output ?? 0).div(1_000_000))
  407. .add(new Decimal(tokens.cache.read).mul(costInfo?.cache?.read ?? 0).div(1_000_000))
  408. .add(new Decimal(tokens.cache.write).mul(costInfo?.cache?.write ?? 0).div(1_000_000))
  409. // TODO: update models.dev to have better pricing model, for now:
  410. // charge reasoning tokens at the same rate as output tokens
  411. .add(new Decimal(tokens.reasoning).mul(costInfo?.output ?? 0).div(1_000_000))
  412. .toNumber(),
  413. ),
  414. tokens,
  415. }
  416. }
  417. export class BusyError extends Schema.TaggedErrorClass<BusyError>()("SessionBusyError", {
  418. sessionID: SessionID,
  419. }) {}
  420. export type NotFound = NotFoundError
  421. export interface Interface {
  422. readonly list: (input?: ListInput) => Effect.Effect<Info[]>
  423. readonly listGlobal: (input?: GlobalListInput) => Effect.Effect<GlobalInfo[]>
  424. readonly create: (input?: {
  425. parentID?: SessionID
  426. title?: string
  427. agent?: string
  428. model?: Schema.Schema.Type<typeof Model>
  429. metadata?: typeof Metadata.Type
  430. permission?: PermissionV1.Ruleset
  431. workspaceID?: WorkspaceV2.ID
  432. }) => Effect.Effect<Info>
  433. readonly fork: (input: { sessionID: SessionID; messageID?: MessageID }) => Effect.Effect<Info, NotFound>
  434. readonly touch: (sessionID: SessionID) => Effect.Effect<void>
  435. readonly get: (id: SessionID) => Effect.Effect<Info, NotFound>
  436. readonly setTitle: (input: { sessionID: SessionID; title: string }) => Effect.Effect<void>
  437. readonly setArchived: (input: { sessionID: SessionID; time?: number }) => Effect.Effect<void>
  438. readonly setMetadata: (input: typeof SetMetadataInput.Type) => Effect.Effect<void>
  439. readonly setPermission: (input: { sessionID: SessionID; permission: PermissionV1.Ruleset }) => Effect.Effect<void>
  440. readonly setRevert: (input: {
  441. sessionID: SessionID
  442. revert: Info["revert"]
  443. summary: Info["summary"]
  444. }) => Effect.Effect<void>
  445. readonly clearRevert: (sessionID: SessionID) => Effect.Effect<void>
  446. readonly setSummary: (input: { sessionID: SessionID; summary: Info["summary"] }) => Effect.Effect<void>
  447. readonly setShare: (input: { sessionID: SessionID; share: Info["share"] }) => Effect.Effect<void>
  448. readonly setWorkspace: (input: { sessionID: SessionID; workspaceID: Info["workspaceID"] }) => Effect.Effect<void>
  449. readonly diff: (sessionID: SessionID) => Effect.Effect<Snapshot.FileDiff[]>
  450. readonly messages: (input: { sessionID: SessionID; limit?: number }) => Effect.Effect<SessionV1.WithParts[], NotFound>
  451. readonly children: (parentID: SessionID) => Effect.Effect<Info[]>
  452. readonly remove: (sessionID: SessionID) => Effect.Effect<void, NotFound>
  453. readonly updateMessage: <T extends SessionV1.Info>(msg: T) => Effect.Effect<T>
  454. readonly removeMessage: (input: { sessionID: SessionID; messageID: MessageID }) => Effect.Effect<MessageID>
  455. readonly removePart: (input: { sessionID: SessionID; messageID: MessageID; partID: PartID }) => Effect.Effect<PartID>
  456. readonly getPart: (input: {
  457. sessionID: SessionID
  458. messageID: MessageID
  459. partID: PartID
  460. }) => Effect.Effect<SessionV1.Part | undefined>
  461. readonly updatePart: <T extends SessionV1.Part>(part: T) => Effect.Effect<T>
  462. readonly updatePartDelta: (input: {
  463. sessionID: SessionID
  464. messageID: MessageID
  465. partID: PartID
  466. field: string
  467. delta: string
  468. }) => Effect.Effect<void>
  469. /** Finds the first message matching the predicate, searching newest-first. */
  470. readonly findMessage: (
  471. sessionID: SessionID,
  472. predicate: (msg: SessionV1.WithParts) => boolean,
  473. ) => Effect.Effect<Option.Option<SessionV1.WithParts>, NotFound>
  474. }
  475. export class Service extends Context.Service<Service, Interface>()("@opencode/Session") {}
  476. export const use = serviceUse(Service)
  477. export type Patch = Omit<Partial<Info>, "time" | "share" | "summary" | "revert" | "permission"> & {
  478. time?: Partial<Info["time"]>
  479. share?: Partial<NonNullable<Info["share"]>> | null
  480. summary?: Info["summary"] | null
  481. revert?: Info["revert"] | null
  482. permission?: Info["permission"] | null
  483. }
  484. export const layer: Layer.Layer<
  485. Service,
  486. never,
  487. BackgroundJob.Service | RuntimeFlags.Service | Database.Service | EventV2Bridge.Service
  488. > = Layer.effect(
  489. Service,
  490. Effect.gen(function* () {
  491. const { db } = yield* Database.Service
  492. const database = yield* Database.Service
  493. const background = yield* BackgroundJob.Service
  494. const events = yield* EventV2Bridge.Service
  495. const flags = yield* RuntimeFlags.Service
  496. const createNext = Effect.fn("Session.createNext")(function* (input: {
  497. id?: SessionID
  498. title?: string
  499. agent?: string
  500. model?: Schema.Schema.Type<typeof Model>
  501. parentID?: SessionID
  502. workspaceID?: WorkspaceV2.ID
  503. directory: string
  504. path?: string
  505. metadata?: typeof Metadata.Type
  506. permission?: PermissionV1.Ruleset
  507. }) {
  508. const ctx = yield* InstanceState.context
  509. const result: Info = {
  510. id: SessionID.descending(input.id),
  511. slug: Slug.create(),
  512. version: InstallationVersion,
  513. projectID: ctx.project.id,
  514. directory: input.directory,
  515. path: input.path,
  516. workspaceID: input.workspaceID,
  517. parentID: input.parentID,
  518. title: input.title ?? (input.parentID ? childTitlePrefix : parentTitlePrefix) + new Date().toISOString(),
  519. agent: input.agent,
  520. model: input.model,
  521. metadata: input.metadata,
  522. permission: input.permission ? [...input.permission] : undefined,
  523. cost: 0,
  524. tokens: EmptyTokens,
  525. time: {
  526. created: Date.now(),
  527. updated: Date.now(),
  528. },
  529. }
  530. log.info("created", result)
  531. yield* events.publish(SessionV1.Event.Created, { sessionID: result.id, info: result })
  532. return result
  533. })
  534. const get = Effect.fn("Session.get")(function* (id: SessionID) {
  535. const row = yield* db.select().from(SessionTable).where(eq(SessionTable.id, id)).get().pipe(Effect.orDie)
  536. if (!row) return yield* Effect.fail(new NotFoundError({ message: `Session not found: ${id}` }))
  537. return fromRow(row)
  538. })
  539. const list = Effect.fn("Session.list")(function* (input?: ListInput) {
  540. const ctx = yield* InstanceState.context
  541. return yield* listByProject(db, {
  542. projectID: ctx.project.id,
  543. experimentalWorkspaces: flags.experimentalWorkspaces,
  544. ...input,
  545. })
  546. })
  547. const listGlobal = Effect.fn("Session.listGlobal")(function* (input?: GlobalListInput) {
  548. const conditions: SQL[] = []
  549. if (input?.directory) conditions.push(eq(SessionTable.directory, input.directory))
  550. if (input?.roots) conditions.push(isNull(SessionTable.parent_id))
  551. if (input?.start) conditions.push(gte(SessionTable.time_updated, input.start))
  552. if (input?.cursor) conditions.push(lt(SessionTable.time_updated, input.cursor))
  553. if (input?.search) conditions.push(like(SessionTable.title, `%${input.search}%`))
  554. if (!input?.archived) conditions.push(isNull(SessionTable.time_archived))
  555. const query =
  556. conditions.length > 0
  557. ? db
  558. .select()
  559. .from(SessionTable)
  560. .where(and(...conditions))
  561. : db.select().from(SessionTable)
  562. const rows = yield* query
  563. .orderBy(desc(SessionTable.time_updated), desc(SessionTable.id))
  564. .limit(input?.limit ?? 100)
  565. .all()
  566. .pipe(Effect.orDie)
  567. const ids = [...new Set(rows.map((row) => row.project_id))]
  568. const projects = new Map<string, ProjectInfo>()
  569. if (ids.length > 0) {
  570. const items = yield* db
  571. .select({ id: ProjectTable.id, name: ProjectTable.name, worktree: ProjectTable.worktree })
  572. .from(ProjectTable)
  573. .where(inArray(ProjectTable.id, ids))
  574. .all()
  575. .pipe(Effect.orDie)
  576. for (const item of items) {
  577. projects.set(item.id, {
  578. id: item.id,
  579. name: item.name ?? undefined,
  580. worktree: item.worktree,
  581. })
  582. }
  583. }
  584. return rows.map((row) => ({ ...fromRow(row), project: projects.get(row.project_id) ?? null }))
  585. })
  586. const children = Effect.fn("Session.children")(function* (parentID: SessionID) {
  587. const rows = yield* db
  588. .select()
  589. .from(SessionTable)
  590. .where(and(eq(SessionTable.parent_id, parentID)))
  591. .all()
  592. .pipe(Effect.orDie)
  593. return rows.map(fromRow)
  594. })
  595. const remove: Interface["remove"] = Effect.fnUntraced(function* (sessionID: SessionID) {
  596. const session = yield* get(sessionID)
  597. try {
  598. // `remove` needs to work in all cases, such as broken sessions that
  599. // run cleanup without instance state.
  600. const hasInstance = yield* InstanceState.directory.pipe(
  601. Effect.as(true),
  602. Effect.catchCause(() => Effect.succeed(false)),
  603. )
  604. if (hasInstance) yield* cancelBackgroundJobs(background, sessionID)
  605. const kids = yield* children(sessionID)
  606. for (const child of kids) {
  607. yield* remove(child.id)
  608. }
  609. yield* events.publish(SessionV1.Event.Deleted, { sessionID, info: session })
  610. yield* events.remove(sessionID)
  611. } catch (e) {
  612. log.error(e)
  613. }
  614. })
  615. const updateMessage = <T extends SessionV1.Info>(msg: T): Effect.Effect<T> =>
  616. Effect.gen(function* () {
  617. yield* events.publish(SessionV1.Event.MessageUpdated, { sessionID: msg.sessionID, info: msg })
  618. return msg
  619. }).pipe(Effect.withSpan("Session.updateMessage"))
  620. const updatePart = <T extends SessionV1.Part>(part: T): Effect.Effect<T> =>
  621. Effect.gen(function* () {
  622. yield* events.publish(SessionV1.Event.PartUpdated, {
  623. sessionID: part.sessionID,
  624. part: structuredClone(part),
  625. time: Date.now(),
  626. })
  627. return part
  628. }).pipe(Effect.withSpan("Session.updatePart"))
  629. const getPart: Interface["getPart"] = Effect.fn("Session.getPart")(function* (input) {
  630. const row = yield* db
  631. .select()
  632. .from(PartTable)
  633. .where(
  634. and(
  635. eq(PartTable.session_id, input.sessionID),
  636. eq(PartTable.message_id, input.messageID),
  637. eq(PartTable.id, input.partID),
  638. ),
  639. )
  640. .get()
  641. .pipe(Effect.orDie)
  642. if (!row) return
  643. return {
  644. ...row.data,
  645. id: row.id,
  646. sessionID: row.session_id,
  647. messageID: row.message_id,
  648. } as SessionV1.Part
  649. })
  650. const create = Effect.fn("Session.create")(function* (input?: {
  651. parentID?: SessionID
  652. title?: string
  653. agent?: string
  654. model?: Schema.Schema.Type<typeof Model>
  655. metadata?: typeof Metadata.Type
  656. permission?: PermissionV1.Ruleset
  657. workspaceID?: WorkspaceV2.ID
  658. }) {
  659. const ctx = yield* InstanceState.context
  660. const workspace = yield* InstanceState.workspaceID
  661. return yield* createNext({
  662. parentID: input?.parentID,
  663. directory: ctx.directory,
  664. path: sessionPath(ctx.worktree, ctx.directory),
  665. title: input?.title,
  666. agent: input?.agent,
  667. model: input?.model,
  668. metadata: input?.metadata,
  669. permission: input?.permission,
  670. workspaceID: input?.workspaceID ?? workspace,
  671. })
  672. })
  673. const fork = Effect.fn("Session.fork")(function* (input: { sessionID: SessionID; messageID?: MessageID }) {
  674. const ctx = yield* InstanceState.context
  675. const original = yield* get(input.sessionID)
  676. const title = getForkedTitle(original.title)
  677. const session = yield* createNext({
  678. directory: ctx.directory,
  679. path: sessionPath(ctx.worktree, ctx.directory),
  680. workspaceID: original.workspaceID,
  681. title,
  682. metadata: structuredClone(original.metadata),
  683. })
  684. const msgs = yield* messages({ sessionID: input.sessionID })
  685. const idMap = new Map<string, MessageID>()
  686. for (const msg of msgs) {
  687. if (input.messageID && msg.info.id >= input.messageID) break
  688. const newID = MessageID.ascending()
  689. idMap.set(msg.info.id, newID)
  690. const parentID = msg.info.role === "assistant" && msg.info.parentID ? idMap.get(msg.info.parentID) : undefined
  691. const cloned = yield* updateMessage({
  692. ...msg.info,
  693. sessionID: session.id,
  694. id: newID,
  695. ...(parentID && { parentID }),
  696. })
  697. for (const part of msg.parts) {
  698. const p: SessionV1.Part = {
  699. ...part,
  700. id: PartID.ascending(),
  701. messageID: cloned.id,
  702. sessionID: session.id,
  703. }
  704. if (p.type === "compaction" && p.tail_start_id) {
  705. p.tail_start_id = idMap.get(p.tail_start_id)
  706. }
  707. yield* updatePart(p)
  708. }
  709. }
  710. return session
  711. })
  712. const patch = (sessionID: SessionID, info: Patch) =>
  713. Effect.gen(function* () {
  714. const current = yield* get(sessionID)
  715. const next = {
  716. ...current,
  717. ...info,
  718. time: info.time ? { ...current.time, ...info.time } : current.time,
  719. share: info.share === null ? undefined : info.share ? { ...current.share, ...info.share } : current.share,
  720. summary: info.summary === null ? undefined : (info.summary ?? current.summary),
  721. revert: info.revert === null ? undefined : (info.revert ?? current.revert),
  722. permission: info.permission === null ? undefined : (info.permission ?? current.permission),
  723. } as Info
  724. yield* events.publish(SessionV1.Event.Updated, { sessionID, info: next })
  725. })
  726. const touch = Effect.fn("Session.touch")(function* (sessionID: SessionID) {
  727. yield* patch(sessionID, { time: { updated: Date.now() } }).pipe(Effect.orDie)
  728. })
  729. const setTitle = Effect.fn("Session.setTitle")(function* (input: { sessionID: SessionID; title: string }) {
  730. yield* patch(input.sessionID, { title: input.title }).pipe(Effect.orDie)
  731. })
  732. const setArchived = Effect.fn("Session.setArchived")(function* (input: { sessionID: SessionID; time?: number }) {
  733. yield* patch(input.sessionID, { time: { archived: input.time } }).pipe(Effect.orDie)
  734. })
  735. const setMetadata = Effect.fn("Session.setMetadata")(function* (input: typeof SetMetadataInput.Type) {
  736. yield* patch(input.sessionID, { metadata: input.metadata, time: { updated: Date.now() } }).pipe(Effect.orDie)
  737. })
  738. const setPermission = Effect.fn("Session.setPermission")(function* (input: {
  739. sessionID: SessionID
  740. permission: PermissionV1.Ruleset
  741. }) {
  742. yield* patch(input.sessionID, { permission: [...input.permission], time: { updated: Date.now() } }).pipe(
  743. Effect.orDie,
  744. )
  745. })
  746. const setRevert = Effect.fn("Session.setRevert")(function* (input: {
  747. sessionID: SessionID
  748. revert: Info["revert"]
  749. summary: Info["summary"]
  750. }) {
  751. yield* patch(input.sessionID, {
  752. summary: input.summary,
  753. time: { updated: Date.now() },
  754. revert: input.revert,
  755. }).pipe(Effect.orDie)
  756. })
  757. const clearRevert = Effect.fn("Session.clearRevert")(function* (sessionID: SessionID) {
  758. yield* patch(sessionID, { time: { updated: Date.now() }, revert: null }).pipe(Effect.orDie)
  759. })
  760. const setSummary = Effect.fn("Session.setSummary")(function* (input: {
  761. sessionID: SessionID
  762. summary: Info["summary"]
  763. }) {
  764. yield* patch(input.sessionID, { time: { updated: Date.now() }, summary: input.summary }).pipe(Effect.orDie)
  765. })
  766. const setShare = Effect.fn("Session.setShare")(function* (input: { sessionID: SessionID; share: Info["share"] }) {
  767. yield* patch(input.sessionID, { share: input.share ?? null, time: { updated: Date.now() } }).pipe(Effect.orDie)
  768. })
  769. const setWorkspace = Effect.fn("Session.setWorkspace")(function* (input: {
  770. sessionID: SessionID
  771. workspaceID: Info["workspaceID"]
  772. }) {
  773. yield* patch(input.sessionID, { workspaceID: input.workspaceID, time: { updated: Date.now() } }).pipe(
  774. Effect.orDie,
  775. )
  776. })
  777. const diff = Effect.fn("Session.diff")(function* (sessionID: SessionID) {
  778. void sessionID
  779. return [] as Snapshot.FileDiff[]
  780. })
  781. const messages: Interface["messages"] = Effect.fn("Session.messages")(function* (input) {
  782. if (input.limit) {
  783. return (yield* MessageV2.page({ sessionID: input.sessionID, limit: input.limit }).pipe(
  784. Effect.provideService(Database.Service, database),
  785. )).items
  786. }
  787. const size = 50
  788. const result = [] as SessionV1.WithParts[]
  789. let before: string | undefined
  790. while (true) {
  791. const page = yield* MessageV2.page({ sessionID: input.sessionID, limit: size, before }).pipe(
  792. Effect.provideService(Database.Service, database),
  793. )
  794. if (page.items.length === 0) break
  795. for (let i = page.items.length - 1; i >= 0; i--) {
  796. const item = page.items[i]
  797. if (item) result.push(item)
  798. }
  799. if (!page.more || !page.cursor) break
  800. before = page.cursor
  801. }
  802. return result.reverse()
  803. })
  804. const removeMessage = Effect.fn("Session.removeMessage")(function* (input: {
  805. sessionID: SessionID
  806. messageID: MessageID
  807. }) {
  808. yield* events.publish(SessionV1.Event.MessageRemoved, {
  809. sessionID: input.sessionID,
  810. messageID: input.messageID,
  811. })
  812. return input.messageID
  813. })
  814. const removePart = Effect.fn("Session.removePart")(function* (input: {
  815. sessionID: SessionID
  816. messageID: MessageID
  817. partID: PartID
  818. }) {
  819. yield* events.publish(SessionV1.Event.PartRemoved, {
  820. sessionID: input.sessionID,
  821. messageID: input.messageID,
  822. partID: input.partID,
  823. })
  824. return input.partID
  825. })
  826. const updatePartDelta = Effect.fnUntraced(function* (input: {
  827. sessionID: SessionID
  828. messageID: MessageID
  829. partID: PartID
  830. field: string
  831. delta: string
  832. }) {
  833. yield* events.publish(MessageV2.Event.PartDelta, input)
  834. })
  835. /** Finds the first message matching the predicate, searching newest-first. */
  836. const findMessage: Interface["findMessage"] = Effect.fn("Session.findMessage")(function* (sessionID, predicate) {
  837. const size = 50
  838. let before: string | undefined
  839. while (true) {
  840. const page = yield* MessageV2.page({ sessionID, limit: size, before }).pipe(
  841. Effect.provideService(Database.Service, database),
  842. )
  843. if (page.items.length === 0) break
  844. for (let i = page.items.length - 1; i >= 0; i--) {
  845. const item = page.items[i]
  846. if (item && predicate(item)) return Option.some(item)
  847. }
  848. if (!page.more || !page.cursor) break
  849. before = page.cursor
  850. }
  851. return Option.none<SessionV1.WithParts>()
  852. })
  853. return Service.of({
  854. list,
  855. listGlobal,
  856. create,
  857. fork,
  858. touch,
  859. get,
  860. setTitle,
  861. setArchived,
  862. setMetadata,
  863. setPermission,
  864. setRevert,
  865. clearRevert,
  866. setSummary,
  867. setShare,
  868. setWorkspace,
  869. diff,
  870. messages,
  871. children,
  872. remove,
  873. updateMessage,
  874. removeMessage,
  875. removePart,
  876. updatePart,
  877. getPart,
  878. updatePartDelta,
  879. findMessage,
  880. })
  881. }),
  882. )
  883. export const defaultLayer = layer.pipe(
  884. Layer.provide(BackgroundJob.defaultLayer),
  885. Layer.provide(Database.defaultLayer),
  886. Layer.provide(EventV2Bridge.defaultLayer),
  887. Layer.provide(SessionExecution.noopLayer),
  888. Layer.provide(SessionV2.defaultLayer),
  889. Layer.provide(RuntimeFlags.defaultLayer),
  890. )
  891. const cancelBackgroundJobs = Effect.fn("Session.cancelBackgroundJobs")(function* (
  892. background: BackgroundJob.Interface,
  893. sessionID: SessionID,
  894. ) {
  895. const jobs = yield* background.list()
  896. yield* Effect.forEach(
  897. jobs.filter((job) => {
  898. if (job.status !== "running") return false
  899. if (job.id === sessionID) return true
  900. if (job.metadata?.sessionId === sessionID) return true
  901. return job.metadata?.parentSessionId === sessionID
  902. }),
  903. (job) => background.cancel(job.id),
  904. { concurrency: "unbounded", discard: true },
  905. )
  906. })
  907. function listByProject(
  908. db: Database.Interface["db"],
  909. input: ListInput & {
  910. projectID: ProjectV2.ID
  911. experimentalWorkspaces: boolean
  912. },
  913. ) {
  914. const conditions = [eq(SessionTable.project_id, input.projectID)]
  915. if (input.workspaceID) {
  916. conditions.push(eq(SessionTable.workspace_id, input.workspaceID))
  917. }
  918. if (input.path !== undefined) {
  919. if (input.path) {
  920. const conds = [
  921. eq(SessionTable.path, input.path),
  922. like(SessionTable.path, sql.param(`${input.path}/%`, SessionTable.path)),
  923. ]
  924. conditions.push(
  925. input.directory
  926. ? or(...conds, and(isNull(SessionTable.path), eq(SessionTable.directory, input.directory))!)!
  927. : or(...conds)!,
  928. )
  929. }
  930. } else if (input.scope !== "project") {
  931. if (input.directory) {
  932. conditions.push(eq(SessionTable.directory, input.directory))
  933. }
  934. }
  935. if (input.roots) {
  936. conditions.push(isNull(SessionTable.parent_id))
  937. }
  938. if (input.start) {
  939. conditions.push(gte(SessionTable.time_updated, input.start))
  940. }
  941. if (input.search) {
  942. conditions.push(like(SessionTable.title, `%${input.search}%`))
  943. }
  944. const limit = input.limit ?? 100
  945. return db
  946. .select()
  947. .from(SessionTable)
  948. .where(and(...conditions))
  949. .orderBy(desc(SessionTable.time_updated))
  950. .limit(limit)
  951. .all()
  952. .pipe(
  953. Effect.orDie,
  954. Effect.map((rows) => rows.map(fromRow)),
  955. )
  956. }
  957. export function* listGlobal(input?: {
  958. directory?: string
  959. roots?: boolean
  960. start?: number
  961. cursor?: number
  962. search?: string
  963. limit?: number
  964. archived?: boolean
  965. }) {
  966. const conditions: SQL[] = []
  967. if (input?.directory) {
  968. conditions.push(eq(SessionTable.directory, input.directory))
  969. }
  970. if (input?.roots) {
  971. conditions.push(isNull(SessionTable.parent_id))
  972. }
  973. if (input?.start) {
  974. conditions.push(gte(SessionTable.time_updated, input.start))
  975. }
  976. if (input?.cursor) {
  977. conditions.push(lt(SessionTable.time_updated, input.cursor))
  978. }
  979. if (input?.search) {
  980. conditions.push(like(SessionTable.title, `%${input.search}%`))
  981. }
  982. if (!input?.archived) {
  983. conditions.push(isNull(SessionTable.time_archived))
  984. }
  985. const limit = input?.limit ?? 100
  986. const rows = runtime.runSync(({ db }) => {
  987. const query =
  988. conditions.length > 0
  989. ? db
  990. .select()
  991. .from(SessionTable)
  992. .where(and(...conditions))
  993. : db.select().from(SessionTable)
  994. return query.orderBy(desc(SessionTable.time_updated), desc(SessionTable.id)).limit(limit).all().pipe(Effect.orDie)
  995. })
  996. const ids = [...new Set(rows.map((row) => row.project_id))]
  997. const projects = new Map<string, ProjectInfo>()
  998. if (ids.length > 0) {
  999. const items = runtime.runSync(({ db }) =>
  1000. db
  1001. .select({ id: ProjectTable.id, name: ProjectTable.name, worktree: ProjectTable.worktree })
  1002. .from(ProjectTable)
  1003. .where(inArray(ProjectTable.id, ids))
  1004. .all()
  1005. .pipe(Effect.orDie),
  1006. )
  1007. for (const item of items) {
  1008. projects.set(item.id, {
  1009. id: item.id,
  1010. name: item.name ?? undefined,
  1011. worktree: item.worktree,
  1012. })
  1013. }
  1014. }
  1015. for (const row of rows) {
  1016. const project = projects.get(row.project_id) ?? null
  1017. yield { ...fromRow(row), project }
  1018. }
  1019. }
  1020. export * as Session from "./session"