diff --git a/apps/example/package.json b/apps/example/package.json index da266153..1c261484 100644 --- a/apps/example/package.json +++ b/apps/example/package.json @@ -38,6 +38,6 @@ "@types/pg": "^8.0.0", "@types/ramda": "^0.28.21", "tsup": "^7.2.0", - "vitest": "^0.34.4" + "vitest": "^0.34.5" } } diff --git a/package.json b/package.json index 524d512e..e5cbbf5e 100644 --- a/package.json +++ b/package.json @@ -18,7 +18,7 @@ "@vitest/coverage-v8": "^0.34.1", "prettier": "latest", "turbo": "latest", - "vitest": "^0.34.4" + "vitest": "^0.34.5" }, "engines": { "node": ">=14.0.0" diff --git a/packages/core/package.json b/packages/core/package.json index fdd79f88..1b0410d4 100644 --- a/packages/core/package.json +++ b/packages/core/package.json @@ -23,7 +23,7 @@ "lru-cache": "^7.14.1", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "^0.34.4", + "vitest": "^0.34.5", "zod": "^3.20.2" }, "dependencies": { diff --git a/packages/dynamo-store/cli/package.json b/packages/dynamo-store/cli/package.json index cde95e0a..30c58b46 100644 --- a/packages/dynamo-store/cli/package.json +++ b/packages/dynamo-store/cli/package.json @@ -8,7 +8,7 @@ "@types/node": "^18.11.18", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "^0.34.4" + "vitest": "^0.34.5" }, "dependencies": { "@aws-sdk/client-dynamodb": "^3.410.0", diff --git a/packages/dynamo-store/dynamo-store/package.json b/packages/dynamo-store/dynamo-store/package.json index d7d3eb58..184fa730 100644 --- a/packages/dynamo-store/dynamo-store/package.json +++ b/packages/dynamo-store/dynamo-store/package.json @@ -19,7 +19,7 @@ "@types/node": "^18.11.18", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "^0.34.4" + "vitest": "^0.34.5" }, "dependencies": { "@aws-sdk/client-dynamodb": "^3.410.0", diff --git a/packages/dynamo-store/indexer/package.json b/packages/dynamo-store/indexer/package.json index 60baf822..f6904a88 100644 --- a/packages/dynamo-store/indexer/package.json +++ b/packages/dynamo-store/indexer/package.json @@ -18,7 +18,7 @@ "@types/node": "^18.11.18", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "0.34.4" + "vitest": "0.34.5" }, "dependencies": { "@aws-sdk/client-dynamodb": "^3.410.0", diff --git a/packages/dynamo-store/source/package.json b/packages/dynamo-store/source/package.json index 33fcc036..38988023 100644 --- a/packages/dynamo-store/source/package.json +++ b/packages/dynamo-store/source/package.json @@ -19,7 +19,7 @@ "p-limit": "^4.0.0", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "^0.34.4" + "vitest": "^0.34.5" }, "dependencies": { "@equinox-js/core": "workspace:*", diff --git a/packages/dynamo-store/source/src/DynamoStoreSource.mts b/packages/dynamo-store/source/src/DynamoStoreSource.mts index ba09f109..d9d36e0d 100644 --- a/packages/dynamo-store/source/src/DynamoStoreSource.mts +++ b/packages/dynamo-store/source/src/DynamoStoreSource.mts @@ -305,14 +305,14 @@ export class DynamoStoreSource { private client: DynamoStoreSourceClient constructor( - private readonly index: AppendsIndex.Reader, - private epochs: AppendsEpoch.Reader.Service, + index: AppendsIndex.Reader, + epochs: AppendsEpoch.Reader.Service, private readonly options: Omit, ) { if (!this.options.categories && !this.options.streamFilter) { - throw new Error("Either categories or categoryFilter must be specified") + throw new Error("Either categories or streamFilter must be specified") } - const sink = new StreamsSink( + const sink = StreamsSink.create( options.handler, options.maxConcurrentStreams, options.maxConcurrentBatches, diff --git a/packages/memory-store/package.json b/packages/memory-store/package.json index 5af37690..d8bcc945 100644 --- a/packages/memory-store/package.json +++ b/packages/memory-store/package.json @@ -18,7 +18,7 @@ "@types/node": "^18.11.18", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "^0.34.4" + "vitest": "^0.34.5" }, "dependencies": { "@equinox-js/core": "workspace:*" diff --git a/packages/message-db/message-db-consumer/package.json b/packages/message-db/message-db-consumer/package.json index a5b57b0b..8d1d74d9 100644 --- a/packages/message-db/message-db-consumer/package.json +++ b/packages/message-db/message-db-consumer/package.json @@ -24,7 +24,7 @@ "@types/pg": "^8.6.6", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "^0.34.4" + "vitest": "^0.34.5" }, "dependencies": { "@equinox-js/core": "workspace:*", diff --git a/packages/message-db/message-db-consumer/src/lib/MessageDbSource.mts b/packages/message-db/message-db-consumer/src/lib/MessageDbSource.mts index e323ae65..b377c026 100644 --- a/packages/message-db/message-db-consumer/src/lib/MessageDbSource.mts +++ b/packages/message-db/message-db-consumer/src/lib/MessageDbSource.mts @@ -91,7 +91,7 @@ export class MessageDbSource { static create(options: Options & { pool: Pool }) { const client = new MessageDbCategoryReader(options.pool) - const sink = new StreamsSink( + const sink = StreamsSink.create( options.handler, options.maxConcurrentStreams, options.maxConcurrentBatches ?? 10, diff --git a/packages/projection-pg/package.json b/packages/projection-pg/package.json index 7e893e68..d3fbb141 100644 --- a/packages/projection-pg/package.json +++ b/packages/projection-pg/package.json @@ -23,7 +23,7 @@ "@types/pg": "^8.6.6", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "^0.34.4" + "vitest": "^0.34.5" }, "dependencies": { "pg": "^8.8.0" diff --git a/packages/propeller/package.json b/packages/propeller/package.json index 56273702..d52641d2 100644 --- a/packages/propeller/package.json +++ b/packages/propeller/package.json @@ -22,7 +22,7 @@ "p-queue": "^7.4.1", "tsup": "^7.2.0", "typescript": "^5.2.2", - "vitest": "^0.34.4" + "vitest": "^0.34.5" }, "dependencies": { "@equinox-js/core": "workspace:*", diff --git a/packages/propeller/src/Queue.test.ts b/packages/propeller/src/Queue.test.ts new file mode 100644 index 00000000..8af78b06 --- /dev/null +++ b/packages/propeller/src/Queue.test.ts @@ -0,0 +1,84 @@ +import { describe, test, expect } from "vitest" +import { Queue, AsyncQueue } from "./Queue.js" + +const toArray = (q: Queue) => { + const result: T[] = [] + let item: T | undefined + while ((item = q.tryGet())) { + result.push(item) + } + return result +} + +const ofArray = (arr: T[]) => { + const q = new Queue() + for (const x of arr) q.add(x) + return q +} + +describe("Queue", () => { + test("Seeking for the last item", () => { + const q = ofArray([1, 2, 3, 4, 5]) + expect(q.tryFind((x) => x === 5)).toEqual(5) + expect(toArray(q)).toEqual([1, 2, 3, 4]) + }) + + test("Seeking for a middle item", () => { + const q = ofArray([1, 2, 3, 4, 5]) + expect(q.tryFind((x) => x === 3)).toEqual(3) + expect(toArray(q)).toEqual([1, 2, 4, 5]) + }) + + test("Seeking for the head item", () => { + const q = ofArray([1, 2, 3, 4, 5]) + expect(q.tryFind((x) => x === 1)).toEqual(1) + expect(toArray(q)).toEqual([2, 3, 4, 5]) + }) + test("Seeking for a non-existent item", () => { + const q = ofArray([1, 2, 3, 4, 5]) + expect(q.tryFind((x) => x === 0)).toEqual(undefined) + expect(toArray(q)).toEqual([1, 2, 3, 4, 5]) + }) +}) + +describe("AsyncQueue", () => { + const true_ = () => true + const ofArray = (arr: T[]) => { + const q = new AsyncQueue() + for (const x of arr) q.add(x) + return q + } + const toArray = async (q: AsyncQueue, signal: AbortSignal) => { + const result: T[] = [] + while (q.size) { + result.push(await q.tryFindAsync(true_, signal)) + } + return result + } + test("Waiting on an item", async () => { + const ctrl = new AbortController() + const q = new AsyncQueue() + const p = q.tryFindAsync(true_, ctrl.signal) + q.add(1) + expect(await p).toEqual(1) + }) + + test("When an item exists it is returned raw", () => { + const ctrl = new AbortController() + const q = new AsyncQueue() + q.add(1) + expect(q.tryFindAsync(true_, ctrl.signal)).toEqual(1) + }) + + test("Seeking an item resolves when a matching item is enqueued", async () => { + const ctrl = new AbortController() + const q = new AsyncQueue() + const p = q.tryFindAsync((x) => x === 3, ctrl.signal) + q.add(1) + q.add(2) + q.add(3) + q.add(4) + expect(await p).toEqual(3) + expect(await toArray(q, ctrl.signal)).toEqual([1, 2, 4]) + }) +}) diff --git a/packages/propeller/src/Queue.ts b/packages/propeller/src/Queue.ts index 888ac086..b181f0d1 100644 --- a/packages/propeller/src/Queue.ts +++ b/packages/propeller/src/Queue.ts @@ -35,15 +35,48 @@ export class Queue { return value } } + + // prev must not meet the predicate + private seek(predicate: (x: T) => boolean, prev: Node) { + let curr = prev.next + + while (curr && !predicate(curr.value)) { + prev = curr + curr = curr.next + } + if (curr) { + prev.next = curr.next + if (this.firstAndLast && this.firstAndLast[1] === curr) this.firstAndLast[1] = prev + return curr.value + } + } + + tryFind(predicate: (x: T) => boolean) { + if (this.firstAndLast) { + const head: Node | undefined = this.firstAndLast[0] + const value = head.value + if (!predicate(value)) return this.seek(predicate, head) + if (this.firstAndLast[0].next) { + this.firstAndLast = [this.firstAndLast[0].next, this.firstAndLast[1]] + } else { + delete this.firstAndLast + } + --this.size + return value + } + } } export class AsyncQueue { private queue = new Queue() - private pendingGets = new Queue<(value: T) => void>() + private pendingGets = new Queue<{ + predicate: (value: T) => boolean + resolve: (value: T) => void + }>() add(value: T) { - const send = this.pendingGets.tryGet() - if (send) return send(value) + const pending = this.pendingGets.tryFind((x) => x.predicate(value)) + if (pending) return pending.resolve(value) this.queue.add(value) } @@ -51,16 +84,19 @@ export class AsyncQueue { return this.queue.size } - tryGetAsync(signal: AbortSignal) { + tryFindAsync(predicate: (x: T) => boolean, signal: AbortSignal): Promise | T { + const value = this.queue.tryFind(predicate) + if (value) return value return new Promise((resolve, reject) => { - const value = this.queue.tryGet() - if (value) return resolve(value) const abort = () => reject(new Error("Aborted")) if (signal.aborted) return abort() signal.addEventListener("abort", abort) - this.pendingGets.add((value) => { - signal.removeEventListener("abort", abort) - resolve(value) + this.pendingGets.add({ + predicate, + resolve(value) { + signal.removeEventListener("abort", abort) + resolve(value) + }, }) }) } diff --git a/packages/propeller/src/StreamsSink.mts b/packages/propeller/src/StreamsSink.mts index 05753ad2..cbdc8674 100644 --- a/packages/propeller/src/StreamsSink.mts +++ b/packages/propeller/src/StreamsSink.mts @@ -21,16 +21,83 @@ class Stream { } } +class QueueWorker { + private stopped = true + constructor( + private readonly tryGetNext: (signal: AbortSignal) => Promise, + private readonly handle: EventHandler, + private readonly tracingAttrs: Attributes, + private readonly onHandled: (stream: Stream) => void, + ) {} + + stop() { + this.stopped = true + } + + async run(signal: AbortSignal): Promise { + this.stopped = false + while (!signal.aborted && !this.stopped) { + const stream = await this.tryGetNext(signal) + await traceHandler(tracer, this.tracingAttrs, stream.name, stream.events, this.handle) + this.onHandled(stream) + } + } +} + +export class BatchLimiter { + private inProgressBatches = 0 + private onReady?: () => void + private onEmpty?: () => void + constructor(private maxConcurrentBatches: number) {} + + waitForCapacity(signal: AbortSignal) { + if (this.inProgressBatches < this.maxConcurrentBatches) return + return new Promise((resolve, reject) => { + const abort = () => { + signal.removeEventListener("abort", abort) + reject(new Error("Aborted")) + } + signal.addEventListener("abort", abort) + this.onReady = () => { + signal.removeEventListener("abort", abort) + resolve() + } + }) + } + + waitForEmpty() { + if (this.inProgressBatches === 0) return + return new Promise((resolve) => { + this.onEmpty = () => { + resolve() + } + }) + } + + startBatch() { + this.inProgressBatches++ + } + + endBatch() { + this.inProgressBatches-- + this.onReady?.() + delete this.onReady + if (this.inProgressBatches === 0) { + this.onEmpty?.() + delete this.onEmpty + } + } +} + export class StreamsSink implements Sink { private queue = new AsyncQueue() private streams = new Map() + private activeStreams = new Set() private batchStreams = new Map>() - private onReady?: () => void - private inProgressBatches = 0 constructor( private readonly handle: EventHandler, private maxConcurrentStreams: number, - private maxConcurrentBatches: number, + private readonly batchLimiter: BatchLimiter, private readonly tracingAttrs: Attributes = {}, ) { this.addTracingAttrs({ "eqx.max_concurrent_streams": maxConcurrentStreams }) @@ -40,7 +107,21 @@ export class StreamsSink implements Sink { Object.assign(this.tracingAttrs, attrs) } - start(signal: AbortSignal) { + private handleStreamCompletion(stream: Stream) { + this.activeStreams.delete(stream.name) + const batches = Array.from(this.batchStreams) + for (const [batch, streams] of this.batchStreams) { + streams.delete(stream) + if (streams.size === 0 && batches[0][0] === batch) { + batches.shift() + batch.onComplete() + this.batchStreams.delete(batch) + this.batchLimiter.endBatch() + } + } + } + + async start(signal: AbortSignal) { // set up a linked abort controller to avoid attaching too many event listeners to the provided signal // node has a default limit of 11 before it emits a warning to the console. We have entirely legitimate reasons // for attaching more than 11 listeners @@ -48,51 +129,36 @@ export class StreamsSink implements Sink { signal.addEventListener("abort", () => ctrl.abort()) const _signal = ctrl.signal // Starts N workers that will process streams in parallel - return new Promise((_resolve, reject) => { - let stopped = false - const aux = async (): Promise => { - if (_signal.aborted || stopped) return - try { - const stream = await this.queue.tryGetAsync(_signal) - this.streams.delete(stream.name) - await traceHandler(tracer, this.tracingAttrs, stream.name, stream.events, this.handle) - for (const [batch, streams] of this.batchStreams) { - streams.delete(stream) - if (streams.size === 0) { - batch.onComplete() - this.batchStreams.delete(batch) - this.inProgressBatches-- - this.onReady?.() - delete this.onReady - } - } - void aux() - } catch (err) { - stopped = true - reject(err) - } - } - for (let i = 0; i < this.maxConcurrentStreams; ++i) aux() - }) - } + const workers = new Array(this.maxConcurrentStreams) - private waitForCapacity(signal: AbortSignal) { - if (this.inProgressBatches < this.maxConcurrentBatches) return - return new Promise((resolve, reject) => { - const abort = () => reject(new Error("Aborted")) - signal.addEventListener("abort", abort) - this.onReady = () => { - signal.removeEventListener("abort", abort) - resolve() - } - }) + for (let i = 0; i < this.maxConcurrentStreams; ++i) { + workers[i] = new QueueWorker( + async (signal) => { + const stream = await this.queue.tryFindAsync( + (stream) => !this.activeStreams.has(stream.name), + signal, + ) + this.activeStreams.add(stream.name) + this.streams.delete(stream.name) + return stream + }, + this.handle, + this.tracingAttrs, + (stream) => this.handleStreamCompletion(stream), + ) + } + + try { + await Promise.all(workers.map((w) => w.run(_signal))) + } catch (err) { + workers.forEach((w) => w.stop()) + throw err + } } - pump(batch: IngesterBatch, signal: AbortSignal): Promise | void { - const p = this.waitForCapacity(signal) - if (p != null) return p.then(() => this.pump(batch, signal)) - this.inProgressBatches++ + async pump(batch: IngesterBatch, signal: AbortSignal): Promise { + this.batchLimiter.startBatch() const streamsInBatch = new Set() this.batchStreams.set(batch, streamsInBatch) for (const [sn, event] of batch.items) { @@ -105,5 +171,16 @@ export class StreamsSink implements Sink { this.queue.add(stream) } } + await this.batchLimiter.waitForCapacity(signal) + } + + static create( + handle: EventHandler, + maxConcurrentStreams: number, + maxConcurrentBatches: number, + tracingAttrs: Attributes = {}, + ) { + const limiter = new BatchLimiter(maxConcurrentBatches) + return new StreamsSink(handle, maxConcurrentStreams, limiter, tracingAttrs) } } diff --git a/packages/propeller/src/StreamsSink.test.mts b/packages/propeller/src/StreamsSink.test.mts index 4c710db8..dc435b22 100644 --- a/packages/propeller/src/StreamsSink.test.mts +++ b/packages/propeller/src/StreamsSink.test.mts @@ -1,5 +1,5 @@ import { describe, test, expect, vi } from "vitest" -import { StreamsSink } from "./StreamsSink.mjs" +import { BatchLimiter, StreamsSink } from "./StreamsSink.mjs" import { ITimelineEvent, StreamId, StreamName } from "@equinox-js/core" import { sleep } from "./Sleep.js" import { IngesterBatch } from "./Types.js" @@ -36,7 +36,8 @@ describe("Concurrency", () => { await new Promise(setImmediate) active.delete(stream) } - const sink = new StreamsSink(handler, concurrency, 1) + const limiter = new BatchLimiter(1) + const sink = new StreamsSink(handler, concurrency, limiter) sink.start(ctrl.signal).catch(() => {}) @@ -45,7 +46,7 @@ describe("Concurrency", () => { ctrl.signal, ) - await new Promise(setImmediate) + await limiter.waitForEmpty() ctrl.abort() expect(maxActive).toBe(concurrency) @@ -60,7 +61,8 @@ test("Correctly merges batches", async () => { invocations++ await new Promise(setImmediate) } - const sink = new StreamsSink(handler, 10, 3) + const limiter = new BatchLimiter(3) + const sink = new StreamsSink(handler, 10, limiter) const checkpoint = vi.fn().mockResolvedValue(undefined) @@ -70,8 +72,7 @@ test("Correctly merges batches", async () => { sink.start(ctrl.signal).catch(() => {}) - // sleep triggers after setImmediate - await sleep(0, ctrl.signal) + await limiter.waitForEmpty() ctrl.abort() expect(invocations).toBe(10) @@ -88,31 +89,75 @@ const mkSingleBatch = ( onComplete: onComplete(checkpoint), }) -test(" Correctly limits in-flight batches", async () => { +test("Correctly limits in-flight batches", async () => { let invocations = 0 const ctrl = new AbortController() async function handler() { invocations++ await new Promise((res) => setTimeout(res, 10)) } - const sink = new StreamsSink(handler, 1, 3) + const limiter = new BatchLimiter(3) + const sink = new StreamsSink(handler, 1, limiter) sink.start(ctrl.signal).catch(() => {}) const completed = vi.fn().mockResolvedValue(undefined) const complete = (n: bigint) => () => completed(n) + // First batch will be immediately picked up + await sink.pump(mkSingleBatch(complete, 0n), ctrl.signal) + // meanwhile we merge two batches together + await sink.pump(mkSingleBatch(complete, 1n), ctrl.signal) + await sink.pump(mkSingleBatch(complete, 2n), ctrl.signal) + + // This batch will be immediately picked up + await sink.pump(mkSingleBatch(complete, 3n), ctrl.signal) + // meanwhile we merge two batches together + await sink.pump(mkSingleBatch(complete, 4n), ctrl.signal) + await sink.pump(mkSingleBatch(complete, 5n), ctrl.signal) + + await limiter.waitForEmpty() + ctrl.abort() + + expect(invocations).toBe(4) + + // onComplete is called in order and for every batch + expect(completed.mock.calls).toEqual([[0n], [1n], [2n], [3n], [4n], [5n]]) +}) + +test("Ensures at-most one handler is per stream", async () => { + let active = 0 + let maxActive = 0 + let invocations = 0 + const ctrl = new AbortController() + async function handler() { + invocations++ + active++ + maxActive = Math.max(maxActive, active) + await new Promise((res) => setTimeout(res, 10)) + active-- + } + const limiter = new BatchLimiter(10) + const sink = new StreamsSink(handler, 100, limiter) + + sink.start(ctrl.signal).catch(() => {}) + + const completed = vi.fn().mockResolvedValue(undefined) + const complete = (n: bigint) => () => completed(n) + await sink.pump(mkSingleBatch(complete, 0n), ctrl.signal) await sink.pump(mkSingleBatch(complete, 1n), ctrl.signal) await sink.pump(mkSingleBatch(complete, 2n), ctrl.signal) await sink.pump(mkSingleBatch(complete, 3n), ctrl.signal) await sink.pump(mkSingleBatch(complete, 4n), ctrl.signal) + await sink.pump(mkSingleBatch(complete, 5n), ctrl.signal) - await sleep(10, ctrl.signal) + await limiter.waitForEmpty() ctrl.abort() - expect(invocations).toBe(3) + // 1 invocation for the first pump, the other 5 are merged + expect(invocations).toBe(2) + expect(maxActive).toBe(1) // onComplete is called in order and for every batch - expect(completed).toHaveBeenCalledTimes(5) - expect(completed.mock.calls).toEqual([[0n], [1n], [2n], [3n], [4n]]) + expect(completed.mock.calls).toEqual([[0n], [1n], [2n], [3n], [4n], [5n]]) }) diff --git a/packages/propeller/src/index.mts b/packages/propeller/src/index.mts index 8f981457..6dadc9f7 100644 --- a/packages/propeller/src/index.mts +++ b/packages/propeller/src/index.mts @@ -1,5 +1,5 @@ export * from "./Checkpoints.js" export * from "./Tracing.js" -export * from "./StreamsSink.mjs" +export { StreamsSink } from "./StreamsSink.mjs" export * from "./Types.js" export * from "./FeedSource.mjs" diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index e96c5416..7d7d9521 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -17,16 +17,16 @@ importers: version: link:packages/dynamo-store/cli '@vitest/coverage-v8': specifier: ^0.34.1 - version: 0.34.1(vitest@0.34.4) + version: 0.34.1(vitest@0.34.5) prettier: specifier: latest version: 3.0.3 turbo: specifier: latest - version: 1.10.13 + version: 1.10.14 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 apps/example: dependencies: @@ -113,8 +113,8 @@ importers: specifier: ^7.2.0 version: 7.2.0(supports-color@9.4.0)(typescript@5.2.2) vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 docs: dependencies: @@ -181,8 +181,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 zod: specifier: ^3.20.2 version: 3.20.2 @@ -237,8 +237,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5(supports-color@9.4.0) packages/dynamo-store/dynamo-store: dependencies: @@ -271,8 +271,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 packages/dynamo-store/indexer: dependencies: @@ -308,8 +308,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: 0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: 0.34.5 + version: 0.34.5 packages/dynamo-store/lambda: dependencies: @@ -376,8 +376,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 packages/memory-store: dependencies: @@ -395,8 +395,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 packages/message-db/message-db: dependencies: @@ -460,8 +460,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 packages/projection-pg: dependencies: @@ -482,8 +482,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 packages/propeller: dependencies: @@ -507,8 +507,8 @@ importers: specifier: ^5.2.2 version: 5.2.2 vitest: - specifier: ^0.34.4 - version: 0.34.4(supports-color@9.4.0) + specifier: ^0.34.5 + version: 0.34.5 packages/test-domain: dependencies: @@ -5943,6 +5943,10 @@ packages: resolution: {integrity: sha512-Y+/1vGBHV/cYk6OI1Na/LHzwnlNCAfU3ZNGrc1LdRe/LAIbdDPTTv/HU3M7yXN448aTVDq3eKRm2cg7iKLb8gw==} dev: true + /@types/node@20.6.3: + resolution: {integrity: sha512-HksnYH4Ljr4VQgEy2lTStbCKv/P590tmPe5HqOnv9Gprffgv5WXAY+Y5Gqniu0GGqeTCUdBnzC3QSrzPkBkAMA==} + dev: true + /@types/parse-json@4.0.0: resolution: {integrity: sha512-//oorEZjL6sbPcKUaCdIGlIUeH26mgzimjBB77G6XRgnDl/L5wOnpyBGRe/Mmf5CVW3PwEBE1NjiMZ/ssFh4wA==} dev: false @@ -6082,7 +6086,7 @@ packages: '@types/yargs-parser': 21.0.0 dev: false - /@vitest/coverage-v8@0.34.1(vitest@0.34.4): + /@vitest/coverage-v8@0.34.1(vitest@0.34.5): resolution: {integrity: sha512-lRgUwjTMr8idXEbUPSNH4jjRZJXJCVY3BqUa+LDXyJVe3pldxYMn/r0HMqatKUGTp0Kyf1j5LfFoY6kRqRp7jw==} peerDependencies: vitest: '>=0.32.0 <1' @@ -6098,47 +6102,47 @@ packages: std-env: 3.3.3 test-exclude: 6.0.0 v8-to-istanbul: 9.1.0 - vitest: 0.34.4(supports-color@9.4.0) + vitest: 0.34.5 transitivePeerDependencies: - supports-color dev: true - /@vitest/expect@0.34.4: - resolution: {integrity: sha512-XlMKX8HyYUqB8dsY8Xxrc64J2Qs9pKMt2Z8vFTL4mBWXJsg4yoALHzJfDWi8h5nkO4Zua4zjqtapQ/IluVkSnA==} + /@vitest/expect@0.34.5: + resolution: {integrity: sha512-/3RBIV9XEH+nRpRMqDJBufKIOQaYUH2X6bt0rKSCW0MfKhXFLYsR5ivHifeajRSTsln0FwJbitxLKHSQz/Xwkw==} dependencies: - '@vitest/spy': 0.34.4 - '@vitest/utils': 0.34.4 + '@vitest/spy': 0.34.5 + '@vitest/utils': 0.34.5 chai: 4.3.8 dev: true - /@vitest/runner@0.34.4: - resolution: {integrity: sha512-hwwdB1StERqUls8oV8YcpmTIpVeJMe4WgYuDongVzixl5hlYLT2G8afhcdADeDeqCaAmZcSgLTLtqkjPQF7x+w==} + /@vitest/runner@0.34.5: + resolution: {integrity: sha512-RDEE3ViVvl7jFSCbnBRyYuu23XxmvRTSZWW6W4M7eC5dOsK75d5LIf6uhE5Fqf809DQ1+9ICZZNxhIolWHU4og==} dependencies: - '@vitest/utils': 0.34.4 + '@vitest/utils': 0.34.5 p-limit: 4.0.0 pathe: 1.1.1 dev: true - /@vitest/snapshot@0.34.4: - resolution: {integrity: sha512-GCsh4coc3YUSL/o+BPUo7lHQbzpdttTxL6f4q0jRx2qVGoYz/cyTRDJHbnwks6TILi6560bVWoBpYC10PuTLHw==} + /@vitest/snapshot@0.34.5: + resolution: {integrity: sha512-+ikwSbhu6z2yOdtKmk/aeoDZ9QPm2g/ZO5rXT58RR9Vmu/kB2MamyDSx77dctqdZfP3Diqv4mbc/yw2kPT8rmA==} dependencies: magic-string: 0.30.3 pathe: 1.1.1 - pretty-format: 29.6.3 + pretty-format: 29.7.0 dev: true - /@vitest/spy@0.34.4: - resolution: {integrity: sha512-PNU+fd7DUPgA3Ya924b1qKuQkonAW6hL7YUjkON3wmBwSTIlhOSpy04SJ0NrRsEbrXgMMj6Morh04BMf8k+w0g==} + /@vitest/spy@0.34.5: + resolution: {integrity: sha512-epsicsfhvBjRjCMOC/3k00mP/TBGQy8/P0DxOFiWyLt55gnZ99dqCfCiAsKO17BWVjn4eZRIjKvcqNmSz8gvmg==} dependencies: tinyspy: 2.1.1 dev: true - /@vitest/utils@0.34.4: - resolution: {integrity: sha512-yR2+5CHhp/K4ySY0Qtd+CAL9f5Yh1aXrKfAT42bq6CtlGPh92jIDDDSg7ydlRow1CP+dys4TrOrbELOyNInHSg==} + /@vitest/utils@0.34.5: + resolution: {integrity: sha512-ur6CmmYQoeHMwmGb0v+qwkwN3yopZuZyf4xt1DBBSGBed8Hf9Gmbm/5dEWqgpLPdRx6Av6jcWXrjcKfkTzg/pw==} dependencies: diff-sequences: 29.6.3 loupe: 2.3.6 - pretty-format: 29.6.3 + pretty-format: 29.7.0 dev: true /@webassemblyjs/ast@1.11.1: @@ -10728,6 +10732,16 @@ packages: nanoid: 3.3.6 picocolors: 1.0.0 source-map-js: 1.0.2 + dev: false + + /postcss@8.4.30: + resolution: {integrity: sha512-7ZEao1g4kd68l97aWG/etQKPKq07us0ieSZ2TnFDk11i0ZfDW2AwKHYU8qv4MZKqN2fdBfg+7q0ES06UA73C1g==} + engines: {node: ^10 || ^12 || >=14} + dependencies: + nanoid: 3.3.6 + picocolors: 1.0.0 + source-map-js: 1.0.2 + dev: true /postgres-array@2.0.0: resolution: {integrity: sha512-VpZrUqU5A69eQyW2c5CA1jtLecCsN2U/bD6VilrFDWq5+5UIEVO7nazS3TEcHf1zuPYO/sqGvUvW62g86RXZuA==} @@ -10765,8 +10779,8 @@ packages: renderkid: 3.0.0 dev: false - /pretty-format@29.6.3: - resolution: {integrity: sha512-ZsBgjVhFAj5KeK+nHfF1305/By3lechHQSMWCTl8iHSbfOm2TN5nHEtFc/+W7fAyUeCs2n5iow72gld4gW0xDw==} + /pretty-format@29.7.0: + resolution: {integrity: sha512-Pdlw/oPxN+aXdmM9R00JVC9WVFoCLTKJvDVLgmJ+qAffBMxsV85l/Lu7sNx4zSzPyoL2euImuEwHhOXdEgNFZQ==} engines: {node: ^14.15.0 || ^16.10.0 || >=18.0.0} dependencies: '@jest/schemas': 29.6.3 @@ -11379,6 +11393,14 @@ packages: fsevents: 2.3.3 dev: true + /rollup@3.29.2: + resolution: {integrity: sha512-CJouHoZ27v6siztc21eEQGo0kIcE5D1gVPA571ez0mMYb25LGYGKnVNXpEj5MGlepmDWGXNjDB5q7uNiPHC11A==} + engines: {node: '>=14.18.0', npm: '>=8.0.0'} + hasBin: true + optionalDependencies: + fsevents: 2.3.3 + dev: true + /rtl-detect@1.0.4: resolution: {integrity: sha512-EBR4I2VDSSYr7PkBmFy04uhycIpDKp+21p/jARYXlCSjQksTBQcJ0HFUPOO79EPPH5JS6VAhiIQbycf0O3JAxQ==} dev: false @@ -12057,8 +12079,8 @@ packages: resolution: {integrity: sha512-lBN9zLN/oAf68o3zNXYrdCt1kP8WsiGW8Oo2ka41b2IM5JL/S1CTyX1rW0mb/zSuJun0ZUrDxx4sqvYS2FWzPA==} dev: false - /tinybench@2.5.0: - resolution: {integrity: sha512-kRwSG8Zx4tjF9ZiyH4bhaebu+EDz1BOx9hOigYHlUW4xxI/wKIUQUqo018UlU4ar6ATPBsaMrdbKZ+tmPdohFA==} + /tinybench@2.5.1: + resolution: {integrity: sha512-65NKvSuAVDP/n4CqH+a9w2kTlLReS9vhsAP06MWx+/89nMinJyB2icyl58RIcqCmIggpojIGeuJGhjU1aGMBSg==} dev: true /tinypool@0.7.0: @@ -12180,64 +12202,64 @@ packages: - ts-node dev: true - /turbo-darwin-64@1.10.13: - resolution: {integrity: sha512-vmngGfa2dlYvX7UFVncsNDMuT4X2KPyPJ2Jj+xvf5nvQnZR/3IeDEGleGVuMi/hRzdinoxwXqgk9flEmAYp0Xw==} + /turbo-darwin-64@1.10.14: + resolution: {integrity: sha512-I8RtFk1b9UILAExPdG/XRgGQz95nmXPE7OiGb6ytjtNIR5/UZBS/xVX/7HYpCdmfriKdVwBKhalCoV4oDvAGEg==} cpu: [x64] os: [darwin] requiresBuild: true dev: true optional: true - /turbo-darwin-arm64@1.10.13: - resolution: {integrity: sha512-eMoJC+k7gIS4i2qL6rKmrIQGP6Wr9nN4odzzgHFngLTMimok2cGLK3qbJs5O5F/XAtEeRAmuxeRnzQwTl/iuAw==} + /turbo-darwin-arm64@1.10.14: + resolution: {integrity: sha512-KAdUWryJi/XX7OD0alOuOa0aJ5TLyd4DNIYkHPHYcM6/d7YAovYvxRNwmx9iv6Vx6IkzTnLeTiUB8zy69QkG9Q==} cpu: [arm64] os: [darwin] requiresBuild: true dev: true optional: true - /turbo-linux-64@1.10.13: - resolution: {integrity: sha512-0CyYmnKTs6kcx7+JRH3nPEqCnzWduM0hj8GP/aodhaIkLNSAGAa+RiYZz6C7IXN+xUVh5rrWTnU2f1SkIy7Gdg==} + /turbo-linux-64@1.10.14: + resolution: {integrity: sha512-BOBzoREC2u4Vgpap/WDxM6wETVqVMRcM8OZw4hWzqCj2bqbQ6L0wxs1LCLWVrghQf93JBQtIGAdFFLyCSBXjWQ==} cpu: [x64] os: [linux] requiresBuild: true dev: true optional: true - /turbo-linux-arm64@1.10.13: - resolution: {integrity: sha512-0iBKviSGQQlh2OjZgBsGjkPXoxvRIxrrLLbLObwJo3sOjIH0loGmVIimGS5E323soMfi/o+sidjk2wU1kFfD7Q==} + /turbo-linux-arm64@1.10.14: + resolution: {integrity: sha512-D8T6XxoTdN5D4V5qE2VZG+/lbZX/89BkAEHzXcsSUTRjrwfMepT3d2z8aT6hxv4yu8EDdooZq/2Bn/vjMI32xw==} cpu: [arm64] os: [linux] requiresBuild: true dev: true optional: true - /turbo-windows-64@1.10.13: - resolution: {integrity: sha512-S5XySRfW2AmnTeY1IT+Jdr6Goq7mxWganVFfrmqU+qqq3Om/nr0GkcUX+KTIo9mPrN0D3p5QViBRzulwB5iuUQ==} + /turbo-windows-64@1.10.14: + resolution: {integrity: sha512-zKNS3c1w4i6432N0cexZ20r/aIhV62g69opUn82FLVs/zk3Ie0GVkSB6h0rqIvMalCp7enIR87LkPSDGz9K4UA==} cpu: [x64] os: [win32] requiresBuild: true dev: true optional: true - /turbo-windows-arm64@1.10.13: - resolution: {integrity: sha512-nKol6+CyiExJIuoIc3exUQPIBjP9nIq5SkMJgJuxsot2hkgGrafAg/izVDRDrRduQcXj2s8LdtxJHvvnbI8hEQ==} + /turbo-windows-arm64@1.10.14: + resolution: {integrity: sha512-rkBwrTPTxNSOUF7of8eVvvM+BkfkhA2OvpHM94if8tVsU+khrjglilp8MTVPHlyS9byfemPAmFN90oRIPB05BA==} cpu: [arm64] os: [win32] requiresBuild: true dev: true optional: true - /turbo@1.10.13: - resolution: {integrity: sha512-vOF5IPytgQPIsgGtT0n2uGZizR2N3kKuPIn4b5p5DdeLoI0BV7uNiydT7eSzdkPRpdXNnO8UwS658VaI4+YSzQ==} + /turbo@1.10.14: + resolution: {integrity: sha512-hr9wDNYcsee+vLkCDIm8qTtwhJ6+UAMJc3nIY6+PNgUTtXcQgHxCq8BGoL7gbABvNWv76CNbK5qL4Lp9G3ZYRA==} hasBin: true optionalDependencies: - turbo-darwin-64: 1.10.13 - turbo-darwin-arm64: 1.10.13 - turbo-linux-64: 1.10.13 - turbo-linux-arm64: 1.10.13 - turbo-windows-64: 1.10.13 - turbo-windows-arm64: 1.10.13 + turbo-darwin-64: 1.10.14 + turbo-darwin-arm64: 1.10.14 + turbo-linux-64: 1.10.14 + turbo-linux-arm64: 1.10.14 + turbo-windows-64: 1.10.14 + turbo-windows-arm64: 1.10.14 dev: true /type-detect@4.0.8: @@ -12561,8 +12583,8 @@ packages: vfile-message: 2.0.4 dev: false - /vite-node@0.34.4(@types/node@18.11.18)(supports-color@9.4.0): - resolution: {integrity: sha512-ho8HtiLc+nsmbwZMw8SlghESEE3KxJNp04F/jPUCLVvaURwt0d+r9LxEqCX5hvrrOQ0GSyxbYr5ZfRYhQ0yVKQ==} + /vite-node@0.34.5(@types/node@20.6.2)(supports-color@9.4.0): + resolution: {integrity: sha512-RNZ+DwbCvDoI5CbCSQSyRyzDTfFvFauvMs6Yq4ObJROKlIKuat1KgSX/Ako5rlDMfVCyMcpMRMTkJBxd6z8YRA==} engines: {node: '>=v14.18.0'} hasBin: true dependencies: @@ -12571,7 +12593,7 @@ packages: mlly: 1.4.2 pathe: 1.1.1 picocolors: 1.0.0 - vite: 4.4.9(@types/node@18.11.18) + vite: 4.4.9(@types/node@20.6.2) transitivePeerDependencies: - '@types/node' - less @@ -12583,7 +12605,29 @@ packages: - terser dev: true - /vite@4.4.9(@types/node@18.11.18): + /vite-node@0.34.5(@types/node@20.6.3): + resolution: {integrity: sha512-RNZ+DwbCvDoI5CbCSQSyRyzDTfFvFauvMs6Yq4ObJROKlIKuat1KgSX/Ako5rlDMfVCyMcpMRMTkJBxd6z8YRA==} + engines: {node: '>=v14.18.0'} + hasBin: true + dependencies: + cac: 6.7.14 + debug: 4.3.4(supports-color@9.4.0) + mlly: 1.4.2 + pathe: 1.1.1 + picocolors: 1.0.0 + vite: 4.4.9(@types/node@20.6.3) + transitivePeerDependencies: + - '@types/node' + - less + - lightningcss + - sass + - stylus + - sugarss + - supports-color + - terser + dev: true + + /vite@4.4.9(@types/node@20.6.2): resolution: {integrity: sha512-2mbUn2LlUmNASWwSCNSJ/EG2HuSRTnVNaydp6vMCm5VIqJsjMfbIWtbH2kDuwUVW5mMUKKZvGPX/rqeqVvv1XA==} engines: {node: ^14.18.0 || >=16.0.0} hasBin: true @@ -12611,16 +12655,52 @@ packages: terser: optional: true dependencies: - '@types/node': 18.11.18 + '@types/node': 20.6.2 esbuild: 0.18.20 - postcss: 8.4.29 - rollup: 3.29.0 + postcss: 8.4.30 + rollup: 3.29.2 optionalDependencies: fsevents: 2.3.3 dev: true - /vitest@0.34.4(supports-color@9.4.0): - resolution: {integrity: sha512-SE/laOsB6995QlbSE6BtkpXDeVNLJc1u2LHRG/OpnN4RsRzM3GQm4nm3PQCK5OBtrsUqnhzLdnT7se3aeNGdlw==} + /vite@4.4.9(@types/node@20.6.3): + resolution: {integrity: sha512-2mbUn2LlUmNASWwSCNSJ/EG2HuSRTnVNaydp6vMCm5VIqJsjMfbIWtbH2kDuwUVW5mMUKKZvGPX/rqeqVvv1XA==} + engines: {node: ^14.18.0 || >=16.0.0} + hasBin: true + peerDependencies: + '@types/node': '>= 14' + less: '*' + lightningcss: ^1.21.0 + sass: '*' + stylus: '*' + sugarss: '*' + terser: ^5.4.0 + peerDependenciesMeta: + '@types/node': + optional: true + less: + optional: true + lightningcss: + optional: true + sass: + optional: true + stylus: + optional: true + sugarss: + optional: true + terser: + optional: true + dependencies: + '@types/node': 20.6.3 + esbuild: 0.18.20 + postcss: 8.4.30 + rollup: 3.29.2 + optionalDependencies: + fsevents: 2.3.3 + dev: true + + /vitest@0.34.5: + resolution: {integrity: sha512-CPI68mmnr2DThSB3frSuE5RLm9wo5wU4fbDrDwWQQB1CWgq9jQVoQwnQSzYAjdoBOPoH2UtXpOgHVge/uScfZg==} engines: {node: '>=v14.18.0'} hasBin: true peerDependencies: @@ -12652,12 +12732,77 @@ packages: dependencies: '@types/chai': 4.3.6 '@types/chai-subset': 1.3.3 - '@types/node': 18.11.18 - '@vitest/expect': 0.34.4 - '@vitest/runner': 0.34.4 - '@vitest/snapshot': 0.34.4 - '@vitest/spy': 0.34.4 - '@vitest/utils': 0.34.4 + '@types/node': 20.6.3 + '@vitest/expect': 0.34.5 + '@vitest/runner': 0.34.5 + '@vitest/snapshot': 0.34.5 + '@vitest/spy': 0.34.5 + '@vitest/utils': 0.34.5 + acorn: 8.10.0 + acorn-walk: 8.2.0 + cac: 6.7.14 + chai: 4.3.8 + debug: 4.3.4(supports-color@9.4.0) + local-pkg: 0.4.3 + magic-string: 0.30.3 + pathe: 1.1.1 + picocolors: 1.0.0 + std-env: 3.4.3 + strip-literal: 1.3.0 + tinybench: 2.5.1 + tinypool: 0.7.0 + vite: 4.4.9(@types/node@20.6.3) + vite-node: 0.34.5(@types/node@20.6.3) + why-is-node-running: 2.2.2 + transitivePeerDependencies: + - less + - lightningcss + - sass + - stylus + - sugarss + - supports-color + - terser + dev: true + + /vitest@0.34.5(supports-color@9.4.0): + resolution: {integrity: sha512-CPI68mmnr2DThSB3frSuE5RLm9wo5wU4fbDrDwWQQB1CWgq9jQVoQwnQSzYAjdoBOPoH2UtXpOgHVge/uScfZg==} + engines: {node: '>=v14.18.0'} + hasBin: true + peerDependencies: + '@edge-runtime/vm': '*' + '@vitest/browser': '*' + '@vitest/ui': '*' + happy-dom: '*' + jsdom: '*' + playwright: '*' + safaridriver: '*' + webdriverio: '*' + peerDependenciesMeta: + '@edge-runtime/vm': + optional: true + '@vitest/browser': + optional: true + '@vitest/ui': + optional: true + happy-dom: + optional: true + jsdom: + optional: true + playwright: + optional: true + safaridriver: + optional: true + webdriverio: + optional: true + dependencies: + '@types/chai': 4.3.6 + '@types/chai-subset': 1.3.3 + '@types/node': 20.6.2 + '@vitest/expect': 0.34.5 + '@vitest/runner': 0.34.5 + '@vitest/snapshot': 0.34.5 + '@vitest/spy': 0.34.5 + '@vitest/utils': 0.34.5 acorn: 8.10.0 acorn-walk: 8.2.0 cac: 6.7.14 @@ -12669,10 +12814,10 @@ packages: picocolors: 1.0.0 std-env: 3.4.3 strip-literal: 1.3.0 - tinybench: 2.5.0 + tinybench: 2.5.1 tinypool: 0.7.0 - vite: 4.4.9(@types/node@18.11.18) - vite-node: 0.34.4(@types/node@18.11.18)(supports-color@9.4.0) + vite: 4.4.9(@types/node@20.6.2) + vite-node: 0.34.5(@types/node@20.6.2)(supports-color@9.4.0) why-is-node-running: 2.2.2 transitivePeerDependencies: - less