event.test.ts 34 KB

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