refactor(api): remove watermark sync surface
This commit is contained in:
parent
4cee5bd824
commit
ff499c3603
25 changed files with 141 additions and 392 deletions
|
|
@ -135,12 +135,6 @@ export interface Interface {
|
|||
readonly after?: number
|
||||
readonly follow?: boolean
|
||||
}) => Stream.Stream<LogItem>
|
||||
/**
|
||||
* Coalescing hint channel: latest committed seq per aggregate, never a
|
||||
* delivery guarantee. Emits `SweepRequired` first on every subscribe and
|
||||
* whenever per-key retention is exceeded. Never fails under backpressure.
|
||||
*/
|
||||
readonly changes: () => Stream.Stream<EventLog.Change>
|
||||
/** Latest committed seq per aggregate. Aggregates without events are absent. */
|
||||
readonly sequences: (aggregateIDs: ReadonlyArray<string>) => Effect.Effect<ReadonlyMap<string, Seq>>
|
||||
/** @deprecated Use `all()` and consume the returned stream. */
|
||||
|
|
@ -183,11 +177,6 @@ export const liveBounded = (
|
|||
|
||||
export interface LayerOptions {
|
||||
readonly beforeAggregateRead?: (aggregateID: string) => Effect.Effect<void>
|
||||
/**
|
||||
* Maximum distinct aggregates buffered per changes subscriber before the
|
||||
* buffer is abandoned and the subscriber is told to sweep.
|
||||
*/
|
||||
readonly changesKeyCapacity?: number
|
||||
/** Maximum durable rows read per page while replaying or tailing an aggregate log. */
|
||||
readonly logReadPageSize?: number
|
||||
}
|
||||
|
|
@ -204,12 +193,6 @@ export const layerWith = (options?: LayerOptions) =>
|
|||
const projectors = new Map<string, Subscriber[]>()
|
||||
// TODO: Bind durable projectors to exact type+version before supporting incompatible historical payloads.
|
||||
const listeners = new Array<Subscriber>()
|
||||
const changesKeyCapacity = options?.changesKeyCapacity ?? 4096
|
||||
const changesSubscribers = new Set<{
|
||||
readonly hints: Map<string, number>
|
||||
sweepRequired: boolean
|
||||
readonly wake: PubSub.PubSub<void>
|
||||
}>()
|
||||
const { db } = yield* Database.Service
|
||||
const logReadPageSize = options?.logReadPageSize ?? 512
|
||||
|
||||
|
|
@ -231,9 +214,6 @@ export const layerWith = (options?: LayerOptions) =>
|
|||
{ discard: true },
|
||||
)
|
||||
yield* Effect.forEach(pubsub.typed.values(), PubSub.shutdown, { discard: true })
|
||||
yield* Effect.forEach(changesSubscribers, (subscriber) => PubSub.shutdown(subscriber.wake), {
|
||||
discard: true,
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
||||
|
|
@ -394,27 +374,6 @@ export const layerWith = (options?: LayerOptions) =>
|
|||
(wake) => PubSub.publish(wake, undefined),
|
||||
{ discard: true },
|
||||
)
|
||||
yield* Effect.forEach(
|
||||
changesSubscribers,
|
||||
(subscriber) =>
|
||||
Effect.sync(() => {
|
||||
// Coalesce to the latest seq per aggregate. Overflowing key
|
||||
// cardinality abandons the buffer instead of dropping hints silently.
|
||||
if (
|
||||
subscriber.hints.size >= changesKeyCapacity &&
|
||||
!subscriber.hints.has(committed.aggregateID)
|
||||
) {
|
||||
subscriber.hints.clear()
|
||||
subscriber.sweepRequired = true
|
||||
} else if (!subscriber.sweepRequired) {
|
||||
subscriber.hints.set(
|
||||
committed.aggregateID,
|
||||
Math.max(subscriber.hints.get(committed.aggregateID) ?? -1, committed.seq),
|
||||
)
|
||||
}
|
||||
}).pipe(Effect.andThen(PubSub.publish(subscriber.wake, undefined)), Effect.asVoid),
|
||||
{ discard: true },
|
||||
)
|
||||
}
|
||||
return committed
|
||||
}),
|
||||
|
|
@ -723,47 +682,6 @@ export const layerWith = (options?: LayerOptions) =>
|
|||
}),
|
||||
)
|
||||
|
||||
const changes = (): Stream.Stream<EventLog.Change> =>
|
||||
Stream.unwrap(
|
||||
Effect.gen(function* () {
|
||||
const wake = yield* PubSub.sliding<void>(1)
|
||||
const subscription = yield* PubSub.subscribe(wake)
|
||||
const subscriber = { hints: new Map<string, number>(), sweepRequired: false, wake }
|
||||
yield* Effect.acquireRelease(
|
||||
Effect.sync(() => changesSubscribers.add(subscriber)),
|
||||
() =>
|
||||
Effect.sync(() => changesSubscribers.delete(subscriber)).pipe(
|
||||
Effect.andThen(PubSub.shutdown(wake)),
|
||||
Effect.asVoid,
|
||||
),
|
||||
)
|
||||
const drain = Effect.sync((): ReadonlyArray<EventLog.Change> => {
|
||||
if (subscriber.sweepRequired) {
|
||||
subscriber.sweepRequired = false
|
||||
subscriber.hints.clear()
|
||||
return [{ type: "log.sweep_required" }]
|
||||
}
|
||||
const hints = Array.from(
|
||||
subscriber.hints,
|
||||
([aggregateID, seq]): EventLog.Change => ({ type: "log.hint", aggregateID, seq: Seq.make(seq) }),
|
||||
)
|
||||
subscriber.hints.clear()
|
||||
return hints
|
||||
})
|
||||
// Hints missed while unsubscribed were never buffered, so every
|
||||
// (re)subscribe starts from the sweep contract.
|
||||
const initial: EventLog.Change = { type: "log.sweep_required" }
|
||||
return Stream.make(initial).pipe(
|
||||
Stream.concat(
|
||||
Stream.fromSubscription(subscription).pipe(
|
||||
Stream.mapEffect(() => drain),
|
||||
Stream.flattenIterable,
|
||||
),
|
||||
),
|
||||
)
|
||||
}),
|
||||
)
|
||||
|
||||
const sequences = (aggregateIDs: ReadonlyArray<string>): Effect.Effect<ReadonlyMap<string, Seq>> => {
|
||||
if (aggregateIDs.length === 0) return Effect.succeed(new Map())
|
||||
return db
|
||||
|
|
@ -798,7 +716,6 @@ export const layerWith = (options?: LayerOptions) =>
|
|||
subscribe,
|
||||
live: streamLive,
|
||||
log,
|
||||
changes,
|
||||
sequences,
|
||||
listen,
|
||||
project,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue