event.test.ts 34 KB

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