event.test.ts 37 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124
  1. import { describe, expect } from "bun:test"
  2. import { Cause, DateTime, Deferred, Effect, Exit, Fiber, Layer, Option, Schema, Stream } from "effect"
  3. import { EventV2 } from "@opencode-ai/core/event"
  4. import { Event } from "@opencode-ai/schema/event"
  5. import { Session } from "@opencode-ai/schema/session"
  6. import { SessionEvent } from "@opencode-ai/schema/session-event"
  7. import { SessionV1 } from "@opencode-ai/schema/session-v1"
  8. import { Database } from "@opencode-ai/core/database/database"
  9. import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
  10. import { LayerNode } from "@opencode-ai/core/effect/layer-node"
  11. import { EventSequenceTable, EventTable } from "@opencode-ai/core/event/sql"
  12. import { Location } from "@opencode-ai/core/location"
  13. import { AbsolutePath } from "@opencode-ai/core/schema"
  14. import { WorkspaceV2 } from "@opencode-ai/core/workspace"
  15. import { eq } from "drizzle-orm"
  16. import { location } from "./fixture/location"
  17. import { testEffect } from "./lib/effect"
  18. const locationLayer = Layer.succeed(
  19. Location.Service,
  20. Location.Service.of(
  21. location({ directory: AbsolutePath.make("project"), workspaceID: WorkspaceV2.ID.make("wrk_test") }),
  22. ),
  23. )
  24. const Message = EventV2.define({
  25. type: "test.message",
  26. schema: {
  27. text: Schema.String,
  28. },
  29. })
  30. const SyncMessage = EventV2.define({
  31. type: "test.sync",
  32. durable: {
  33. version: 1,
  34. aggregate: "id",
  35. },
  36. schema: {
  37. id: Schema.String,
  38. text: Schema.String,
  39. },
  40. })
  41. const SyncSent = EventV2.define({
  42. type: "test.sent",
  43. durable: {
  44. version: 1,
  45. aggregate: "messageID",
  46. },
  47. schema: {
  48. messageID: Schema.String,
  49. text: Schema.String,
  50. },
  51. })
  52. const GlobalMessage = EventV2.define({
  53. type: "test.global",
  54. schema: {
  55. text: Schema.String,
  56. },
  57. })
  58. const VersionedMessage = EventV2.define({
  59. type: "test.versioned",
  60. durable: {
  61. version: 2,
  62. aggregate: "id",
  63. },
  64. schema: {
  65. id: Schema.String,
  66. text: Schema.String,
  67. },
  68. })
  69. const DurableMessage = SessionV1.Event.MessageRemoved
  70. const durableData = (sessionID: Session.ID, text: string) => ({
  71. sessionID,
  72. messageID: SessionV1.MessageID.ascending(`msg_${text}`),
  73. })
  74. const it = testEffect(
  75. AppNodeBuilder.build(LayerNode.group([Database.node, EventV2.node, Location.node]), [[Location.node, locationLayer]]),
  76. )
  77. const itWithoutLocation = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, EventV2.node])))
  78. describe("EventV2", () => {
  79. it.effect("publishes events with the current location", () =>
  80. Effect.gen(function* () {
  81. const events = yield* EventV2.Service
  82. const fiber = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  83. yield* Effect.yieldNow
  84. const event = yield* events.publish(Message, { text: "hello" })
  85. const received = Array.from(yield* Fiber.join(fiber))
  86. expect(received).toEqual([event])
  87. expect(event.type).toBe("test.message")
  88. expect(event).not.toHaveProperty("version")
  89. expect(event.data).toEqual({ text: "hello" })
  90. expect(event.location).toEqual({
  91. directory: AbsolutePath.make("project"),
  92. workspaceID: WorkspaceV2.ID.make("wrk_test"),
  93. })
  94. }),
  95. )
  96. itWithoutLocation.effect("omits location when no location is available", () =>
  97. Effect.gen(function* () {
  98. const events = yield* EventV2.Service
  99. const event = yield* events.publish(GlobalMessage, { text: "hello" })
  100. expect(event).not.toHaveProperty("location")
  101. expect(event.type).toBe("test.global")
  102. }),
  103. )
  104. it.effect("publishes definition version", () =>
  105. Effect.gen(function* () {
  106. const events = yield* EventV2.Service
  107. const event = yield* events.publish(VersionedMessage, { id: "one", text: "hello" })
  108. expect(event.type).toBe("test.versioned")
  109. expect(event.durable?.version).toBe(2)
  110. }),
  111. )
  112. it.effect("selects the latest durable definition independent of declaration order", () =>
  113. Effect.sync(() => {
  114. const latest = EventV2.define({
  115. type: "test.out-of-order",
  116. durable: { version: 2, aggregate: "id" },
  117. schema: { id: Schema.String },
  118. })
  119. const historical = EventV2.define({
  120. type: "test.out-of-order",
  121. durable: { version: 1, aggregate: "id" },
  122. schema: { id: Schema.String },
  123. })
  124. expect(Event.latest([latest, historical]).get("test.out-of-order")).toBe(latest)
  125. expect(Event.latest([historical, latest]).get("test.out-of-order")).toBe(latest)
  126. }),
  127. )
  128. it.effect("publishes to typed and wildcard subscriptions", () =>
  129. Effect.gen(function* () {
  130. const events = yield* EventV2.Service
  131. const typed = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  132. const wildcard = yield* events.all().pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  133. yield* Effect.yieldNow
  134. const event = yield* events.publish(Message, { text: "hello" })
  135. expect(Array.from(yield* Fiber.join(typed))).toEqual([event])
  136. expect(Array.from(yield* Fiber.join(wildcard))).toEqual([event])
  137. }),
  138. )
  139. it.effect("runs projectors inline", () =>
  140. Effect.gen(function* () {
  141. const events = yield* EventV2.Service
  142. const received = new Array<EventV2.Payload>()
  143. yield* events.project(SyncMessage, (event) =>
  144. Effect.sync(() => {
  145. received.push(event)
  146. }),
  147. )
  148. const event = yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  149. yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" })
  150. expect(received[0]).toEqual(event)
  151. expect(received[1]?.data).toEqual({ id: "one", text: "after unsubscribe" })
  152. }),
  153. )
  154. it.effect("commits local operational state inside a new durable event transaction", () =>
  155. Effect.gen(function* () {
  156. const events = yield* EventV2.Service
  157. const received = new Array<string>()
  158. const aggregateID = EventV2.ID.create()
  159. yield* events.project(SyncMessage, () => Effect.sync(() => received.push("projector")))
  160. yield* events.publish(
  161. SyncMessage,
  162. { id: aggregateID, text: "hello" },
  163. { commit: (seq) => Effect.sync(() => received.push(`commit:${seq}`)) },
  164. )
  165. expect(received).toEqual(["projector", "commit:0"])
  166. }),
  167. )
  168. it.effect("rolls back the durable event and projector when the local commit fails", () =>
  169. Effect.gen(function* () {
  170. const events = yield* EventV2.Service
  171. const { db } = yield* Database.Service
  172. const aggregateID = EventV2.ID.create()
  173. yield* db.run("CREATE TABLE IF NOT EXISTS event_commit_probe (value text NOT NULL)")
  174. yield* db.run("DELETE FROM event_commit_probe")
  175. yield* events.project(SyncMessage, () =>
  176. db.run("INSERT INTO event_commit_probe (value) VALUES ('projected')").pipe(Effect.orDie, Effect.asVoid),
  177. )
  178. const exit = yield* events
  179. .publish(SyncMessage, { id: aggregateID, text: "hello" }, { commit: () => Effect.die("commit failed") })
  180. .pipe(Effect.exit)
  181. expect(String(exit)).toContain("commit failed")
  182. expect(yield* db.all("SELECT value FROM event_commit_probe")).toEqual([])
  183. expect(yield* db.select().from(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).all()).toEqual([])
  184. expect(
  185. yield* db.select().from(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).all(),
  186. ).toEqual([])
  187. }),
  188. )
  189. it.effect("rejects local commit hooks on live-only events", () =>
  190. Effect.gen(function* () {
  191. const events = yield* EventV2.Service
  192. const exit = yield* events.publish(Message, { text: "hello" }, { commit: () => Effect.void }).pipe(Effect.exit)
  193. expect(String(exit)).toContain("Local commit hooks require a durable event")
  194. }),
  195. )
  196. it.effect("runs projectors before publishing to streams", () =>
  197. Effect.gen(function* () {
  198. const events = yield* EventV2.Service
  199. const received = new Array<string>()
  200. const fiber = yield* events.all().pipe(
  201. Stream.take(1),
  202. Stream.runForEach(() => Effect.sync(() => received.push("stream"))),
  203. Effect.forkScoped,
  204. )
  205. yield* events.project(SyncMessage, (event) =>
  206. Effect.sync(() => {
  207. received.push(event.type)
  208. }),
  209. )
  210. yield* Effect.yieldNow
  211. yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  212. yield* Fiber.join(fiber)
  213. expect(received).toEqual([SyncMessage.type, "stream"])
  214. }),
  215. )
  216. it.effect("runs listeners inline after projectors", () =>
  217. Effect.gen(function* () {
  218. const events = yield* EventV2.Service
  219. const received = new Array<string>()
  220. yield* events.project(SyncMessage, () =>
  221. Effect.sync(() => {
  222. received.push("projector")
  223. }),
  224. )
  225. const unsubscribe = yield* events.listen(() =>
  226. Effect.sync(() => {
  227. received.push("listener")
  228. }),
  229. )
  230. yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  231. yield* unsubscribe
  232. yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" })
  233. expect(received).toEqual(["projector", "listener", "projector"])
  234. }),
  235. )
  236. it.effect("isolates observer defects after durable events commit", () =>
  237. Effect.gen(function* () {
  238. const events = yield* EventV2.Service
  239. const received = new Array<string>()
  240. yield* events.listen(() => {
  241. throw new Error("listener defect")
  242. })
  243. yield* events.listen((event) =>
  244. Effect.sync(() => {
  245. received.push(event.type)
  246. }),
  247. )
  248. const event = yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  249. expect(received).toEqual([SyncMessage.type])
  250. expect(event.durable?.seq).toBeNumber()
  251. }),
  252. )
  253. it.effect("notifies global listeners only after a durable event is committed", () =>
  254. Effect.gen(function* () {
  255. const events = yield* EventV2.Service
  256. const { db } = yield* Database.Service
  257. const aggregateID = EventV2.ID.create()
  258. const observed = new Array<{ id: string; seq: number }>()
  259. yield* events.listen((event) =>
  260. event.type !== SyncMessage.type
  261. ? Effect.void
  262. : db
  263. .select({ id: EventTable.id, seq: EventTable.seq })
  264. .from(EventTable)
  265. .where(eq(EventTable.id, event.id))
  266. .get()
  267. .pipe(
  268. Effect.orDie,
  269. Effect.tap((row) =>
  270. Effect.sync(() => {
  271. if (row) observed.push(row)
  272. }),
  273. ),
  274. Effect.asVoid,
  275. ),
  276. )
  277. const event = yield* events.publish(SyncMessage, { id: aggregateID, text: "committed" })
  278. if (!event.durable) throw new Error("Expected durable event metadata")
  279. expect(observed).toEqual([{ id: event.id, seq: event.durable.seq }])
  280. }),
  281. )
  282. it.effect("ends only an overflowing bounded subscriber without blocking other listeners", () =>
  283. Effect.gen(function* () {
  284. const events = yield* EventV2.Service
  285. const consuming = yield* Deferred.make<void>()
  286. const release = yield* Deferred.make<void>()
  287. const slowStream = yield* EventV2.allBounded(events, 1)
  288. const fastStream = yield* EventV2.allBounded(events, 8)
  289. const slow = yield* slowStream.pipe(
  290. Stream.runForEach(() => Deferred.succeed(consuming, undefined).pipe(Effect.andThen(Deferred.await(release)))),
  291. Effect.forkScoped,
  292. )
  293. const fast = yield* fastStream.pipe(Stream.take(4), Stream.runCollect, Effect.forkScoped)
  294. yield* events.publish(Message, { text: "one" })
  295. yield* Deferred.await(consuming)
  296. yield* events.publish(Message, { text: "two" })
  297. yield* events.publish(Message, { text: "overflow" })
  298. const last = yield* events.publish(Message, { text: "still delivered" })
  299. yield* Deferred.succeed(release, undefined)
  300. const slowExit = yield* Fiber.await(slow)
  301. expect(Exit.findErrorOption(slowExit).pipe(Option.getOrUndefined)).toBeInstanceOf(EventV2.SubscriberOverflowError)
  302. expect(Array.from(yield* Fiber.join(fast))).toEqual([
  303. expect.objectContaining({ data: { text: "one" } }),
  304. expect.objectContaining({ data: { text: "two" } }),
  305. expect.objectContaining({ data: { text: "overflow" } }),
  306. last,
  307. ])
  308. }),
  309. )
  310. it.effect("preserves observer interruption", () =>
  311. Effect.gen(function* () {
  312. const events = yield* EventV2.Service
  313. const { db } = yield* Database.Service
  314. yield* events.listen(() => Effect.interrupt)
  315. const exit = yield* events.publish(SyncMessage, { id: "interrupted", text: "hello" }).pipe(Effect.exit)
  316. const committed = yield* db
  317. .select({ id: EventTable.id })
  318. .from(EventTable)
  319. .where(eq(EventTable.aggregate_id, "interrupted"))
  320. .get()
  321. .pipe(Effect.orDie)
  322. expect(Exit.isFailure(exit) && Cause.hasInterrupts(exit.cause)).toBeTrue()
  323. expect(committed).toBeDefined()
  324. }),
  325. )
  326. it.effect("keeps live-only listener defects fail-fast", () =>
  327. Effect.gen(function* () {
  328. const events = yield* EventV2.Service
  329. const defect = new Error("listener defect")
  330. yield* events.listen(() => Effect.die(defect))
  331. expect(yield* events.publish(Message, { text: "hello" }).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect)
  332. }),
  333. )
  334. it.effect("inserts durable event rows on publish", () =>
  335. Effect.gen(function* () {
  336. const events = yield* EventV2.Service
  337. const { db } = yield* Database.Service
  338. const aggregateID = EventV2.ID.create()
  339. yield* events.publish(SyncMessage, { id: aggregateID, text: "first" })
  340. const rows = yield* db
  341. .select()
  342. .from(EventTable)
  343. .where(eq(EventTable.aggregate_id, aggregateID))
  344. .all()
  345. .pipe(Effect.orDie)
  346. expect(rows).toHaveLength(1)
  347. expect(rows[0]?.type).toBe(EventV2.versionedType(SyncMessage.type, 1))
  348. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  349. }),
  350. )
  351. it.effect("increments durable event seq per aggregate", () =>
  352. Effect.gen(function* () {
  353. const events = yield* EventV2.Service
  354. const { db } = yield* Database.Service
  355. const aggregateID = EventV2.ID.create()
  356. yield* events.publish(SyncMessage, { id: aggregateID, text: "first" })
  357. yield* events.publish(SyncMessage, { id: aggregateID, text: "second" })
  358. const rows = yield* db
  359. .select()
  360. .from(EventTable)
  361. .where(eq(EventTable.aggregate_id, aggregateID))
  362. .all()
  363. .pipe(Effect.orDie)
  364. expect(rows.map((row) => row.seq)).toEqual([0, 1])
  365. }),
  366. )
  367. it.effect("replays durable aggregate events after a sequence and tails new events", () =>
  368. Effect.gen(function* () {
  369. const events = yield* EventV2.Service
  370. const aggregateID = Session.ID.create()
  371. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  372. yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
  373. const fiber = yield* events
  374. .durable({ aggregateID, after: 0 })
  375. .pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
  376. yield* Effect.yieldNow
  377. yield* events.publish(DurableMessage, durableData(aggregateID, "two"))
  378. expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
  379. [1, durableData(aggregateID, "one")],
  380. [2, durableData(aggregateID, "two")],
  381. ])
  382. }),
  383. )
  384. it.effect("catches durable aggregate events published during replay handoff", () =>
  385. Effect.gen(function* () {
  386. const events = yield* EventV2.Service
  387. const aggregateID = Session.ID.create()
  388. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  389. const fiber = yield* events.durable({ aggregateID }).pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
  390. yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
  391. expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
  392. [0, durableData(aggregateID, "zero")],
  393. [1, durableData(aggregateID, "one")],
  394. ])
  395. }),
  396. )
  397. it.effect("retains a durable wake committed while historical replay is paused", () =>
  398. Effect.gen(function* () {
  399. const readStarted = yield* Deferred.make<void>()
  400. const continueRead = yield* Deferred.make<void>()
  401. let pause = true
  402. const eventLayer = EventV2.layerWith({
  403. beforeAggregateRead: () =>
  404. pause
  405. ? Deferred.succeed(readStarted, undefined).pipe(Effect.andThen(Deferred.await(continueRead)))
  406. : Effect.void,
  407. }).pipe(Layer.provide(LayerNode.compile(Database.node)))
  408. yield* Effect.gen(function* () {
  409. const events = yield* EventV2.Service
  410. const aggregateID = Session.ID.create()
  411. const fiber = yield* events.durable({ aggregateID }).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  412. yield* Deferred.await(readStarted)
  413. pause = false
  414. yield* events.publish(DurableMessage, durableData(aggregateID, "during handoff"))
  415. yield* Deferred.succeed(continueRead, undefined)
  416. expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
  417. [0, durableData(aggregateID, "during handoff")],
  418. ])
  419. }).pipe(Effect.provide(Layer.merge(LayerNode.compile(Database.node), eventLayer)))
  420. }),
  421. )
  422. it.effect("coalesces durable aggregate wakes while draining every committed event", () =>
  423. Effect.gen(function* () {
  424. const events = yield* EventV2.Service
  425. const aggregateID = Session.ID.create()
  426. const count = 64
  427. const fiber = yield* events
  428. .durable({ aggregateID })
  429. .pipe(Stream.take(count), Stream.runCollect, Effect.forkScoped)
  430. yield* Effect.yieldNow
  431. for (let index = 0; index < count; index++) {
  432. yield* events.publish(DurableMessage, durableData(aggregateID, String(index)))
  433. }
  434. expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual(
  435. Array.from({ length: count }, (_, index) => [index, durableData(aggregateID, String(index))]),
  436. )
  437. }),
  438. )
  439. it.effect("omits live-only events from durable aggregate streams", () =>
  440. Effect.gen(function* () {
  441. const events = yield* EventV2.Service
  442. const aggregateID = Session.ID.create()
  443. const fiber = yield* events.durable({ aggregateID }).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  444. yield* Effect.yieldNow
  445. yield* events.publish(Message, { text: "live only" })
  446. yield* events.publish(DurableMessage, durableData(aggregateID, "durable"))
  447. expect(Array.from(yield* Fiber.join(fiber)).map((event) => event.type)).toEqual([DurableMessage.type])
  448. }),
  449. )
  450. it.effect("uses custom sync aggregate field", () =>
  451. Effect.gen(function* () {
  452. const events = yield* EventV2.Service
  453. const { db } = yield* Database.Service
  454. const aggregateID = EventV2.ID.create()
  455. yield* events.publish(SyncSent, { messageID: aggregateID, text: "sent" })
  456. const rows = yield* db
  457. .select()
  458. .from(EventTable)
  459. .where(eq(EventTable.aggregate_id, aggregateID))
  460. .all()
  461. .pipe(Effect.orDie)
  462. expect(rows).toHaveLength(1)
  463. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  464. }),
  465. )
  466. it.effect("replays durable events through projectors", () =>
  467. Effect.gen(function* () {
  468. const events = yield* EventV2.Service
  469. const received = new Array<EventV2.Payload>()
  470. yield* events.project(DurableMessage, (event) =>
  471. Effect.sync(() => {
  472. received.push(event)
  473. }),
  474. )
  475. const aggregateID = Session.ID.create()
  476. yield* events.replay({
  477. id: EventV2.ID.create(),
  478. type: EventV2.versionedType(DurableMessage.type, 1),
  479. seq: 0,
  480. aggregateID,
  481. data: durableData(aggregateID, "hello"),
  482. })
  483. expect(received[0]?.type).toBe(DurableMessage.type)
  484. expect(received[0]?.data).toEqual(durableData(aggregateID, "hello"))
  485. }),
  486. )
  487. it.effect("replay inserts external event rows", () =>
  488. Effect.gen(function* () {
  489. const events = yield* EventV2.Service
  490. const { db } = yield* Database.Service
  491. const aggregateID = Session.ID.create()
  492. yield* events.replay({
  493. id: EventV2.ID.create(),
  494. type: EventV2.versionedType(DurableMessage.type, 1),
  495. seq: 0,
  496. aggregateID,
  497. data: durableData(aggregateID, "replayed"),
  498. })
  499. const rows = yield* db
  500. .select()
  501. .from(EventTable)
  502. .where(eq(EventTable.aggregate_id, aggregateID))
  503. .all()
  504. .pipe(Effect.orDie)
  505. expect(rows).toHaveLength(1)
  506. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  507. }),
  508. )
  509. it.effect(
  510. "replay rejects an envelope aggregate that differs from its payload without mutating the payload aggregate",
  511. () =>
  512. Effect.gen(function* () {
  513. const events = yield* EventV2.Service
  514. const { db } = yield* Database.Service
  515. const envelopeAggregateID = Session.ID.create()
  516. const payloadAggregateID = Session.ID.create()
  517. const received = new Array<EventV2.Payload>()
  518. yield* events.publish(DurableMessage, durableData(payloadAggregateID, "seed"))
  519. yield* events.project(DurableMessage, (event) =>
  520. Effect.sync(() => {
  521. received.push(event)
  522. }),
  523. )
  524. const exit = yield* events
  525. .replay({
  526. id: EventV2.ID.create(),
  527. type: EventV2.versionedType(DurableMessage.type, 1),
  528. seq: 1,
  529. aggregateID: envelopeAggregateID,
  530. data: durableData(payloadAggregateID, "replayed"),
  531. })
  532. .pipe(Effect.exit)
  533. const rows = yield* db
  534. .select()
  535. .from(EventTable)
  536. .where(eq(EventTable.aggregate_id, payloadAggregateID))
  537. .all()
  538. .pipe(Effect.orDie)
  539. const sequence = yield* db
  540. .select({ seq: EventSequenceTable.seq })
  541. .from(EventSequenceTable)
  542. .where(eq(EventSequenceTable.aggregate_id, payloadAggregateID))
  543. .get()
  544. .pipe(Effect.orDie)
  545. expect(String(exit)).toContain("Aggregate mismatch")
  546. expect(received).toHaveLength(0)
  547. expect(rows).toHaveLength(1)
  548. expect(sequence).toEqual({ seq: 0 })
  549. }),
  550. )
  551. it.effect("replay defects on sequence mismatch", () =>
  552. Effect.gen(function* () {
  553. const events = yield* EventV2.Service
  554. const aggregateID = Session.ID.create()
  555. yield* events.replay({
  556. id: EventV2.ID.create(),
  557. type: EventV2.versionedType(DurableMessage.type, 1),
  558. seq: 0,
  559. aggregateID,
  560. data: durableData(aggregateID, "first"),
  561. })
  562. const exit = yield* events
  563. .replay({
  564. id: EventV2.ID.create(),
  565. type: EventV2.versionedType(DurableMessage.type, 1),
  566. seq: 5,
  567. aggregateID,
  568. data: durableData(aggregateID, "bad"),
  569. })
  570. .pipe(Effect.exit)
  571. expect(String(exit)).toContain("Sequence mismatch")
  572. }),
  573. )
  574. it.effect("replay decodes synchronized transformed values before projection", () =>
  575. Effect.gen(function* () {
  576. const events = yield* EventV2.Service
  577. const aggregateID = Session.ID.create()
  578. const received = new Array<typeof SessionEvent.ContextUpdated.Type>()
  579. yield* events.project(SessionEvent.ContextUpdated, (event) =>
  580. Effect.sync(() => {
  581. received.push(event)
  582. }),
  583. )
  584. yield* events.replay({
  585. id: EventV2.ID.create(),
  586. type: EventV2.versionedType(SessionEvent.ContextUpdated.type, 1),
  587. seq: 0,
  588. aggregateID,
  589. data: { sessionID: aggregateID, messageID: "msg_context", timestamp: 0, text: "context" },
  590. })
  591. expect(received[0]?.data.timestamp).toEqual(DateTime.makeUnsafe(0))
  592. }),
  593. )
  594. it.effect("replay defects on unknown event type", () =>
  595. Effect.gen(function* () {
  596. const events = yield* EventV2.Service
  597. const exit = yield* events
  598. .replay({
  599. id: EventV2.ID.create(),
  600. type: "unknown.event.1",
  601. seq: 0,
  602. aggregateID: EventV2.ID.create(),
  603. data: {},
  604. })
  605. .pipe(Effect.exit)
  606. expect(String(exit)).toContain("Unknown durable event type")
  607. }),
  608. )
  609. it.effect("replayAll validates contiguous aggregate events", () =>
  610. Effect.gen(function* () {
  611. const events = yield* EventV2.Service
  612. const aggregateID = Session.ID.create()
  613. const source = yield* events.replayAll([
  614. {
  615. id: EventV2.ID.create(),
  616. type: EventV2.versionedType(DurableMessage.type, 1),
  617. seq: 0,
  618. aggregateID,
  619. data: durableData(aggregateID, "one"),
  620. },
  621. {
  622. id: EventV2.ID.create(),
  623. type: EventV2.versionedType(DurableMessage.type, 1),
  624. seq: 1,
  625. aggregateID,
  626. data: durableData(aggregateID, "two"),
  627. },
  628. ])
  629. expect(source).toBe(aggregateID)
  630. }),
  631. )
  632. it.effect("replayAll accepts later chunks after the first batch", () =>
  633. Effect.gen(function* () {
  634. const events = yield* EventV2.Service
  635. const { db } = yield* Database.Service
  636. const aggregateID = Session.ID.create()
  637. const one = yield* events.replayAll([
  638. {
  639. id: EventV2.ID.create(),
  640. type: EventV2.versionedType(DurableMessage.type, 1),
  641. seq: 0,
  642. aggregateID,
  643. data: durableData(aggregateID, "one"),
  644. },
  645. {
  646. id: EventV2.ID.create(),
  647. type: EventV2.versionedType(DurableMessage.type, 1),
  648. seq: 1,
  649. aggregateID,
  650. data: durableData(aggregateID, "two"),
  651. },
  652. ])
  653. const two = yield* events.replayAll([
  654. {
  655. id: EventV2.ID.create(),
  656. type: EventV2.versionedType(DurableMessage.type, 1),
  657. seq: 2,
  658. aggregateID,
  659. data: durableData(aggregateID, "three"),
  660. },
  661. {
  662. id: EventV2.ID.create(),
  663. type: EventV2.versionedType(DurableMessage.type, 1),
  664. seq: 3,
  665. aggregateID,
  666. data: durableData(aggregateID, "four"),
  667. },
  668. ])
  669. const rows = yield* db
  670. .select()
  671. .from(EventTable)
  672. .where(eq(EventTable.aggregate_id, aggregateID))
  673. .all()
  674. .pipe(Effect.orDie)
  675. expect(one).toBe(aggregateID)
  676. expect(two).toBe(aggregateID)
  677. expect(rows.map((row) => row.seq)).toEqual([0, 1, 2, 3])
  678. }),
  679. )
  680. it.effect("claim fences replay owners", () =>
  681. Effect.gen(function* () {
  682. const events = yield* EventV2.Service
  683. const received = new Array<EventV2.Payload>()
  684. const aggregateID = Session.ID.create()
  685. yield* events.publish(DurableMessage, durableData(aggregateID, "seed"))
  686. yield* events.claim(aggregateID, "owner-a")
  687. yield* events.project(DurableMessage, (event) =>
  688. Effect.sync(() => {
  689. received.push(event)
  690. }),
  691. )
  692. yield* events.replay(
  693. {
  694. id: EventV2.ID.create(),
  695. type: EventV2.versionedType(DurableMessage.type, 1),
  696. seq: 1,
  697. aggregateID,
  698. data: durableData(aggregateID, "ignored"),
  699. },
  700. { ownerID: "owner-b" },
  701. )
  702. expect(received).toHaveLength(0)
  703. }),
  704. )
  705. it.effect("strict owner fences exact replay", () =>
  706. Effect.gen(function* () {
  707. const events = yield* EventV2.Service
  708. const aggregateID = Session.ID.create()
  709. const id = EventV2.ID.create()
  710. const replayed = {
  711. id,
  712. type: EventV2.versionedType(DurableMessage.type, 1),
  713. seq: 0,
  714. aggregateID,
  715. data: durableData(aggregateID, "owned"),
  716. }
  717. yield* events.replay(replayed, { ownerID: "owner-a" })
  718. const exit = yield* events.replay(replayed, { ownerID: "owner-b", strictOwner: true }).pipe(Effect.exit)
  719. expect(String(exit)).toContain("Replay owner mismatch")
  720. }),
  721. )
  722. it.effect("exact replay claims an unowned aggregate", () =>
  723. Effect.gen(function* () {
  724. const events = yield* EventV2.Service
  725. const { db } = yield* Database.Service
  726. const aggregateID = Session.ID.create()
  727. const published = yield* events.publish(DurableMessage, durableData(aggregateID, "owned"))
  728. const replayed = {
  729. id: published.id,
  730. type: EventV2.versionedType(DurableMessage.type, 1),
  731. seq: published.durable!.seq,
  732. aggregateID,
  733. data: published.data,
  734. }
  735. yield* events.replay(replayed, { ownerID: "owner-a", strictOwner: true })
  736. const row = yield* db
  737. .select({ ownerID: EventSequenceTable.owner_id })
  738. .from(EventSequenceTable)
  739. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  740. .get()
  741. .pipe(Effect.orDie)
  742. expect(row?.ownerID).toBe("owner-a")
  743. const exit = yield* events
  744. .replay(
  745. { ...replayed, id: EventV2.ID.create(), seq: 1, data: durableData(aggregateID, "conflict") },
  746. { ownerID: "owner-b", strictOwner: true },
  747. )
  748. .pipe(Effect.exit)
  749. expect(String(exit)).toContain("Replay owner mismatch")
  750. }),
  751. )
  752. it.effect("replay with owner claims an unowned sequence", () =>
  753. Effect.gen(function* () {
  754. const events = yield* EventV2.Service
  755. const { db } = yield* Database.Service
  756. const aggregateID = Session.ID.create()
  757. yield* events.replay(
  758. {
  759. id: EventV2.ID.create(),
  760. type: EventV2.versionedType(DurableMessage.type, 1),
  761. seq: 0,
  762. aggregateID,
  763. data: durableData(aggregateID, "owned"),
  764. },
  765. { ownerID: "owner-1" },
  766. )
  767. const row = yield* db
  768. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  769. .from(EventSequenceTable)
  770. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  771. .get()
  772. .pipe(Effect.orDie)
  773. expect(row).toEqual({ seq: 0, ownerID: "owner-1" })
  774. }),
  775. )
  776. it.effect("replay claims an existing unowned sequence before fencing a different owner", () =>
  777. Effect.gen(function* () {
  778. const events = yield* EventV2.Service
  779. const { db } = yield* Database.Service
  780. const aggregateID = Session.ID.create()
  781. yield* events.publish(DurableMessage, durableData(aggregateID, "local"))
  782. yield* events.replay(
  783. {
  784. id: EventV2.ID.create(),
  785. type: EventV2.versionedType(DurableMessage.type, 1),
  786. seq: 1,
  787. aggregateID,
  788. data: durableData(aggregateID, "claimed"),
  789. },
  790. { ownerID: "owner-1" },
  791. )
  792. yield* events.replay(
  793. {
  794. id: EventV2.ID.create(),
  795. type: EventV2.versionedType(DurableMessage.type, 1),
  796. seq: 2,
  797. aggregateID,
  798. data: durableData(aggregateID, "fenced"),
  799. },
  800. { ownerID: "owner-2" },
  801. )
  802. const rows = yield* db
  803. .select()
  804. .from(EventTable)
  805. .where(eq(EventTable.aggregate_id, aggregateID))
  806. .all()
  807. .pipe(Effect.orDie)
  808. const sequence = yield* db
  809. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  810. .from(EventSequenceTable)
  811. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  812. .get()
  813. .pipe(Effect.orDie)
  814. expect(rows.map((row) => row.seq)).toEqual([0, 1])
  815. expect(sequence).toEqual({ seq: 1, ownerID: "owner-1" })
  816. }),
  817. )
  818. it.effect("strict replay rejects an owner conflict instead of silently skipping it", () =>
  819. Effect.gen(function* () {
  820. const events = yield* EventV2.Service
  821. const aggregateID = Session.ID.create()
  822. yield* events.replay(
  823. {
  824. id: EventV2.ID.create(),
  825. type: EventV2.versionedType(DurableMessage.type, 1),
  826. seq: 0,
  827. aggregateID,
  828. data: durableData(aggregateID, "claimed"),
  829. },
  830. { ownerID: "owner-1" },
  831. )
  832. const exit = yield* events
  833. .replay(
  834. {
  835. id: EventV2.ID.create(),
  836. type: EventV2.versionedType(DurableMessage.type, 1),
  837. seq: 1,
  838. aggregateID,
  839. data: durableData(aggregateID, "conflict"),
  840. },
  841. { ownerID: "owner-2", strictOwner: true },
  842. )
  843. .pipe(Effect.exit)
  844. expect(String(exit)).toContain("Replay owner mismatch")
  845. }),
  846. )
  847. it.effect("publishes accepted replay with its durable sequence and suppresses stale replay", () =>
  848. Effect.gen(function* () {
  849. const events = yield* EventV2.Service
  850. const received = new Array<EventV2.Payload>()
  851. const aggregateID = Session.ID.create()
  852. yield* events.listen((event) => Effect.sync(() => received.push(event)))
  853. const replayed = {
  854. id: EventV2.ID.create(),
  855. type: EventV2.versionedType(DurableMessage.type, 1),
  856. seq: 0,
  857. aggregateID,
  858. data: durableData(aggregateID, "replayed"),
  859. }
  860. yield* events.replay(replayed, { publish: true })
  861. yield* events.replay(replayed, { publish: true })
  862. expect(received).toMatchObject([{ id: replayed.id, durable: { seq: 0, version: 1 }, data: replayed.data }])
  863. }),
  864. )
  865. it.effect("rejects divergent stale replay without publishing it", () =>
  866. Effect.gen(function* () {
  867. const events = yield* EventV2.Service
  868. const received = new Array<EventV2.Payload>()
  869. const aggregateID = Session.ID.create()
  870. const replayed = {
  871. id: EventV2.ID.create(),
  872. type: EventV2.versionedType(DurableMessage.type, 1),
  873. seq: 0,
  874. aggregateID,
  875. data: durableData(aggregateID, "original"),
  876. }
  877. yield* events.listen((event) => Effect.sync(() => received.push(event)))
  878. yield* events.replay(replayed, { publish: true })
  879. const exit = yield* events
  880. .replay({ ...replayed, data: durableData(aggregateID, "divergent") }, { publish: true })
  881. .pipe(Effect.exit)
  882. expect(String(exit)).toContain("Replay diverged")
  883. expect(received).toHaveLength(1)
  884. }),
  885. )
  886. it.effect("rejects an event ID reused at another aggregate position", () =>
  887. Effect.gen(function* () {
  888. const events = yield* EventV2.Service
  889. const aggregateID = Session.ID.create()
  890. const id = EventV2.ID.create()
  891. yield* events.replay({
  892. id,
  893. type: EventV2.versionedType(DurableMessage.type, 1),
  894. seq: 0,
  895. aggregateID,
  896. data: durableData(aggregateID, "first"),
  897. })
  898. const exit = yield* events
  899. .replay({
  900. id,
  901. type: EventV2.versionedType(DurableMessage.type, 1),
  902. seq: 1,
  903. aggregateID,
  904. data: durableData(aggregateID, "second"),
  905. })
  906. .pipe(Effect.exit)
  907. expect(String(exit)).toContain(`Event ${id} already exists`)
  908. }),
  909. )
  910. it.effect("replay from a different owner leaves claimed sequence unchanged", () =>
  911. Effect.gen(function* () {
  912. const events = yield* EventV2.Service
  913. const { db } = yield* Database.Service
  914. const aggregateID = Session.ID.create()
  915. const received = new Array<EventV2.Payload>()
  916. yield* events.listen((event) => Effect.sync(() => received.push(event)))
  917. yield* events.replay(
  918. {
  919. id: EventV2.ID.create(),
  920. type: EventV2.versionedType(DurableMessage.type, 1),
  921. seq: 0,
  922. aggregateID,
  923. data: durableData(aggregateID, "first"),
  924. },
  925. { ownerID: "owner-1" },
  926. )
  927. yield* events.replay(
  928. {
  929. id: EventV2.ID.create(),
  930. type: EventV2.versionedType(DurableMessage.type, 1),
  931. seq: 1,
  932. aggregateID,
  933. data: durableData(aggregateID, "ignored"),
  934. },
  935. { ownerID: "owner-2", publish: true },
  936. )
  937. const rows = yield* db
  938. .select()
  939. .from(EventTable)
  940. .where(eq(EventTable.aggregate_id, aggregateID))
  941. .all()
  942. .pipe(Effect.orDie)
  943. const sequence = yield* db
  944. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  945. .from(EventSequenceTable)
  946. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  947. .get()
  948. .pipe(Effect.orDie)
  949. expect(rows).toHaveLength(1)
  950. expect(sequence).toEqual({ seq: 0, ownerID: "owner-1" })
  951. expect(received).toHaveLength(0)
  952. }),
  953. )
  954. it.effect("claim updates the event sequence owner", () =>
  955. Effect.gen(function* () {
  956. const events = yield* EventV2.Service
  957. const { db } = yield* Database.Service
  958. const aggregateID = EventV2.ID.create()
  959. yield* events.publish(SyncMessage, { id: aggregateID, text: "claimed" })
  960. yield* events.claim(aggregateID, "owner-1")
  961. yield* events.claim(aggregateID, "owner-2")
  962. const row = yield* db
  963. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  964. .from(EventSequenceTable)
  965. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  966. .get()
  967. .pipe(Effect.orDie)
  968. expect(row).toEqual({ seq: 0, ownerID: "owner-2" })
  969. }),
  970. )
  971. it.effect("remove clears durable event sequence", () =>
  972. Effect.gen(function* () {
  973. const events = yield* EventV2.Service
  974. const received = new Array<EventV2.Payload>()
  975. const aggregateID = Session.ID.create()
  976. yield* events.publish(DurableMessage, durableData(aggregateID, "seed"))
  977. yield* events.remove(aggregateID)
  978. yield* events.project(DurableMessage, (event) =>
  979. Effect.sync(() => {
  980. received.push(event)
  981. }),
  982. )
  983. yield* events.replay({
  984. id: EventV2.ID.create(),
  985. type: EventV2.versionedType(DurableMessage.type, 1),
  986. seq: 0,
  987. aggregateID,
  988. data: durableData(aggregateID, "replayed"),
  989. })
  990. expect(received[0]?.data).toEqual(durableData(aggregateID, "replayed"))
  991. }),
  992. )
  993. })