event.test.ts 37 KB

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