event.test.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576
  1. import { describe, expect } from "bun:test"
  2. import { 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 { eq } from "drizzle-orm"
  9. import { location } from "./fixture/location"
  10. import { testEffect } from "./lib/effect"
  11. const locationLayer = Layer.succeed(
  12. Location.Service,
  13. Location.Service.of(location({ directory: AbsolutePath.make("project"), workspaceID: "workspace" })),
  14. )
  15. const eventLayer = Layer.mergeAll(EventV2.defaultLayer, Database.defaultLayer)
  16. const it = testEffect(eventLayer.pipe(Layer.provideMerge(locationLayer)))
  17. const itWithoutLocation = testEffect(eventLayer)
  18. const Message = EventV2.define({
  19. type: "test.message",
  20. schema: {
  21. text: Schema.String,
  22. },
  23. })
  24. const SyncMessage = EventV2.define({
  25. type: "test.sync",
  26. sync: {
  27. version: 1,
  28. aggregate: "id",
  29. },
  30. schema: {
  31. id: Schema.String,
  32. text: Schema.String,
  33. },
  34. })
  35. const SyncSent = EventV2.define({
  36. type: "test.sent",
  37. sync: {
  38. version: 1,
  39. aggregate: "messageID",
  40. },
  41. schema: {
  42. messageID: Schema.String,
  43. text: Schema.String,
  44. },
  45. })
  46. const GlobalMessage = EventV2.define({
  47. type: "test.global",
  48. schema: {
  49. text: Schema.String,
  50. },
  51. })
  52. const VersionedMessage = EventV2.define({
  53. type: "test.versioned",
  54. sync: {
  55. version: 2,
  56. aggregate: "id",
  57. },
  58. schema: {
  59. id: Schema.String,
  60. text: Schema.String,
  61. },
  62. })
  63. describe("EventV2", () => {
  64. it.effect("publishes events with the current location", () =>
  65. Effect.gen(function* () {
  66. const events = yield* EventV2.Service
  67. const fiber = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  68. yield* Effect.yieldNow
  69. const event = yield* events.publish(Message, { text: "hello" })
  70. const received = Array.from(yield* Fiber.join(fiber))
  71. expect(received).toEqual([event])
  72. expect(event.type).toBe("test.message")
  73. expect(event).not.toHaveProperty("version")
  74. expect(event.data).toEqual({ text: "hello" })
  75. expect(event.location).toEqual({ directory: AbsolutePath.make("project"), workspaceID: "workspace" })
  76. }),
  77. )
  78. itWithoutLocation.effect("omits location when no location is available", () =>
  79. Effect.gen(function* () {
  80. const events = yield* EventV2.Service
  81. const event = yield* events.publish(GlobalMessage, { text: "hello" })
  82. expect(event).not.toHaveProperty("location")
  83. expect(event.type).toBe("test.global")
  84. }),
  85. )
  86. it.effect("publishes definition version", () =>
  87. Effect.gen(function* () {
  88. const events = yield* EventV2.Service
  89. const event = yield* events.publish(VersionedMessage, { id: "one", text: "hello" })
  90. expect(event.type).toBe("test.versioned")
  91. expect(event.version).toBe(2)
  92. }),
  93. )
  94. it.effect("stores definitions in the exported registry", () =>
  95. Effect.sync(() => {
  96. expect(EventV2.registry.get(Message.type)).toBe(Message)
  97. }),
  98. )
  99. it.effect("keeps the latest sync definition in the registry", () =>
  100. Effect.sync(() => {
  101. const latest = EventV2.define({
  102. type: "test.out-of-order",
  103. sync: { version: 2, aggregate: "id" },
  104. schema: { id: Schema.String },
  105. })
  106. EventV2.define({
  107. type: "test.out-of-order",
  108. sync: { version: 1, aggregate: "id" },
  109. schema: { id: Schema.String },
  110. })
  111. expect(EventV2.registry.get("test.out-of-order")).toBe(latest)
  112. }),
  113. )
  114. it.effect("publishes to typed and wildcard subscriptions", () =>
  115. Effect.gen(function* () {
  116. const events = yield* EventV2.Service
  117. const typed = yield* events.subscribe(Message).pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  118. const wildcard = yield* events.all().pipe(Stream.take(1), Stream.runCollect, Effect.forkScoped)
  119. yield* Effect.yieldNow
  120. const event = yield* events.publish(Message, { text: "hello" })
  121. expect(Array.from(yield* Fiber.join(typed))).toEqual([event])
  122. expect(Array.from(yield* Fiber.join(wildcard))).toEqual([event])
  123. }),
  124. )
  125. it.effect("runs projectors inline", () =>
  126. Effect.gen(function* () {
  127. const events = yield* EventV2.Service
  128. const received = new Array<EventV2.Payload>()
  129. yield* events.project(SyncMessage, (event) =>
  130. Effect.sync(() => {
  131. received.push(event)
  132. }),
  133. )
  134. const event = yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  135. yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" })
  136. expect(received[0]).toEqual(event)
  137. expect(received[1]?.data).toEqual({ id: "one", text: "after unsubscribe" })
  138. }),
  139. )
  140. it.effect("runs projectors before publishing to streams", () =>
  141. Effect.gen(function* () {
  142. const events = yield* EventV2.Service
  143. const received = new Array<string>()
  144. const fiber = yield* events.all().pipe(
  145. Stream.take(1),
  146. Stream.runForEach(() => Effect.sync(() => received.push("stream"))),
  147. Effect.forkScoped,
  148. )
  149. yield* events.project(SyncMessage, (event) =>
  150. Effect.sync(() => {
  151. received.push(event.type)
  152. }),
  153. )
  154. yield* Effect.yieldNow
  155. yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  156. yield* Fiber.join(fiber)
  157. expect(received).toEqual([SyncMessage.type, "stream"])
  158. }),
  159. )
  160. it.effect("runs listeners inline after projectors", () =>
  161. Effect.gen(function* () {
  162. const events = yield* EventV2.Service
  163. const received = new Array<string>()
  164. yield* events.project(SyncMessage, () =>
  165. Effect.sync(() => {
  166. received.push("projector")
  167. }),
  168. )
  169. const unsubscribe = yield* events.listen(() =>
  170. Effect.sync(() => {
  171. received.push("listener")
  172. }),
  173. )
  174. yield* events.publish(SyncMessage, { id: "one", text: "hello" })
  175. yield* unsubscribe
  176. yield* events.publish(SyncMessage, { id: "one", text: "after unsubscribe" })
  177. expect(received).toEqual(["projector", "listener", "projector"])
  178. }),
  179. )
  180. it.effect("inserts sync event rows on publish", () =>
  181. Effect.gen(function* () {
  182. const events = yield* EventV2.Service
  183. const { db } = yield* Database.Service
  184. const aggregateID = EventV2.ID.create()
  185. yield* events.publish(SyncMessage, { id: aggregateID, text: "first" })
  186. const rows = yield* db
  187. .select()
  188. .from(EventTable)
  189. .where(eq(EventTable.aggregate_id, aggregateID))
  190. .all()
  191. .pipe(Effect.orDie)
  192. expect(rows).toHaveLength(1)
  193. expect(rows[0]?.type).toBe(EventV2.versionedType(SyncMessage.type, 1))
  194. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  195. }),
  196. )
  197. it.effect("increments sync event seq per aggregate", () =>
  198. Effect.gen(function* () {
  199. const events = yield* EventV2.Service
  200. const { db } = yield* Database.Service
  201. const aggregateID = EventV2.ID.create()
  202. yield* events.publish(SyncMessage, { id: aggregateID, text: "first" })
  203. yield* events.publish(SyncMessage, { id: aggregateID, text: "second" })
  204. const rows = yield* db
  205. .select()
  206. .from(EventTable)
  207. .where(eq(EventTable.aggregate_id, aggregateID))
  208. .all()
  209. .pipe(Effect.orDie)
  210. expect(rows.map((row) => row.seq)).toEqual([0, 1])
  211. }),
  212. )
  213. it.effect("uses custom sync aggregate field", () =>
  214. Effect.gen(function* () {
  215. const events = yield* EventV2.Service
  216. const { db } = yield* Database.Service
  217. const aggregateID = EventV2.ID.create()
  218. yield* events.publish(SyncSent, { messageID: aggregateID, text: "sent" })
  219. const rows = yield* db
  220. .select()
  221. .from(EventTable)
  222. .where(eq(EventTable.aggregate_id, aggregateID))
  223. .all()
  224. .pipe(Effect.orDie)
  225. expect(rows).toHaveLength(1)
  226. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  227. }),
  228. )
  229. it.effect("replays sync events through projectors", () =>
  230. Effect.gen(function* () {
  231. const events = yield* EventV2.Service
  232. const received = new Array<EventV2.Payload>()
  233. yield* events.project(SyncMessage, (event) =>
  234. Effect.sync(() => {
  235. received.push(event)
  236. }),
  237. )
  238. const aggregateID = EventV2.ID.create()
  239. yield* events.replay({
  240. id: EventV2.ID.create(),
  241. type: EventV2.versionedType(SyncMessage.type, 1),
  242. seq: 0,
  243. aggregateID,
  244. data: { id: aggregateID, text: "hello" },
  245. })
  246. expect(received[0]?.type).toBe(SyncMessage.type)
  247. expect(received[0]?.data).toEqual({ id: aggregateID, text: "hello" })
  248. }),
  249. )
  250. it.effect("replay inserts external event rows", () =>
  251. Effect.gen(function* () {
  252. const events = yield* EventV2.Service
  253. const { db } = yield* Database.Service
  254. const aggregateID = EventV2.ID.create()
  255. yield* events.replay({
  256. id: EventV2.ID.create(),
  257. type: EventV2.versionedType(SyncMessage.type, 1),
  258. seq: 0,
  259. aggregateID,
  260. data: { id: aggregateID, text: "replayed" },
  261. })
  262. const rows = yield* db
  263. .select()
  264. .from(EventTable)
  265. .where(eq(EventTable.aggregate_id, aggregateID))
  266. .all()
  267. .pipe(Effect.orDie)
  268. expect(rows).toHaveLength(1)
  269. expect(rows[0]?.aggregate_id).toBe(aggregateID)
  270. }),
  271. )
  272. it.effect("replay defects on sequence mismatch", () =>
  273. Effect.gen(function* () {
  274. const events = yield* EventV2.Service
  275. const aggregateID = EventV2.ID.create()
  276. yield* events.replay({
  277. id: EventV2.ID.create(),
  278. type: EventV2.versionedType(SyncMessage.type, 1),
  279. seq: 0,
  280. aggregateID,
  281. data: { id: aggregateID, text: "first" },
  282. })
  283. const exit = yield* events
  284. .replay({
  285. id: EventV2.ID.create(),
  286. type: EventV2.versionedType(SyncMessage.type, 1),
  287. seq: 5,
  288. aggregateID,
  289. data: { id: aggregateID, text: "bad" },
  290. })
  291. .pipe(Effect.exit)
  292. expect(String(exit)).toContain("Sequence mismatch")
  293. }),
  294. )
  295. it.effect("replay defects on unknown event type", () =>
  296. Effect.gen(function* () {
  297. const events = yield* EventV2.Service
  298. const exit = yield* events
  299. .replay({
  300. id: EventV2.ID.create(),
  301. type: "unknown.event.1",
  302. seq: 0,
  303. aggregateID: EventV2.ID.create(),
  304. data: {},
  305. })
  306. .pipe(Effect.exit)
  307. expect(String(exit)).toContain("Unknown sync event type")
  308. }),
  309. )
  310. it.effect("replayAll validates contiguous aggregate events", () =>
  311. Effect.gen(function* () {
  312. const events = yield* EventV2.Service
  313. const aggregateID = EventV2.ID.create()
  314. const source = yield* events.replayAll([
  315. {
  316. id: EventV2.ID.create(),
  317. type: EventV2.versionedType(SyncMessage.type, 1),
  318. seq: 0,
  319. aggregateID,
  320. data: { id: aggregateID, text: "one" },
  321. },
  322. {
  323. id: EventV2.ID.create(),
  324. type: EventV2.versionedType(SyncMessage.type, 1),
  325. seq: 1,
  326. aggregateID,
  327. data: { id: aggregateID, text: "two" },
  328. },
  329. ])
  330. expect(source).toBe(aggregateID)
  331. }),
  332. )
  333. it.effect("replayAll accepts later chunks after the first batch", () =>
  334. Effect.gen(function* () {
  335. const events = yield* EventV2.Service
  336. const { db } = yield* Database.Service
  337. const aggregateID = EventV2.ID.create()
  338. const one = yield* events.replayAll([
  339. {
  340. id: EventV2.ID.create(),
  341. type: EventV2.versionedType(SyncMessage.type, 1),
  342. seq: 0,
  343. aggregateID,
  344. data: { id: aggregateID, text: "one" },
  345. },
  346. {
  347. id: EventV2.ID.create(),
  348. type: EventV2.versionedType(SyncMessage.type, 1),
  349. seq: 1,
  350. aggregateID,
  351. data: { id: aggregateID, text: "two" },
  352. },
  353. ])
  354. const two = yield* events.replayAll([
  355. {
  356. id: EventV2.ID.create(),
  357. type: EventV2.versionedType(SyncMessage.type, 1),
  358. seq: 2,
  359. aggregateID,
  360. data: { id: aggregateID, text: "three" },
  361. },
  362. {
  363. id: EventV2.ID.create(),
  364. type: EventV2.versionedType(SyncMessage.type, 1),
  365. seq: 3,
  366. aggregateID,
  367. data: { id: aggregateID, text: "four" },
  368. },
  369. ])
  370. const rows = yield* db
  371. .select()
  372. .from(EventTable)
  373. .where(eq(EventTable.aggregate_id, aggregateID))
  374. .all()
  375. .pipe(Effect.orDie)
  376. expect(one).toBe(aggregateID)
  377. expect(two).toBe(aggregateID)
  378. expect(rows.map((row) => row.seq)).toEqual([0, 1, 2, 3])
  379. }),
  380. )
  381. it.effect("claim fences replay owners", () =>
  382. Effect.gen(function* () {
  383. const events = yield* EventV2.Service
  384. const received = new Array<EventV2.Payload>()
  385. const aggregateID = EventV2.ID.create()
  386. yield* events.publish(SyncMessage, { id: aggregateID, text: "seed" })
  387. yield* events.claim(aggregateID, "owner-a")
  388. yield* events.project(SyncMessage, (event) =>
  389. Effect.sync(() => {
  390. received.push(event)
  391. }),
  392. )
  393. yield* events.replay(
  394. {
  395. id: EventV2.ID.create(),
  396. type: EventV2.versionedType(SyncMessage.type, 1),
  397. seq: 1,
  398. aggregateID,
  399. data: { id: aggregateID, text: "ignored" },
  400. },
  401. { ownerID: "owner-b" },
  402. )
  403. expect(received).toHaveLength(0)
  404. }),
  405. )
  406. it.effect("replay with owner claims an unowned sequence", () =>
  407. Effect.gen(function* () {
  408. const events = yield* EventV2.Service
  409. const { db } = yield* Database.Service
  410. const aggregateID = EventV2.ID.create()
  411. yield* events.replay(
  412. {
  413. id: EventV2.ID.create(),
  414. type: EventV2.versionedType(SyncMessage.type, 1),
  415. seq: 0,
  416. aggregateID,
  417. data: { id: aggregateID, text: "owned" },
  418. },
  419. { ownerID: "owner-1" },
  420. )
  421. const row = yield* db
  422. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  423. .from(EventSequenceTable)
  424. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  425. .get()
  426. .pipe(Effect.orDie)
  427. expect(row).toEqual({ seq: 0, ownerID: "owner-1" })
  428. }),
  429. )
  430. it.effect("replay from a different owner leaves claimed sequence unchanged", () =>
  431. Effect.gen(function* () {
  432. const events = yield* EventV2.Service
  433. const { db } = yield* Database.Service
  434. const aggregateID = EventV2.ID.create()
  435. yield* events.replay(
  436. {
  437. id: EventV2.ID.create(),
  438. type: EventV2.versionedType(SyncMessage.type, 1),
  439. seq: 0,
  440. aggregateID,
  441. data: { id: aggregateID, text: "first" },
  442. },
  443. { ownerID: "owner-1" },
  444. )
  445. yield* events.replay(
  446. {
  447. id: EventV2.ID.create(),
  448. type: EventV2.versionedType(SyncMessage.type, 1),
  449. seq: 1,
  450. aggregateID,
  451. data: { id: aggregateID, text: "ignored" },
  452. },
  453. { ownerID: "owner-2" },
  454. )
  455. const rows = yield* db
  456. .select()
  457. .from(EventTable)
  458. .where(eq(EventTable.aggregate_id, aggregateID))
  459. .all()
  460. .pipe(Effect.orDie)
  461. const sequence = yield* db
  462. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  463. .from(EventSequenceTable)
  464. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  465. .get()
  466. .pipe(Effect.orDie)
  467. expect(rows).toHaveLength(1)
  468. expect(sequence).toEqual({ seq: 0, ownerID: "owner-1" })
  469. }),
  470. )
  471. it.effect("claim updates the event sequence owner", () =>
  472. Effect.gen(function* () {
  473. const events = yield* EventV2.Service
  474. const { db } = yield* Database.Service
  475. const aggregateID = EventV2.ID.create()
  476. yield* events.publish(SyncMessage, { id: aggregateID, text: "claimed" })
  477. yield* events.claim(aggregateID, "owner-1")
  478. yield* events.claim(aggregateID, "owner-2")
  479. const row = yield* db
  480. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  481. .from(EventSequenceTable)
  482. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  483. .get()
  484. .pipe(Effect.orDie)
  485. expect(row).toEqual({ seq: 0, ownerID: "owner-2" })
  486. }),
  487. )
  488. it.effect("remove clears sync event sequence", () =>
  489. Effect.gen(function* () {
  490. const events = yield* EventV2.Service
  491. const received = new Array<EventV2.Payload>()
  492. const aggregateID = EventV2.ID.create()
  493. yield* events.publish(SyncMessage, { id: aggregateID, text: "seed" })
  494. yield* events.remove(aggregateID)
  495. yield* events.project(SyncMessage, (event) =>
  496. Effect.sync(() => {
  497. received.push(event)
  498. }),
  499. )
  500. yield* events.replay({
  501. id: EventV2.ID.create(),
  502. type: EventV2.versionedType(SyncMessage.type, 1),
  503. seq: 0,
  504. aggregateID,
  505. data: { id: aggregateID, text: "replayed" },
  506. })
  507. expect(received[0]?.data).toEqual({ id: aggregateID, text: "replayed" })
  508. }),
  509. )
  510. })