event.test.ts 29 KB

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