|
|
@@ -146,7 +146,10 @@ export interface Interface {
|
|
|
) => Effect.Effect<Payload<D>>
|
|
|
readonly subscribe: <D extends Definition>(definition: D) => Stream.Stream<Payload<D>>
|
|
|
readonly all: () => Stream.Stream<Payload>
|
|
|
- readonly aggregateEvents: (input: { readonly aggregateID: string; readonly after?: Cursor }) => Stream.Stream<CursorEvent>
|
|
|
+ readonly aggregateEvents: (input: {
|
|
|
+ readonly aggregateID: string
|
|
|
+ readonly after?: Cursor
|
|
|
+ }) => Stream.Stream<CursorEvent>
|
|
|
readonly sync: (handler: Sync) => Effect.Effect<Unsubscribe>
|
|
|
readonly listen: (listener: Listener) => Effect.Effect<Unsubscribe>
|
|
|
readonly beforeCommit: (guard: CommitGuard) => Effect.Effect<void>
|
|
|
@@ -169,394 +172,434 @@ export interface LayerOptions {
|
|
|
readonly beforeAggregateRead?: (aggregateID: string) => Effect.Effect<void>
|
|
|
}
|
|
|
|
|
|
-export const layerWith = (options?: LayerOptions) => Layer.effect(
|
|
|
- Service,
|
|
|
- Effect.gen(function* () {
|
|
|
- const all = yield* PubSub.unbounded<Payload>()
|
|
|
- const synchronized = new Map<string, Set<PubSub.PubSub<void>>>()
|
|
|
- const typed = new Map<string, PubSub.PubSub<Payload>>()
|
|
|
- const projectors = new Map<string, AnyProjector[]>()
|
|
|
- const commitGuards = new Array<CommitGuard>()
|
|
|
- const listeners = new Array<Listener>()
|
|
|
- const syncHandlers = new Array<Sync>()
|
|
|
- const { db } = yield* Database.Service
|
|
|
-
|
|
|
- const getOrCreate = (definition: Definition) =>
|
|
|
- Effect.gen(function* () {
|
|
|
- const existing = typed.get(definition.type)
|
|
|
- if (existing) return existing
|
|
|
- const pubsub = yield* PubSub.unbounded<Payload>()
|
|
|
- typed.set(definition.type, pubsub)
|
|
|
- return pubsub
|
|
|
- })
|
|
|
+export const layerWith = (options?: LayerOptions) =>
|
|
|
+ Layer.effect(
|
|
|
+ Service,
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const all = yield* PubSub.unbounded<Payload>()
|
|
|
+ const synchronized = new Map<string, Set<PubSub.PubSub<void>>>()
|
|
|
+ const typed = new Map<string, PubSub.PubSub<Payload>>()
|
|
|
+ const projectors = new Map<string, AnyProjector[]>()
|
|
|
+ const commitGuards = new Array<CommitGuard>()
|
|
|
+ const listeners = new Array<Listener>()
|
|
|
+ const syncHandlers = new Array<Sync>()
|
|
|
+ const { db } = yield* Database.Service
|
|
|
+
|
|
|
+ const getOrCreate = (definition: Definition) =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const existing = typed.get(definition.type)
|
|
|
+ if (existing) return existing
|
|
|
+ const pubsub = yield* PubSub.unbounded<Payload>()
|
|
|
+ typed.set(definition.type, pubsub)
|
|
|
+ return pubsub
|
|
|
+ })
|
|
|
|
|
|
- yield* Effect.addFinalizer(() =>
|
|
|
- Effect.gen(function* () {
|
|
|
- yield* PubSub.shutdown(all)
|
|
|
- yield* Effect.forEach(synchronized.values(), (pubsubs) =>
|
|
|
- Effect.forEach(pubsubs, PubSub.shutdown, { discard: true }),
|
|
|
- { discard: true })
|
|
|
- yield* Effect.forEach(typed.values(), PubSub.shutdown, { discard: true })
|
|
|
- }),
|
|
|
- )
|
|
|
+ yield* Effect.addFinalizer(() =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ yield* PubSub.shutdown(all)
|
|
|
+ yield* Effect.forEach(
|
|
|
+ synchronized.values(),
|
|
|
+ (pubsubs) => Effect.forEach(pubsubs, PubSub.shutdown, { discard: true }),
|
|
|
+ { discard: true },
|
|
|
+ )
|
|
|
+ yield* Effect.forEach(typed.values(), PubSub.shutdown, { discard: true })
|
|
|
+ }),
|
|
|
+ )
|
|
|
|
|
|
- function commitSyncEvent(
|
|
|
- event: Payload,
|
|
|
- input?: { readonly seq: number; readonly aggregateID: string; readonly ownerID?: string; readonly strictOwner?: boolean },
|
|
|
- ) {
|
|
|
- return Effect.gen(function* () {
|
|
|
- const definition = registry.get(event.type)
|
|
|
- const sync = definition?.sync
|
|
|
- if (sync) {
|
|
|
- if (event.version !== sync.version) {
|
|
|
- yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
- type: event.type,
|
|
|
- message: `Expected event version ${sync.version}, got ${event.version}`,
|
|
|
- }),
|
|
|
- )
|
|
|
- }
|
|
|
- const aggregateID = (event.data as Record<string, unknown>)[sync.aggregate]
|
|
|
- if (typeof aggregateID !== "string") {
|
|
|
- yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
- type: event.type,
|
|
|
- message: `Expected string aggregate field ${sync.aggregate}`,
|
|
|
- }),
|
|
|
- )
|
|
|
- } else {
|
|
|
- if (input && input.aggregateID !== aggregateID) {
|
|
|
+ function commitSyncEvent(
|
|
|
+ event: Payload,
|
|
|
+ input?: {
|
|
|
+ readonly seq: number
|
|
|
+ readonly aggregateID: string
|
|
|
+ readonly ownerID?: string
|
|
|
+ readonly strictOwner?: boolean
|
|
|
+ },
|
|
|
+ ) {
|
|
|
+ return Effect.gen(function* () {
|
|
|
+ const definition = registry.get(event.type)
|
|
|
+ const sync = definition?.sync
|
|
|
+ if (sync) {
|
|
|
+ if (event.version !== sync.version) {
|
|
|
yield* Effect.die(
|
|
|
new InvalidSyncEventError({
|
|
|
type: event.type,
|
|
|
- message: `Aggregate mismatch: expected ${input.aggregateID}, got ${aggregateID}`,
|
|
|
+ message: `Expected event version ${sync.version}, got ${event.version}`,
|
|
|
}),
|
|
|
)
|
|
|
}
|
|
|
- const list = projectors.get(event.type) ?? []
|
|
|
- return yield* Effect.uninterruptible(
|
|
|
- Effect.gen(function* () {
|
|
|
- const committed = yield* db
|
|
|
- .transaction(
|
|
|
- () =>
|
|
|
- Effect.gen(function* () {
|
|
|
- const row = yield* db
|
|
|
- .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
|
|
|
- .from(EventSequenceTable)
|
|
|
- .where(eq(EventSequenceTable.aggregate_id, aggregateID))
|
|
|
- .get()
|
|
|
- .pipe(Effect.orDie)
|
|
|
- const latest = row?.seq ?? -1
|
|
|
- if (input && input.seq <= latest) return
|
|
|
- if (input && row?.ownerID && row.ownerID !== input.ownerID) {
|
|
|
- if (input.strictOwner) {
|
|
|
+ const aggregateID = (event.data as Record<string, unknown>)[sync.aggregate]
|
|
|
+ if (typeof aggregateID !== "string") {
|
|
|
+ yield* Effect.die(
|
|
|
+ new InvalidSyncEventError({
|
|
|
+ type: event.type,
|
|
|
+ message: `Expected string aggregate field ${sync.aggregate}`,
|
|
|
+ }),
|
|
|
+ )
|
|
|
+ } else {
|
|
|
+ if (input && input.aggregateID !== aggregateID) {
|
|
|
+ yield* Effect.die(
|
|
|
+ new InvalidSyncEventError({
|
|
|
+ type: event.type,
|
|
|
+ message: `Aggregate mismatch: expected ${input.aggregateID}, got ${aggregateID}`,
|
|
|
+ }),
|
|
|
+ )
|
|
|
+ }
|
|
|
+ const list = projectors.get(event.type) ?? []
|
|
|
+ return yield* Effect.uninterruptible(
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const committed = yield* db
|
|
|
+ .transaction(
|
|
|
+ () =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const row = yield* db
|
|
|
+ .select({ seq: EventSequenceTable.seq, ownerID: EventSequenceTable.owner_id })
|
|
|
+ .from(EventSequenceTable)
|
|
|
+ .where(eq(EventSequenceTable.aggregate_id, aggregateID))
|
|
|
+ .get()
|
|
|
+ .pipe(Effect.orDie)
|
|
|
+ const latest = row?.seq ?? -1
|
|
|
+ if (input && input.seq <= latest) return
|
|
|
+ if (input && row?.ownerID && row.ownerID !== input.ownerID) {
|
|
|
+ if (input.strictOwner) {
|
|
|
+ yield* Effect.die(
|
|
|
+ new InvalidSyncEventError({
|
|
|
+ type: event.type,
|
|
|
+ message: `Replay owner mismatch for aggregate ${aggregateID}: expected ${row.ownerID}, got ${input.ownerID ?? "none"}`,
|
|
|
+ }),
|
|
|
+ )
|
|
|
+ }
|
|
|
+ return
|
|
|
+ }
|
|
|
+ const seq = input?.seq ?? latest + 1
|
|
|
+ if (input && seq !== latest + 1) {
|
|
|
yield* Effect.die(
|
|
|
new InvalidSyncEventError({
|
|
|
type: event.type,
|
|
|
- message: `Replay owner mismatch for aggregate ${aggregateID}: expected ${row.ownerID}, got ${input.ownerID ?? "none"}`,
|
|
|
+ message: `Sequence mismatch for aggregate ${aggregateID}: expected ${latest + 1}, got ${seq}`,
|
|
|
}),
|
|
|
)
|
|
|
}
|
|
|
- return
|
|
|
- }
|
|
|
- const seq = input?.seq ?? latest + 1
|
|
|
- if (input && seq !== latest + 1) {
|
|
|
- yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
- type: event.type,
|
|
|
- message: `Sequence mismatch for aggregate ${aggregateID}: expected ${latest + 1}, got ${seq}`,
|
|
|
- }),
|
|
|
- )
|
|
|
- }
|
|
|
- for (const guard of commitGuards) {
|
|
|
- yield* guard(event)
|
|
|
- }
|
|
|
- for (const projector of list) {
|
|
|
- yield* projector({ ...event, seq } as Payload)
|
|
|
- }
|
|
|
- const encoded = syncRegistry.get(versionedType(definition.type, sync.version))!.encode(event.data)
|
|
|
- yield* db
|
|
|
- .insert(EventSequenceTable)
|
|
|
- .values([{ aggregate_id: aggregateID, seq, owner_id: input?.ownerID }])
|
|
|
- .onConflictDoUpdate({
|
|
|
- target: EventSequenceTable.aggregate_id,
|
|
|
- set: { seq, ...(input?.ownerID && row?.ownerID == null ? { owner_id: input.ownerID } : {}) },
|
|
|
- })
|
|
|
- .run()
|
|
|
- .pipe(Effect.orDie)
|
|
|
- yield* db
|
|
|
- .insert(EventTable)
|
|
|
- .values([
|
|
|
- {
|
|
|
- id: event.id,
|
|
|
- aggregate_id: aggregateID,
|
|
|
- seq,
|
|
|
- type: versionedType(definition.type, sync.version),
|
|
|
- data: encoded as Record<string, unknown>,
|
|
|
- },
|
|
|
- ])
|
|
|
- .run()
|
|
|
- .pipe(Effect.orDie)
|
|
|
- return { aggregateID, seq }
|
|
|
- }),
|
|
|
- { behavior: "immediate" },
|
|
|
- )
|
|
|
- .pipe(Effect.orDie)
|
|
|
- if (committed) {
|
|
|
- yield* Effect.forEach(
|
|
|
- synchronized.get(committed.aggregateID) ?? [],
|
|
|
- (pubsub) => PubSub.publish(pubsub, undefined),
|
|
|
- { discard: true },
|
|
|
- )
|
|
|
- }
|
|
|
- return committed
|
|
|
- }),
|
|
|
- )
|
|
|
+ for (const guard of commitGuards) {
|
|
|
+ yield* guard(event)
|
|
|
+ }
|
|
|
+ for (const projector of list) {
|
|
|
+ yield* projector({ ...event, seq } as Payload)
|
|
|
+ }
|
|
|
+ const encoded = syncRegistry
|
|
|
+ .get(versionedType(definition.type, sync.version))!
|
|
|
+ .encode(event.data)
|
|
|
+ yield* db
|
|
|
+ .insert(EventSequenceTable)
|
|
|
+ .values([{ aggregate_id: aggregateID, seq, owner_id: input?.ownerID }])
|
|
|
+ .onConflictDoUpdate({
|
|
|
+ target: EventSequenceTable.aggregate_id,
|
|
|
+ set: {
|
|
|
+ seq,
|
|
|
+ ...(input?.ownerID && row?.ownerID == null ? { owner_id: input.ownerID } : {}),
|
|
|
+ },
|
|
|
+ })
|
|
|
+ .run()
|
|
|
+ .pipe(Effect.orDie)
|
|
|
+ yield* db
|
|
|
+ .insert(EventTable)
|
|
|
+ .values([
|
|
|
+ {
|
|
|
+ id: event.id,
|
|
|
+ aggregate_id: aggregateID,
|
|
|
+ seq,
|
|
|
+ type: versionedType(definition.type, sync.version),
|
|
|
+ data: encoded as Record<string, unknown>,
|
|
|
+ },
|
|
|
+ ])
|
|
|
+ .run()
|
|
|
+ .pipe(Effect.orDie)
|
|
|
+ return { aggregateID, seq }
|
|
|
+ }),
|
|
|
+ { behavior: "immediate" },
|
|
|
+ )
|
|
|
+ .pipe(Effect.orDie)
|
|
|
+ if (committed) {
|
|
|
+ yield* Effect.forEach(
|
|
|
+ synchronized.get(committed.aggregateID) ?? [],
|
|
|
+ (pubsub) => PubSub.publish(pubsub, undefined),
|
|
|
+ { discard: true },
|
|
|
+ )
|
|
|
+ }
|
|
|
+ return committed
|
|
|
+ }),
|
|
|
+ )
|
|
|
+ }
|
|
|
}
|
|
|
- }
|
|
|
- })
|
|
|
- }
|
|
|
-
|
|
|
- function publishEvent<D extends Definition>(event: Payload<D>) {
|
|
|
- return Effect.gen(function* () {
|
|
|
- const durable = registry.get(event.type)?.sync !== undefined
|
|
|
- if (durable) {
|
|
|
- for (const sync of syncHandlers) {
|
|
|
- yield* sync(event as Payload)
|
|
|
+ })
|
|
|
+ }
|
|
|
+
|
|
|
+ function publishEvent<D extends Definition>(event: Payload<D>) {
|
|
|
+ return Effect.gen(function* () {
|
|
|
+ const durable = registry.get(event.type)?.sync !== undefined
|
|
|
+ if (durable) {
|
|
|
+ for (const sync of syncHandlers) {
|
|
|
+ yield* sync(event as Payload)
|
|
|
+ }
|
|
|
+ const committed = yield* commitSyncEvent(event as Payload)
|
|
|
+ if (committed) event = { ...event, seq: committed.seq }
|
|
|
}
|
|
|
- const committed = yield* commitSyncEvent(event as Payload)
|
|
|
- if (committed) event = { ...event, seq: committed.seq }
|
|
|
- }
|
|
|
- for (const listener of listeners) {
|
|
|
- yield* listener(event as Payload)
|
|
|
- }
|
|
|
- const pubsub = typed.get(event.type)
|
|
|
- if (pubsub) yield* PubSub.publish(pubsub, event as Payload)
|
|
|
- yield* PubSub.publish(all, event as Payload)
|
|
|
- return event
|
|
|
- })
|
|
|
- }
|
|
|
-
|
|
|
- function publish<D extends Definition>(definition: D, data: Data<D>, options?: PublishOptions) {
|
|
|
- return Effect.gen(function* () {
|
|
|
- const serviceLocation = Option.getOrUndefined(yield* Effect.serviceOption(Location.Service))
|
|
|
- const location =
|
|
|
- options?.location ??
|
|
|
- (serviceLocation
|
|
|
- ? { directory: serviceLocation.directory, workspaceID: serviceLocation.workspaceID }
|
|
|
- : undefined)
|
|
|
- return yield* publishEvent({
|
|
|
- id: options?.id ?? ID.create(),
|
|
|
- ...(options?.metadata ? { metadata: options.metadata } : {}),
|
|
|
- type: definition.type,
|
|
|
- ...(definition.sync === undefined ? {} : { version: definition.sync.version }),
|
|
|
- ...(location ? { location } : {}),
|
|
|
- data,
|
|
|
- } as Payload<D>)
|
|
|
- })
|
|
|
- }
|
|
|
+ for (const listener of listeners) {
|
|
|
+ yield* listener(event as Payload)
|
|
|
+ }
|
|
|
+ const pubsub = typed.get(event.type)
|
|
|
+ if (pubsub) yield* PubSub.publish(pubsub, event as Payload)
|
|
|
+ yield* PubSub.publish(all, event as Payload)
|
|
|
+ return event
|
|
|
+ })
|
|
|
+ }
|
|
|
|
|
|
- function replay(event: SerializedEvent, options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean }) {
|
|
|
- return Effect.gen(function* () {
|
|
|
- const definition = syncRegistry.get(event.type)
|
|
|
- if (!definition) {
|
|
|
- yield* Effect.die(
|
|
|
- new InvalidSyncEventError({ type: event.type, message: `Unknown sync event type ${event.type}` }),
|
|
|
- )
|
|
|
- } else {
|
|
|
- const payload = {
|
|
|
- id: event.id,
|
|
|
+ function publish<D extends Definition>(definition: D, data: Data<D>, options?: PublishOptions) {
|
|
|
+ return Effect.gen(function* () {
|
|
|
+ const serviceLocation = Option.getOrUndefined(yield* Effect.serviceOption(Location.Service))
|
|
|
+ const location =
|
|
|
+ options?.location ??
|
|
|
+ (serviceLocation
|
|
|
+ ? { directory: serviceLocation.directory, workspaceID: serviceLocation.workspaceID }
|
|
|
+ : undefined)
|
|
|
+ return yield* publishEvent({
|
|
|
+ id: options?.id ?? ID.create(),
|
|
|
+ ...(options?.metadata ? { metadata: options.metadata } : {}),
|
|
|
type: definition.type,
|
|
|
- version: definition.sync.version,
|
|
|
- data: definition.decode(event.data),
|
|
|
- } as Payload
|
|
|
- const committed = yield* commitSyncEvent(payload, { seq: event.seq, aggregateID: event.aggregateID, ownerID: options?.ownerID, strictOwner: options?.strictOwner })
|
|
|
- if (committed && options?.publish) {
|
|
|
- const published = { ...payload, seq: committed.seq }
|
|
|
- for (const listener of listeners) {
|
|
|
- yield* listener(published)
|
|
|
+ ...(definition.sync === undefined ? {} : { version: definition.sync.version }),
|
|
|
+ ...(location ? { location } : {}),
|
|
|
+ data,
|
|
|
+ } as Payload<D>)
|
|
|
+ })
|
|
|
+ }
|
|
|
+
|
|
|
+ function replay(
|
|
|
+ event: SerializedEvent,
|
|
|
+ options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
|
|
|
+ ) {
|
|
|
+ return Effect.gen(function* () {
|
|
|
+ const definition = syncRegistry.get(event.type)
|
|
|
+ if (!definition) {
|
|
|
+ yield* Effect.die(
|
|
|
+ new InvalidSyncEventError({ type: event.type, message: `Unknown sync event type ${event.type}` }),
|
|
|
+ )
|
|
|
+ } else {
|
|
|
+ const payload = {
|
|
|
+ id: event.id,
|
|
|
+ type: definition.type,
|
|
|
+ version: definition.sync.version,
|
|
|
+ data: definition.decode(event.data),
|
|
|
+ } as Payload
|
|
|
+ const committed = yield* commitSyncEvent(payload, {
|
|
|
+ seq: event.seq,
|
|
|
+ aggregateID: event.aggregateID,
|
|
|
+ ownerID: options?.ownerID,
|
|
|
+ strictOwner: options?.strictOwner,
|
|
|
+ })
|
|
|
+ if (committed && options?.publish) {
|
|
|
+ const published = { ...payload, seq: committed.seq }
|
|
|
+ for (const listener of listeners) {
|
|
|
+ yield* listener(published)
|
|
|
+ }
|
|
|
+ const pubsub = typed.get(payload.type)
|
|
|
+ if (pubsub) yield* PubSub.publish(pubsub, published)
|
|
|
+ yield* PubSub.publish(all, published)
|
|
|
}
|
|
|
- const pubsub = typed.get(payload.type)
|
|
|
- if (pubsub) yield* PubSub.publish(pubsub, published)
|
|
|
- yield* PubSub.publish(all, published)
|
|
|
}
|
|
|
- }
|
|
|
- })
|
|
|
- }
|
|
|
-
|
|
|
- function replayAll(events: SerializedEvent[], options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean }) {
|
|
|
- return Effect.gen(function* () {
|
|
|
- const source = events[0]?.aggregateID
|
|
|
- if (!source) return undefined
|
|
|
- if (events.some((event) => event.aggregateID !== source)) {
|
|
|
- yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
- type: events[0]?.type ?? "unknown",
|
|
|
- message: "Replay events must belong to the same aggregate",
|
|
|
- }),
|
|
|
- )
|
|
|
- }
|
|
|
- const start = events[0]?.seq ?? 0
|
|
|
- for (const [index, event] of events.entries()) {
|
|
|
- const seq = start + index
|
|
|
- if (event.seq !== seq) {
|
|
|
+ })
|
|
|
+ }
|
|
|
+
|
|
|
+ function replayAll(
|
|
|
+ events: SerializedEvent[],
|
|
|
+ options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
|
|
|
+ ) {
|
|
|
+ return Effect.gen(function* () {
|
|
|
+ const source = events[0]?.aggregateID
|
|
|
+ if (!source) return undefined
|
|
|
+ if (events.some((event) => event.aggregateID !== source)) {
|
|
|
yield* Effect.die(
|
|
|
new InvalidSyncEventError({
|
|
|
- type: event.type,
|
|
|
- message: `Replay sequence mismatch at index ${index}: expected ${seq}, got ${event.seq}`,
|
|
|
+ type: events[0]?.type ?? "unknown",
|
|
|
+ message: "Replay events must belong to the same aggregate",
|
|
|
}),
|
|
|
)
|
|
|
}
|
|
|
- }
|
|
|
- for (const event of events) {
|
|
|
- yield* replay(event, options)
|
|
|
- }
|
|
|
- return source
|
|
|
- })
|
|
|
- }
|
|
|
+ const start = events[0]?.seq ?? 0
|
|
|
+ for (const [index, event] of events.entries()) {
|
|
|
+ const seq = start + index
|
|
|
+ if (event.seq !== seq) {
|
|
|
+ yield* Effect.die(
|
|
|
+ new InvalidSyncEventError({
|
|
|
+ type: event.type,
|
|
|
+ message: `Replay sequence mismatch at index ${index}: expected ${seq}, got ${event.seq}`,
|
|
|
+ }),
|
|
|
+ )
|
|
|
+ }
|
|
|
+ }
|
|
|
+ for (const event of events) {
|
|
|
+ yield* replay(event, options)
|
|
|
+ }
|
|
|
+ return source
|
|
|
+ })
|
|
|
+ }
|
|
|
|
|
|
- function remove(aggregateID: string) {
|
|
|
- return db
|
|
|
- .transaction(() =>
|
|
|
- Effect.gen(function* () {
|
|
|
- yield* db.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).run()
|
|
|
- yield* db.delete(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).run()
|
|
|
- }),
|
|
|
+ function remove(aggregateID: string) {
|
|
|
+ return db
|
|
|
+ .transaction(() =>
|
|
|
+ Effect.gen(function* () {
|
|
|
+ yield* db.delete(EventSequenceTable).where(eq(EventSequenceTable.aggregate_id, aggregateID)).run()
|
|
|
+ yield* db.delete(EventTable).where(eq(EventTable.aggregate_id, aggregateID)).run()
|
|
|
+ }),
|
|
|
+ )
|
|
|
+ .pipe(Effect.orDie)
|
|
|
+ }
|
|
|
+
|
|
|
+ function claim(aggregateID: string, ownerID: string) {
|
|
|
+ return db
|
|
|
+ .update(EventSequenceTable)
|
|
|
+ .set({ owner_id: ownerID })
|
|
|
+ .where(eq(EventSequenceTable.aggregate_id, aggregateID))
|
|
|
+ .run()
|
|
|
+ .pipe(Effect.orDie)
|
|
|
+ }
|
|
|
+
|
|
|
+ const subscribe = <D extends Definition>(definition: D): Stream.Stream<Payload<D>> =>
|
|
|
+ Stream.unwrap(getOrCreate(definition).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub)))).pipe(
|
|
|
+ Stream.map((event) => event as Payload<D>),
|
|
|
)
|
|
|
- .pipe(Effect.orDie)
|
|
|
- }
|
|
|
-
|
|
|
- function claim(aggregateID: string, ownerID: string) {
|
|
|
- return db
|
|
|
- .update(EventSequenceTable)
|
|
|
- .set({ owner_id: ownerID })
|
|
|
- .where(eq(EventSequenceTable.aggregate_id, aggregateID))
|
|
|
- .run()
|
|
|
- .pipe(Effect.orDie)
|
|
|
- }
|
|
|
-
|
|
|
- const subscribe = <D extends Definition>(definition: D): Stream.Stream<Payload<D>> =>
|
|
|
- Stream.unwrap(getOrCreate(definition).pipe(Effect.map((pubsub) => Stream.fromPubSub(pubsub)))).pipe(
|
|
|
- Stream.map((event) => event as Payload<D>),
|
|
|
- )
|
|
|
|
|
|
- const streamAll = (): Stream.Stream<Payload> => Stream.fromPubSub(all)
|
|
|
+ const streamAll = (): Stream.Stream<Payload> => Stream.fromPubSub(all)
|
|
|
|
|
|
- const decodeSerializedEvent = (event: SerializedEvent): CursorEvent => {
|
|
|
- const definition = syncRegistry.get(event.type)
|
|
|
- if (!definition) {
|
|
|
- throw new InvalidSyncEventError({ type: event.type, message: `Unknown sync event type ${event.type}` })
|
|
|
- }
|
|
|
- return {
|
|
|
- cursor: Cursor.make(event.seq),
|
|
|
- event: {
|
|
|
- id: event.id,
|
|
|
- type: definition.type,
|
|
|
- version: definition.sync.version,
|
|
|
- seq: event.seq,
|
|
|
- data: definition.decode(event.data),
|
|
|
- },
|
|
|
+ const decodeSerializedEvent = (event: SerializedEvent): CursorEvent => {
|
|
|
+ const definition = syncRegistry.get(event.type)
|
|
|
+ if (!definition) {
|
|
|
+ throw new InvalidSyncEventError({ type: event.type, message: `Unknown sync event type ${event.type}` })
|
|
|
+ }
|
|
|
+ return {
|
|
|
+ cursor: Cursor.make(event.seq),
|
|
|
+ event: {
|
|
|
+ id: event.id,
|
|
|
+ type: definition.type,
|
|
|
+ version: definition.sync.version,
|
|
|
+ seq: event.seq,
|
|
|
+ data: definition.decode(event.data),
|
|
|
+ },
|
|
|
+ }
|
|
|
}
|
|
|
- }
|
|
|
-
|
|
|
- const readAfter = (aggregateID: string, after: number) =>
|
|
|
- (options?.beforeAggregateRead?.(aggregateID) ?? Effect.void).pipe(
|
|
|
- Effect.andThen(
|
|
|
- db
|
|
|
- .select()
|
|
|
- .from(EventTable)
|
|
|
- .where(and(eq(EventTable.aggregate_id, aggregateID), gt(EventTable.seq, after)))
|
|
|
- .orderBy(asc(EventTable.seq))
|
|
|
- .all(),
|
|
|
- ),
|
|
|
- Effect.orDie,
|
|
|
- Effect.map((rows) =>
|
|
|
- rows.map((event) =>
|
|
|
- decodeSerializedEvent({
|
|
|
- id: event.id,
|
|
|
- aggregateID: event.aggregate_id,
|
|
|
- seq: event.seq,
|
|
|
- type: event.type,
|
|
|
- data: event.data,
|
|
|
- }),
|
|
|
- ),
|
|
|
- ),
|
|
|
- )
|
|
|
|
|
|
- const subscribeSynchronized = (aggregateID: string) =>
|
|
|
- Effect.gen(function* () {
|
|
|
- const pubsub = yield* PubSub.sliding<void>(1)
|
|
|
- const subscription = yield* PubSub.subscribe(pubsub)
|
|
|
- yield* Effect.acquireRelease(
|
|
|
- Effect.sync(() => {
|
|
|
- const pubsubs = synchronized.get(aggregateID) ?? new Set()
|
|
|
- pubsubs.add(pubsub)
|
|
|
- synchronized.set(aggregateID, pubsubs)
|
|
|
- }),
|
|
|
- () =>
|
|
|
- Effect.sync(() => {
|
|
|
- const pubsubs = synchronized.get(aggregateID)
|
|
|
- pubsubs?.delete(pubsub)
|
|
|
- if (pubsubs?.size === 0) synchronized.delete(aggregateID)
|
|
|
- }).pipe(Effect.andThen(PubSub.shutdown(pubsub))),
|
|
|
+ const readAfter = (aggregateID: string, after: number) =>
|
|
|
+ (options?.beforeAggregateRead?.(aggregateID) ?? Effect.void).pipe(
|
|
|
+ Effect.andThen(
|
|
|
+ db
|
|
|
+ .select()
|
|
|
+ .from(EventTable)
|
|
|
+ .where(and(eq(EventTable.aggregate_id, aggregateID), gt(EventTable.seq, after)))
|
|
|
+ .orderBy(asc(EventTable.seq))
|
|
|
+ .all(),
|
|
|
+ ),
|
|
|
+ Effect.orDie,
|
|
|
+ Effect.map((rows) =>
|
|
|
+ rows.map((event) =>
|
|
|
+ decodeSerializedEvent({
|
|
|
+ id: event.id,
|
|
|
+ aggregateID: event.aggregate_id,
|
|
|
+ seq: event.seq,
|
|
|
+ type: event.type,
|
|
|
+ data: event.data,
|
|
|
+ }),
|
|
|
+ ),
|
|
|
+ ),
|
|
|
)
|
|
|
- return subscription
|
|
|
- })
|
|
|
|
|
|
- const streamEvents = (input: { readonly aggregateID: string; readonly after?: Cursor }): Stream.Stream<CursorEvent> =>
|
|
|
- Stream.unwrap(
|
|
|
+ const subscribeSynchronized = (aggregateID: string) =>
|
|
|
Effect.gen(function* () {
|
|
|
- const synchronized = yield* subscribeSynchronized(input.aggregateID)
|
|
|
- let cursor = input.after ?? -1
|
|
|
- const read = Effect.suspend(() => readAfter(input.aggregateID, cursor)).pipe(
|
|
|
- Effect.tap((events) =>
|
|
|
+ const pubsub = yield* PubSub.sliding<void>(1)
|
|
|
+ const subscription = yield* PubSub.subscribe(pubsub)
|
|
|
+ yield* Effect.acquireRelease(
|
|
|
+ Effect.sync(() => {
|
|
|
+ const pubsubs = synchronized.get(aggregateID) ?? new Set()
|
|
|
+ pubsubs.add(pubsub)
|
|
|
+ synchronized.set(aggregateID, pubsubs)
|
|
|
+ }),
|
|
|
+ () =>
|
|
|
Effect.sync(() => {
|
|
|
- cursor = events.at(-1)?.cursor ?? cursor
|
|
|
- }),
|
|
|
- ),
|
|
|
+ const pubsubs = synchronized.get(aggregateID)
|
|
|
+ pubsubs?.delete(pubsub)
|
|
|
+ if (pubsubs?.size === 0) synchronized.delete(aggregateID)
|
|
|
+ }).pipe(Effect.andThen(PubSub.shutdown(pubsub))),
|
|
|
)
|
|
|
- const historical = yield* read
|
|
|
- const live = Stream.fromSubscription(synchronized).pipe(
|
|
|
- Stream.mapEffect(() => read),
|
|
|
- Stream.flattenIterable,
|
|
|
- )
|
|
|
- return Stream.concat(Stream.fromIterable(historical), live)
|
|
|
- }),
|
|
|
- )
|
|
|
+ return subscription
|
|
|
+ })
|
|
|
+
|
|
|
+ const streamEvents = (input: {
|
|
|
+ readonly aggregateID: string
|
|
|
+ readonly after?: Cursor
|
|
|
+ }): Stream.Stream<CursorEvent> =>
|
|
|
+ Stream.unwrap(
|
|
|
+ Effect.gen(function* () {
|
|
|
+ const synchronized = yield* subscribeSynchronized(input.aggregateID)
|
|
|
+ let cursor = input.after ?? -1
|
|
|
+ const read = Effect.suspend(() => readAfter(input.aggregateID, cursor)).pipe(
|
|
|
+ Effect.tap((events) =>
|
|
|
+ Effect.sync(() => {
|
|
|
+ cursor = events.at(-1)?.cursor ?? cursor
|
|
|
+ }),
|
|
|
+ ),
|
|
|
+ )
|
|
|
+ const historical = yield* read
|
|
|
+ const live = Stream.fromSubscription(synchronized).pipe(
|
|
|
+ Stream.mapEffect(() => read),
|
|
|
+ Stream.flattenIterable,
|
|
|
+ )
|
|
|
+ return Stream.concat(Stream.fromIterable(historical), live)
|
|
|
+ }),
|
|
|
+ )
|
|
|
|
|
|
- const listen = (listener: Listener): Effect.Effect<Unsubscribe> =>
|
|
|
- Effect.sync(() => {
|
|
|
- listeners.push(listener)
|
|
|
- return Effect.sync(() => {
|
|
|
- const index = listeners.indexOf(listener)
|
|
|
- if (index >= 0) listeners.splice(index, 1)
|
|
|
+ const listen = (listener: Listener): Effect.Effect<Unsubscribe> =>
|
|
|
+ Effect.sync(() => {
|
|
|
+ listeners.push(listener)
|
|
|
+ return Effect.sync(() => {
|
|
|
+ const index = listeners.indexOf(listener)
|
|
|
+ if (index >= 0) listeners.splice(index, 1)
|
|
|
+ })
|
|
|
})
|
|
|
- })
|
|
|
|
|
|
- const sync = (handler: Sync): Effect.Effect<Unsubscribe> =>
|
|
|
- Effect.sync(() => {
|
|
|
- syncHandlers.push(handler)
|
|
|
- return Effect.sync(() => {
|
|
|
- const index = syncHandlers.indexOf(handler)
|
|
|
- if (index >= 0) syncHandlers.splice(index, 1)
|
|
|
+ const sync = (handler: Sync): Effect.Effect<Unsubscribe> =>
|
|
|
+ Effect.sync(() => {
|
|
|
+ syncHandlers.push(handler)
|
|
|
+ return Effect.sync(() => {
|
|
|
+ const index = syncHandlers.indexOf(handler)
|
|
|
+ if (index >= 0) syncHandlers.splice(index, 1)
|
|
|
+ })
|
|
|
})
|
|
|
- })
|
|
|
|
|
|
- const beforeCommit = (guard: CommitGuard): Effect.Effect<void> =>
|
|
|
- Effect.sync(() => {
|
|
|
- commitGuards.push(guard)
|
|
|
- })
|
|
|
+ const beforeCommit = (guard: CommitGuard): Effect.Effect<void> =>
|
|
|
+ Effect.sync(() => {
|
|
|
+ commitGuards.push(guard)
|
|
|
+ })
|
|
|
|
|
|
- const project = <D extends Definition>(definition: D, projector: Projector<D>): Effect.Effect<void> =>
|
|
|
- Effect.sync(() => {
|
|
|
- const list = projectors.get(definition.type) ?? []
|
|
|
- list.push((event) => projector(event as Payload<D>))
|
|
|
- projectors.set(definition.type, list)
|
|
|
- })
|
|
|
+ const project = <D extends Definition>(definition: D, projector: Projector<D>): Effect.Effect<void> =>
|
|
|
+ Effect.sync(() => {
|
|
|
+ const list = projectors.get(definition.type) ?? []
|
|
|
+ list.push((event) => projector(event as Payload<D>))
|
|
|
+ projectors.set(definition.type, list)
|
|
|
+ })
|
|
|
|
|
|
- return Service.of({ publish, subscribe, all: streamAll, aggregateEvents: streamEvents, sync, listen, beforeCommit, project, replay, replayAll, remove, claim })
|
|
|
- }),
|
|
|
-)
|
|
|
+ return Service.of({
|
|
|
+ publish,
|
|
|
+ subscribe,
|
|
|
+ all: streamAll,
|
|
|
+ aggregateEvents: streamEvents,
|
|
|
+ sync,
|
|
|
+ listen,
|
|
|
+ beforeCommit,
|
|
|
+ project,
|
|
|
+ replay,
|
|
|
+ replayAll,
|
|
|
+ remove,
|
|
|
+ claim,
|
|
|
+ })
|
|
|
+ }),
|
|
|
+ )
|
|
|
|
|
|
export const layer = layerWith()
|
|
|
|