session-create.test.ts 32 KB

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