session.test.ts 9.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253
  1. import { describe, expect } from "bun:test"
  2. import { SessionV1 } from "@kirincode-ai/core/v1/session"
  3. import { EventV2 } from "@kirincode-ai/core/event"
  4. import { SessionProjector } from "@kirincode-ai/core/session/projector"
  5. import { Deferred, Effect, Exit, Layer } from "effect"
  6. import { Session as SessionNs } from "@/session/session"
  7. import { MessageV2 } from "../../src/session/message-v2"
  8. import { MessageID, PartID, type SessionID } from "../../src/session/schema"
  9. import { CrossSpawnSpawner } from "@kirincode-ai/core/cross-spawn-spawner"
  10. import { provideInstance, tmpdirScoped } from "../fixture/fixture"
  11. import { testEffect } from "../lib/effect"
  12. import { RuntimeFlags } from "@/effect/runtime-flags"
  13. import { EventV2Bridge } from "@/event-v2-bridge"
  14. import { GlobalBus } from "@/bus/global"
  15. import { AppNodeBuilder } from "@kirincode-ai/core/effect/app-node-builder"
  16. import { LayerNode } from "@kirincode-ai/core/effect/layer-node"
  17. import { InstanceStore } from "@/project/instance-store"
  18. import { InstanceBootstrap } from "@/project/bootstrap"
  19. const it = testEffect(
  20. AppNodeBuilder.build(
  21. LayerNode.group([
  22. SessionNs.node,
  23. EventV2Bridge.node,
  24. SessionProjector.node,
  25. CrossSpawnSpawner.node,
  26. InstanceStore.node,
  27. ]),
  28. [
  29. [RuntimeFlags.node, RuntimeFlags.layer({ experimentalWorkspaces: false })],
  30. [
  31. InstanceBootstrap.node,
  32. Layer.succeed(InstanceBootstrap.Service, InstanceBootstrap.Service.of({ run: Effect.void })),
  33. ],
  34. ],
  35. ),
  36. )
  37. const awaitDeferred = <T>(deferred: Deferred.Deferred<T>, message: string) =>
  38. Effect.race(
  39. Deferred.await(deferred),
  40. Effect.sleep("2 seconds").pipe(Effect.flatMap(() => Effect.fail(new Error(message)))),
  41. )
  42. const remove = (id: SessionID) => SessionNs.use.remove(id)
  43. describe("session.created event", () => {
  44. it.instance("should emit session.created event when session is created", () =>
  45. Effect.gen(function* () {
  46. const session = yield* SessionNs.Service
  47. const events = yield* EventV2Bridge.Service
  48. const received = yield* Deferred.make<SessionNs.Info>()
  49. const unsub = yield* events.listen((event) => {
  50. if (event.type === SessionNs.Event.Created.type)
  51. Deferred.doneUnsafe(
  52. received,
  53. Effect.succeed((event.data as typeof SessionNs.Event.Created.data.Type).info as SessionNs.Info),
  54. )
  55. return Effect.void
  56. })
  57. yield* Effect.addFinalizer(() => unsub)
  58. const info = yield* session.create({})
  59. const receivedInfo = yield* awaitDeferred(received, "timed out waiting for session.created")
  60. expect(receivedInfo.id).toBe(info.id)
  61. expect(receivedInfo.projectID).toBe(info.projectID)
  62. expect(receivedInfo.directory).toBe(info.directory)
  63. expect(receivedInfo.path).toBe(info.path)
  64. expect(receivedInfo.title).toBe(info.title)
  65. yield* session.remove(info.id)
  66. }),
  67. )
  68. it.instance("session.created event should be emitted before session.updated", () =>
  69. Effect.gen(function* () {
  70. const session = yield* SessionNs.Service
  71. const source = yield* EventV2Bridge.Service
  72. const events: string[] = []
  73. const received = yield* Deferred.make<string[]>()
  74. const push = (event: string) => {
  75. events.push(event)
  76. if (events.includes("created") && events.includes("updated")) {
  77. Deferred.doneUnsafe(received, Effect.succeed(events))
  78. }
  79. }
  80. const unsubscribe = yield* source.listen((event) => {
  81. if (event.type === SessionNs.Event.Created.type) push("created")
  82. if (event.type === SessionNs.Event.Updated.type) push("updated")
  83. return Effect.void
  84. })
  85. yield* Effect.addFinalizer(() => unsubscribe)
  86. const info = yield* session.create({})
  87. yield* session.setTitle({ sessionID: info.id, title: "updated" })
  88. const receivedEvents = yield* awaitDeferred(received, "timed out waiting for session created/updated events")
  89. expect(receivedEvents).toContain("created")
  90. expect(receivedEvents).toContain("updated")
  91. expect(receivedEvents.indexOf("created")).toBeLessThan(receivedEvents.indexOf("updated"))
  92. yield* session.remove(info.id)
  93. }),
  94. )
  95. it.instance("emits legacy global sync payload", () =>
  96. Effect.gen(function* () {
  97. const session = yield* SessionNs.Service
  98. const received = yield* Deferred.make<{ syncEvent: EventV2.SerializedEvent }>()
  99. const listener = (event: { payload: { type?: string; syncEvent?: EventV2.SerializedEvent } }) => {
  100. if (event.payload.type === "sync" && event.payload.syncEvent)
  101. Deferred.doneUnsafe(received, Effect.succeed({ syncEvent: event.payload.syncEvent }))
  102. }
  103. GlobalBus.on("event", listener)
  104. yield* Effect.addFinalizer(() => Effect.sync(() => GlobalBus.off("event", listener)))
  105. const info = yield* session.create({})
  106. const event = yield* awaitDeferred(received, "timed out waiting for legacy global sync event")
  107. expect(event.syncEvent).toMatchObject({
  108. type: EventV2.versionedType(SessionNs.Event.Created.type, 1),
  109. seq: 0,
  110. aggregateID: info.id,
  111. data: { sessionID: info.id },
  112. })
  113. yield* session.remove(info.id)
  114. }),
  115. )
  116. })
  117. describe("step-finish token propagation via event", () => {
  118. it.instance(
  119. "non-zero tokens propagate through PartUpdated event",
  120. () =>
  121. Effect.gen(function* () {
  122. const session = yield* SessionNs.Service
  123. const events = yield* EventV2Bridge.Service
  124. const info = yield* session.create({})
  125. const messageID = MessageID.ascending()
  126. yield* session.updateMessage({
  127. id: messageID,
  128. sessionID: info.id,
  129. role: "user",
  130. time: { created: Date.now() },
  131. agent: "user",
  132. model: { providerID: "test", modelID: "test" },
  133. tools: {},
  134. mode: "",
  135. } as unknown as SessionV1.Info)
  136. // Event subscribers receive readonly Schema.Type payloads; `SessionV1.Part`
  137. // is the mutable domain type. Cast bridges the two — safe because the
  138. // test only reads the value afterwards.
  139. const received = yield* Deferred.make<SessionV1.Part>()
  140. const unsub = yield* events.listen((event) => {
  141. if (event.type === MessageV2.Event.PartUpdated.type)
  142. Deferred.doneUnsafe(
  143. received,
  144. Effect.succeed((event.data as typeof MessageV2.Event.PartUpdated.data.Type).part as SessionV1.Part),
  145. )
  146. return Effect.void
  147. })
  148. yield* Effect.addFinalizer(() => unsub)
  149. const tokens = {
  150. total: 1500,
  151. input: 500,
  152. output: 800,
  153. reasoning: 200,
  154. cache: { read: 100, write: 50 },
  155. }
  156. const partInput = {
  157. id: PartID.ascending(),
  158. messageID,
  159. sessionID: info.id,
  160. type: "step-finish" as const,
  161. reason: "stop",
  162. cost: 0.005,
  163. tokens,
  164. }
  165. yield* session.updatePart(partInput)
  166. const receivedPart = yield* awaitDeferred(received, "timed out waiting for message.part.updated")
  167. expect(receivedPart.type).toBe("step-finish")
  168. const finish = receivedPart as SessionV1.StepFinishPart
  169. expect(finish.tokens.input).toBe(500)
  170. expect(finish.tokens.output).toBe(800)
  171. expect(finish.tokens.reasoning).toBe(200)
  172. expect(finish.tokens.total).toBe(1500)
  173. expect(finish.tokens.cache.read).toBe(100)
  174. expect(finish.tokens.cache.write).toBe(50)
  175. expect(finish.cost).toBe(0.005)
  176. expect(receivedPart).not.toBe(partInput)
  177. yield* session.remove(info.id)
  178. }),
  179. { timeout: 30000 },
  180. )
  181. })
  182. describe("Session", () => {
  183. it.live("remove works without an instance", () =>
  184. Effect.gen(function* () {
  185. const session = yield* SessionNs.Service
  186. const dir = yield* tmpdirScoped({ git: true })
  187. const info = yield* provideInstance(dir)(session.create({ title: "remove-without-instance" }))
  188. const removeExit = yield* remove(info.id).pipe(Effect.exit)
  189. expect(Exit.isSuccess(removeExit)).toBe(true)
  190. const getExit = yield* session.get(info.id).pipe(Effect.exit)
  191. expect(Exit.isFailure(getExit)).toBe(true)
  192. }),
  193. )
  194. it.instance("persists metadata and copies it on fork by default", () =>
  195. Effect.gen(function* () {
  196. const session = yield* SessionNs.Service
  197. const meta = { source: "sdk", trace: { id: "abc" } }
  198. const created = yield* Effect.acquireRelease(session.create({ title: "with-meta", metadata: meta }), (info) =>
  199. session.remove(info.id).pipe(Effect.ignore),
  200. )
  201. const saved = yield* session.get(created.id)
  202. const fork = yield* Effect.acquireRelease(session.fork({ sessionID: created.id }), (info) =>
  203. session.remove(info.id).pipe(Effect.ignore),
  204. )
  205. expect(saved.metadata).toEqual(meta)
  206. expect(fork.metadata).toEqual(meta)
  207. expect(fork.metadata).not.toBe(meta)
  208. }),
  209. )
  210. it.instance("omits metadata when not provided", () =>
  211. Effect.gen(function* () {
  212. const session = yield* SessionNs.Service
  213. const created = yield* Effect.acquireRelease(session.create({ title: "empty-meta" }), (info) =>
  214. session.remove(info.id).pipe(Effect.ignore),
  215. )
  216. const saved = yield* session.get(created.id)
  217. expect(created.metadata).toBeUndefined()
  218. expect(saved.metadata).toBeUndefined()
  219. }),
  220. )
  221. })