session-create.test.ts 35 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901
  1. import { describe, expect } from "bun:test"
  2. import { $ } from "bun"
  3. import fs from "fs/promises"
  4. import path from "path"
  5. import { DateTime, Effect, Layer, Stream } from "effect"
  6. import { Money } from "@opencode-ai/schema/money"
  7. import { Agent } from "@opencode-ai/core/agent"
  8. import { asc, eq } from "drizzle-orm"
  9. import { Database } from "@opencode-ai/core/database/database"
  10. import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
  11. import { LayerNode } from "@opencode-ai/util/effect/layer-node"
  12. import { Hash } from "@opencode-ai/util/hash"
  13. import { Bus } from "@opencode-ai/core/bus"
  14. import { EventTable } from "@opencode-ai/core/event/sql"
  15. import { Location } from "@opencode-ai/core/location"
  16. import { Model } from "@opencode-ai/core/model"
  17. import { Project } from "@opencode-ai/core/project"
  18. import { ProjectTable } from "@opencode-ai/core/project/sql"
  19. import { Provider } from "@opencode-ai/core/provider"
  20. import { AbsolutePath, RelativePath } from "@opencode-ai/core/schema"
  21. import { Session } from "@opencode-ai/core/session"
  22. import { SessionMessage } from "@opencode-ai/core/session/message"
  23. import { SessionProjector } from "@opencode-ai/core/session/projector"
  24. import { SessionExecution } from "@opencode-ai/core/session/execution"
  25. import { SessionInbox } from "@opencode-ai/core/session/inbox"
  26. import { SessionEvent } from "@opencode-ai/core/session/event"
  27. import { SessionTable } from "@opencode-ai/core/session/sql"
  28. import { SessionStore } from "@opencode-ai/core/session/store"
  29. import { SessionTransfer } from "@opencode-ai/core/session/transfer"
  30. import { Workspace } from "@opencode-ai/core/workspace"
  31. import { testEffect } from "./lib/effect"
  32. import { globalProjectLayer } from "./lib/project"
  33. import { tmpdir } from "./fixture/tmpdir"
  34. const it = testEffect(
  35. AppNodeBuilder.build(
  36. LayerNode.group([
  37. Database.node,
  38. Bus.node,
  39. SessionProjector.node,
  40. SessionStore.node,
  41. Session.node,
  42. SessionTransfer.node,
  43. ]),
  44. [
  45. [Bus.node, Bus.configured({ persist: true })],
  46. [Project.node, globalProjectLayer],
  47. [SessionExecution.node, SessionExecution.noopLayer],
  48. ],
  49. ),
  50. )
  51. const liveIt = testEffect(
  52. AppNodeBuilder.build(
  53. LayerNode.group([Database.node, Bus.node, Project.node, SessionProjector.node, SessionStore.node, Session.node]),
  54. [
  55. [Bus.node, Bus.configured({ persist: true })],
  56. [SessionExecution.node, SessionExecution.noopLayer],
  57. ],
  58. ),
  59. )
  60. const location = Location.Ref.make({ directory: AbsolutePath.make("/project") })
  61. const id = Session.ID.create()
  62. /** Public session events from a `log` read, without synced markers. */
  63. const logEvents = (session: Session.Interface, sessionID: Session.ID, follow?: boolean) =>
  64. session
  65. .log({ sessionID, follow })
  66. .pipe(Stream.filter((item): item is SessionEvent.DurableEvent => !Bus.isSynced(item)))
  67. const assertCreateInputTypes = (session: Session.Interface) => {
  68. // @ts-expect-error location or parentID is required.
  69. session.create({})
  70. // @ts-expect-error child sessions inherit their parent's location.
  71. session.create({ parentID: Session.ID.create(), location })
  72. }
  73. void assertCreateInputTypes
  74. function withTmp<A, E, R>(f: (directory: string) => Effect.Effect<A, E, R>) {
  75. return Effect.acquireRelease(
  76. Effect.promise(() => tmpdir()),
  77. (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
  78. ).pipe(Effect.flatMap((tmp) => f(tmp.path)))
  79. }
  80. describe("Session.create", () => {
  81. liveIt.live("follows the directory's project identity established after creation", () =>
  82. withTmp((directory) =>
  83. Effect.gen(function* () {
  84. const session = yield* Session.Service
  85. const projects = yield* Project.Service
  86. const { db } = yield* Database.Service
  87. const ref = Location.Ref.make({ directory: AbsolutePath.make(directory) })
  88. const nested = Location.Ref.make({ directory: AbsolutePath.make(path.join(directory, "packages", "app")) })
  89. const created = yield* session.create({ location: ref, title: "Before git" })
  90. const child = yield* session.create({ location: nested, title: "Nested before git" })
  91. const originalUpdated = created.time.updated
  92. yield* Effect.promise(async () => {
  93. await $`git init -q`.cwd(directory)
  94. await $`git config user.email test@example.com`.cwd(directory)
  95. await $`git config user.name Test`.cwd(directory)
  96. await fs.writeFile(path.join(directory, "README.md"), "test\n")
  97. await $`git add README.md`.cwd(directory)
  98. await $`git commit -qm initial`.cwd(directory)
  99. await $`git remote add origin git@github.com:owner/adopted.git`.cwd(directory)
  100. })
  101. const project = yield* projects.resolve(ref.directory)
  102. const repeat = yield* projects.resolve(ref.directory)
  103. const adopted = yield* session.get(created.id)
  104. const nestedAdopted = yield* session.get(child.id)
  105. const page = yield* session.list({ project: project.id })
  106. const log = Array.from(yield* Stream.runCollect(logEvents(session, created.id)))
  107. expect(created.projectID).toBe(Project.ID.global)
  108. expect(project.id).toBe(Project.ID.make(Hash.fast("git-remote:github.com/owner/adopted")))
  109. expect(repeat.id).toBe(project.id)
  110. expect(page.data.map((item) => item.id)).toEqual(expect.arrayContaining([created.id, child.id]))
  111. expect(adopted).toMatchObject({
  112. projectID: project.id,
  113. location: ref,
  114. subpath: undefined,
  115. time: { updated: originalUpdated },
  116. })
  117. expect(nestedAdopted).toMatchObject({
  118. projectID: project.id,
  119. location: nested,
  120. subpath: RelativePath.make("packages/app"),
  121. })
  122. // Adoption is a project-domain fact; the session log records nothing new.
  123. expect(log.map((event) => event.type)).toEqual(["session.created"])
  124. expect(yield* session.messages({ sessionID: created.id })).toEqual([])
  125. // Repeated resolution announces the directory's identity exactly once.
  126. const announced = yield* db
  127. .select({ type: EventTable.type })
  128. .from(EventTable)
  129. .where(eq(EventTable.aggregate_id, project.id))
  130. .all()
  131. .pipe(Effect.orDie)
  132. expect(announced.map((event) => event.type)).toEqual(["worktree.resolved.1"])
  133. }),
  134. ),
  135. )
  136. it.effect("persists a missing title until one is generated or supplied", () =>
  137. Effect.gen(function* () {
  138. const session = yield* Session.Service
  139. const { db } = yield* Database.Service
  140. const created = yield* session.create({ location })
  141. const row = yield* db.select().from(SessionTable).where(eq(SessionTable.id, created.id)).get().pipe(Effect.orDie)
  142. const event = yield* db
  143. .select({ data: EventTable.data })
  144. .from(EventTable)
  145. .where(eq(EventTable.aggregate_id, created.id))
  146. .get()
  147. .pipe(Effect.orDie)
  148. expect(created.title).toBeUndefined()
  149. expect(row?.title).toBeNull()
  150. expect(event?.data).not.toHaveProperty("title")
  151. expect((yield* session.create({ location, title: "Explicit title" })).title).toBe("Explicit title")
  152. }),
  153. )
  154. it.effect("creates a fresh projected session when the ID is omitted", () =>
  155. Effect.gen(function* () {
  156. const session = yield* Session.Service
  157. const first = yield* session.create({ location })
  158. const second = yield* session.create({ location })
  159. expect(second.id).not.toBe(first.id)
  160. expect((yield* session.list()).data).toHaveLength(2)
  161. }),
  162. )
  163. it.effect("returns the original session when the ID is retried", () =>
  164. Effect.gen(function* () {
  165. const session = yield* Session.Service
  166. const input = { id, location }
  167. const first = yield* session.create(input)
  168. const retried = yield* session.create(input)
  169. expect(retried).toEqual(first)
  170. expect((yield* session.list()).data).toEqual([first])
  171. }),
  172. )
  173. it.effect("stores supplied immutable create attributes", () =>
  174. Effect.gen(function* () {
  175. const session = yield* Session.Service
  176. const workspaceID = Workspace.ID.make("wrk_test")
  177. const model = Model.Ref.make({
  178. id: Model.ID.make("sonnet"),
  179. providerID: Provider.ID.anthropic,
  180. variant: Model.VariantID.make("fast"),
  181. })
  182. expect(
  183. yield* session.create({
  184. location: Location.Ref.make({ directory: location.directory, workspaceID }),
  185. agent: Agent.ID.make("build"),
  186. model,
  187. }),
  188. ).toMatchObject({ location: { directory: location.directory, workspaceID }, agent: "build", model })
  189. }),
  190. )
  191. it.effect("inherits location from an existing parent when omitted", () =>
  192. Effect.gen(function* () {
  193. const session = yield* Session.Service
  194. const parent = yield* session.create({ location })
  195. const child = yield* session.create({ parentID: parent.id, title: "child" })
  196. expect(child).toMatchObject({ parentID: parent.id, location })
  197. }),
  198. )
  199. it.effect("rejects child creation when the parent does not exist", () =>
  200. Effect.gen(function* () {
  201. const session = yield* Session.Service
  202. const missing = Session.ID.create()
  203. expect(yield* Effect.flip(session.create({ parentID: missing, title: "child" }))).toEqual(
  204. new Session.NotFoundError({ sessionID: missing }),
  205. )
  206. }),
  207. )
  208. it.effect("filters root sessions before applying the page limit", () =>
  209. Effect.gen(function* () {
  210. const session = yield* Session.Service
  211. const { db } = yield* Database.Service
  212. const staleRoot = yield* session.create({ location, title: "stale root" })
  213. const root = yield* session.create({ location, title: "root" })
  214. const children = yield* Effect.forEach(Array.from({ length: 60 }), (_, index) =>
  215. session.create({ parentID: root.id, title: `child ${index}` }),
  216. )
  217. yield* Effect.forEach(children, (item, index) =>
  218. db
  219. .update(SessionTable)
  220. .set({ time_created: index + 100, time_updated: index + 20_000 })
  221. .where(eq(SessionTable.id, item.id))
  222. .run(),
  223. )
  224. yield* db
  225. .update(SessionTable)
  226. .set({ time_created: 2, time_updated: 5_000 })
  227. .where(eq(SessionTable.id, staleRoot.id))
  228. .run()
  229. yield* db
  230. .update(SessionTable)
  231. .set({ time_created: 1, time_updated: 10_000 })
  232. .where(eq(SessionTable.id, root.id))
  233. .run()
  234. const page = yield* session.list({ directory: location.directory, parentID: null, limit: 1, order: "desc" })
  235. expect(page.data.map((item) => item.id)).toEqual([root.id])
  236. }),
  237. )
  238. it.effect("orders sessions by their latest prompt", () =>
  239. Effect.gen(function* () {
  240. const session = yield* Session.Service
  241. const { db } = yield* Database.Service
  242. const active = yield* session.create({ location, title: "active" })
  243. const newer = yield* session.create({ location, title: "newer" })
  244. yield* db
  245. .update(SessionTable)
  246. .set({ time_created: -2, time_updated: -2 })
  247. .where(eq(SessionTable.id, active.id))
  248. .run()
  249. yield* db
  250. .update(SessionTable)
  251. .set({ time_created: -1, time_updated: -1 })
  252. .where(eq(SessionTable.id, newer.id))
  253. .run()
  254. yield* session.prompt({ sessionID: active.id, text: "continue", resume: false })
  255. expect((yield* session.list()).data.map((item) => item.id)).toEqual([active.id, newer.id])
  256. }),
  257. )
  258. it.effect("filters direct child sessions by parent ID", () =>
  259. Effect.gen(function* () {
  260. const session = yield* Session.Service
  261. const parent = yield* session.create({ location, title: "parent" })
  262. const child = yield* session.create({ parentID: parent.id, title: "child" })
  263. yield* session.create({ location, title: "other root" })
  264. const page = yield* session.list({ parentID: parent.id })
  265. expect(page.data.map((item) => item.id)).toEqual([child.id])
  266. }),
  267. )
  268. it.effect("filters project sessions by subpath", () =>
  269. Effect.gen(function* () {
  270. const session = yield* Session.Service
  271. const { db } = yield* Database.Service
  272. const root = yield* session.create({ location, title: "root" })
  273. const nested = yield* session.create({ location, title: "nested" })
  274. yield* db.update(SessionTable).set({ path: "packages/tui" }).where(eq(SessionTable.id, nested.id)).run()
  275. const page = yield* session.list({
  276. project: Project.ID.global,
  277. subpath: RelativePath.make("packages/tui"),
  278. parentID: null,
  279. })
  280. expect(page.data.map((item) => item.id)).toEqual([nested.id])
  281. expect(page.data.map((item) => item.id)).not.toContain(root.id)
  282. }),
  283. )
  284. it.effect("forks a session by replaying a durable fork event into copied projected rows", () =>
  285. Effect.gen(function* () {
  286. const session = yield* Session.Service
  287. const bus = yield* Bus.Service
  288. const { db } = yield* Database.Service
  289. const parent = yield* session.create({ location, title: "Parent" })
  290. const admitted = yield* session.prompt({
  291. sessionID: parent.id,
  292. text: "First",
  293. resume: false,
  294. })
  295. yield* SessionInbox.promote(db, bus, parent.id, "steer")
  296. yield* session.synthetic({ sessionID: parent.id, text: "parent note", resume: false })
  297. yield* SessionInbox.promote(db, bus, parent.id, "steer")
  298. const forked = yield* session.fork({ sessionID: parent.id, boundary: { type: "through" } })
  299. const parentContext = yield* session.context(parent.id)
  300. const forkContext = yield* session.context(forked.id)
  301. const history = Array.from(yield* Stream.runCollect(logEvents(session, forked.id)))
  302. expect(forked).toMatchObject({ title: "Parent (fork #1)", fork: { sessionID: parent.id } })
  303. expect(forked.parentID).toBeUndefined()
  304. expect(forkContext).toMatchObject([
  305. { type: "user", text: "First" },
  306. { type: "synthetic", text: "parent note" },
  307. ])
  308. expect(forkContext.map((message) => message.id)).not.toEqual(parentContext.map((message) => message.id))
  309. expect(history).toHaveLength(1)
  310. expect(history[0]).toMatchObject({
  311. type: "session.forked",
  312. durable: { seq: 0 },
  313. data: { sessionID: forked.id, parentID: parent.id },
  314. })
  315. expect(yield* SessionInbox.find(db, forkContext[0].id)).toBeUndefined()
  316. expect(yield* SessionInbox.find(db, forkContext[1].id)).toBeUndefined()
  317. expect(
  318. yield* session.prompt({ id: forkContext[0].id, sessionID: forked.id, text: "First", resume: false }),
  319. ).toMatchObject({ id: forkContext[0].id, type: "user", payload: { text: "First" } })
  320. yield* session.prompt({
  321. sessionID: parent.id,
  322. text: "Parent changed",
  323. resume: false,
  324. })
  325. yield* SessionInbox.promote(db, bus, parent.id, "steer")
  326. yield* session.prompt({
  327. sessionID: forked.id,
  328. text: "Child continues",
  329. resume: false,
  330. })
  331. yield* SessionInbox.promote(db, bus, forked.id, "steer")
  332. expect((yield* session.context(parent.id)).map((message) => message.type)).toEqual(["user", "synthetic", "user"])
  333. expect((yield* session.context(forked.id)).map((message) => message.type)).toEqual(["user", "synthetic", "user"])
  334. expect((yield* session.context(forked.id)).at(-1)).toMatchObject({ text: "Child continues" })
  335. expect(
  336. Array.from(yield* Stream.runCollect(logEvents(session, forked.id))).map(
  337. (event): number | undefined => event.durable?.seq,
  338. ),
  339. ).toEqual([0, 5, 6])
  340. expect(yield* SessionInbox.find(db, admitted.id)).toBeUndefined()
  341. }),
  342. )
  343. it.effect("keeps a fork untitled when its parent is untitled", () =>
  344. Effect.gen(function* () {
  345. const session = yield* Session.Service
  346. const bus = yield* Bus.Service
  347. const { db } = yield* Database.Service
  348. const parent = yield* session.create({ location })
  349. yield* session.prompt({ sessionID: parent.id, text: "First", resume: false })
  350. yield* SessionInbox.promote(db, bus, parent.id, "steer")
  351. const forked = yield* session.fork({ sessionID: parent.id, boundary: { type: "through" } })
  352. const row = yield* db.select().from(SessionTable).where(eq(SessionTable.id, forked.id)).get().pipe(Effect.orDie)
  353. expect(forked.title).toBeUndefined()
  354. expect(row?.title).toBeNull()
  355. }),
  356. )
  357. it.effect("rejects forking an empty session", () =>
  358. Effect.gen(function* () {
  359. const session = yield* Session.Service
  360. const parent = yield* session.create({ location })
  361. expect(
  362. yield* session.fork({ sessionID: parent.id, boundary: { type: "through" } }).pipe(Effect.flip),
  363. ).toMatchObject({ _tag: "Session.ForkEmptyError", sessionID: parent.id })
  364. }),
  365. )
  366. it.effect("forks before the selected boundary message", () =>
  367. Effect.gen(function* () {
  368. const session = yield* Session.Service
  369. const bus = yield* Bus.Service
  370. const { db } = yield* Database.Service
  371. const parent = yield* session.create({ location })
  372. const first = yield* session.prompt({
  373. sessionID: parent.id,
  374. text: "First",
  375. resume: false,
  376. })
  377. yield* SessionInbox.promote(db, bus, parent.id, "steer")
  378. const second = yield* session.prompt({
  379. sessionID: parent.id,
  380. text: "Second",
  381. resume: false,
  382. })
  383. yield* SessionInbox.promote(db, bus, parent.id, "steer")
  384. const assistantMessageID = SessionMessage.ID.create()
  385. const model = Model.Ref.make({ id: Model.ID.make("model"), providerID: Provider.ID.make("provider") })
  386. yield* bus.publish(SessionEvent.Step.Started, {
  387. sessionID: parent.id,
  388. assistantMessageID,
  389. agent: Agent.ID.make("build"),
  390. model,
  391. })
  392. yield* bus.publish(SessionEvent.Step.Ended, {
  393. sessionID: parent.id,
  394. assistantMessageID,
  395. finish: "stop",
  396. cost: Money.USD.make(0.75),
  397. tokens: { input: 6, output: 3, reasoning: 1, cache: { read: 2, write: 1 } },
  398. })
  399. const forked = yield* session.fork({
  400. sessionID: parent.id,
  401. boundary: { type: "before", messageID: second.id },
  402. })
  403. const beforeFirst = yield* session.fork({
  404. sessionID: parent.id,
  405. boundary: { type: "before", messageID: first.id },
  406. })
  407. const complete = yield* session.fork({ sessionID: parent.id, boundary: { type: "through" } })
  408. const context = yield* session.context(forked.id)
  409. const history = Array.from(yield* Stream.runCollect(logEvents(session, forked.id)))
  410. expect(forked.fork).toEqual({
  411. sessionID: parent.id,
  412. boundary: { type: "before", messageID: second.id },
  413. })
  414. expect(context).toMatchObject([{ text: "First" }])
  415. expect(context[0]?.id).not.toBe(first.id)
  416. expect(history[0]).toMatchObject({
  417. data: { boundary: { type: "before", messageID: second.id } },
  418. })
  419. expect(forked).toMatchObject({ cost: 0, tokens: { input: 0, output: 0, reasoning: 0 } })
  420. expect(yield* session.context(beforeFirst.id)).toEqual([])
  421. expect(beforeFirst).toMatchObject({ cost: 0, tokens: { input: 0, output: 0, reasoning: 0 } })
  422. expect(complete).toMatchObject({
  423. cost: 0,
  424. tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
  425. })
  426. }),
  427. )
  428. it.effect("returns the existing Session when one ID is reused with different create arguments", () =>
  429. Effect.gen(function* () {
  430. const session = yield* Session.Service
  431. const created = yield* session.create({ id, location })
  432. const changed = [
  433. { id, location: Location.Ref.make({ directory: AbsolutePath.make("/other") }) },
  434. { id, location, agent: Agent.ID.make("build") },
  435. {
  436. id,
  437. location,
  438. model: Model.Ref.make({ id: Model.ID.make("sonnet"), providerID: Provider.ID.anthropic }),
  439. },
  440. ]
  441. for (const input of changed) {
  442. expect(yield* session.create(input)).toEqual(created)
  443. }
  444. expect((yield* session.list()).data).toHaveLength(1)
  445. }),
  446. )
  447. it.effect("returns one recorded session to concurrent exact retries", () =>
  448. Effect.gen(function* () {
  449. const session = yield* Session.Service
  450. const input = { id, location }
  451. const created = yield* Effect.all([session.create(input), session.create(input)], { concurrency: "unbounded" })
  452. expect(created[1]).toEqual(created[0])
  453. expect((yield* session.list()).data).toEqual([created[0]])
  454. }),
  455. )
  456. it.effect("returns the current Session projection after updates", () =>
  457. Effect.gen(function* () {
  458. const session = yield* Session.Service
  459. const { db } = yield* Database.Service
  460. const input = { id, location }
  461. const created = yield* session.create(input)
  462. yield* db.update(SessionTable).set({ agent: "build" }).where(eq(SessionTable.id, id)).run().pipe(Effect.orDie)
  463. expect(yield* session.create(input)).toMatchObject({ id: created.id, agent: "build" })
  464. }),
  465. )
  466. it.effect("persists creation through the current created event", () =>
  467. Effect.gen(function* () {
  468. const session = yield* Session.Service
  469. const { db } = yield* Database.Service
  470. const created = yield* session.create({ location })
  471. expect(
  472. yield* db.select().from(EventTable).where(eq(EventTable.aggregate_id, created.id)).all().pipe(Effect.orDie),
  473. ).toMatchObject([{ type: Bus.versionedType(SessionEvent.Created.type, 1) }])
  474. }),
  475. )
  476. it.effect("persists caller-ID creation through the existing created event", () =>
  477. Effect.gen(function* () {
  478. const session = yield* Session.Service
  479. const { db } = yield* Database.Service
  480. const created = yield* session.create({ id, location })
  481. expect(
  482. yield* db.select().from(EventTable).where(eq(EventTable.aggregate_id, created.id)).get().pipe(Effect.orDie),
  483. ).toMatchObject({
  484. data: { sessionID: id },
  485. })
  486. }),
  487. )
  488. it.effect("includes current creation rows in the Session event stream", () =>
  489. Effect.gen(function* () {
  490. const session = yield* Session.Service
  491. const bus = yield* Bus.Service
  492. const { db } = yield* Database.Service
  493. const created = yield* session.create({ location })
  494. yield* session.prompt({
  495. sessionID: created.id,
  496. text: "Hello",
  497. resume: false,
  498. })
  499. yield* SessionInbox.promote(db, bus, created.id, "steer")
  500. expect(
  501. Array.from(yield* logEvents(session, created.id, true).pipe(Stream.take(3), Stream.runCollect)),
  502. ).toMatchObject([
  503. { durable: { seq: 0 }, type: "session.created" },
  504. {
  505. durable: { seq: 1 },
  506. type: "session.inbox.enqueued",
  507. data: {
  508. inboxID: expect.any(String),
  509. item: { type: "user", payload: { text: "Hello" }, delivery: "steer" },
  510. },
  511. },
  512. { durable: { seq: 2 }, type: "session.inbox.delivered" },
  513. ])
  514. }),
  515. )
  516. it.effect("replays one prompt lifecycle into a fresh target database", () =>
  517. Effect.gen(function* () {
  518. const session = yield* Session.Service
  519. const sourceEvents = yield* Bus.Service
  520. const sourceDb = (yield* Database.Service).db
  521. const created = yield* session.create({ id: Session.ID.make("ses_fresh_target_replay"), location })
  522. const admitted = yield* session.prompt({
  523. sessionID: created.id,
  524. text: "Replay lifecycle",
  525. resume: false,
  526. })
  527. yield* SessionInbox.promote(sourceDb, sourceEvents, created.id, "steer")
  528. const serialized = (yield* sourceDb
  529. .select()
  530. .from(EventTable)
  531. .where(eq(EventTable.aggregate_id, created.id))
  532. .orderBy(asc(EventTable.seq))
  533. .all()
  534. .pipe(Effect.orDie)).map((event) => ({
  535. id: event.id,
  536. created: event.created,
  537. aggregateID: event.aggregate_id,
  538. seq: event.seq,
  539. type: event.type,
  540. data: event.data,
  541. }))
  542. const tmp = yield* Effect.acquireRelease(
  543. Effect.promise(() => tmpdir()),
  544. (tmp) => Effect.promise(() => tmp[Symbol.asyncDispose]()),
  545. )
  546. const targetLayer = AppNodeBuilder.build(
  547. LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node]),
  548. [
  549. [Database.node, Database.configured({ path: path.join(tmp.path, "target.sqlite") })],
  550. [Bus.node, Bus.configured({ persist: true })],
  551. ],
  552. )
  553. yield* Effect.gen(function* () {
  554. const db = (yield* Database.Service).db
  555. const bus = yield* Bus.Service
  556. const store = yield* SessionStore.Service
  557. yield* db
  558. .insert(ProjectTable)
  559. .values({ id: Project.ID.global, worktree: location.directory, sandboxes: [] })
  560. .run()
  561. .pipe(Effect.orDie)
  562. expect(yield* store.get(created.id)).toBeUndefined()
  563. yield* Effect.forEach(serialized.slice(0, 2), (event) => bus.replay(event), { discard: true })
  564. expect(yield* SessionInbox.find(db, admitted.id)).toMatchObject({
  565. id: admitted.id,
  566. sessionID: created.id,
  567. type: "user",
  568. payload: { text: "Replay lifecycle" },
  569. delivery: "steer",
  570. })
  571. expect(yield* store.context(created.id)).toEqual([])
  572. yield* Effect.forEach(serialized.slice(2), (event) => bus.replay(event), { discard: true })
  573. expect(yield* SessionInbox.find(db, admitted.id)).toBeUndefined()
  574. expect(yield* store.context(created.id)).toMatchObject([
  575. { id: admitted.id, type: "user", text: "Replay lifecycle" },
  576. ])
  577. expect(
  578. (yield* db
  579. .select()
  580. .from(EventTable)
  581. .where(eq(EventTable.aggregate_id, created.id))
  582. .orderBy(asc(EventTable.seq))
  583. .all()
  584. .pipe(Effect.orDie)).map((event) => [event.seq, event.type]),
  585. ).toEqual([
  586. [0, Bus.versionedType(SessionEvent.Created.type, 1)],
  587. [1, Bus.versionedType(SessionEvent.InboxEnqueued.type, 1)],
  588. [2, Bus.versionedType(SessionEvent.InboxDelivered.type, 1)],
  589. ])
  590. }).pipe(Effect.provide(Layer.fresh(targetLayer)))
  591. }),
  592. )
  593. it.effect("does not mask unrelated created projector defects", () =>
  594. Effect.gen(function* () {
  595. const session = yield* Session.Service
  596. const event = yield* Bus.Service
  597. const defect = new Error("unrelated projector defect")
  598. yield* event.project(SessionEvent.Created, () => Effect.die(defect))
  599. expect(yield* session.create({ id, location }).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect)
  600. }),
  601. )
  602. it.live("runs a shell command and projects the started/ended shell message", () =>
  603. withTmp((directory) =>
  604. Effect.gen(function* () {
  605. const session = yield* Session.Service
  606. const created = yield* session.create({
  607. location: Location.Ref.make({ directory: AbsolutePath.make(directory) }),
  608. })
  609. yield* session.shell({ sessionID: created.id, command: "echo hello" })
  610. const messages = yield* session.messages({ sessionID: created.id, order: "asc" })
  611. const shell = messages.find((message): message is SessionMessage.Shell => message.type === "shell")
  612. expect(shell).toMatchObject({ type: "shell", command: "echo hello", status: "exited", exit: 0 })
  613. expect(shell?.output?.output).toContain("hello")
  614. expect(shell?.output?.truncated).toBe(false)
  615. expect(shell?.time.completed).toBeDefined()
  616. }),
  617. ),
  618. )
  619. it.live("still emits shell ended for a failing command", () =>
  620. withTmp((directory) =>
  621. Effect.gen(function* () {
  622. const session = yield* Session.Service
  623. const created = yield* session.create({
  624. location: Location.Ref.make({ directory: AbsolutePath.make(directory) }),
  625. })
  626. yield* session.shell({ sessionID: created.id, command: "false" })
  627. const messages = yield* session.messages({ sessionID: created.id, order: "asc" })
  628. const shell = messages.find((message): message is SessionMessage.Shell => message.type === "shell")
  629. expect(shell).toMatchObject({ type: "shell", command: "false", status: "exited" })
  630. expect(shell?.exit).not.toBe(0)
  631. expect(shell?.time.completed).toBeDefined()
  632. }),
  633. ),
  634. )
  635. it.effect("switches the selected agent through the durable Session event", () =>
  636. Effect.gen(function* () {
  637. const session = yield* Session.Service
  638. const created = yield* session.create({ location, agent: Agent.ID.make("build") })
  639. yield* session.switchAgent({ sessionID: created.id, agent: Agent.ID.make("plan") })
  640. expect(yield* session.get(created.id)).toMatchObject({ agent: "plan" })
  641. expect(
  642. Array.from(yield* logEvents(session, created.id, true).pipe(Stream.drop(1), Stream.take(1), Stream.runCollect)),
  643. ).toMatchObject([{ type: "session.agent.selected", data: { agent: "plan", previous: "build" } }])
  644. expect(yield* session.messages({ sessionID: created.id, order: "asc" })).toMatchObject([
  645. { type: "agent-switched", agent: "plan", previous: "build" },
  646. ])
  647. }),
  648. )
  649. it.effect("rejects an agent switch for a missing Session", () =>
  650. Effect.gen(function* () {
  651. const session = yield* Session.Service
  652. const missing = Session.ID.make("ses_missing_agent_switch")
  653. expect(
  654. yield* session.switchAgent({ sessionID: missing, agent: Agent.ID.make("plan") }).pipe(
  655. Effect.flip,
  656. Effect.map((error) => error._tag),
  657. ),
  658. ).toBe("Session.NotFoundError")
  659. }),
  660. )
  661. it.effect("switches the selected model through the durable Session event", () =>
  662. Effect.gen(function* () {
  663. const session = yield* Session.Service
  664. const previous = Model.Ref.make({
  665. id: Model.ID.make("haiku"),
  666. providerID: Provider.ID.anthropic,
  667. variant: Model.VariantID.make("default"),
  668. })
  669. const created = yield* session.create({ location, model: previous })
  670. const model = Model.Ref.make({
  671. id: Model.ID.make("sonnet"),
  672. providerID: Provider.ID.anthropic,
  673. variant: Model.VariantID.make("high"),
  674. })
  675. yield* session.switchModel({ sessionID: created.id, model })
  676. expect(yield* session.get(created.id)).toMatchObject({ model })
  677. const bus = Array.from(
  678. yield* logEvents(session, created.id, true).pipe(Stream.drop(1), Stream.take(1), Stream.runCollect),
  679. )
  680. expect(bus).toMatchObject([{ type: "session.model.selected" }])
  681. expect(bus[0]?.data).toEqual({ sessionID: created.id, model, previous })
  682. expect(yield* session.messages({ sessionID: created.id, order: "asc" })).toMatchObject([
  683. { type: "model-switched", model, previous },
  684. ])
  685. }),
  686. )
  687. it.effect("ignores a model switch when the selected model is unchanged", () =>
  688. Effect.gen(function* () {
  689. const session = yield* Session.Service
  690. const created = yield* session.create({ location })
  691. const model = Model.Ref.make({ id: Model.ID.make("sonnet"), providerID: Provider.ID.anthropic })
  692. yield* session.switchModel({ sessionID: created.id, model })
  693. yield* session.switchModel({ sessionID: created.id, model })
  694. const { db } = yield* Database.Service
  695. expect(
  696. yield* db.select().from(EventTable).where(eq(EventTable.aggregate_id, created.id)).all().pipe(Effect.orDie),
  697. ).toHaveLength(2)
  698. expect(yield* session.get(created.id)).toMatchObject({ model })
  699. }),
  700. )
  701. it.effect("treats an omitted variant as the default variant", () =>
  702. Effect.gen(function* () {
  703. const session = yield* Session.Service
  704. const model = Model.Ref.make({ id: Model.ID.make("sonnet"), providerID: Provider.ID.anthropic })
  705. const created = yield* session.create({ location, model })
  706. yield* session.switchModel({
  707. sessionID: created.id,
  708. model: Model.Ref.make({ ...model, variant: Model.VariantID.make("default") }),
  709. })
  710. const { db } = yield* Database.Service
  711. expect(
  712. yield* db.select().from(EventTable).where(eq(EventTable.aggregate_id, created.id)).all().pipe(Effect.orDie),
  713. ).toHaveLength(1)
  714. }),
  715. )
  716. it.effect("rejects a model switch for a missing Session", () =>
  717. Effect.gen(function* () {
  718. const session = yield* Session.Service
  719. const missing = Session.ID.make("ses_missing_model_switch")
  720. expect(
  721. yield* session
  722. .switchModel({
  723. sessionID: missing,
  724. model: Model.Ref.make({ id: Model.ID.make("sonnet"), providerID: Provider.ID.anthropic }),
  725. })
  726. .pipe(
  727. Effect.flip,
  728. Effect.map((error) => error._tag),
  729. ),
  730. ).toBe("Session.NotFoundError")
  731. }),
  732. )
  733. })
  734. describe("SessionTransfer", () => {
  735. it.effect("imports projected messages and reserves their aggregate sequence", () =>
  736. Effect.gen(function* () {
  737. const session = yield* Session.Service
  738. const transfer = yield* SessionTransfer.Service
  739. const bus = yield* Bus.Service
  740. const { db } = yield* Database.Service
  741. const template = yield* session.create({ location, title: "Exported" })
  742. const sessionID = Session.ID.create()
  743. const sourceMessageID = SessionMessage.ID.create()
  744. const errorMessageID = SessionMessage.ID.create()
  745. const imported = yield* transfer.import({
  746. data: {
  747. info: { ...template, id: sessionID },
  748. messages: [
  749. {
  750. id: sourceMessageID,
  751. type: "user",
  752. text: "Imported message",
  753. time: { created: DateTime.makeUnsafe(100) },
  754. },
  755. {
  756. id: errorMessageID,
  757. type: "compaction",
  758. status: "failed",
  759. reason: "manual",
  760. error: { type: "test_error", message: "Original error" },
  761. time: { created: DateTime.makeUnsafe(101) },
  762. },
  763. ],
  764. },
  765. location,
  766. })
  767. const messages = yield* session.messages({ sessionID, order: "asc" })
  768. expect(imported).toMatchObject({ id: sessionID, title: "Exported", location })
  769. expect(messages).toMatchObject([
  770. { id: sourceMessageID, type: "user", text: "Imported message" },
  771. { id: errorMessageID, type: "compaction", error: { type: "test_error", message: "Original error" } },
  772. ])
  773. expect(yield* Bus.latestSequence(db, sessionID)).toBe(2)
  774. expect((yield* transfer.export({ sessionID })).messages).toEqual(messages)
  775. expect((yield* transfer.export({ sessionID, sanitize: true })).messages).toMatchObject([
  776. { id: sourceMessageID, text: `[redacted:text:${sourceMessageID}]` },
  777. { id: errorMessageID, error: { type: "test_error", message: "Original error" } },
  778. ])
  779. yield* session.prompt({ sessionID, text: "Continue", resume: false })
  780. yield* SessionInbox.promote(db, bus, sessionID, "steer")
  781. expect((yield* session.messages({ sessionID, order: "asc" })).map((message) => message.type)).toEqual([
  782. "user",
  783. "compaction",
  784. "user",
  785. ])
  786. expect(yield* Bus.latestSequence(db, sessionID)).toBe(4)
  787. }),
  788. )
  789. it.effect("rejects an existing session ID without changing its transcript", () =>
  790. Effect.gen(function* () {
  791. const session = yield* Session.Service
  792. const transfer = yield* SessionTransfer.Service
  793. const existing = yield* session.create({ location, title: "Existing" })
  794. const exit = yield* Effect.exit(transfer.import({ data: { info: existing, messages: [] }, location }))
  795. expect(exit._tag).toBe("Failure")
  796. expect((yield* session.get(existing.id)).title).toBe("Existing")
  797. expect(yield* session.messages({ sessionID: existing.id })).toEqual([])
  798. }),
  799. )
  800. })