event.test.ts 35 KB

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