|
@@ -5,7 +5,7 @@ import { and, asc, eq, gt } from "drizzle-orm"
|
|
|
import { Database } from "./database/database"
|
|
import { Database } from "./database/database"
|
|
|
import { EventSequenceTable, EventTable } from "./event/sql"
|
|
import { EventSequenceTable, EventTable } from "./event/sql"
|
|
|
import { Location } from "./location"
|
|
import { Location } from "./location"
|
|
|
-import { externalID, type ExternalID, NonNegativeInt, withStatics } from "./schema"
|
|
|
|
|
|
|
+import { externalID, type ExternalID, withStatics } from "./schema"
|
|
|
import { Identifier } from "./util/identifier"
|
|
import { Identifier } from "./util/identifier"
|
|
|
import { LayerNode } from "./effect/layer-node"
|
|
import { LayerNode } from "./effect/layer-node"
|
|
|
import { isDeepStrictEqual } from "node:util"
|
|
import { isDeepStrictEqual } from "node:util"
|
|
@@ -19,16 +19,9 @@ export const ID = Schema.String.check(Schema.isStartsWith("evt_")).pipe(
|
|
|
)
|
|
)
|
|
|
export type ID = typeof ID.Type
|
|
export type ID = typeof ID.Type
|
|
|
|
|
|
|
|
-/**
|
|
|
|
|
- * Durable aggregate continuation position for embedded replay streams.
|
|
|
|
|
- * TODO: Decide whether a future HTTP / SDK surface should expose an opaque cursor instead.
|
|
|
|
|
- */
|
|
|
|
|
-export const Cursor = NonNegativeInt.pipe(Schema.brand("EventV2.Cursor"))
|
|
|
|
|
-export type Cursor = typeof Cursor.Type
|
|
|
|
|
-
|
|
|
|
|
export type Definition<Type extends string = string, DataSchema extends Schema.Top = Schema.Top> = {
|
|
export type Definition<Type extends string = string, DataSchema extends Schema.Top = Schema.Top> = {
|
|
|
readonly type: Type
|
|
readonly type: Type
|
|
|
- readonly sync?: {
|
|
|
|
|
|
|
+ readonly durable?: {
|
|
|
readonly version: number
|
|
readonly version: number
|
|
|
readonly aggregate: string
|
|
readonly aggregate: string
|
|
|
}
|
|
}
|
|
@@ -41,20 +34,16 @@ export type Payload<D extends Definition = Definition> = {
|
|
|
readonly id: ID
|
|
readonly id: ID
|
|
|
readonly type: D["type"]
|
|
readonly type: D["type"]
|
|
|
readonly data: Data<D>
|
|
readonly data: Data<D>
|
|
|
- /** Durable aggregate order, populated while synchronized events are projected. */
|
|
|
|
|
- readonly seq?: number
|
|
|
|
|
- readonly version?: number
|
|
|
|
|
|
|
+ readonly durable?: {
|
|
|
|
|
+ readonly aggregateID: string
|
|
|
|
|
+ readonly seq: number
|
|
|
|
|
+ readonly version: number
|
|
|
|
|
+ }
|
|
|
readonly location?: Location.Ref
|
|
readonly location?: Location.Ref
|
|
|
readonly metadata?: Record<string, unknown>
|
|
readonly metadata?: Record<string, unknown>
|
|
|
- /** Internal replay marker for projectors that own non-replicated operational state. */
|
|
|
|
|
- readonly replay?: boolean
|
|
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
-export type Projector<D extends Definition = Definition> = (event: Payload<D>) => Effect.Effect<void>
|
|
|
|
|
-type AnyProjector = (event: Payload) => Effect.Effect<void>
|
|
|
|
|
-export type CommitGuard = (event: Payload) => Effect.Effect<void>
|
|
|
|
|
-export type Listener = (event: Payload) => Effect.Effect<void>
|
|
|
|
|
-export type Sync = (event: Payload) => Effect.Effect<void>
|
|
|
|
|
|
|
+export type Subscriber<D extends Definition = Definition> = (event: Payload<D>) => Effect.Effect<void>
|
|
|
export type Unsubscribe = Effect.Effect<void>
|
|
export type Unsubscribe = Effect.Effect<void>
|
|
|
|
|
|
|
|
export type SerializedEvent = {
|
|
export type SerializedEvent = {
|
|
@@ -65,13 +54,8 @@ export type SerializedEvent = {
|
|
|
readonly data: Record<string, unknown>
|
|
readonly data: Record<string, unknown>
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
-export type CursorEvent<E extends Payload = Payload> = {
|
|
|
|
|
- readonly cursor: Cursor
|
|
|
|
|
- readonly event: E
|
|
|
|
|
-}
|
|
|
|
|
-
|
|
|
|
|
-export class InvalidSyncEventError extends Schema.TaggedErrorClass<InvalidSyncEventError>()(
|
|
|
|
|
- "EventV2.InvalidSyncEvent",
|
|
|
|
|
|
|
+export class InvalidDurableEventError extends Schema.TaggedErrorClass<InvalidDurableEventError>()(
|
|
|
|
|
+ "EventV2.InvalidDurableEvent",
|
|
|
{
|
|
{
|
|
|
type: Schema.String,
|
|
type: Schema.String,
|
|
|
message: Schema.String,
|
|
message: Schema.String,
|
|
@@ -83,19 +67,11 @@ export function versionedType(type: string, version: number) {
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
export const registry = new Map<string, Definition>()
|
|
export const registry = new Map<string, Definition>()
|
|
|
-type SyncDefinition = Definition & {
|
|
|
|
|
- readonly sync: NonNullable<Definition["sync"]>
|
|
|
|
|
- readonly encode: (data: unknown) => unknown
|
|
|
|
|
- readonly decode: (data: unknown) => unknown
|
|
|
|
|
-}
|
|
|
|
|
-const syncRegistry = new Map<string, SyncDefinition>()
|
|
|
|
|
-
|
|
|
|
|
-// Synchronized events cross a JSON boundary, so their data schemas must encode and decode without services.
|
|
|
|
|
-const syncCodec = (definition: Definition) => definition.data as Schema.Codec<unknown, unknown, never, never>
|
|
|
|
|
|
|
+const durableRegistry = new Map<string, Definition>()
|
|
|
|
|
|
|
|
export function define<const Type extends string, Fields extends Schema.Struct.Fields>(input: {
|
|
export function define<const Type extends string, Fields extends Schema.Struct.Fields>(input: {
|
|
|
readonly type: Type
|
|
readonly type: Type
|
|
|
- readonly sync?: {
|
|
|
|
|
|
|
+ readonly durable?: {
|
|
|
readonly version: number
|
|
readonly version: number
|
|
|
readonly aggregate: string
|
|
readonly aggregate: string
|
|
|
}
|
|
}
|
|
@@ -106,28 +82,26 @@ export function define<const Type extends string, Fields extends Schema.Struct.F
|
|
|
id: ID,
|
|
id: ID,
|
|
|
metadata: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)),
|
|
metadata: Schema.optional(Schema.Record(Schema.String, Schema.Unknown)),
|
|
|
type: Schema.Literal(input.type),
|
|
type: Schema.Literal(input.type),
|
|
|
- version: Schema.optional(Schema.Number),
|
|
|
|
|
|
|
+ durable: Schema.optional(Schema.Struct({ aggregateID: Schema.String, seq: Schema.Number, version: Schema.Number })),
|
|
|
location: Schema.optional(Location.Ref),
|
|
location: Schema.optional(Location.Ref),
|
|
|
data: Data,
|
|
data: Data,
|
|
|
}).annotate({ identifier: input.type })
|
|
}).annotate({ identifier: input.type })
|
|
|
|
|
|
|
|
const definition = Object.assign(Payload, {
|
|
const definition = Object.assign(Payload, {
|
|
|
type: input.type,
|
|
type: input.type,
|
|
|
- ...(input.sync === undefined ? {} : { sync: input.sync }),
|
|
|
|
|
|
|
+ ...(input.durable === undefined ? {} : { durable: input.durable }),
|
|
|
data: Data,
|
|
data: Data,
|
|
|
})
|
|
})
|
|
|
const existing = registry.get(input.type)
|
|
const existing = registry.get(input.type)
|
|
|
- if (input.sync === undefined || existing?.sync === undefined || input.sync.version >= existing.sync.version) {
|
|
|
|
|
|
|
+ if (
|
|
|
|
|
+ input.durable === undefined ||
|
|
|
|
|
+ existing?.durable === undefined ||
|
|
|
|
|
+ input.durable.version >= existing.durable.version
|
|
|
|
|
+ ) {
|
|
|
registry.set(input.type, definition)
|
|
registry.set(input.type, definition)
|
|
|
}
|
|
}
|
|
|
- if (input.sync)
|
|
|
|
|
- syncRegistry.set(
|
|
|
|
|
- versionedType(input.type, input.sync.version),
|
|
|
|
|
- Object.assign(definition, {
|
|
|
|
|
- encode: Schema.encodeUnknownSync(syncCodec(definition)),
|
|
|
|
|
- decode: Schema.decodeUnknownSync(syncCodec(definition)),
|
|
|
|
|
- }) as SyncDefinition,
|
|
|
|
|
- )
|
|
|
|
|
|
|
+ if (input.durable)
|
|
|
|
|
+ durableRegistry.set(versionedType(input.type, input.durable.version), definition)
|
|
|
return definition as Schema.Schema<Payload<Definition<Type, Schema.Struct<Fields>>>> &
|
|
return definition as Schema.Schema<Payload<Definition<Type, Schema.Struct<Fields>>>> &
|
|
|
Definition<Type, Schema.Struct<Fields>>
|
|
Definition<Type, Schema.Struct<Fields>>
|
|
|
}
|
|
}
|
|
@@ -140,7 +114,7 @@ export interface PublishOptions {
|
|
|
readonly id?: ID
|
|
readonly id?: ID
|
|
|
readonly metadata?: Record<string, unknown>
|
|
readonly metadata?: Record<string, unknown>
|
|
|
readonly location?: Location.Ref
|
|
readonly location?: Location.Ref
|
|
|
- /** Local operational projection committed atomically with a new synchronized event. Not replayed or serialized. */
|
|
|
|
|
|
|
+ /** Local operational projection committed atomically with a new durable event. Not replayed or serialized. */
|
|
|
readonly commit?: (seq: number) => Effect.Effect<void>
|
|
readonly commit?: (seq: number) => Effect.Effect<void>
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -152,14 +126,13 @@ export interface Interface {
|
|
|
) => Effect.Effect<Payload<D>>
|
|
) => Effect.Effect<Payload<D>>
|
|
|
readonly subscribe: <D extends Definition>(definition: D) => Stream.Stream<Payload<D>>
|
|
readonly subscribe: <D extends Definition>(definition: D) => Stream.Stream<Payload<D>>
|
|
|
readonly all: () => Stream.Stream<Payload>
|
|
readonly all: () => Stream.Stream<Payload>
|
|
|
- readonly aggregateEvents: (input: {
|
|
|
|
|
|
|
+ readonly durable: (input: {
|
|
|
readonly aggregateID: string
|
|
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>
|
|
|
|
|
- readonly project: <D extends Definition>(definition: D, projector: Projector<D>) => Effect.Effect<void>
|
|
|
|
|
|
|
+ readonly after?: number
|
|
|
|
|
+ }) => Stream.Stream<Payload>
|
|
|
|
|
+ /** @deprecated Use `all()` and consume the returned stream. */
|
|
|
|
|
+ readonly listen: (listener: Subscriber) => Effect.Effect<Unsubscribe>
|
|
|
|
|
+ readonly project: <D extends Definition>(definition: D, projector: Subscriber<D>) => Effect.Effect<void>
|
|
|
readonly replay: (
|
|
readonly replay: (
|
|
|
event: SerializedEvent,
|
|
event: SerializedEvent,
|
|
|
options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
|
|
options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
|
|
@@ -182,37 +155,37 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
Layer.effect(
|
|
Layer.effect(
|
|
|
Service,
|
|
Service,
|
|
|
Effect.gen(function* () {
|
|
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 pubsub = {
|
|
|
|
|
+ all: yield* PubSub.unbounded<Payload>(),
|
|
|
|
|
+ durable: new Map<string, Set<PubSub.PubSub<void>>>(),
|
|
|
|
|
+ typed: new Map<string, PubSub.PubSub<Payload>>(),
|
|
|
|
|
+ }
|
|
|
|
|
+ const projectors = new Map<string, Subscriber[]>()
|
|
|
|
|
+ const listeners = new Array<Subscriber>()
|
|
|
const { db } = yield* Database.Service
|
|
const { db } = yield* Database.Service
|
|
|
|
|
|
|
|
const getOrCreate = (definition: Definition) =>
|
|
const getOrCreate = (definition: Definition) =>
|
|
|
Effect.gen(function* () {
|
|
Effect.gen(function* () {
|
|
|
- const existing = typed.get(definition.type)
|
|
|
|
|
|
|
+ const existing = pubsub.typed.get(definition.type)
|
|
|
if (existing) return existing
|
|
if (existing) return existing
|
|
|
- const pubsub = yield* PubSub.unbounded<Payload>()
|
|
|
|
|
- typed.set(definition.type, pubsub)
|
|
|
|
|
- return pubsub
|
|
|
|
|
|
|
+ const created = yield* PubSub.unbounded<Payload>()
|
|
|
|
|
+ pubsub.typed.set(definition.type, created)
|
|
|
|
|
+ return created
|
|
|
})
|
|
})
|
|
|
|
|
|
|
|
yield* Effect.addFinalizer(() =>
|
|
yield* Effect.addFinalizer(() =>
|
|
|
Effect.gen(function* () {
|
|
Effect.gen(function* () {
|
|
|
- yield* PubSub.shutdown(all)
|
|
|
|
|
|
|
+ yield* PubSub.shutdown(pubsub.all)
|
|
|
yield* Effect.forEach(
|
|
yield* Effect.forEach(
|
|
|
- synchronized.values(),
|
|
|
|
|
|
|
+ pubsub.durable.values(),
|
|
|
(pubsubs) => Effect.forEach(pubsubs, PubSub.shutdown, { discard: true }),
|
|
(pubsubs) => Effect.forEach(pubsubs, PubSub.shutdown, { discard: true }),
|
|
|
{ discard: true },
|
|
{ discard: true },
|
|
|
)
|
|
)
|
|
|
- yield* Effect.forEach(typed.values(), PubSub.shutdown, { discard: true })
|
|
|
|
|
|
|
+ yield* Effect.forEach(pubsub.typed.values(), PubSub.shutdown, { discard: true })
|
|
|
}),
|
|
}),
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
- function commitSyncEvent(
|
|
|
|
|
|
|
+ function commitDurableEvent(
|
|
|
event: Payload,
|
|
event: Payload,
|
|
|
input?: {
|
|
input?: {
|
|
|
readonly seq: number
|
|
readonly seq: number
|
|
@@ -224,28 +197,20 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
) {
|
|
) {
|
|
|
return Effect.gen(function* () {
|
|
return Effect.gen(function* () {
|
|
|
const definition = registry.get(event.type)
|
|
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]
|
|
|
|
|
|
|
+ const durable = definition?.durable
|
|
|
|
|
+ if (durable) {
|
|
|
|
|
+ const aggregateID = (event.data as Record<string, unknown>)[durable.aggregate]
|
|
|
if (typeof aggregateID !== "string") {
|
|
if (typeof aggregateID !== "string") {
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: event.type,
|
|
type: event.type,
|
|
|
- message: `Expected string aggregate field ${sync.aggregate}`,
|
|
|
|
|
|
|
+ message: `Expected string aggregate field ${durable.aggregate}`,
|
|
|
}),
|
|
}),
|
|
|
)
|
|
)
|
|
|
} else {
|
|
} else {
|
|
|
if (input && input.aggregateID !== aggregateID) {
|
|
if (input && input.aggregateID !== aggregateID) {
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: event.type,
|
|
type: event.type,
|
|
|
message: `Aggregate mismatch: expected ${input.aggregateID}, got ${aggregateID}`,
|
|
message: `Aggregate mismatch: expected ${input.aggregateID}, got ${aggregateID}`,
|
|
|
}),
|
|
}),
|
|
@@ -265,12 +230,12 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
.get()
|
|
.get()
|
|
|
.pipe(Effect.orDie)
|
|
.pipe(Effect.orDie)
|
|
|
const latest = row?.seq ?? -1
|
|
const latest = row?.seq ?? -1
|
|
|
- const encoded = syncRegistry
|
|
|
|
|
- .get(versionedType(definition.type, sync.version))!
|
|
|
|
|
- .encode(event.data) as Record<string, unknown>
|
|
|
|
|
|
|
+ const encoded = Schema.encodeUnknownSync(
|
|
|
|
|
+ definition.data as Schema.Codec<unknown, unknown, never, never>,
|
|
|
|
|
+ )(event.data) as Record<string, unknown>
|
|
|
if (input?.strictOwner && row?.ownerID && row.ownerID !== input.ownerID) {
|
|
if (input?.strictOwner && row?.ownerID && row.ownerID !== input.ownerID) {
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: event.type,
|
|
type: event.type,
|
|
|
message: `Replay owner mismatch for aggregate ${aggregateID}: expected ${row.ownerID}, got ${input.ownerID ?? "none"}`,
|
|
message: `Replay owner mismatch for aggregate ${aggregateID}: expected ${row.ownerID}, got ${input.ownerID ?? "none"}`,
|
|
|
}),
|
|
}),
|
|
@@ -285,7 +250,7 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
.pipe(Effect.orDie)
|
|
.pipe(Effect.orDie)
|
|
|
if (
|
|
if (
|
|
|
stored?.id === event.id &&
|
|
stored?.id === event.id &&
|
|
|
- stored.type === versionedType(definition.type, sync.version) &&
|
|
|
|
|
|
|
+ stored.type === versionedType(definition.type, durable.version) &&
|
|
|
isDeepStrictEqual(stored.data, encoded)
|
|
isDeepStrictEqual(stored.data, encoded)
|
|
|
) {
|
|
) {
|
|
|
if (input.ownerID && row?.ownerID == null) {
|
|
if (input.ownerID && row?.ownerID == null) {
|
|
@@ -299,7 +264,7 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
return
|
|
return
|
|
|
}
|
|
}
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: event.type,
|
|
type: event.type,
|
|
|
message: `Replay diverged at aggregate ${aggregateID} sequence ${input.seq}`,
|
|
message: `Replay diverged at aggregate ${aggregateID} sequence ${input.seq}`,
|
|
|
}),
|
|
}),
|
|
@@ -311,7 +276,7 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
const seq = input?.seq ?? latest + 1
|
|
const seq = input?.seq ?? latest + 1
|
|
|
if (input && seq !== latest + 1) {
|
|
if (input && seq !== latest + 1) {
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: event.type,
|
|
type: event.type,
|
|
|
message: `Sequence mismatch for aggregate ${aggregateID}: expected ${latest + 1}, got ${seq}`,
|
|
message: `Sequence mismatch for aggregate ${aggregateID}: expected ${latest + 1}, got ${seq}`,
|
|
|
}),
|
|
}),
|
|
@@ -325,16 +290,17 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
.pipe(Effect.orDie)
|
|
.pipe(Effect.orDie)
|
|
|
if (stored)
|
|
if (stored)
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: event.type,
|
|
type: event.type,
|
|
|
message: `Event ${event.id} already exists at aggregate ${stored.aggregateID} sequence ${stored.seq}`,
|
|
message: `Event ${event.id} already exists at aggregate ${stored.aggregateID} sequence ${stored.seq}`,
|
|
|
}),
|
|
}),
|
|
|
)
|
|
)
|
|
|
- for (const guard of commitGuards) {
|
|
|
|
|
- yield* guard(event)
|
|
|
|
|
- }
|
|
|
|
|
|
|
+ const committed = {
|
|
|
|
|
+ ...event,
|
|
|
|
|
+ durable: { aggregateID, seq, version: durable.version },
|
|
|
|
|
+ } as Payload
|
|
|
for (const projector of list) {
|
|
for (const projector of list) {
|
|
|
- yield* projector({ ...event, seq } as Payload)
|
|
|
|
|
|
|
+ yield* projector(committed)
|
|
|
}
|
|
}
|
|
|
if (commit) yield* commit(seq)
|
|
if (commit) yield* commit(seq)
|
|
|
yield* db
|
|
yield* db
|
|
@@ -356,7 +322,7 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
id: event.id,
|
|
id: event.id,
|
|
|
aggregate_id: aggregateID,
|
|
aggregate_id: aggregateID,
|
|
|
seq,
|
|
seq,
|
|
|
- type: versionedType(definition.type, sync.version),
|
|
|
|
|
|
|
+ type: versionedType(definition.type, durable.version),
|
|
|
data: encoded,
|
|
data: encoded,
|
|
|
},
|
|
},
|
|
|
])
|
|
])
|
|
@@ -369,8 +335,8 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
.pipe(Effect.orDie)
|
|
.pipe(Effect.orDie)
|
|
|
if (committed) {
|
|
if (committed) {
|
|
|
yield* Effect.forEach(
|
|
yield* Effect.forEach(
|
|
|
- synchronized.get(committed.aggregateID) ?? [],
|
|
|
|
|
- (pubsub) => PubSub.publish(pubsub, undefined),
|
|
|
|
|
|
|
+ pubsub.durable.get(committed.aggregateID) ?? [],
|
|
|
|
|
+ (wake) => PubSub.publish(wake, undefined),
|
|
|
{ discard: true },
|
|
{ discard: true },
|
|
|
)
|
|
)
|
|
|
}
|
|
}
|
|
@@ -384,19 +350,25 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
|
|
|
|
|
function publishEvent<D extends Definition>(event: Payload<D>, commit?: PublishOptions["commit"]) {
|
|
function publishEvent<D extends Definition>(event: Payload<D>, commit?: PublishOptions["commit"]) {
|
|
|
return Effect.gen(function* () {
|
|
return Effect.gen(function* () {
|
|
|
- const durable = registry.get(event.type)?.sync !== undefined
|
|
|
|
|
- if (!durable && commit)
|
|
|
|
|
|
|
+ const definition = registry.get(event.type)
|
|
|
|
|
+ if (!definition?.durable && commit)
|
|
|
return yield* Effect.die(
|
|
return yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: event.type,
|
|
type: event.type,
|
|
|
- message: "Local commit hooks require a synchronized event",
|
|
|
|
|
|
|
+ message: "Local commit hooks require a durable event",
|
|
|
}),
|
|
}),
|
|
|
)
|
|
)
|
|
|
- if (durable) {
|
|
|
|
|
- const committed = yield* commitSyncEvent(event as Payload, undefined, commit)
|
|
|
|
|
|
|
+ if (definition?.durable) {
|
|
|
|
|
+ const committed = yield* commitDurableEvent(event as Payload, undefined, commit)
|
|
|
if (committed) {
|
|
if (committed) {
|
|
|
- event = { ...event, seq: committed.seq }
|
|
|
|
|
- yield* Effect.forEach(syncHandlers, (sync) => observe(event as Payload, "sync", sync), { discard: true })
|
|
|
|
|
|
|
+ event = {
|
|
|
|
|
+ ...event,
|
|
|
|
|
+ durable: {
|
|
|
|
|
+ aggregateID: committed.aggregateID,
|
|
|
|
|
+ seq: committed.seq,
|
|
|
|
|
+ version: definition.durable.version,
|
|
|
|
|
+ },
|
|
|
|
|
+ }
|
|
|
yield* notify(event as Payload, true)
|
|
yield* notify(event as Payload, true)
|
|
|
return event
|
|
return event
|
|
|
}
|
|
}
|
|
@@ -406,12 +378,12 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
})
|
|
})
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- const observe = (event: Payload, kind: "sync" | "listener", observer: (event: Payload) => Effect.Effect<void>) =>
|
|
|
|
|
|
|
+ const observe = (event: Payload, observer: (event: Payload) => Effect.Effect<void>) =>
|
|
|
Effect.suspend(() => observer(event)).pipe(
|
|
Effect.suspend(() => observer(event)).pipe(
|
|
|
Effect.catchCauseIf(
|
|
Effect.catchCauseIf(
|
|
|
(cause) => !Cause.hasInterrupts(cause),
|
|
(cause) => !Cause.hasInterrupts(cause),
|
|
|
(cause) =>
|
|
(cause) =>
|
|
|
- Effect.logError("Event observer failed", { eventID: event.id, eventType: event.type, kind, cause }),
|
|
|
|
|
|
|
+ Effect.logError("Event listener failed", { eventID: event.id, eventType: event.type, cause }),
|
|
|
),
|
|
),
|
|
|
)
|
|
)
|
|
|
|
|
|
|
@@ -419,12 +391,12 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
return Effect.gen(function* () {
|
|
return Effect.gen(function* () {
|
|
|
yield* Effect.forEach(
|
|
yield* Effect.forEach(
|
|
|
listeners,
|
|
listeners,
|
|
|
- (listener) => (isolateListeners ? observe(event, "listener", listener) : listener(event)),
|
|
|
|
|
|
|
+ (listener) => (isolateListeners ? observe(event, listener) : listener(event)),
|
|
|
{ discard: true },
|
|
{ discard: true },
|
|
|
)
|
|
)
|
|
|
- const pubsub = typed.get(event.type)
|
|
|
|
|
- if (pubsub) yield* PubSub.publish(pubsub, event)
|
|
|
|
|
- yield* PubSub.publish(all, event)
|
|
|
|
|
|
|
+ const typed = pubsub.typed.get(event.type)
|
|
|
|
|
+ if (typed) yield* PubSub.publish(typed, event)
|
|
|
|
|
+ yield* PubSub.publish(pubsub.all, event)
|
|
|
})
|
|
})
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -441,7 +413,6 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
id: options?.id ?? ID.create(),
|
|
id: options?.id ?? ID.create(),
|
|
|
...(options?.metadata ? { metadata: options.metadata } : {}),
|
|
...(options?.metadata ? { metadata: options.metadata } : {}),
|
|
|
type: definition.type,
|
|
type: definition.type,
|
|
|
- ...(definition.sync === undefined ? {} : { version: definition.sync.version }),
|
|
|
|
|
...(location ? { location } : {}),
|
|
...(location ? { location } : {}),
|
|
|
data,
|
|
data,
|
|
|
} as Payload<D>,
|
|
} as Payload<D>,
|
|
@@ -455,27 +426,37 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
|
|
options?: { readonly publish?: boolean; readonly ownerID?: string; readonly strictOwner?: boolean },
|
|
|
) {
|
|
) {
|
|
|
return Effect.gen(function* () {
|
|
return Effect.gen(function* () {
|
|
|
- const definition = syncRegistry.get(event.type)
|
|
|
|
|
- if (!definition) {
|
|
|
|
|
|
|
+ const definition = durableRegistry.get(event.type)
|
|
|
|
|
+ if (!definition?.durable) {
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({ type: event.type, message: `Unknown sync event type ${event.type}` }),
|
|
|
|
|
|
|
+ new InvalidDurableEventError({ type: event.type, message: `Unknown durable event type ${event.type}` }),
|
|
|
)
|
|
)
|
|
|
} else {
|
|
} else {
|
|
|
const payload = {
|
|
const payload = {
|
|
|
id: event.id,
|
|
id: event.id,
|
|
|
type: definition.type,
|
|
type: definition.type,
|
|
|
- version: definition.sync.version,
|
|
|
|
|
- data: definition.decode(event.data),
|
|
|
|
|
- replay: true,
|
|
|
|
|
|
|
+ data: Schema.decodeUnknownSync(
|
|
|
|
|
+ definition.data as Schema.Codec<unknown, unknown, never, never>,
|
|
|
|
|
+ )(event.data),
|
|
|
} as Payload
|
|
} as Payload
|
|
|
- const committed = yield* commitSyncEvent(payload, {
|
|
|
|
|
|
|
+ const committed = yield* commitDurableEvent(payload, {
|
|
|
seq: event.seq,
|
|
seq: event.seq,
|
|
|
aggregateID: event.aggregateID,
|
|
aggregateID: event.aggregateID,
|
|
|
ownerID: options?.ownerID,
|
|
ownerID: options?.ownerID,
|
|
|
strictOwner: options?.strictOwner,
|
|
strictOwner: options?.strictOwner,
|
|
|
})
|
|
})
|
|
|
if (committed && options?.publish) {
|
|
if (committed && options?.publish) {
|
|
|
- yield* notify({ ...payload, seq: committed.seq }, true)
|
|
|
|
|
|
|
+ yield* notify(
|
|
|
|
|
+ {
|
|
|
|
|
+ ...payload,
|
|
|
|
|
+ durable: {
|
|
|
|
|
+ aggregateID: committed.aggregateID,
|
|
|
|
|
+ seq: committed.seq,
|
|
|
|
|
+ version: definition.durable.version,
|
|
|
|
|
+ },
|
|
|
|
|
+ },
|
|
|
|
|
+ true,
|
|
|
|
|
+ )
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
})
|
|
})
|
|
@@ -490,7 +471,7 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
if (!source) return undefined
|
|
if (!source) return undefined
|
|
|
if (events.some((event) => event.aggregateID !== source)) {
|
|
if (events.some((event) => event.aggregateID !== source)) {
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: events[0]?.type ?? "unknown",
|
|
type: events[0]?.type ?? "unknown",
|
|
|
message: "Replay events must belong to the same aggregate",
|
|
message: "Replay events must belong to the same aggregate",
|
|
|
}),
|
|
}),
|
|
@@ -501,7 +482,7 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
const seq = start + index
|
|
const seq = start + index
|
|
|
if (event.seq !== seq) {
|
|
if (event.seq !== seq) {
|
|
|
yield* Effect.die(
|
|
yield* Effect.die(
|
|
|
- new InvalidSyncEventError({
|
|
|
|
|
|
|
+ new InvalidDurableEventError({
|
|
|
type: event.type,
|
|
type: event.type,
|
|
|
message: `Replay sequence mismatch at index ${index}: expected ${seq}, got ${event.seq}`,
|
|
message: `Replay sequence mismatch at index ${index}: expected ${seq}, got ${event.seq}`,
|
|
|
}),
|
|
}),
|
|
@@ -540,22 +521,18 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
Stream.map((event) => event as Payload<D>),
|
|
Stream.map((event) => event as Payload<D>),
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
- const streamAll = (): Stream.Stream<Payload> => Stream.fromPubSub(all)
|
|
|
|
|
|
|
+ const streamAll = (): Stream.Stream<Payload> => Stream.fromPubSub(pubsub.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}` })
|
|
|
|
|
|
|
+ const decodeSerializedEvent = (event: SerializedEvent): Payload => {
|
|
|
|
|
+ const definition = durableRegistry.get(event.type)
|
|
|
|
|
+ if (!definition?.durable) {
|
|
|
|
|
+ throw new InvalidDurableEventError({ type: event.type, message: `Unknown durable event type ${event.type}` })
|
|
|
}
|
|
}
|
|
|
return {
|
|
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),
|
|
|
|
|
- },
|
|
|
|
|
|
|
+ id: event.id,
|
|
|
|
|
+ type: definition.type,
|
|
|
|
|
+ durable: { aggregateID: event.aggregateID, seq: event.seq, version: definition.durable.version },
|
|
|
|
|
+ data: Schema.decodeUnknownSync(definition.data as Schema.Codec<unknown, unknown, never, never>)(event.data),
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -583,43 +560,43 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
),
|
|
),
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
- const subscribeSynchronized = (aggregateID: string) =>
|
|
|
|
|
|
|
+ const subscribeDurable = (aggregateID: string) =>
|
|
|
Effect.gen(function* () {
|
|
Effect.gen(function* () {
|
|
|
- const pubsub = yield* PubSub.sliding<void>(1)
|
|
|
|
|
- const subscription = yield* PubSub.subscribe(pubsub)
|
|
|
|
|
|
|
+ const wake = yield* PubSub.sliding<void>(1)
|
|
|
|
|
+ const subscription = yield* PubSub.subscribe(wake)
|
|
|
yield* Effect.acquireRelease(
|
|
yield* Effect.acquireRelease(
|
|
|
Effect.sync(() => {
|
|
Effect.sync(() => {
|
|
|
- const pubsubs = synchronized.get(aggregateID) ?? new Set()
|
|
|
|
|
- pubsubs.add(pubsub)
|
|
|
|
|
- synchronized.set(aggregateID, pubsubs)
|
|
|
|
|
|
|
+ const wakes = pubsub.durable.get(aggregateID) ?? new Set()
|
|
|
|
|
+ wakes.add(wake)
|
|
|
|
|
+ pubsub.durable.set(aggregateID, wakes)
|
|
|
}),
|
|
}),
|
|
|
() =>
|
|
() =>
|
|
|
Effect.sync(() => {
|
|
Effect.sync(() => {
|
|
|
- const pubsubs = synchronized.get(aggregateID)
|
|
|
|
|
- pubsubs?.delete(pubsub)
|
|
|
|
|
- if (pubsubs?.size === 0) synchronized.delete(aggregateID)
|
|
|
|
|
- }).pipe(Effect.andThen(PubSub.shutdown(pubsub))),
|
|
|
|
|
|
|
+ const wakes = pubsub.durable.get(aggregateID)
|
|
|
|
|
+ wakes?.delete(wake)
|
|
|
|
|
+ if (wakes?.size === 0) pubsub.durable.delete(aggregateID)
|
|
|
|
|
+ }).pipe(Effect.andThen(PubSub.shutdown(wake))),
|
|
|
)
|
|
)
|
|
|
return subscription
|
|
return subscription
|
|
|
})
|
|
})
|
|
|
|
|
|
|
|
- const streamEvents = (input: {
|
|
|
|
|
|
|
+ const durable = (input: {
|
|
|
readonly aggregateID: string
|
|
readonly aggregateID: string
|
|
|
- readonly after?: Cursor
|
|
|
|
|
- }): Stream.Stream<CursorEvent> =>
|
|
|
|
|
|
|
+ readonly after?: number
|
|
|
|
|
+ }): Stream.Stream<Payload> =>
|
|
|
Stream.unwrap(
|
|
Stream.unwrap(
|
|
|
Effect.gen(function* () {
|
|
Effect.gen(function* () {
|
|
|
- const synchronized = yield* subscribeSynchronized(input.aggregateID)
|
|
|
|
|
- let cursor = input.after ?? -1
|
|
|
|
|
- const read = Effect.suspend(() => readAfter(input.aggregateID, cursor)).pipe(
|
|
|
|
|
|
|
+ const wakes = yield* subscribeDurable(input.aggregateID)
|
|
|
|
|
+ let sequence = input.after ?? -1
|
|
|
|
|
+ const read = Effect.suspend(() => readAfter(input.aggregateID, sequence)).pipe(
|
|
|
Effect.tap((events) =>
|
|
Effect.tap((events) =>
|
|
|
Effect.sync(() => {
|
|
Effect.sync(() => {
|
|
|
- cursor = events.at(-1)?.cursor ?? cursor
|
|
|
|
|
|
|
+ sequence = events.at(-1)?.durable?.seq ?? sequence
|
|
|
}),
|
|
}),
|
|
|
),
|
|
),
|
|
|
)
|
|
)
|
|
|
const historical = yield* read
|
|
const historical = yield* read
|
|
|
- const live = Stream.fromSubscription(synchronized).pipe(
|
|
|
|
|
|
|
+ const live = Stream.fromSubscription(wakes).pipe(
|
|
|
Stream.mapEffect(() => read),
|
|
Stream.mapEffect(() => read),
|
|
|
Stream.flattenIterable,
|
|
Stream.flattenIterable,
|
|
|
)
|
|
)
|
|
@@ -627,7 +604,7 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
}),
|
|
}),
|
|
|
)
|
|
)
|
|
|
|
|
|
|
|
- const listen = (listener: Listener): Effect.Effect<Unsubscribe> =>
|
|
|
|
|
|
|
+ const listen = (listener: Subscriber): Effect.Effect<Unsubscribe> =>
|
|
|
Effect.sync(() => {
|
|
Effect.sync(() => {
|
|
|
listeners.push(listener)
|
|
listeners.push(listener)
|
|
|
return Effect.sync(() => {
|
|
return Effect.sync(() => {
|
|
@@ -636,21 +613,7 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
})
|
|
})
|
|
|
})
|
|
})
|
|
|
|
|
|
|
|
- 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 project = <D extends Definition>(definition: D, projector: Projector<D>): Effect.Effect<void> =>
|
|
|
|
|
|
|
+ const project = <D extends Definition>(definition: D, projector: Subscriber<D>): Effect.Effect<void> =>
|
|
|
Effect.sync(() => {
|
|
Effect.sync(() => {
|
|
|
const list = projectors.get(definition.type) ?? []
|
|
const list = projectors.get(definition.type) ?? []
|
|
|
list.push((event) => projector(event as Payload<D>))
|
|
list.push((event) => projector(event as Payload<D>))
|
|
@@ -661,10 +624,8 @@ export const layerWith = (options?: LayerOptions) =>
|
|
|
publish,
|
|
publish,
|
|
|
subscribe,
|
|
subscribe,
|
|
|
all: streamAll,
|
|
all: streamAll,
|
|
|
- aggregateEvents: streamEvents,
|
|
|
|
|
- sync,
|
|
|
|
|
|
|
+ durable,
|
|
|
listen,
|
|
listen,
|
|
|
- beforeCommit,
|
|
|
|
|
project,
|
|
project,
|
|
|
replay,
|
|
replay,
|
|
|
replayAll,
|
|
replayAll,
|