event.test.ts 35 KB

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