event.test.ts 37 KB

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