session-log.test.ts 5.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134
  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. commit: () => Effect.void,
  26. }),
  27. )
  28. const it = testEffect(
  29. AppNodeBuilder.build(
  30. LayerNode.group([Database.node, Bus.node, SessionProjector.node, SessionStore.node, Session.node]),
  31. [
  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 bus = yield* Bus.Service
  43. const created = yield* session.create({ location })
  44. yield* session.rename({ sessionID: created.id, title: "session.renamed" })
  45. const items = Array.from(yield* Stream.runCollect(session.log({ sessionID: created.id })))
  46. const watermark = (yield* bus.sequences([created.id])).get(created.id)
  47. // Session creation commits a non-public durable event, so the marker's
  48. // seq covers more of the aggregate than the public events emitted.
  49. expect(items.map((item) => item.type)).toEqual(["session.renamed", "log.synced"])
  50. expect(items.at(-1)).toEqual({ type: "log.synced", aggregateID: created.id, seq: watermark })
  51. }),
  52. )
  53. it.effect("continues with live public events when following", () =>
  54. Effect.gen(function* () {
  55. const session = yield* Session.Service
  56. const created = yield* session.create({ location })
  57. const fiber = yield* session
  58. .log({ sessionID: created.id, follow: true })
  59. .pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
  60. yield* Effect.yieldNow
  61. yield* session.rename({ sessionID: created.id, title: "renamed live" })
  62. const items = Array.from(yield* Fiber.join(fiber))
  63. expect(items.map((item) => item.type)).toEqual(["log.synced", "session.renamed"])
  64. }),
  65. )
  66. it.effect("fails with NotFound for an unknown session", () =>
  67. Effect.gen(function* () {
  68. const session = yield* Session.Service
  69. const error = yield* Effect.flip(Stream.runCollect(session.log({ sessionID: Session.ID.create() })))
  70. expect(error._tag).toBe("Session.NotFoundError")
  71. }),
  72. )
  73. it.effect("reads across undecodable gaps in aggregate order and marks the true log position", () =>
  74. Effect.gen(function* () {
  75. const GapEvent = Bus.durable({
  76. type: "test.session.log.gap",
  77. durable: { aggregate: "sessionID", version: 1 },
  78. schema: { sessionID: Session.ID, value: Schema.String },
  79. })
  80. const session = yield* Session.Service
  81. const bus = yield* Bus.Service
  82. const created = yield* session.create({ location })
  83. yield* session.switchAgent({ sessionID: created.id, agent: Agent.ID.make("one") })
  84. // Not in the durable manifest, so reads must skip it without failing.
  85. yield* bus.publish(GapEvent, { sessionID: created.id, value: "filtered" })
  86. yield* session.switchAgent({ sessionID: created.id, agent: Agent.ID.make("two") })
  87. yield* session.switchAgent({ sessionID: created.id, agent: Agent.ID.make("three") })
  88. const items = Array.from(yield* Stream.runCollect(session.log({ sessionID: created.id, after: 1 })))
  89. expect(
  90. items.map((item): number | string | undefined => (Bus.isSynced(item) ? item.type : item.durable?.seq)),
  91. ).toEqual([3, 4, "log.synced"])
  92. expect(items.at(-1)).toEqual({ type: "log.synced", aggregateID: created.id, seq: Event.Seq.make(4) })
  93. }),
  94. )
  95. it.effect("completes with a bare synced marker for a migrated Session with no event sequence", () =>
  96. Effect.gen(function* () {
  97. const db = (yield* Database.Service).db
  98. const session = yield* Session.Service
  99. const sessionID = Session.ID.make("ses_empty_log")
  100. yield* db
  101. .insert(ProjectTable)
  102. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  103. .onConflictDoNothing()
  104. .run()
  105. yield* db
  106. .insert(SessionTable)
  107. .values({
  108. id: sessionID,
  109. project_id: Project.ID.global,
  110. slug: "empty-log",
  111. directory: "/project",
  112. title: "Empty log",
  113. version: "test",
  114. })
  115. .run()
  116. const items = Array.from(yield* Stream.runCollect(session.log({ sessionID })))
  117. expect(items).toEqual([{ type: "log.synced", aggregateID: sessionID }])
  118. }),
  119. )
  120. })