index.test.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350
  1. import { describe, expect, beforeEach, afterEach, afterAll } from "bun:test"
  2. import { provideTmpdirInstance } from "../fixture/fixture"
  3. import { Effect, Layer, Schema } from "effect"
  4. import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
  5. import { Bus } from "../../src/bus"
  6. import { SyncEvent } from "../../src/sync"
  7. import { Database, eq } from "@/storage/db"
  8. import { EventSequenceTable, EventTable } from "../../src/sync/event.sql"
  9. import { MessageID } from "../../src/session/schema"
  10. import { Flag } from "@opencode-ai/core/flag/flag"
  11. import { initProjectors } from "../../src/server/projectors"
  12. import { testEffect } from "../lib/effect"
  13. const original = Flag.OPENCODE_EXPERIMENTAL_WORKSPACES
  14. const it = testEffect(Layer.mergeAll(SyncEvent.defaultLayer, CrossSpawnSpawner.defaultLayer))
  15. beforeEach(() => {
  16. Database.close()
  17. Flag.OPENCODE_EXPERIMENTAL_WORKSPACES = true
  18. })
  19. afterEach(() => {
  20. Flag.OPENCODE_EXPERIMENTAL_WORKSPACES = original
  21. })
  22. describe("SyncEvent", () => {
  23. function setup() {
  24. SyncEvent.reset()
  25. const Created = SyncEvent.define({
  26. type: "item.created",
  27. version: 1,
  28. aggregate: "id",
  29. schema: Schema.Struct({ id: Schema.String, name: Schema.String }),
  30. })
  31. const Sent = SyncEvent.define({
  32. type: "item.sent",
  33. version: 1,
  34. aggregate: "item_id",
  35. schema: Schema.Struct({ item_id: Schema.String, to: Schema.String }),
  36. })
  37. SyncEvent.init({
  38. projectors: [SyncEvent.project(Created, () => {}), SyncEvent.project(Sent, () => {})],
  39. })
  40. return { Created, Sent }
  41. }
  42. function expectDefect<A, E, R>(effect: Effect.Effect<A, E, R>, pattern: RegExp) {
  43. return Effect.gen(function* () {
  44. const exit = yield* Effect.exit(effect)
  45. if (exit._tag === "Success") throw new Error("Expected effect to fail")
  46. expect(String(exit.cause)).toMatch(pattern)
  47. })
  48. }
  49. afterAll(() => {
  50. SyncEvent.reset()
  51. initProjectors()
  52. })
  53. describe("run", () => {
  54. it.live(
  55. "inserts event row",
  56. provideTmpdirInstance(() =>
  57. Effect.gen(function* () {
  58. const { Created } = setup()
  59. yield* SyncEvent.use.run(Created, { id: "evt_1", name: "first" })
  60. const rows = Database.use((db) => db.select().from(EventTable).all())
  61. expect(rows).toHaveLength(1)
  62. expect(rows[0].type).toBe("item.created.1")
  63. expect(rows[0].aggregate_id).toBe("evt_1")
  64. }),
  65. ),
  66. )
  67. it.live(
  68. "increments seq per aggregate",
  69. provideTmpdirInstance(() =>
  70. Effect.gen(function* () {
  71. const { Created } = setup()
  72. yield* SyncEvent.use.run(Created, { id: "evt_1", name: "first" })
  73. yield* SyncEvent.use.run(Created, { id: "evt_1", name: "second" })
  74. const rows = Database.use((db) => db.select().from(EventTable).all())
  75. expect(rows).toHaveLength(2)
  76. expect(rows[1].seq).toBe(rows[0].seq + 1)
  77. }),
  78. ),
  79. )
  80. it.live(
  81. "uses custom aggregate field from agg()",
  82. provideTmpdirInstance(() =>
  83. Effect.gen(function* () {
  84. const { Sent } = setup()
  85. yield* SyncEvent.use.run(Sent, { item_id: "evt_1", to: "james" })
  86. const rows = Database.use((db) => db.select().from(EventTable).all())
  87. expect(rows).toHaveLength(1)
  88. expect(rows[0].aggregate_id).toBe("evt_1")
  89. }),
  90. ),
  91. )
  92. it.live(
  93. "emits events",
  94. provideTmpdirInstance(() =>
  95. Effect.gen(function* () {
  96. const { Created } = setup()
  97. const events: Array<{
  98. type: string
  99. properties: { id: string; name: string }
  100. }> = []
  101. let resolve = () => {}
  102. const received = new Promise<void>((done) => {
  103. resolve = done
  104. })
  105. const dispose = Bus.subscribeAll((event) => {
  106. events.push(event)
  107. resolve()
  108. })
  109. try {
  110. yield* SyncEvent.use.run(Created, { id: "evt_1", name: "test" })
  111. yield* Effect.promise(() => received)
  112. expect(events).toHaveLength(1)
  113. expect(events[0]).toMatchObject({
  114. type: "item.created",
  115. properties: {
  116. id: "evt_1",
  117. name: "test",
  118. },
  119. })
  120. } finally {
  121. dispose()
  122. }
  123. }),
  124. ),
  125. )
  126. })
  127. describe("replay", () => {
  128. it.live(
  129. "inserts event from external payload",
  130. provideTmpdirInstance(() =>
  131. Effect.gen(function* () {
  132. const id = MessageID.ascending()
  133. yield* SyncEvent.use.replay({
  134. id: "evt_1",
  135. type: "item.created.1",
  136. seq: 0,
  137. aggregateID: id,
  138. data: { id, name: "replayed" },
  139. })
  140. const rows = Database.use((db) => db.select().from(EventTable).all())
  141. expect(rows).toHaveLength(1)
  142. expect(rows[0].aggregate_id).toBe(id)
  143. }),
  144. ),
  145. )
  146. it.live(
  147. "throws on sequence mismatch",
  148. provideTmpdirInstance(() =>
  149. Effect.gen(function* () {
  150. const id = MessageID.ascending()
  151. yield* SyncEvent.use.replay({
  152. id: "evt_1",
  153. type: "item.created.1",
  154. seq: 0,
  155. aggregateID: id,
  156. data: { id, name: "first" },
  157. })
  158. yield* expectDefect(
  159. SyncEvent.use.replay({
  160. id: "evt_1",
  161. type: "item.created.1",
  162. seq: 5,
  163. aggregateID: id,
  164. data: { id, name: "bad" },
  165. }),
  166. /Sequence mismatch/,
  167. )
  168. }),
  169. ),
  170. )
  171. it.live(
  172. "throws on unknown event type",
  173. provideTmpdirInstance(() =>
  174. Effect.gen(function* () {
  175. yield* expectDefect(
  176. SyncEvent.use.replay({
  177. id: "evt_1",
  178. type: "unknown.event.1",
  179. seq: 0,
  180. aggregateID: "x",
  181. data: {},
  182. }),
  183. /Unknown event type/,
  184. )
  185. }),
  186. ),
  187. )
  188. it.live(
  189. "replayAll accepts later chunks after the first batch",
  190. provideTmpdirInstance(() =>
  191. Effect.gen(function* () {
  192. const { Created } = setup()
  193. const id = MessageID.ascending()
  194. const one = yield* SyncEvent.use.replayAll([
  195. {
  196. id: "evt_1",
  197. type: SyncEvent.versionedType(Created.type, Created.version),
  198. seq: 0,
  199. aggregateID: id,
  200. data: { id, name: "first" },
  201. },
  202. {
  203. id: "evt_2",
  204. type: SyncEvent.versionedType(Created.type, Created.version),
  205. seq: 1,
  206. aggregateID: id,
  207. data: { id, name: "second" },
  208. },
  209. ])
  210. const two = yield* SyncEvent.use.replayAll([
  211. {
  212. id: "evt_3",
  213. type: SyncEvent.versionedType(Created.type, Created.version),
  214. seq: 2,
  215. aggregateID: id,
  216. data: { id, name: "third" },
  217. },
  218. {
  219. id: "evt_4",
  220. type: SyncEvent.versionedType(Created.type, Created.version),
  221. seq: 3,
  222. aggregateID: id,
  223. data: { id, name: "fourth" },
  224. },
  225. ])
  226. expect(one).toBe(id)
  227. expect(two).toBe(id)
  228. const rows = Database.use((db) => db.select().from(EventTable).all())
  229. expect(rows.map((row) => row.seq)).toEqual([0, 1, 2, 3])
  230. }),
  231. ),
  232. )
  233. it.live(
  234. "claims unowned event sequence on replay with ownerID",
  235. provideTmpdirInstance(() =>
  236. Effect.gen(function* () {
  237. const { Created } = setup()
  238. const id = MessageID.ascending()
  239. yield* SyncEvent.use.replay(
  240. {
  241. id: "evt_1",
  242. type: SyncEvent.versionedType(Created.type, Created.version),
  243. seq: 0,
  244. aggregateID: id,
  245. data: { id, name: "owned" },
  246. },
  247. { publish: false, ownerID: "owner-1" },
  248. )
  249. const row = Database.use((db) =>
  250. db
  251. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  252. .from(EventSequenceTable)
  253. .get(),
  254. )
  255. expect(row).toEqual({ seq: 0, ownerID: "owner-1" })
  256. }),
  257. ),
  258. )
  259. it.live(
  260. "ignores replay from a different owner after sequence is claimed",
  261. provideTmpdirInstance(() =>
  262. Effect.gen(function* () {
  263. const { Created } = setup()
  264. const id = MessageID.ascending()
  265. yield* SyncEvent.use.replay(
  266. {
  267. id: "evt_1",
  268. type: SyncEvent.versionedType(Created.type, Created.version),
  269. seq: 0,
  270. aggregateID: id,
  271. data: { id, name: "first" },
  272. },
  273. { publish: false, ownerID: "owner-1" },
  274. )
  275. yield* SyncEvent.use.replay(
  276. {
  277. id: "evt_2",
  278. type: SyncEvent.versionedType(Created.type, Created.version),
  279. seq: 1,
  280. aggregateID: id,
  281. data: { id, name: "ignored" },
  282. },
  283. { publish: false, ownerID: "owner-2" },
  284. )
  285. const events = Database.use((db) => db.select().from(EventTable).all())
  286. const sequence = Database.use((db) =>
  287. db
  288. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  289. .from(EventSequenceTable)
  290. .get(),
  291. )
  292. expect(events).toHaveLength(1)
  293. expect(events[0].id).toBe("evt_1")
  294. expect(sequence).toEqual({ seq: 0, ownerID: "owner-1" })
  295. }),
  296. ),
  297. )
  298. it.live(
  299. "claim updates the event sequence owner",
  300. provideTmpdirInstance(() =>
  301. Effect.gen(function* () {
  302. const { Created } = setup()
  303. const id = MessageID.ascending()
  304. yield* SyncEvent.use.run(Created, { id, name: "claimed" }, { publish: false })
  305. yield* SyncEvent.use.claim(id, "owner-1")
  306. yield* SyncEvent.use.claim(id, "owner-2")
  307. const row = Database.use((db) =>
  308. db
  309. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  310. .from(EventSequenceTable)
  311. .where(eq(EventSequenceTable.aggregate_id, id))
  312. .get(),
  313. )
  314. expect(row).toEqual({ seq: 0, ownerID: "owner-2" })
  315. }),
  316. ),
  317. )
  318. })
  319. })