session-log.test.ts 5.6 KB

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