session.ts 35 KB

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