event.ts 25 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663
  1. export * as EventV2 from "./event"
  2. import { Cause, Context, Effect, Layer, Option, PubSub, Schema, Stream } from "effect"
  3. import { and, asc, eq, gt } from "drizzle-orm"
  4. import { Database } from "./database/database"
  5. import { EventSequenceTable, EventTable } from "./event/sql"
  6. import { Location } from "./location"
  7. import { externalID, type ExternalID, NonNegativeInt, withStatics } from "./schema"
  8. import { Identifier } from "./util/identifier"
  9. import { isDeepStrictEqual } from "node:util"
  10. export const ID = Schema.String.check(Schema.isStartsWith("evt_")).pipe(
  11. Schema.brand("Event.ID"),
  12. withStatics((schema) => ({
  13. create: () => schema.make("evt_" + Identifier.ascending()),
  14. fromExternal: (input: ExternalID) => schema.make(externalID("evt", input)),
  15. })),
  16. )
  17. export type ID = typeof ID.Type
  18. /**
  19. * Durable aggregate continuation position for embedded replay streams.
  20. * TODO: Decide whether a future HTTP / SDK surface should expose an opaque cursor instead.
  21. */
  22. export const Cursor = NonNegativeInt.pipe(Schema.brand("EventV2.Cursor"))
  23. export type Cursor = typeof Cursor.Type
  24. export type Definition<Type extends string = string, DataSchema extends Schema.Top = Schema.Top> = {
  25. readonly type: Type
  26. readonly sync?: {
  27. readonly version: number
  28. readonly aggregate: string
  29. }
  30. readonly data: DataSchema
  31. }
  32. export type Data<D extends Definition> = Schema.Schema.Type<D["data"]>
  33. export type Payload<D extends Definition = Definition> = {
  34. readonly id: ID
  35. readonly type: D["type"]
  36. readonly data: Data<D>
  37. /** Durable aggregate order, populated while synchronized events are projected. */
  38. readonly seq?: number
  39. readonly version?: number
  40. readonly location?: Location.Ref
  41. readonly metadata?: Record<string, unknown>
  42. }
  43. export type Projector<D extends Definition = Definition> = (event: Payload<D>) => Effect.Effect<void>
  44. type AnyProjector = (event: Payload) => Effect.Effect<void>
  45. export type CommitGuard = (event: Payload) => Effect.Effect<void>
  46. export type Listener = (event: Payload) => Effect.Effect<void>
  47. export type Sync = (event: Payload) => Effect.Effect<void>
  48. export type Unsubscribe = Effect.Effect<void>
  49. export type SerializedEvent = {
  50. readonly id: ID
  51. readonly type: string
  52. readonly seq: number
  53. readonly aggregateID: string
  54. readonly data: Record<string, unknown>
  55. }
  56. export type CursorEvent<E extends Payload = Payload> = {
  57. readonly cursor: Cursor
  58. readonly event: E
  59. }
  60. export class InvalidSyncEventError extends Schema.TaggedErrorClass<InvalidSyncEventError>()(
  61. "EventV2.InvalidSyncEvent",
  62. {
  63. type: Schema.String,
  64. message: Schema.String,
  65. },
  66. ) {}
  67. export function versionedType(type: string, version: number) {
  68. return `${type}.${version}`
  69. }
  70. export const registry = new Map<string, Definition>()
  71. type SyncDefinition = Definition & {
  72. readonly sync: NonNullable<Definition["sync"]>
  73. readonly encode: (data: unknown) => unknown
  74. readonly decode: (data: unknown) => unknown
  75. }
  76. const syncRegistry = new Map<string, SyncDefinition>()
  77. // Synchronized events cross a JSON boundary, so their data schemas must encode and decode without services.
  78. const syncCodec = (definition: Definition) => definition.data as Schema.Codec<unknown, unknown, never, never>
  79. export function define<const Type extends string, Fields extends Schema.Struct.Fields>(input: {
  80. readonly type: Type
  81. readonly sync?: {
  82. readonly version: number
  83. readonly aggregate: string
  84. }
  85. readonly schema: Fields
  86. }): Schema.Schema<Payload<Definition<Type, Schema.Struct<Fields>>>> & Definition<Type, Schema.Struct<Fields>> {
  87. const Data = Schema.Struct(input.schema)
  88. const Payload = Schema.Struct({
  89. id: ID,
  90. metadata: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)),
  91. type: Schema.Literal(input.type),
  92. version: Schema.optional(Schema.Number),
  93. location: Schema.optional(Location.Ref),
  94. data: Data,
  95. }).annotate({ identifier: input.type })
  96. const definition = Object.assign(Payload, {
  97. type: input.type,
  98. ...(input.sync === undefined ? {} : { sync: input.sync }),
  99. data: Data,
  100. })
  101. const existing = registry.get(input.type)
  102. if (input.sync === undefined || existing?.sync === undefined || input.sync.version >= existing.sync.version) {
  103. registry.set(input.type, definition)
  104. }
  105. if (input.sync)
  106. syncRegistry.set(
  107. versionedType(input.type, input.sync.version),
  108. Object.assign(definition, {
  109. encode: Schema.encodeUnknownSync(syncCodec(definition)),
  110. decode: Schema.decodeUnknownSync(syncCodec(definition)),
  111. }) as SyncDefinition,
  112. )
  113. return definition as Schema.Schema<Payload<Definition<Type, Schema.Struct<Fields>>>> &
  114. Definition<Type, Schema.Struct<Fields>>
  115. }
  116. export function definitions() {
  117. return registry.values().toArray()
  118. }
  119. export interface PublishOptions {
  120. readonly id?: ID
  121. readonly metadata?: Record<string, unknown>
  122. readonly location?: Location.Ref
  123. }
  124. export interface Interface {
  125. readonly publish: <D extends Definition>(
  126. definition: D,
  127. data: Data<D>,
  128. options?: PublishOptions,
  129. ) => Effect.Effect<Payload<D>>
  130. readonly subscribe: <D extends Definition>(definition: D) => Stream.Stream<Payload<D>>
  131. readonly all: () => Stream.Stream<Payload>
  132. readonly aggregateEvents: (input: {
  133. readonly aggregateID: string
  134. readonly after?: Cursor
  135. }) => Stream.Stream<CursorEvent>
  136. readonly sync: (handler: Sync) => Effect.Effect<Unsubscribe>
  137. readonly listen: (listener: Listener) => Effect.Effect<Unsubscribe>
  138. readonly beforeCommit: (guard: CommitGuard) => Effect.Effect<void>
  139. readonly project: <D extends Definition>(definition: D, projector: Projector<D>) => Effect.Effect<void>
  140. readonly replay: (
  141. event: SerializedEvent,
  142. options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
  143. ) => Effect.Effect<void>
  144. readonly replayAll: (
  145. events: SerializedEvent[],
  146. options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
  147. ) => Effect.Effect<string | undefined>
  148. readonly remove: (aggregateID: string) => Effect.Effect<void>
  149. readonly claim: (aggregateID: string, ownerID: string) => Effect.Effect<void>
  150. }
  151. export class Service extends Context.Service<Service, Interface>()("@opencode/Event") {}
  152. export interface LayerOptions {
  153. readonly beforeAggregateRead?: (aggregateID: string) => Effect.Effect<void>
  154. }
  155. export const layerWith = (options?: LayerOptions) =>
  156. Layer.effect(
  157. Service,
  158. Effect.gen(function* () {
  159. const all = yield* PubSub.unbounded<Payload>()
  160. const synchronized = new Map<string, Set<PubSub.PubSub<void>>>()
  161. const typed = new Map<string, PubSub.PubSub<Payload>>()
  162. const projectors = new Map<string, AnyProjector[]>()
  163. const commitGuards = new Array<CommitGuard>()
  164. const listeners = new Array<Listener>()
  165. const syncHandlers = new Array<Sync>()
  166. const { db } = yield* Database.Service
  167. const getOrCreate = (definition: Definition) =>
  168. Effect.gen(function* () {
  169. const existing = typed.get(definition.type)
  170. if (existing) return existing
  171. const pubsub = yield* PubSub.unbounded<Payload>()
  172. typed.set(definition.type, pubsub)
  173. return pubsub
  174. })
  175. yield* Effect.addFinalizer(() =>
  176. Effect.gen(function* () {
  177. yield* PubSub.shutdown(all)
  178. yield* Effect.forEach(
  179. synchronized.values(),
  180. (pubsubs) => Effect.forEach(pubsubs, PubSub.shutdown, { discard: true }),
  181. { discard: true },
  182. )
  183. yield* Effect.forEach(typed.values(), PubSub.shutdown, { discard: true })
  184. }),
  185. )
  186. function commitSyncEvent(
  187. event: Payload,
  188. input?: {
  189. readonly seq: number
  190. readonly aggregateID: string
  191. readonly ownerID?: string
  192. readonly strictOwner?: boolean
  193. },
  194. ) {
  195. return Effect.gen(function* () {
  196. const definition = registry.get(event.type)
  197. const sync = definition?.sync
  198. if (sync) {
  199. if (event.version !== sync.version) {
  200. yield* Effect.die(
  201. new InvalidSyncEventError({
  202. type: event.type,
  203. message: `Expected event version ${sync.version}, got ${event.version}`,
  204. }),
  205. )
  206. }
  207. const aggregateID = (event.data as Record<string, unknown>)[sync.aggregate]
  208. if (typeof aggregateID !== "string") {
  209. yield* Effect.die(
  210. new InvalidSyncEventError({
  211. type: event.type,
  212. message: `Expected string aggregate field ${sync.aggregate}`,
  213. }),
  214. )
  215. } else {
  216. if (input && input.aggregateID !== aggregateID) {
  217. yield* Effect.die(
  218. new InvalidSyncEventError({
  219. type: event.type,
  220. message: `Aggregate mismatch: expected ${input.aggregateID}, got ${aggregateID}`,
  221. }),
  222. )
  223. }
  224. const list = projectors.get(event.type) ?? []
  225. return yield* Effect.uninterruptible(
  226. Effect.gen(function* () {
  227. const committed = yield* db
  228. .transaction(
  229. () =>
  230. Effect.gen(function* () {
  231. const row = yield* db
  232. .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
  233. .from(EventSequenceTable)
  234. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  235. .get()
  236. .pipe(Effect.orDie)
  237. const latest = row?.seq ?? -1
  238. const encoded = syncRegistry
  239. .get(versionedType(definition.type, sync.version))!
  240. .encode(event.data) as Record<string, unknown>
  241. if (input?.strictOwner && row?.ownerID && row.ownerID !== input.ownerID) {
  242. yield* Effect.die(
  243. new InvalidSyncEventError({
  244. type: event.type,
  245. message: `Replay owner mismatch for aggregate ${aggregateID}: expected ${row.ownerID}, got ${input.ownerID ?? "none"}`,
  246. }),
  247. )
  248. }
  249. if (input && input.seq <= latest) {
  250. const stored = yield* db
  251. .select()
  252. .from(EventTable)
  253. .where(and(eq(EventTable.aggregate_id, aggregateID), eq(EventTable.seq, input.seq)))
  254. .get()
  255. .pipe(Effect.orDie)
  256. if (
  257. stored?.id === event.id &&
  258. stored.type === versionedType(definition.type, sync.version) &&
  259. isDeepStrictEqual(stored.data, encoded)
  260. ) {
  261. if (input.ownerID && row?.ownerID == null) {
  262. yield* db
  263. .update(EventSequenceTable)
  264. .set({ owner_id: input.ownerID })
  265. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  266. .run()
  267. .pipe(Effect.orDie)
  268. }
  269. return
  270. }
  271. yield* Effect.die(
  272. new InvalidSyncEventError({
  273. type: event.type,
  274. message: `Replay diverged at aggregate ${aggregateID} sequence ${input.seq}`,
  275. }),
  276. )
  277. }
  278. if (input && row?.ownerID && row.ownerID !== input.ownerID) {
  279. return
  280. }
  281. const seq = input?.seq ?? latest + 1
  282. if (input && seq !== latest + 1) {
  283. yield* Effect.die(
  284. new InvalidSyncEventError({
  285. type: event.type,
  286. message: `Sequence mismatch for aggregate ${aggregateID}: expected ${latest + 1}, got ${seq}`,
  287. }),
  288. )
  289. }
  290. const stored = yield* db
  291. .select({ aggregateID: EventTable.aggregate_id, seq: EventTable.seq })
  292. .from(EventTable)
  293. .where(eq(EventTable.id, event.id))
  294. .get()
  295. .pipe(Effect.orDie)
  296. if (stored)
  297. yield* Effect.die(
  298. new InvalidSyncEventError({
  299. type: event.type,
  300. message: `Event ${event.id} already exists at aggregate ${stored.aggregateID} sequence ${stored.seq}`,
  301. }),
  302. )
  303. for (const guard of commitGuards) {
  304. yield* guard(event)
  305. }
  306. for (const projector of list) {
  307. yield* projector({ ...event, seq } as Payload)
  308. }
  309. yield* db
  310. .insert(EventSequenceTable)
  311. .values([{ aggregate_id: aggregateID, seq, owner_id: input?.ownerID }])
  312. .onConflictDoUpdate({
  313. target: EventSequenceTable.aggregate_id,
  314. set: {
  315. seq,
  316. ...(input?.ownerID && row?.ownerID == null ? { owner_id: input.ownerID } : {}),
  317. },
  318. })
  319. .run()
  320. .pipe(Effect.orDie)
  321. yield* db
  322. .insert(EventTable)
  323. .values([
  324. {
  325. id: event.id,
  326. aggregate_id: aggregateID,
  327. seq,
  328. type: versionedType(definition.type, sync.version),
  329. data: encoded,
  330. },
  331. ])
  332. .run()
  333. .pipe(Effect.orDie)
  334. return { aggregateID, seq }
  335. }),
  336. { behavior: "immediate" },
  337. )
  338. .pipe(Effect.orDie)
  339. if (committed) {
  340. yield* Effect.forEach(
  341. synchronized.get(committed.aggregateID) ?? [],
  342. (pubsub) => PubSub.publish(pubsub, undefined),
  343. { discard: true },
  344. )
  345. }
  346. return committed
  347. }),
  348. )
  349. }
  350. }
  351. })
  352. }
  353. function publishEvent<D extends Definition>(event: Payload<D>) {
  354. return Effect.gen(function* () {
  355. const durable = registry.get(event.type)?.sync !== undefined
  356. if (durable) {
  357. const committed = yield* commitSyncEvent(event as Payload)
  358. if (committed) {
  359. event = { ...event, seq: committed.seq }
  360. yield* Effect.forEach(syncHandlers, (sync) => observe(event as Payload, "sync", sync), { discard: true })
  361. yield* notify(event as Payload, true)
  362. return event
  363. }
  364. }
  365. yield* notify(event as Payload, false)
  366. return event
  367. })
  368. }
  369. const observe = (event: Payload, kind: "sync" | "listener", observer: (event: Payload) => Effect.Effect<void>) =>
  370. Effect.suspend(() => observer(event)).pipe(
  371. Effect.catchCauseIf(
  372. (cause) => !Cause.hasInterrupts(cause),
  373. (cause) =>
  374. Effect.logError("Event observer failed").pipe(
  375. Effect.annotateLogs({ eventID: event.id, eventType: event.type, kind, cause }),
  376. ),
  377. ),
  378. )
  379. function notify(event: Payload, isolateListeners: boolean) {
  380. return Effect.gen(function* () {
  381. yield* Effect.forEach(
  382. listeners,
  383. (listener) => (isolateListeners ? observe(event, "listener", listener) : listener(event)),
  384. { discard: true },
  385. )
  386. const pubsub = typed.get(event.type)
  387. if (pubsub) yield* PubSub.publish(pubsub, event)
  388. yield* PubSub.publish(all, event)
  389. })
  390. }
  391. function publish<D extends Definition>(definition: D, data: Data<D>, options?: PublishOptions) {
  392. return Effect.gen(function* () {
  393. const serviceLocation = Option.getOrUndefined(yield* Effect.serviceOption(Location.Service))
  394. const location =
  395. options?.location ??
  396. (serviceLocation
  397. ? { directory: serviceLocation.directory, workspaceID: serviceLocation.workspaceID }
  398. : undefined)
  399. return yield* publishEvent({
  400. id: options?.id ?? ID.create(),
  401. ...(options?.metadata ? { metadata: options.metadata } : {}),
  402. type: definition.type,
  403. ...(definition.sync === undefined ? {} : { version: definition.sync.version }),
  404. ...(location ? { location } : {}),
  405. data,
  406. } as Payload<D>)
  407. })
  408. }
  409. function replay(
  410. event: SerializedEvent,
  411. options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
  412. ) {
  413. return Effect.gen(function* () {
  414. const definition = syncRegistry.get(event.type)
  415. if (!definition) {
  416. yield* Effect.die(
  417. new InvalidSyncEventError({ type: event.type, message: `Unknown sync event type ${event.type}` }),
  418. )
  419. } else {
  420. const payload = {
  421. id: event.id,
  422. type: definition.type,
  423. version: definition.sync.version,
  424. data: definition.decode(event.data),
  425. } as Payload
  426. const committed = yield* commitSyncEvent(payload, {
  427. seq: event.seq,
  428. aggregateID: event.aggregateID,
  429. ownerID: options?.ownerID,
  430. strictOwner: options?.strictOwner,
  431. })
  432. if (committed && options?.publish) {
  433. yield* notify({ ...payload, seq: committed.seq }, true)
  434. }
  435. }
  436. })
  437. }
  438. function replayAll(
  439. events: SerializedEvent[],
  440. options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
  441. ) {
  442. return Effect.gen(function* () {
  443. const source = events[0]?.aggregateID
  444. if (!source) return undefined
  445. if (events.some((event) => event.aggregateID !== source)) {
  446. yield* Effect.die(
  447. new InvalidSyncEventError({
  448. type: events[0]?.type ?? "unknown",
  449. message: "Replay events must belong to the same aggregate",
  450. }),
  451. )
  452. }
  453. const start = events[0]?.seq ?? 0
  454. for (const [index, event] of events.entries()) {
  455. const seq = start + index
  456. if (event.seq !== seq) {
  457. yield* Effect.die(
  458. new InvalidSyncEventError({
  459. type: event.type,
  460. message: `Replay sequence mismatch at index ${index}: expected ${seq}, got ${event.seq}`,
  461. }),
  462. )
  463. }
  464. }
  465. for (const event of events) {
  466. yield* replay(event, options)
  467. }
  468. return source
  469. })
  470. }
  471. function remove(aggregateID: string) {
  472. return db
  473. .transaction(() =>
  474. Effect.gen(function* () {
  475. yield* db.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).run()
  476. yield* db.delete(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).run()
  477. }),
  478. )
  479. .pipe(Effect.orDie)
  480. }
  481. function claim(aggregateID: string, ownerID: string) {
  482. return db
  483. .update(EventSequenceTable)
  484. .set({ owner_id: ownerID })
  485. .where(eq(EventSequenceTable.aggregate_id, aggregateID))
  486. .run()
  487. .pipe(Effect.orDie)
  488. }
  489. const subscribe = <D extends Definition>(definition: D): Stream.Stream<Payload<D>> =>
  490. Stream.unwrap(getOrCreate(definition).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub)))).pipe(
  491. Stream.map((event) => event as Payload<D>),
  492. )
  493. const streamAll = (): Stream.Stream<Payload> => Stream.fromPubSub(all)
  494. const decodeSerializedEvent = (event: SerializedEvent): CursorEvent => {
  495. const definition = syncRegistry.get(event.type)
  496. if (!definition) {
  497. throw new InvalidSyncEventError({ type: event.type, message: `Unknown sync event type ${event.type}` })
  498. }
  499. return {
  500. cursor: Cursor.make(event.seq),
  501. event: {
  502. id: event.id,
  503. type: definition.type,
  504. version: definition.sync.version,
  505. seq: event.seq,
  506. data: definition.decode(event.data),
  507. },
  508. }
  509. }
  510. const readAfter = (aggregateID: string, after: number) =>
  511. (options?.beforeAggregateRead?.(aggregateID) ?? Effect.void).pipe(
  512. Effect.andThen(
  513. db
  514. .select()
  515. .from(EventTable)
  516. .where(and(eq(EventTable.aggregate_id, aggregateID), gt(EventTable.seq, after)))
  517. .orderBy(asc(EventTable.seq))
  518. .all(),
  519. ),
  520. Effect.orDie,
  521. Effect.map((rows) =>
  522. rows.map((event) =>
  523. decodeSerializedEvent({
  524. id: event.id,
  525. aggregateID: event.aggregate_id,
  526. seq: event.seq,
  527. type: event.type,
  528. data: event.data,
  529. }),
  530. ),
  531. ),
  532. )
  533. const subscribeSynchronized = (aggregateID: string) =>
  534. Effect.gen(function* () {
  535. const pubsub = yield* PubSub.sliding<void>(1)
  536. const subscription = yield* PubSub.subscribe(pubsub)
  537. yield* Effect.acquireRelease(
  538. Effect.sync(() => {
  539. const pubsubs = synchronized.get(aggregateID) ?? new Set()
  540. pubsubs.add(pubsub)
  541. synchronized.set(aggregateID, pubsubs)
  542. }),
  543. () =>
  544. Effect.sync(() => {
  545. const pubsubs = synchronized.get(aggregateID)
  546. pubsubs?.delete(pubsub)
  547. if (pubsubs?.size === 0) synchronized.delete(aggregateID)
  548. }).pipe(Effect.andThen(PubSub.shutdown(pubsub))),
  549. )
  550. return subscription
  551. })
  552. const streamEvents = (input: {
  553. readonly aggregateID: string
  554. readonly after?: Cursor
  555. }): Stream.Stream<CursorEvent> =>
  556. Stream.unwrap(
  557. Effect.gen(function* () {
  558. const synchronized = yield* subscribeSynchronized(input.aggregateID)
  559. let cursor = input.after ?? -1
  560. const read = Effect.suspend(() => readAfter(input.aggregateID, cursor)).pipe(
  561. Effect.tap((events) =>
  562. Effect.sync(() => {
  563. cursor = events.at(-1)?.cursor ?? cursor
  564. }),
  565. ),
  566. )
  567. const historical = yield* read
  568. const live = Stream.fromSubscription(synchronized).pipe(
  569. Stream.mapEffect(() => read),
  570. Stream.flattenIterable,
  571. )
  572. return Stream.concat(Stream.fromIterable(historical), live)
  573. }),
  574. )
  575. const listen = (listener: Listener): Effect.Effect<Unsubscribe> =>
  576. Effect.sync(() => {
  577. listeners.push(listener)
  578. return Effect.sync(() => {
  579. const index = listeners.indexOf(listener)
  580. if (index >= 0) listeners.splice(index, 1)
  581. })
  582. })
  583. const sync = (handler: Sync): Effect.Effect<Unsubscribe> =>
  584. Effect.sync(() => {
  585. syncHandlers.push(handler)
  586. return Effect.sync(() => {
  587. const index = syncHandlers.indexOf(handler)
  588. if (index >= 0) syncHandlers.splice(index, 1)
  589. })
  590. })
  591. const beforeCommit = (guard: CommitGuard): Effect.Effect<void> =>
  592. Effect.sync(() => {
  593. commitGuards.push(guard)
  594. })
  595. const project = <D extends Definition>(definition: D, projector: Projector<D>): Effect.Effect<void> =>
  596. Effect.sync(() => {
  597. const list = projectors.get(definition.type) ?? []
  598. list.push((event) => projector(event as Payload<D>))
  599. projectors.set(definition.type, list)
  600. })
  601. return Service.of({
  602. publish,
  603. subscribe,
  604. all: streamAll,
  605. aggregateEvents: streamEvents,
  606. sync,
  607. listen,
  608. beforeCommit,
  609. project,
  610. replay,
  611. replayAll,
  612. remove,
  613. claim,
  614. })
  615. }),
  616. )
  617. export const layer = layerWith()
  618. export const defaultLayer = layer.pipe(Layer.provide(Database.defaultLayer))