session-log.test.ts 5.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130
  1. import { describe, expect } from "bun:test"
  2. import { Effect, Fiber, Layer, Schema, Stream } from "effect"
  3. import { Database } from "@opencode-ai/core/database/database"
  4. import { Agent } from "@opencode-ai/core/agent"
  5. import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
  6. import { LayerNode } from "@opencode-ai/util/effect/layer-node"
  7. import { Bus } from "@opencode-ai/core/bus"
  8. import { Event } from "@opencode-ai/schema/event"
  9. import { Location } from "@opencode-ai/core/location"
  10. import { Project } from "@opencode-ai/core/project"
  11. import { ProjectTable } from "@opencode-ai/core/project/sql"
  12. import { AbsolutePath } from "@opencode-ai/core/schema"
  13. import { Session } from "@opencode-ai/core/session"
  14. import { SessionProjector } from "@opencode-ai/core/session/projector"
  15. import { SessionExecution } from "@opencode-ai/core/session/execution"
  16. import { SessionStore } from "@opencode-ai/core/session/store"
  17. import { SessionTable } from "@opencode-ai/core/session/sql"
  18. import { testEffect } from "./lib/effect"
  19. const projects = Layer.succeed(
  20. Project.Service,
  21. Project.Service.of({
  22. list: () => Effect.succeed([]),
  23. resolve: (directory) => Effect.succeed({ id: Project.ID.global, directory, canonical: directory }),
  24. directories: () => Effect.succeed([]),
  25. }),
  26. )
  27. const it = testEffect(
  28. AppNodeBuilder.build(
  29. LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node, Session.node]),
  30. [
  31. [Bus.node, Bus.configured({ persist: true })],
  32. [Project.node, projects],
  33. [SessionExecution.node, SessionExecution.noopLayer],
  34. ],
  35. ),
  36. )
  37. const location = Location.Ref.make({ directory: AbsolutePath.make("/project") })
  38. describe("Session.log", () => {
  39. it.effect("replays public session events and marks synced at the aggregate watermark", () =>
  40. Effect.gen(function* () {
  41. const session = yield* Session.Service
  42. const created = yield* session.create({ location })
  43. yield* session.rename({ sessionID: created.id, title: "session.renamed" })
  44. const items = Array.from(yield* Stream.runCollect(session.log({ sessionID: created.id })))
  45. expect(items.map((item) => item.type)).toEqual(["session.created", "session.renamed", "log.synced"])
  46. expect(items.at(-1)).toEqual({ type: "log.synced", aggregateID: created.id, seq: Event.Seq.make(1) })
  47. }),
  48. )
  49. it.effect("continues with live public events when following", () =>
  50. Effect.gen(function* () {
  51. const session = yield* Session.Service
  52. const created = yield* session.create({ location })
  53. const fiber = yield* session
  54. .log({ sessionID: created.id, after: Event.Seq.make(0), follow: true })
  55. .pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
  56. yield* Effect.yieldNow
  57. yield* session.rename({ sessionID: created.id, title: "renamed live" })
  58. const items = Array.from(yield* Fiber.join(fiber))
  59. expect(items.map((item) => item.type)).toEqual(["log.synced", "session.renamed"])
  60. }),
  61. )
  62. it.effect("fails with NotFound for an unknown session", () =>
  63. Effect.gen(function* () {
  64. const session = yield* Session.Service
  65. const error = yield* Effect.flip(Stream.runCollect(session.log({ sessionID: Session.ID.create() })))
  66. expect(error._tag).toBe("Session.NotFoundError")
  67. }),
  68. )
  69. it.effect("reads across undecodable gaps in aggregate order and marks the true log position", () =>
  70. Effect.gen(function* () {
  71. const GapEvent = Bus.durable({
  72. type: "test.session.log.gap",
  73. durable: { aggregate: "sessionID", version: 1 },
  74. schema: { sessionID: Session.ID, value: Schema.String },
  75. })
  76. const session = yield* Session.Service
  77. const bus = yield* Bus.Service
  78. const created = yield* session.create({ location })
  79. yield* session.switchAgent({ sessionID: created.id, agent: Agent.ID.make("one") })
  80. // Not in the durable manifest, so reads must skip it without failing.
  81. yield* bus.publish(GapEvent, { sessionID: created.id, value: "filtered" })
  82. yield* session.switchAgent({ sessionID: created.id, agent: Agent.ID.make("two") })
  83. yield* session.switchAgent({ sessionID: created.id, agent: Agent.ID.make("three") })
  84. const items = Array.from(yield* Stream.runCollect(session.log({ sessionID: created.id, after: 1 })))
  85. expect(
  86. items.map((item): number | string | undefined => (Bus.isSynced(item) ? item.type : item.durable?.seq)),
  87. ).toEqual([3, 4, "log.synced"])
  88. expect(items.at(-1)).toEqual({ type: "log.synced", aggregateID: created.id, seq: Event.Seq.make(4) })
  89. }),
  90. )
  91. it.effect("completes with a bare synced marker for a migrated Session with no event sequence", () =>
  92. Effect.gen(function* () {
  93. const db = (yield* Database.Service).db
  94. const session = yield* Session.Service
  95. const sessionID = Session.ID.make("ses_empty_log")
  96. yield* db
  97. .insert(ProjectTable)
  98. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  99. .onConflictDoNothing()
  100. .run()
  101. yield* db
  102. .insert(SessionTable)
  103. .values({
  104. id: sessionID,
  105. project_id: Project.ID.global,
  106. slug: "empty-log",
  107. directory: "/project",
  108. title: "Empty log",
  109. version: "test",
  110. })
  111. .run()
  112. const items = Array.from(yield* Stream.runCollect(session.log({ sessionID })))
  113. expect(items).toEqual([{ type: "log.synced", aggregateID: sessionID }])
  114. }),
  115. )
  116. })