event.test.ts 48 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368
  1. import { describe, expect } from "bun:test"
  2. import { Cause, DateTime, Deferred, Effect, Exit, Fiber, Layer, Option, Ref, Schema, Stream } from "effect"
  3. import { EventV2 } from "@opencode-ai/core/event"
  4. import { Event } from "@opencode-ai/schema/event"
  5. import { EventManifest } from "@opencode-ai/schema/event-manifest"
  6. import { McpEvent } from "@opencode-ai/schema/mcp-event"
  7. import { Session } from "@opencode-ai/schema/session"
  8. import { SessionEvent } from "@opencode-ai/schema/session-event"
  9. import { SessionV1 } from "@opencode-ai/schema/session-v1"
  10. import { Database } from "@opencode-ai/core/database/database"
  11. import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
  12. import { LayerNode } from "@opencode-ai/core/effect/layer-node"
  13. import { EventSequenceTable, EventTable } from "@opencode-ai/core/event/sql"
  14. import { Location } from "@opencode-ai/core/location"
  15. import { AbsolutePath } from "@opencode-ai/core/schema"
  16. import { WorkspaceV2 } from "@opencode-ai/core/workspace"
  17. import { eq } from "drizzle-orm"
  18. import { location } from "./fixture/location"
  19. import { testEffect } from "./lib/effect"
  20. const locationLayer = Layer.succeed(
  21. Location.Service,
  22. Location.Service.of(
  23. location({ directory: AbsolutePath.make("project"), workspaceID: WorkspaceV2.ID.make("wrk_test") }),
  24. ),
  25. )
  26. const Message = EventV2.ephemeral({
  27. type: "test.message",
  28. schema: {
  29. text: Schema.String,
  30. },
  31. })
  32. const SyncMessage = EventV2.durable({
  33. type: "test.sync",
  34. durable: {
  35. version: 1,
  36. aggregate: "id",
  37. },
  38. schema: {
  39. id: Schema.String,
  40. text: Schema.String,
  41. },
  42. })
  43. const SyncSent = EventV2.durable({
  44. type: "test.sent",
  45. durable: {
  46. version: 1,
  47. aggregate: "messageID",
  48. },
  49. schema: {
  50. messageID: Schema.String,
  51. text: Schema.String,
  52. },
  53. })
  54. const VersionedMessageV1 = EventV2.durable({
  55. type: "test.versioned",
  56. durable: { version: 1, aggregate: "id" },
  57. schema: { id: Schema.String },
  58. })
  59. const VersionedMessageV2 = EventV2.durable({
  60. type: "test.versioned",
  61. durable: { version: 2, aggregate: "id" },
  62. schema: { id: Schema.String },
  63. })
  64. const GlobalMessage = EventV2.ephemeral({
  65. type: "test.global",
  66. schema: {
  67. text: Schema.String,
  68. },
  69. })
  70. const CountMessage = EventV2.ephemeral({
  71. type: "test.count",
  72. schema: {
  73. count: Schema.Number,
  74. },
  75. })
  76. const VersionedMessage = EventV2.durable({
  77. type: "test.versioned",
  78. durable: {
  79. version: 2,
  80. aggregate: "id",
  81. },
  82. schema: {
  83. id: Schema.String,
  84. text: Schema.String,
  85. },
  86. })
  87. const DurableMessage = SessionV1.Event.MessageRemoved
  88. const durableData = (sessionID: Session.ID, text: string) => ({
  89. sessionID,
  90. messageID: SessionV1.MessageID.ascending(`msg_${text}`),
  91. })
  92. /** Followed log read without markers: the old `durable` stream shape. */
  93. const tail = (events: EventV2.Interface, input: { aggregateID: string; after?: number }) =>
  94. events.log({ ...input, follow: true }).pipe(Stream.filter((item): item is EventV2.Payload => !EventV2.isSynced(item)))
  95. const it = testEffect(
  96. AppNodeBuilder.build(LayerNode.group([Database.node, EventV2.node, Location.node]), [[Location.node, locationLayer]]),
  97. )
  98. const itWithoutLocation = testEffect(AppNodeBuilder.build(LayerNode.group([Database.node, EventV2.node])))
  99. describe("EventV2", () => {
  100. it.effect("subscribes to multiple event definitions with a discriminated payload union", () =>
  101. Effect.gen(function* () {
  102. const events = yield* EventV2.Service
  103. // @ts-expect-error multi-definition subscriptions require at least one definition
  104. events.subscribe([])
  105. const fiber = yield* events
  106. .subscribe([Message, CountMessage])
  107. .pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
  108. yield* Effect.yieldNow
  109. yield* events.publish(Message, { text: "hello" })
  110. yield* events.publish(CountMessage, { count: 2 })
  111. const received = Array.from(yield* Fiber.join(fiber)).map((event) =>
  112. event.type === "test.message" ? event.data.text : event.data.count,
  113. )
  114. expect(received).toEqual(["hello", 2])
  115. }),
  116. )
  117. it.effect("publishes events with the current location", () =>
  118. Effect.gen(function* () {
  119. const events = yield* EventV2.Service
  120. const fiber = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  121. yield* Effect.yieldNow
  122. const event = yield* events.publish(Message, { text: "hello" })
  123. const received = Array.from(yield* Fiber.join(fiber))
  124. expect(received).toEqual([event])
  125. expect(event.type).toBe("test.message")
  126. expect(event).not.toHaveProperty("version")
  127. expect(event.data).toEqual({ text: "hello" })
  128. expect(event.location).toEqual({
  129. directory: AbsolutePath.make("project"),
  130. workspaceID: WorkspaceV2.ID.make("wrk_test"),
  131. })
  132. }),
  133. )
  134. itWithoutLocation.effect("omits location when no location is available", () =>
  135. Effect.gen(function* () {
  136. const events = yield* EventV2.Service
  137. const event = yield* events.publish(GlobalMessage, { text: "hello" })
  138. expect(event).not.toHaveProperty("location")
  139. expect(event.type).toBe("test.global")
  140. }),
  141. )
  142. it.effect("publishes definition version", () =>
  143. Effect.gen(function* () {
  144. const events = yield* EventV2.Service
  145. const event = yield* events.publish(VersionedMessage, { id: "one", text: "hello" })
  146. expect(event.type).toBe("test.versioned")
  147. expect(event.durable?.version).toBe(EventV2.Version.make(2))
  148. }),
  149. )
  150. it.effect("selects the latest durable definition independent of declaration order", () =>
  151. Effect.sync(() => {
  152. const latest = EventV2.durable({
  153. type: "test.out-of-order",
  154. durable: { version: 2, aggregate: "id" },
  155. schema: { id: Schema.String },
  156. })
  157. const historical = EventV2.durable({
  158. type: "test.out-of-order",
  159. durable: { version: 1, aggregate: "id" },
  160. schema: { id: Schema.String },
  161. })
  162. expect(Event.latest([latest, historical]).get("test.out-of-order")).toBe(latest)
  163. expect(Event.latest([historical, latest]).get("test.out-of-order")).toBe(latest)
  164. }),
  165. )
  166. it.effect("publishes to typed and wildcard subscriptions", () =>
  167. Effect.gen(function* () {
  168. const events = yield* EventV2.Service
  169. const typed = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  170. const wildcard = yield* events.subscribe().pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  171. yield* Effect.yieldNow
  172. const event = yield* events.publish(Message, { text: "hello" })
  173. expect(Array.from(yield* Fiber.join(typed))).toEqual([event])
  174. expect(Array.from(yield* Fiber.join(wildcard))).toEqual([event])
  175. }),
  176. )
  177. it.effect("runs projectors inline", () =>
  178. Effect.gen(function* () {
  179. const events = yield* EventV2.Service
  180. const received = new Array<EventV2.Payload>()
  181. yield* events.project(SyncMessage, (event) =>
  182. Effect.sync(() => {
  183. received.push(event)
  184. }),
  185. )
  186. const event = yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  187. yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" })
  188. expect(received[0]).toEqual(event)
  189. expect(received[1]?.data).toEqual({ id: "one", text: "after unsubscribe" })
  190. }),
  191. )
  192. it.effect("commits local operational state inside a new durable event transaction", () =>
  193. Effect.gen(function* () {
  194. const events = yield* EventV2.Service
  195. const received = new Array<string>()
  196. const aggregateID = EventV2.ID.create()
  197. yield* events.project(SyncMessage, () => Effect.sync(() => received.push("projector")))
  198. yield* events.publish(
  199. SyncMessage,
  200. { id: aggregateID, text: "hello" },
  201. { commit: (seq) => Effect.sync(() => received.push(`commit:${seq}`)) },
  202. )
  203. expect(received).toEqual(["projector", "commit:0"])
  204. }),
  205. )
  206. it.effect("rolls back the durable event and projector when the local commit fails", () =>
  207. Effect.gen(function* () {
  208. const events = yield* EventV2.Service
  209. const { db } = yield* Database.Service
  210. const aggregateID = EventV2.ID.create()
  211. yield* db.run("CREATE TABLE IF NOT EXISTS event_commit_probe (value text NOT NULL)")
  212. yield* db.run("DELETE FROM event_commit_probe")
  213. yield* events.project(SyncMessage, () =>
  214. db.run("INSERT INTO event_commit_probe (value) VALUES ('projected')").pipe(Effect.orDie, Effect.asVoid),
  215. )
  216. const exit = yield* events
  217. .publish(SyncMessage, { id: aggregateID, text: "hello" }, { commit: () => Effect.die("commit failed") })
  218. .pipe(Effect.exit)
  219. expect(String(exit)).toContain("commit failed")
  220. expect(yield* db.all("SELECT value FROM event_commit_probe")).toEqual([])
  221. expect(yield* db.select().from(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).all()).toEqual([])
  222. expect(
  223. yield* db.select().from(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).all(),
  224. ).toEqual([])
  225. }),
  226. )
  227. it.effect("rejects local commit hooks on live-only events", () =>
  228. Effect.gen(function* () {
  229. const events = yield* EventV2.Service
  230. const exit = yield* events.publish(Message, { text: "hello" }, { commit: () => Effect.void }).pipe(Effect.exit)
  231. expect(String(exit)).toContain("Local commit hooks require a durable event")
  232. }),
  233. )
  234. it.effect("runs projectors before publishing to streams", () =>
  235. Effect.gen(function* () {
  236. const events = yield* EventV2.Service
  237. const received = new Array<string>()
  238. const fiber = yield* events.subscribe().pipe(
  239. Stream.take(1),
  240. Stream.runForEach(() => Effect.sync(() => received.push("stream"))),
  241. Effect.forkScoped,
  242. )
  243. yield* events.project(SyncMessage, (event) =>
  244. Effect.sync(() => {
  245. received.push(event.type)
  246. }),
  247. )
  248. yield* Effect.yieldNow
  249. yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  250. yield* Fiber.join(fiber)
  251. expect(received).toEqual([SyncMessage.type, "stream"])
  252. }),
  253. )
  254. it.effect("runs listeners inline after projectors", () =>
  255. Effect.gen(function* () {
  256. const events = yield* EventV2.Service
  257. const received = new Array<string>()
  258. yield* events.project(SyncMessage, () =>
  259. Effect.sync(() => {
  260. received.push("projector")
  261. }),
  262. )
  263. const unsubscribe = yield* events.listen(() =>
  264. Effect.sync(() => {
  265. received.push("listener")
  266. }),
  267. )
  268. yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  269. yield* unsubscribe
  270. yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" })
  271. expect(received).toEqual(["projector", "listener", "projector"])
  272. }),
  273. )
  274. it.effect("isolates observer defects after durable events commit", () =>
  275. Effect.gen(function* () {
  276. const events = yield* EventV2.Service
  277. const received = new Array<string>()
  278. yield* events.listen(() => {
  279. throw new Error("listener defect")
  280. })
  281. yield* events.listen((event) =>
  282. Effect.sync(() => {
  283. received.push(event.type)
  284. }),
  285. )
  286. const event = yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  287. expect(received).toEqual([SyncMessage.type])
  288. expect(event.durable?.seq).toBeNumber()
  289. }),
  290. )
  291. it.effect("notifies global listeners only after a durable event is committed", () =>
  292. Effect.gen(function* () {
  293. const events = yield* EventV2.Service
  294. const { db } = yield* Database.Service
  295. const aggregateID = EventV2.ID.create()
  296. const observed = new Array<{ id: string; seq: number }>()
  297. yield* events.listen((event) =>
  298. event.type !== SyncMessage.type
  299. ? Effect.void
  300. : db
  301. .select({ id: EventTable.id, seq: EventTable.seq })
  302. .from(EventTable)
  303. .where(eq(EventTable.id, event.id))
  304. .get()
  305. .pipe(
  306. Effect.orDie,
  307. Effect.tap((row) =>
  308. Effect.sync(() => {
  309. if (row) observed.push(row)
  310. }),
  311. ),
  312. Effect.asVoid,
  313. ),
  314. )
  315. const event = yield* events.publish(SyncMessage, { id: aggregateID, text: "committed" })
  316. if (!event.durable) throw new Error("Expected durable event metadata")
  317. expect(observed).toEqual([{ id: event.id, seq: event.durable.seq }])
  318. }),
  319. )
  320. it.effect("ends only an overflowing bounded subscriber without blocking other listeners", () =>
  321. Effect.gen(function* () {
  322. const events = yield* EventV2.Service
  323. const consuming = yield* Deferred.make<void>()
  324. const release = yield* Deferred.make<void>()
  325. const slowStream = yield* EventV2.liveBounded(events, { capacity: 1 })
  326. const fastStream = yield* EventV2.liveBounded(events, { capacity: 8 })
  327. const slow = yield* slowStream.pipe(
  328. Stream.runForEach(() => Deferred.succeed(consuming, undefined).pipe(Effect.andThen(Deferred.await(release)))),
  329. Effect.forkScoped,
  330. )
  331. const fast = yield* fastStream.pipe(Stream.take(4), Stream.runCollect, Effect.forkScoped)
  332. yield* events.publish(Message, { text: "one" })
  333. yield* Deferred.await(consuming)
  334. yield* events.publish(Message, { text: "two" })
  335. yield* events.publish(Message, { text: "overflow" })
  336. const last = yield* events.publish(Message, { text: "still delivered" })
  337. yield* Deferred.succeed(release, undefined)
  338. const slowExit = yield* Fiber.await(slow)
  339. expect(Exit.findErrorOption(slowExit).pipe(Option.getOrUndefined)).toBeInstanceOf(EventV2.SubscriberOverflowError)
  340. expect(Array.from(yield* Fiber.join(fast))).toEqual([
  341. expect.objectContaining({ data: { text: "one" } }),
  342. expect.objectContaining({ data: { text: "two" } }),
  343. expect.objectContaining({ data: { text: "overflow" } }),
  344. last,
  345. ])
  346. }),
  347. )
  348. it.effect("filters internal events before they enter a bounded server stream", () =>
  349. Effect.gen(function* () {
  350. const events = yield* EventV2.Service
  351. const stream = yield* EventV2.liveBounded(events, { capacity: 1, accept: EventManifest.isServer })
  352. const received = yield* stream.pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  353. yield* events.publish(McpEvent.ToolsChanged, { server: "one" })
  354. yield* events.publish(McpEvent.ToolsChanged, { server: "two" })
  355. const published = yield* events.publish(McpEvent.StatusChanged, { server: "example" })
  356. expect(Array.from(yield* Fiber.join(received))).toEqual([published])
  357. }),
  358. )
  359. it.effect("preserves observer interruption", () =>
  360. Effect.gen(function* () {
  361. const events = yield* EventV2.Service
  362. const { db } = yield* Database.Service
  363. yield* events.listen(() => Effect.interrupt)
  364. const exit = yield* events.publish(SyncMessage, { id: "interrupted", text: "hello" }).pipe(Effect.exit)
  365. const committed = yield* db
  366. .select({ id: EventTable.id })
  367. .from(EventTable)
  368. .where(eq(EventTable.aggregate_id, "interrupted"))
  369. .get()
  370. .pipe(Effect.orDie)
  371. expect(Exit.isFailure(exit) && Cause.hasInterrupts(exit.cause)).toBeTrue()
  372. expect(committed).toBeDefined()
  373. }),
  374. )
  375. it.effect("keeps live-only listener defects fail-fast", () =>
  376. Effect.gen(function* () {
  377. const events = yield* EventV2.Service
  378. const defect = new Error("listener defect")
  379. yield* events.listen(() => Effect.die(defect))
  380. expect(yield* events.publish(Message, { text: "hello" }).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect)
  381. }),
  382. )
  383. it.effect("inserts durable event rows on publish", () =>
  384. Effect.gen(function* () {
  385. const events = yield* EventV2.Service
  386. const { db } = yield* Database.Service
  387. const aggregateID = EventV2.ID.create()
  388. yield* events.publish(SyncMessage, { id: aggregateID, text: "first" })
  389. const rows = yield* db
  390. .select()
  391. .from(EventTable)
  392. .where(eq(EventTable.aggregate_id, aggregateID))
  393. .all()
  394. .pipe(Effect.orDie)
  395. expect(rows).toHaveLength(1)
  396. expect(rows[0]?.type).toBe(EventV2.versionedType(SyncMessage.type, 1))
  397. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  398. }),
  399. )
  400. it.effect("increments durable event seq per aggregate", () =>
  401. Effect.gen(function* () {
  402. const events = yield* EventV2.Service
  403. const { db } = yield* Database.Service
  404. const aggregateID = EventV2.ID.create()
  405. yield* events.publish(SyncMessage, { id: aggregateID, text: "first" })
  406. yield* events.publish(SyncMessage, { id: aggregateID, text: "second" })
  407. const rows = yield* db
  408. .select()
  409. .from(EventTable)
  410. .where(eq(EventTable.aggregate_id, aggregateID))
  411. .all()
  412. .pipe(Effect.orDie)
  413. expect(rows.map((row) => row.seq)).toEqual([0, 1])
  414. }),
  415. )
  416. it.effect("replays durable aggregate events after a sequence and tails new events", () =>
  417. Effect.gen(function* () {
  418. const events = yield* EventV2.Service
  419. const aggregateID = Session.ID.create()
  420. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  421. yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
  422. const fiber = yield* tail(events, { aggregateID, after: 0 }).pipe(
  423. Stream.take(2),
  424. Stream.runCollect,
  425. Effect.forkScoped,
  426. )
  427. yield* Effect.yieldNow
  428. yield* events.publish(DurableMessage, durableData(aggregateID, "two"))
  429. expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
  430. [1, durableData(aggregateID, "one")],
  431. [2, durableData(aggregateID, "two")],
  432. ])
  433. }),
  434. )
  435. it.effect("catches durable aggregate events published during replay handoff", () =>
  436. Effect.gen(function* () {
  437. const events = yield* EventV2.Service
  438. const aggregateID = Session.ID.create()
  439. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  440. const fiber = yield* tail(events, { aggregateID }).pipe(Stream.take(2), Stream.runCollect, Effect.forkScoped)
  441. yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
  442. expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
  443. [0, durableData(aggregateID, "zero")],
  444. [1, durableData(aggregateID, "one")],
  445. ])
  446. }),
  447. )
  448. it.effect("retains a durable wake committed while historical replay is paused", () =>
  449. Effect.gen(function* () {
  450. const readStarted = yield* Deferred.make<void>()
  451. const continueRead = yield* Deferred.make<void>()
  452. let pause = true
  453. const eventLayer = EventV2.layerWith({
  454. beforeAggregateRead: () =>
  455. pause
  456. ? Deferred.succeed(readStarted, undefined).pipe(Effect.andThen(Deferred.await(continueRead)))
  457. : Effect.void,
  458. }).pipe(Layer.provide(LayerNode.compile(Database.node)))
  459. yield* Effect.gen(function* () {
  460. const events = yield* EventV2.Service
  461. const aggregateID = Session.ID.create()
  462. const fiber = yield* tail(events, { aggregateID }).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  463. yield* Deferred.await(readStarted)
  464. pause = false
  465. yield* events.publish(DurableMessage, durableData(aggregateID, "during handoff"))
  466. yield* Deferred.succeed(continueRead, undefined)
  467. expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual([
  468. [0, durableData(aggregateID, "during handoff")],
  469. ])
  470. }).pipe(Effect.provide(Layer.merge(LayerNode.compile(Database.node), eventLayer)))
  471. }),
  472. )
  473. it.effect("coalesces durable aggregate wakes while draining every committed event", () =>
  474. Effect.gen(function* () {
  475. const events = yield* EventV2.Service
  476. const aggregateID = Session.ID.create()
  477. const count = 64
  478. const fiber = yield* tail(events, { aggregateID }).pipe(Stream.take(count), Stream.runCollect, Effect.forkScoped)
  479. yield* Effect.yieldNow
  480. for (let index = 0; index < count; index++) {
  481. yield* events.publish(DurableMessage, durableData(aggregateID, String(index)))
  482. }
  483. expect(Array.from(yield* Fiber.join(fiber)).map((event) => [event.durable?.seq, event.data])).toEqual(
  484. Array.from({ length: count }, (_, index) => [index, durableData(aggregateID, String(index))]),
  485. )
  486. }),
  487. )
  488. it.effect("omits live-only events from durable aggregate streams", () =>
  489. Effect.gen(function* () {
  490. const events = yield* EventV2.Service
  491. const aggregateID = Session.ID.create()
  492. const fiber = yield* tail(events, { aggregateID }).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  493. yield* Effect.yieldNow
  494. yield* events.publish(Message, { text: "live only" })
  495. yield* events.publish(DurableMessage, durableData(aggregateID, "durable"))
  496. expect(Array.from(yield* Fiber.join(fiber)).map((event) => event.type)).toEqual([DurableMessage.type])
  497. }),
  498. )
  499. it.effect("uses custom sync aggregate field", () =>
  500. Effect.gen(function* () {
  501. const events = yield* EventV2.Service
  502. const { db } = yield* Database.Service
  503. const aggregateID = EventV2.ID.create()
  504. yield* events.publish(SyncSent, { messageID: aggregateID, text: "sent" })
  505. const rows = yield* db
  506. .select()
  507. .from(EventTable)
  508. .where(eq(EventTable.aggregate_id, aggregateID))
  509. .all()
  510. .pipe(Effect.orDie)
  511. expect(rows).toHaveLength(1)
  512. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  513. }),
  514. )
  515. it.effect("replays durable events through projectors", () =>
  516. Effect.gen(function* () {
  517. const events = yield* EventV2.Service
  518. const received = new Array<EventV2.Payload>()
  519. yield* events.project(DurableMessage, (event) =>
  520. Effect.sync(() => {
  521. received.push(event)
  522. }),
  523. )
  524. const aggregateID = Session.ID.create()
  525. yield* events.replay({
  526. id: EventV2.ID.create(),
  527. created: DateTime.makeUnsafe(0),
  528. type: EventV2.versionedType(DurableMessage.type, 1),
  529. seq: 0,
  530. aggregateID,
  531. data: durableData(aggregateID, "hello"),
  532. })
  533. expect(received[0]?.type).toBe(DurableMessage.type)
  534. expect(received[0]?.data).toEqual(durableData(aggregateID, "hello"))
  535. }),
  536. )
  537. it.effect("replay inserts external event rows", () =>
  538. Effect.gen(function* () {
  539. const events = yield* EventV2.Service
  540. const { db } = yield* Database.Service
  541. const aggregateID = Session.ID.create()
  542. yield* events.replay({
  543. id: EventV2.ID.create(),
  544. created: DateTime.makeUnsafe(0),
  545. type: EventV2.versionedType(DurableMessage.type, 1),
  546. seq: 0,
  547. aggregateID,
  548. data: durableData(aggregateID, "replayed"),
  549. })
  550. const rows = yield* db
  551. .select()
  552. .from(EventTable)
  553. .where(eq(EventTable.aggregate_id, aggregateID))
  554. .all()
  555. .pipe(Effect.orDie)
  556. expect(rows).toHaveLength(1)
  557. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  558. }),
  559. )
  560. it.effect(
  561. "replay rejects an envelope aggregate that differs from its payload without mutating the payload aggregate",
  562. () =>
  563. Effect.gen(function* () {
  564. const events = yield* EventV2.Service
  565. const { db } = yield* Database.Service
  566. const envelopeAggregateID = Session.ID.create()
  567. const payloadAggregateID = Session.ID.create()
  568. const received = new Array<EventV2.Payload>()
  569. yield* events.publish(DurableMessage, durableData(payloadAggregateID, "seed"))
  570. yield* events.project(DurableMessage, (event) =>
  571. Effect.sync(() => {
  572. received.push(event)
  573. }),
  574. )
  575. const exit = yield* events
  576. .replay({
  577. id: EventV2.ID.create(),
  578. created: DateTime.makeUnsafe(0),
  579. type: EventV2.versionedType(DurableMessage.type, 1),
  580. seq: 1,
  581. aggregateID: envelopeAggregateID,
  582. data: durableData(payloadAggregateID, "replayed"),
  583. })
  584. .pipe(Effect.exit)
  585. const rows = yield* db
  586. .select()
  587. .from(EventTable)
  588. .where(eq(EventTable.aggregate_id, payloadAggregateID))
  589. .all()
  590. .pipe(Effect.orDie)
  591. const sequence = yield* db
  592. .select({ seq: EventSequenceTable.seq })
  593. .from(EventSequenceTable)
  594. .where(eq(EventSequenceTable.aggregate_id, payloadAggregateID))
  595. .get()
  596. .pipe(Effect.orDie)
  597. expect(String(exit)).toContain("Aggregate mismatch")
  598. expect(received).toHaveLength(0)
  599. expect(rows).toHaveLength(1)
  600. expect(sequence).toEqual({ seq: 0 })
  601. }),
  602. )
  603. it.effect("replay defects on sequence mismatch", () =>
  604. Effect.gen(function* () {
  605. const events = yield* EventV2.Service
  606. const aggregateID = Session.ID.create()
  607. yield* events.replay({
  608. id: EventV2.ID.create(),
  609. created: DateTime.makeUnsafe(0),
  610. type: EventV2.versionedType(DurableMessage.type, 1),
  611. seq: 0,
  612. aggregateID,
  613. data: durableData(aggregateID, "first"),
  614. })
  615. const exit = yield* events
  616. .replay({
  617. id: EventV2.ID.create(),
  618. created: DateTime.makeUnsafe(0),
  619. type: EventV2.versionedType(DurableMessage.type, 1),
  620. seq: 5,
  621. aggregateID,
  622. data: durableData(aggregateID, "bad"),
  623. })
  624. .pipe(Effect.exit)
  625. expect(String(exit)).toContain("Sequence mismatch")
  626. }),
  627. )
  628. it.effect("replay decodes synchronized transformed values before projection", () =>
  629. Effect.gen(function* () {
  630. const events = yield* EventV2.Service
  631. const aggregateID = Session.ID.create()
  632. const received = new Array<typeof SessionEvent.InstructionsUpdated.Type>()
  633. yield* events.project(SessionEvent.InstructionsUpdated, (event) =>
  634. Effect.sync(() => {
  635. received.push(event)
  636. }),
  637. )
  638. yield* events.replay({
  639. id: EventV2.ID.create(),
  640. created: DateTime.makeUnsafe(0),
  641. type: EventV2.versionedType(SessionEvent.InstructionsUpdated.type, 2),
  642. seq: 0,
  643. aggregateID,
  644. data: { sessionID: aggregateID, delta: { "core/context": "0".repeat(64) } },
  645. })
  646. expect(received[0]?.created).toEqual(DateTime.makeUnsafe(0))
  647. }),
  648. )
  649. it.effect("dispatches durable projectors by exact event version", () =>
  650. Effect.gen(function* () {
  651. const events = yield* EventV2.Service
  652. const aggregateID = Session.ID.create()
  653. const received = new Array<typeof VersionedMessageV2.Type>()
  654. yield* events.project(VersionedMessageV2, (event) =>
  655. Effect.sync(() => {
  656. received.push(event)
  657. }),
  658. )
  659. yield* events.publish(VersionedMessageV1, { id: aggregateID })
  660. yield* events.publish(VersionedMessageV2, { id: aggregateID })
  661. expect(received).toHaveLength(1)
  662. expect(received[0]?.durable.version).toBe(EventV2.Version.make(2))
  663. }),
  664. )
  665. it.effect("replay defects on unknown event type", () =>
  666. Effect.gen(function* () {
  667. const events = yield* EventV2.Service
  668. const exit = yield* events
  669. .replay({
  670. id: EventV2.ID.create(),
  671. created: DateTime.makeUnsafe(0),
  672. type: "unknown.event.1",
  673. seq: 0,
  674. aggregateID: EventV2.ID.create(),
  675. data: {},
  676. })
  677. .pipe(Effect.exit)
  678. expect(String(exit)).toContain("Unknown durable event type")
  679. }),
  680. )
  681. it.effect("replayAll validates contiguous aggregate events", () =>
  682. Effect.gen(function* () {
  683. const events = yield* EventV2.Service
  684. const aggregateID = Session.ID.create()
  685. const source = yield* events.replayAll([
  686. {
  687. id: EventV2.ID.create(),
  688. created: DateTime.makeUnsafe(0),
  689. type: EventV2.versionedType(DurableMessage.type, 1),
  690. seq: 0,
  691. aggregateID,
  692. data: durableData(aggregateID, "one"),
  693. },
  694. {
  695. id: EventV2.ID.create(),
  696. created: DateTime.makeUnsafe(0),
  697. type: EventV2.versionedType(DurableMessage.type, 1),
  698. seq: 1,
  699. aggregateID,
  700. data: durableData(aggregateID, "two"),
  701. },
  702. ])
  703. expect(source).toBe(aggregateID)
  704. }),
  705. )
  706. it.effect("replayAll accepts later chunks after the first batch", () =>
  707. Effect.gen(function* () {
  708. const events = yield* EventV2.Service
  709. const { db } = yield* Database.Service
  710. const aggregateID = Session.ID.create()
  711. const one = yield* events.replayAll([
  712. {
  713. id: EventV2.ID.create(),
  714. created: DateTime.makeUnsafe(0),
  715. type: EventV2.versionedType(DurableMessage.type, 1),
  716. seq: 0,
  717. aggregateID,
  718. data: durableData(aggregateID, "one"),
  719. },
  720. {
  721. id: EventV2.ID.create(),
  722. created: DateTime.makeUnsafe(0),
  723. type: EventV2.versionedType(DurableMessage.type, 1),
  724. seq: 1,
  725. aggregateID,
  726. data: durableData(aggregateID, "two"),
  727. },
  728. ])
  729. const two = yield* events.replayAll([
  730. {
  731. id: EventV2.ID.create(),
  732. created: DateTime.makeUnsafe(0),
  733. type: EventV2.versionedType(DurableMessage.type, 1),
  734. seq: 2,
  735. aggregateID,
  736. data: durableData(aggregateID, "three"),
  737. },
  738. {
  739. id: EventV2.ID.create(),
  740. created: DateTime.makeUnsafe(0),
  741. type: EventV2.versionedType(DurableMessage.type, 1),
  742. seq: 3,
  743. aggregateID,
  744. data: durableData(aggregateID, "four"),
  745. },
  746. ])
  747. const rows = yield* db
  748. .select()
  749. .from(EventTable)
  750. .where(eq(EventTable.aggregate_id, aggregateID))
  751. .all()
  752. .pipe(Effect.orDie)
  753. expect(one).toBe(aggregateID)
  754. expect(two).toBe(aggregateID)
  755. expect(rows.map((row) => row.seq)).toEqual([0, 1, 2, 3])
  756. }),
  757. )
  758. it.effect("claim fences replay owners", () =>
  759. Effect.gen(function* () {
  760. const events = yield* EventV2.Service
  761. const received = new Array<EventV2.Payload>()
  762. const aggregateID = Session.ID.create()
  763. yield* events.publish(DurableMessage, durableData(aggregateID, "seed"))
  764. yield* events.claim(aggregateID, "owner-a")
  765. yield* events.project(DurableMessage, (event) =>
  766. Effect.sync(() => {
  767. received.push(event)
  768. }),
  769. )
  770. yield* events.replay(
  771. {
  772. id: EventV2.ID.create(),
  773. created: DateTime.makeUnsafe(0),
  774. type: EventV2.versionedType(DurableMessage.type, 1),
  775. seq: 1,
  776. aggregateID,
  777. data: durableData(aggregateID, "ignored"),
  778. },
  779. { ownerID: "owner-b" },
  780. )
  781. expect(received).toHaveLength(0)
  782. }),
  783. )
  784. it.effect("strict owner fences exact replay", () =>
  785. Effect.gen(function* () {
  786. const events = yield* EventV2.Service
  787. const aggregateID = Session.ID.create()
  788. const id = EventV2.ID.create()
  789. const replayed = {
  790. id,
  791. created: DateTime.makeUnsafe(0),
  792. type: EventV2.versionedType(DurableMessage.type, 1),
  793. seq: 0,
  794. aggregateID,
  795. data: durableData(aggregateID, "owned"),
  796. }
  797. yield* events.replay(replayed, { ownerID: "owner-a" })
  798. const exit = yield* events.replay(replayed, { ownerID: "owner-b", strictOwner: true }).pipe(Effect.exit)
  799. expect(String(exit)).toContain("Replay owner mismatch")
  800. }),
  801. )
  802. it.effect("exact replay claims an unowned aggregate", () =>
  803. Effect.gen(function* () {
  804. const events = yield* EventV2.Service
  805. const { db } = yield* Database.Service
  806. const aggregateID = Session.ID.create()
  807. const published = yield* events.publish(DurableMessage, durableData(aggregateID, "owned"))
  808. const replayed = {
  809. id: published.id,
  810. created: published.created,
  811. type: EventV2.versionedType(DurableMessage.type, 1),
  812. seq: published.durable!.seq,
  813. aggregateID,
  814. data: published.data,
  815. }
  816. yield* events.replay(replayed, { ownerID: "owner-a", strictOwner: true })
  817. const row = yield* db
  818. .select({ ownerID: EventSequenceTable.owner_id })
  819. .from(EventSequenceTable)
  820. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  821. .get()
  822. .pipe(Effect.orDie)
  823. expect(row?.ownerID).toBe("owner-a")
  824. const exit = yield* events
  825. .replay(
  826. { ...replayed, id: EventV2.ID.create(), seq: 1, data: durableData(aggregateID, "conflict") },
  827. { ownerID: "owner-b", strictOwner: true },
  828. )
  829. .pipe(Effect.exit)
  830. expect(String(exit)).toContain("Replay owner mismatch")
  831. }),
  832. )
  833. it.effect("replay with owner claims an unowned sequence", () =>
  834. Effect.gen(function* () {
  835. const events = yield* EventV2.Service
  836. const { db } = yield* Database.Service
  837. const aggregateID = Session.ID.create()
  838. yield* events.replay(
  839. {
  840. id: EventV2.ID.create(),
  841. created: DateTime.makeUnsafe(0),
  842. type: EventV2.versionedType(DurableMessage.type, 1),
  843. seq: 0,
  844. aggregateID,
  845. data: durableData(aggregateID, "owned"),
  846. },
  847. { ownerID: "owner-1" },
  848. )
  849. const row = yield* db
  850. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  851. .from(EventSequenceTable)
  852. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  853. .get()
  854. .pipe(Effect.orDie)
  855. expect(row).toEqual({ seq: 0, ownerID: "owner-1" })
  856. }),
  857. )
  858. it.effect("replay claims an existing unowned sequence before fencing a different owner", () =>
  859. Effect.gen(function* () {
  860. const events = yield* EventV2.Service
  861. const { db } = yield* Database.Service
  862. const aggregateID = Session.ID.create()
  863. yield* events.publish(DurableMessage, durableData(aggregateID, "local"))
  864. yield* events.replay(
  865. {
  866. id: EventV2.ID.create(),
  867. created: DateTime.makeUnsafe(0),
  868. type: EventV2.versionedType(DurableMessage.type, 1),
  869. seq: 1,
  870. aggregateID,
  871. data: durableData(aggregateID, "claimed"),
  872. },
  873. { ownerID: "owner-1" },
  874. )
  875. yield* events.replay(
  876. {
  877. id: EventV2.ID.create(),
  878. created: DateTime.makeUnsafe(0),
  879. type: EventV2.versionedType(DurableMessage.type, 1),
  880. seq: 2,
  881. aggregateID,
  882. data: durableData(aggregateID, "fenced"),
  883. },
  884. { ownerID: "owner-2" },
  885. )
  886. const rows = yield* db
  887. .select()
  888. .from(EventTable)
  889. .where(eq(EventTable.aggregate_id, aggregateID))
  890. .all()
  891. .pipe(Effect.orDie)
  892. const sequence = yield* db
  893. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  894. .from(EventSequenceTable)
  895. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  896. .get()
  897. .pipe(Effect.orDie)
  898. expect(rows.map((row) => row.seq)).toEqual([0, 1])
  899. expect(sequence).toEqual({ seq: 1, ownerID: "owner-1" })
  900. }),
  901. )
  902. it.effect("strict replay rejects an owner conflict instead of silently skipping it", () =>
  903. Effect.gen(function* () {
  904. const events = yield* EventV2.Service
  905. const aggregateID = Session.ID.create()
  906. yield* events.replay(
  907. {
  908. id: EventV2.ID.create(),
  909. created: DateTime.makeUnsafe(0),
  910. type: EventV2.versionedType(DurableMessage.type, 1),
  911. seq: 0,
  912. aggregateID,
  913. data: durableData(aggregateID, "claimed"),
  914. },
  915. { ownerID: "owner-1" },
  916. )
  917. const exit = yield* events
  918. .replay(
  919. {
  920. id: EventV2.ID.create(),
  921. created: DateTime.makeUnsafe(0),
  922. type: EventV2.versionedType(DurableMessage.type, 1),
  923. seq: 1,
  924. aggregateID,
  925. data: durableData(aggregateID, "conflict"),
  926. },
  927. { ownerID: "owner-2", strictOwner: true },
  928. )
  929. .pipe(Effect.exit)
  930. expect(String(exit)).toContain("Replay owner mismatch")
  931. }),
  932. )
  933. it.effect("publishes accepted replay with its durable sequence and suppresses stale replay", () =>
  934. Effect.gen(function* () {
  935. const events = yield* EventV2.Service
  936. const received = new Array<EventV2.Payload>()
  937. const aggregateID = Session.ID.create()
  938. yield* events.listen((event) => Effect.sync(() => received.push(event)))
  939. const replayed = {
  940. id: EventV2.ID.create(),
  941. created: DateTime.makeUnsafe(0),
  942. type: EventV2.versionedType(DurableMessage.type, 1),
  943. seq: 0,
  944. aggregateID,
  945. data: durableData(aggregateID, "replayed"),
  946. }
  947. yield* events.replay(replayed, { publish: true })
  948. yield* events.replay(replayed, { publish: true })
  949. expect(received).toMatchObject([{ id: replayed.id, durable: { seq: 0, version: 1 }, data: replayed.data }])
  950. }),
  951. )
  952. it.effect("rejects divergent stale replay without publishing it", () =>
  953. Effect.gen(function* () {
  954. const events = yield* EventV2.Service
  955. const received = new Array<EventV2.Payload>()
  956. const aggregateID = Session.ID.create()
  957. const replayed = {
  958. id: EventV2.ID.create(),
  959. created: DateTime.makeUnsafe(0),
  960. type: EventV2.versionedType(DurableMessage.type, 1),
  961. seq: 0,
  962. aggregateID,
  963. data: durableData(aggregateID, "original"),
  964. }
  965. yield* events.listen((event) => Effect.sync(() => received.push(event)))
  966. yield* events.replay(replayed, { publish: true })
  967. const exit = yield* events
  968. .replay({ ...replayed, data: durableData(aggregateID, "divergent") }, { publish: true })
  969. .pipe(Effect.exit)
  970. expect(String(exit)).toContain("Replay diverged")
  971. expect(received).toHaveLength(1)
  972. }),
  973. )
  974. it.effect("rejects an event ID reused at another aggregate position", () =>
  975. Effect.gen(function* () {
  976. const events = yield* EventV2.Service
  977. const aggregateID = Session.ID.create()
  978. const id = EventV2.ID.create()
  979. yield* events.replay({
  980. id,
  981. created: DateTime.makeUnsafe(0),
  982. type: EventV2.versionedType(DurableMessage.type, 1),
  983. seq: 0,
  984. aggregateID,
  985. data: durableData(aggregateID, "first"),
  986. })
  987. const exit = yield* events
  988. .replay({
  989. id,
  990. created: DateTime.makeUnsafe(0),
  991. type: EventV2.versionedType(DurableMessage.type, 1),
  992. seq: 1,
  993. aggregateID,
  994. data: durableData(aggregateID, "second"),
  995. })
  996. .pipe(Effect.exit)
  997. expect(String(exit)).toContain(`Event ${id} already exists`)
  998. }),
  999. )
  1000. it.effect("replay from a different owner leaves claimed sequence unchanged", () =>
  1001. Effect.gen(function* () {
  1002. const events = yield* EventV2.Service
  1003. const { db } = yield* Database.Service
  1004. const aggregateID = Session.ID.create()
  1005. const received = new Array<EventV2.Payload>()
  1006. yield* events.listen((event) => Effect.sync(() => received.push(event)))
  1007. yield* events.replay(
  1008. {
  1009. id: EventV2.ID.create(),
  1010. created: DateTime.makeUnsafe(0),
  1011. type: EventV2.versionedType(DurableMessage.type, 1),
  1012. seq: 0,
  1013. aggregateID,
  1014. data: durableData(aggregateID, "first"),
  1015. },
  1016. { ownerID: "owner-1" },
  1017. )
  1018. yield* events.replay(
  1019. {
  1020. id: EventV2.ID.create(),
  1021. created: DateTime.makeUnsafe(0),
  1022. type: EventV2.versionedType(DurableMessage.type, 1),
  1023. seq: 1,
  1024. aggregateID,
  1025. data: durableData(aggregateID, "ignored"),
  1026. },
  1027. { ownerID: "owner-2", publish: true },
  1028. )
  1029. const rows = yield* db
  1030. .select()
  1031. .from(EventTable)
  1032. .where(eq(EventTable.aggregate_id, aggregateID))
  1033. .all()
  1034. .pipe(Effect.orDie)
  1035. const sequence = yield* db
  1036. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  1037. .from(EventSequenceTable)
  1038. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  1039. .get()
  1040. .pipe(Effect.orDie)
  1041. expect(rows).toHaveLength(1)
  1042. expect(sequence).toEqual({ seq: 0, ownerID: "owner-1" })
  1043. expect(received).toHaveLength(0)
  1044. }),
  1045. )
  1046. it.effect("claim updates the event sequence owner", () =>
  1047. Effect.gen(function* () {
  1048. const events = yield* EventV2.Service
  1049. const { db } = yield* Database.Service
  1050. const aggregateID = EventV2.ID.create()
  1051. yield* events.publish(SyncMessage, { id: aggregateID, text: "claimed" })
  1052. yield* events.claim(aggregateID, "owner-1")
  1053. yield* events.claim(aggregateID, "owner-2")
  1054. const row = yield* db
  1055. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  1056. .from(EventSequenceTable)
  1057. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  1058. .get()
  1059. .pipe(Effect.orDie)
  1060. expect(row).toEqual({ seq: 0, ownerID: "owner-2" })
  1061. }),
  1062. )
  1063. it.effect("remove clears durable event sequence", () =>
  1064. Effect.gen(function* () {
  1065. const events = yield* EventV2.Service
  1066. const received = new Array<EventV2.Payload>()
  1067. const aggregateID = Session.ID.create()
  1068. yield* events.publish(DurableMessage, durableData(aggregateID, "seed"))
  1069. yield* events.remove(aggregateID)
  1070. yield* events.project(DurableMessage, (event) =>
  1071. Effect.sync(() => {
  1072. received.push(event)
  1073. }),
  1074. )
  1075. yield* events.replay({
  1076. id: EventV2.ID.create(),
  1077. created: DateTime.makeUnsafe(0),
  1078. type: EventV2.versionedType(DurableMessage.type, 1),
  1079. seq: 0,
  1080. aggregateID,
  1081. data: durableData(aggregateID, "replayed"),
  1082. })
  1083. expect(received[0]?.data).toEqual(durableData(aggregateID, "replayed"))
  1084. }),
  1085. )
  1086. it.effect("log without follow replays events and completes with a synced marker", () =>
  1087. Effect.gen(function* () {
  1088. const events = yield* EventV2.Service
  1089. const aggregateID = Session.ID.create()
  1090. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  1091. yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
  1092. const items = Array.from(yield* Stream.runCollect(events.log({ aggregateID })))
  1093. expect(items.map((item) => (EventV2.isSynced(item) ? item.type : item.durable?.seq))).toEqual([
  1094. EventV2.Seq.make(0),
  1095. EventV2.Seq.make(1),
  1096. "log.synced",
  1097. ])
  1098. expect(items.at(-1)).toEqual({ type: "log.synced", aggregateID, seq: EventV2.Seq.make(1) })
  1099. }),
  1100. )
  1101. it.effect("log synced marker omits seq for an empty log and keeps the cursor otherwise", () =>
  1102. Effect.gen(function* () {
  1103. const events = yield* EventV2.Service
  1104. const aggregateID = Session.ID.create()
  1105. const empty = Array.from(yield* Stream.runCollect(events.log({ aggregateID })))
  1106. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  1107. const drained = Array.from(yield* Stream.runCollect(events.log({ aggregateID, after: 0 })))
  1108. expect(empty).toEqual([{ type: "log.synced", aggregateID }])
  1109. expect(empty[0]).not.toHaveProperty("seq")
  1110. expect(drained).toEqual([{ type: "log.synced", aggregateID, seq: EventV2.Seq.make(0) }])
  1111. }),
  1112. )
  1113. it.effect("log with follow emits the synced marker at the replay-to-live boundary", () =>
  1114. Effect.gen(function* () {
  1115. const events = yield* EventV2.Service
  1116. const aggregateID = Session.ID.create()
  1117. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  1118. const fiber = yield* events
  1119. .log({ aggregateID, follow: true })
  1120. .pipe(Stream.take(3), Stream.runCollect, Effect.forkScoped)
  1121. yield* Effect.yieldNow
  1122. yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
  1123. const items = Array.from(yield* Fiber.join(fiber))
  1124. expect(items.map((item) => (EventV2.isSynced(item) ? item : item.durable?.seq))).toEqual([
  1125. EventV2.Seq.make(0),
  1126. { type: "log.synced", aggregateID, seq: EventV2.Seq.make(0) },
  1127. EventV2.Seq.make(1),
  1128. ])
  1129. }),
  1130. )
  1131. it.effect("log replays across configured read pages", () =>
  1132. Effect.gen(function* () {
  1133. const eventLayer = EventV2.layerWith({ logReadPageSize: 2 }).pipe(Layer.provide(LayerNode.compile(Database.node)))
  1134. yield* Effect.gen(function* () {
  1135. const events = yield* EventV2.Service
  1136. const aggregateID = Session.ID.create()
  1137. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  1138. yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
  1139. yield* events.publish(DurableMessage, durableData(aggregateID, "two"))
  1140. yield* events.publish(DurableMessage, durableData(aggregateID, "three"))
  1141. yield* events.publish(DurableMessage, durableData(aggregateID, "four"))
  1142. const items = Array.from(yield* Stream.runCollect(events.log({ aggregateID })))
  1143. expect(items.map((item) => (EventV2.isSynced(item) ? item.type : item.durable?.seq))).toEqual([
  1144. EventV2.Seq.make(0),
  1145. EventV2.Seq.make(1),
  1146. EventV2.Seq.make(2),
  1147. EventV2.Seq.make(3),
  1148. EventV2.Seq.make(4),
  1149. "log.synced",
  1150. ])
  1151. expect(items.at(-1)).toEqual({ type: "log.synced", aggregateID, seq: EventV2.Seq.make(4) })
  1152. }).pipe(Effect.provide(Layer.merge(LayerNode.compile(Database.node), eventLayer)))
  1153. }),
  1154. )
  1155. it.effect("log with follow emits events committed during replay after the synced marker", () =>
  1156. Effect.gen(function* () {
  1157. const readStarted = yield* Deferred.make<void>()
  1158. const releaseRead = yield* Deferred.make<void>()
  1159. const firstRead = yield* Ref.make(true)
  1160. const eventLayer = EventV2.layerWith({
  1161. beforeAggregateRead: () =>
  1162. Ref.getAndSet(firstRead, false).pipe(
  1163. Effect.flatMap((shouldBlock) => {
  1164. if (!shouldBlock) return Effect.void
  1165. return Deferred.succeed(readStarted, undefined).pipe(Effect.andThen(Deferred.await(releaseRead)))
  1166. }),
  1167. ),
  1168. }).pipe(Layer.provide(LayerNode.compile(Database.node)))
  1169. yield* Effect.gen(function* () {
  1170. const events = yield* EventV2.Service
  1171. const aggregateID = Session.ID.create()
  1172. yield* events.publish(DurableMessage, durableData(aggregateID, "zero"))
  1173. const fiber = yield* events
  1174. .log({ aggregateID, follow: true })
  1175. .pipe(Stream.take(3), Stream.runCollect, Effect.forkScoped)
  1176. yield* Deferred.await(readStarted)
  1177. yield* events.publish(DurableMessage, durableData(aggregateID, "one"))
  1178. yield* Deferred.succeed(releaseRead, undefined)
  1179. const items = Array.from(yield* Fiber.join(fiber))
  1180. expect(items.map((item) => (EventV2.isSynced(item) ? item : item.durable?.seq))).toEqual([
  1181. EventV2.Seq.make(0),
  1182. { type: "log.synced", aggregateID, seq: EventV2.Seq.make(0) },
  1183. EventV2.Seq.make(1),
  1184. ])
  1185. }).pipe(Effect.provide(Layer.merge(LayerNode.compile(Database.node), eventLayer)))
  1186. }),
  1187. )
  1188. it.effect("sequences returns the latest committed seq per aggregate and omits unknown aggregates", () =>
  1189. Effect.gen(function* () {
  1190. const events = yield* EventV2.Service
  1191. const first = Session.ID.create()
  1192. const second = Session.ID.create()
  1193. yield* events.publish(DurableMessage, durableData(first, "zero"))
  1194. yield* events.publish(DurableMessage, durableData(first, "one"))
  1195. yield* events.publish(DurableMessage, durableData(second, "zero"))
  1196. const sequences = yield* events.sequences([first, second, Session.ID.create()])
  1197. expect(sequences).toEqual(
  1198. new Map([
  1199. [first, EventV2.Seq.make(1)],
  1200. [second, EventV2.Seq.make(0)],
  1201. ]),
  1202. )
  1203. expect(yield* events.sequences([])).toEqual(new Map())
  1204. }),
  1205. )
  1206. })