diff --git a/CHANGELOG.md b/CHANGELOG.md index 8c66c20..9e31c38 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -6,6 +6,11 @@ The format follows Keep a Changelog and the package uses semantic versioning. ## Unreleased +- Applied every supported lower and upper bound from a conjunctive range query + before loading secondary-index candidates. +- Bounded parallel secondary-index uploads to three operations, overlapped + independent immutable Snapshot or Trie writes, and kept HEAD as the final + awaited compare-and-swap. - Started the mutable-HEAD freshness interval when a successful fetch or conditional revalidation completes, preventing slow 304 responses from arriving already expired. diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index dda019e..8e5289e 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -110,9 +110,11 @@ private. gzip-compressed when useful, and encrypted. 5. Every configured secondary index and declared covering projection is updated or rebuilt. -6. New immutable document and index objects are created. -7. HEAD publishes the document root and all active index references with one - ETag compare-and-swap. +6. New immutable document and index objects are created. Secondary-index + uploads use a concurrency bound of three and overlap independent Snapshot + or Trie object writes. +7. After every immutable upload completes, HEAD publishes the document root + and all active index references with one ETag compare-and-swap. 8. The response includes the new HEAD, changed immutable objects, and document. 9. The writing tab updates its cache and broadcasts the bundle to other tabs. diff --git a/docs/BENCHMARKS.md b/docs/BENCHMARKS.md index f9c1268..a7875c4 100644 --- a/docs/BENCHMARKS.md +++ b/docs/BENCHMARKS.md @@ -183,10 +183,13 @@ The range returned 25 documents from covering index fields. | Medium | 473.98 ms | 1,612.96 ms | 510.99 ms | 1,919.19 ms | 1,500 | | Large | 1,044.55 ms | 2,115.30 ms | 1,026.00 ms | 2,018.80 ms | 7,500 | -Current range planning selected every value above the lower bound before -applying the upper bound, so the 25-document large result evaluated 7,500 -index candidates. Covering fields still kept the operation to two network -reads and avoided loading document pages. +At the benchmarked commit, range planning selected every value above the +lower bound before applying the upper bound, so the 25-document large result +evaluated 7,500 index candidates. The current planner applies both bounds and +the equivalent regression test selects 25 candidates. The latency table +remains the historical measurement and has not been relabelled as a +post-correction benchmark. Covering fields kept the measured operation to two +network reads and avoided loading document pages. ## Full scans diff --git a/docs/QUERIES-INDEXES.md b/docs/QUERIES-INDEXES.md index 956c31e..e712911 100644 --- a/docs/QUERIES-INDEXES.md +++ b/docs/QUERIES-INDEXES.md @@ -189,7 +189,9 @@ defineIndex( ``` Range indexes contain exactly one field. Equality indexes can contain up to -four fields. +four fields. For conjunctive range queries, every supported `eq`, `lt`, `lte`, +`gt`, and `gte` comparison on the indexed field is applied before document +candidates are loaded. Only scalar string, number, boolean, or null values are indexed. Arrays and objects remain available to bounded local filtering. diff --git a/src/engines/content-trie.ts b/src/engines/content-trie.ts index 25f795b..472d1a7 100644 --- a/src/engines/content-trie.ts +++ b/src/engines/content-trie.ts @@ -9,11 +9,13 @@ import { type StoredObject, } from "../core.js"; import { + createAsyncOperationLimiter, createDictionary, decodeJson, encodeJson, isPreconditionFailure, ownValue, + type AsyncOperationLimiter, validateName, } from "../shared-utils.js"; import { @@ -67,6 +69,8 @@ type PreparedSecondaryIndex = { reference: SecondaryIndexReference; }; +const MAX_PARALLEL_INDEX_WRITES = 3; + export class ContentAddressedTrieEngine implements DatabaseEngine { readonly name = "content-addressed-trie"; private casRetries = 0; @@ -547,108 +551,29 @@ export class ContentAddressedTrieEngine implements DatabaseEngine { head.state, collapsedChanges, ); - const root = - head.state.rootHash === null - ? this.emptyRoot() - : await this.readNode( - normalized, - head.state.rootHash, - "root", - ); - const nextRoot: TrieRootNode = { - kind: "root", - children: { ...root.children }, - }; - - const byBranch = groupUpdates(updates); - const changedBranches = await Promise.all( - [...byBranch].map(async ([first, byLeaf]) => { - const currentBranchHash = root.children[first]; - const currentBranch = currentBranchHash - ? await this.readNode( - normalized, - currentBranchHash, - "branch", - ) - : this.emptyBranch(); - const nextBranch: TrieBranchNode = { - kind: "branch", - children: { ...currentBranch.children }, - leafMetadata: createDictionary( - currentBranch.leafMetadata, - ), - }; - - const changedLeaves = await Promise.all( - [...byLeaf].map(async ([second, leafUpdates]) => { - const currentLeafHash = - currentBranch.children[second]; - const currentLeaf = currentLeafHash - ? await this.readNode( - normalized, - currentLeafHash, - "leaf", - ) - : this.emptyLeaf(); - const nextLeaf: TrieLeafNode = { - kind: "leaf", - documents: createDictionary( - currentLeaf.documents, - ), - }; - for (const update of leafUpdates) { - if (update.document === null) { - delete nextLeaf.documents[update.id]; - } else { - nextLeaf.documents[update.id] = update.document; - } - } - if (Object.keys(nextLeaf.documents).length === 0) { - return [second, null] as const; - } - return [ - second, - await this.writeLeafNode(normalized, nextLeaf), - ] as const; - }), - ); - for (const [second, leafResult] of changedLeaves) { - if (leafResult === null) { - delete nextBranch.children[second]; - delete nextBranch.leafMetadata?.[second]; - } else { - nextBranch.children[second] = leafResult.hash; - nextBranch.leafMetadata ??= - createDictionary(); - nextBranch.leafMetadata[second] = - leafResult.metadata; - } - } - if (Object.keys(nextBranch.children).length === 0) { - return [first, null] as const; - } - - return [ - first, - await this.writeNode(normalized, nextBranch), - ] as const; - }), + const limitIndexWrite = createAsyncOperationLimiter( + MAX_PARALLEL_INDEX_WRITES, ); - for (const [first, branchHash] of changedBranches) { - if (branchHash === null) { - delete nextRoot.children[first]; - } else { - nextRoot.children[first] = branchHash; - } + const [indexCommit, treeCommit] = + await Promise.allSettled([ + this.commitIndexes( + preparedIndexes, + limitIndexWrite, + ), + this.commitTreeChanges( + normalized, + head.state, + updates, + ), + ]); + if (indexCommit.status === "rejected") { + throw indexCommit.reason; } - - const rootHash = - Object.keys(nextRoot.children).length === 0 - ? null - : await this.writeNode(normalized, nextRoot); - const indexes = await this.commitIndexes( - preparedIndexes, - ); + if (treeCommit.status === "rejected") { + throw treeCommit.reason; + } + const indexes = indexCommit.value; + const rootHash = treeCommit.value; const nextHead: TrieHead = { revision: head.state.revision + 1, rootHash, @@ -685,6 +610,121 @@ export class ContentAddressedTrieEngine implements DatabaseEngine { ); } + private async commitTreeChanges( + collection: string, + head: TrieHead, + updates: TrieUpdate[], + ): Promise { + const root = + head.rootHash === null + ? this.emptyRoot() + : await this.readNode( + collection, + head.rootHash, + "root", + ); + const nextRoot: TrieRootNode = { + kind: "root", + children: { ...root.children }, + }; + + const byBranch = groupUpdates(updates); + const changedBranches = await Promise.all( + [...byBranch].map(async ([first, byLeaf]) => { + const currentBranchHash = root.children[first]; + const currentBranch = currentBranchHash + ? await this.readNode( + collection, + currentBranchHash, + "branch", + ) + : this.emptyBranch(); + const nextBranch: TrieBranchNode = { + kind: "branch", + children: { ...currentBranch.children }, + leafMetadata: createDictionary( + currentBranch.leafMetadata, + ), + }; + + const changedLeaves = await Promise.all( + [...byLeaf].map(async ([second, leafUpdates]) => { + const currentLeafHash = + currentBranch.children[second]; + const currentLeaf = currentLeafHash + ? await this.readNode( + collection, + currentLeafHash, + "leaf", + ) + : this.emptyLeaf(); + const nextLeaf: TrieLeafNode = { + kind: "leaf", + documents: createDictionary( + currentLeaf.documents, + ), + }; + for (const update of leafUpdates) { + if (update.document === null) { + delete nextLeaf.documents[update.id]; + } else { + nextLeaf.documents[update.id] = + update.document; + } + } + if (Object.keys(nextLeaf.documents).length === 0) { + return [second, null] as const; + } + return [ + second, + await this.writeLeafNode( + collection, + nextLeaf, + ), + ] as const; + }), + ); + for (const [second, leafResult] of changedLeaves) { + if (leafResult === null) { + delete nextBranch.children[second]; + delete nextBranch.leafMetadata?.[second]; + } else { + nextBranch.children[second] = leafResult.hash; + nextBranch.leafMetadata ??= + createDictionary(); + nextBranch.leafMetadata[second] = + leafResult.metadata; + } + } + if (Object.keys(nextBranch.children).length === 0) { + return [first, null] as const; + } + + return [ + first, + await this.writeNode( + collection, + nextBranch, + ), + ] as const; + }), + ); + for (const [first, branchHash] of changedBranches) { + if (branchHash === null) { + delete nextRoot.children[first]; + } else { + nextRoot.children[first] = branchHash; + } + } + + return Object.keys(nextRoot.children).length === 0 + ? null + : this.writeNode( + collection, + nextRoot, + ); + } + async compact(collection: string): Promise { if (!this.allowQuiescentGarbageCollection) { return; @@ -1016,22 +1056,46 @@ export class ContentAddressedTrieEngine implements DatabaseEngine { private async commitIndexes( prepared: PreparedSecondaryIndex[], + limitWrite: AsyncOperationLimiter = + createAsyncOperationLimiter( + MAX_PARALLEL_INDEX_WRITES, + ), ): Promise { const references = createDictionary(); - for (const index of prepared) { - try { - await this.store.put( - index.key, - index.bytes, - { ifNoneMatch: true }, - ); - } catch (error) { - if (!isPreconditionFailure(error)) { - throw error; - } + const committed = await Promise.allSettled( + prepared.map((index) => + limitWrite(async () => { + try { + await this.store.put( + index.key, + index.bytes, + { ifNoneMatch: true }, + ); + } catch (error) { + if (!isPreconditionFailure(error)) { + throw error; + } + } + return index; + }), + ), + ); + const failure = committed.find( + ( + result, + ): result is PromiseRejectedResult => + result.status === "rejected", + ); + if (failure) { + throw failure.reason; + } + for (const result of committed) { + if (result.status === "fulfilled") { + const index = result.value; + references[index.name] = + index.reference; } - references[index.name] = index.reference; } return references; } @@ -1250,9 +1314,13 @@ export class ContentAddressedTrieEngine implements DatabaseEngine { const bytes = encodeJson(node as unknown as JsonValue); const hash = await this.addressNode(bytes); try { - await this.store.put(this.nodeKey(collection, hash), bytes, { - ifNoneMatch: true, - }); + await this.store.put( + this.nodeKey(collection, hash), + bytes, + { + ifNoneMatch: true, + }, + ); this.nodesCreated += 1; } catch (error) { if (!isPreconditionFailure(error)) { @@ -1273,9 +1341,13 @@ export class ContentAddressedTrieEngine implements DatabaseEngine { const bytes = encodeJson(leaf as unknown as JsonValue); const hash = await this.addressNode(bytes); try { - await this.store.put(this.nodeKey(collection, hash), bytes, { - ifNoneMatch: true, - }); + await this.store.put( + this.nodeKey(collection, hash), + bytes, + { + ifNoneMatch: true, + }, + ); this.nodesCreated += 1; } catch (error) { if (!isPreconditionFailure(error)) { diff --git a/src/engines/immutable-snapshot.ts b/src/engines/immutable-snapshot.ts index 83de84b..81208a0 100644 --- a/src/engines/immutable-snapshot.ts +++ b/src/engines/immutable-snapshot.ts @@ -9,11 +9,13 @@ import { type StoredObject, } from "../core.js"; import { + createAsyncOperationLimiter, createDictionary, decodeJson, encodeJson, isPreconditionFailure, ownValue, + type AsyncOperationLimiter, validateName, } from "../shared-utils.js"; import { @@ -54,6 +56,8 @@ type PreparedSecondaryIndex = { reference: SecondaryIndexReference; }; +const MAX_IMMUTABLE_WRITE_CONCURRENCY = 3; + export class ImmutableSnapshotEngine implements DatabaseEngine { readonly name = "immutable-snapshot"; private casRetries = 0; @@ -487,13 +491,32 @@ export class ImmutableSnapshotEngine implements DatabaseEngine { const snapshotHash = Object.keys(documents).length === 0 ? null - : await this.writeSnapshot( - normalized, - pageBytes, - ); - const indexes = await this.commitIndexes( - preparedIndexes, + : await this.addressSnapshot(pageBytes); + const limitWrite = createAsyncOperationLimiter( + MAX_IMMUTABLE_WRITE_CONCURRENCY, ); + const [snapshotCommit, indexCommit] = + await Promise.allSettled([ + snapshotHash === null + ? Promise.resolve() + : this.commitSnapshot( + normalized, + snapshotHash, + pageBytes, + limitWrite, + ), + this.commitIndexes( + preparedIndexes, + limitWrite, + ), + ]); + if (snapshotCommit.status === "rejected") { + throw snapshotCommit.reason; + } + if (indexCommit.status === "rejected") { + throw indexCommit.reason; + } + const indexes = indexCommit.value; const nextHead: SnapshotHead = { revision: loaded.head.state.revision + 1, snapshotHash, @@ -567,23 +590,29 @@ export class ImmutableSnapshotEngine implements DatabaseEngine { }; } - private async writeSnapshot( + private async commitSnapshot( collection: string, + hash: string, bytes: Uint8Array, - ): Promise { - const hash = await this.addressSnapshot(bytes); - try { - await this.store.put(snapshotPageKey(collection, hash), bytes, { - ifNoneMatch: true, - }); - this.snapshotsCreated += 1; - } catch (error) { - if (!isPreconditionFailure(error)) { - throw error; + limitWrite: AsyncOperationLimiter, + ): Promise { + await limitWrite(async () => { + try { + await this.store.put( + snapshotPageKey(collection, hash), + bytes, + { + ifNoneMatch: true, + }, + ); + this.snapshotsCreated += 1; + } catch (error) { + if (!isPreconditionFailure(error)) { + throw error; + } + this.reusedSnapshots += 1; } - this.reusedSnapshots += 1; - } - return hash; + }); } private async prepareIndexes( @@ -619,22 +648,46 @@ export class ImmutableSnapshotEngine implements DatabaseEngine { private async commitIndexes( prepared: PreparedSecondaryIndex[], + limitWrite: AsyncOperationLimiter = + createAsyncOperationLimiter( + MAX_IMMUTABLE_WRITE_CONCURRENCY, + ), ): Promise { const references = createDictionary(); - for (const index of prepared) { - try { - await this.store.put( - index.key, - index.bytes, - { ifNoneMatch: true }, - ); - } catch (error) { - if (!isPreconditionFailure(error)) { - throw error; - } + const committed = await Promise.allSettled( + prepared.map((index) => + limitWrite(async () => { + try { + await this.store.put( + index.key, + index.bytes, + { ifNoneMatch: true }, + ); + } catch (error) { + if (!isPreconditionFailure(error)) { + throw error; + } + } + return index; + }), + ), + ); + const failure = committed.find( + ( + result, + ): result is PromiseRejectedResult => + result.status === "rejected", + ); + if (failure) { + throw failure.reason; + } + for (const result of committed) { + if (result.status === "fulfilled") { + const index = result.value; + references[index.name] = + index.reference; } - references[index.name] = index.reference; } return references; } diff --git a/src/secondary-index.ts b/src/secondary-index.ts index 0335e29..185e786 100644 --- a/src/secondary-index.ts +++ b/src/secondary-index.ts @@ -300,7 +300,7 @@ export function planSecondaryIndex( } continue; } - const comparison = comparisons.find( + const matched = comparisons.filter( (candidate) => candidate.field === definition.fields[0] && ["eq", "lt", "lte", "gt", "gte"].includes( @@ -308,10 +308,10 @@ export function planSecondaryIndex( ) && isJsonPrimitive(candidate.value), ); - if (comparison) { + if (matched.length > 0) { return { definition, - comparisons: [comparison], + comparisons: matched, }; } } @@ -346,16 +346,21 @@ export function idsFromSecondaryIndex( ?.ids ?? [] ); } - const comparison = plan.comparisons[0]!; - if (!isJsonPrimitive(comparison.value)) { + if ( + !plan.comparisons.every((comparison) => + isJsonPrimitive(comparison.value), + ) + ) { return []; } return page.entries .filter((entry) => - compareIndexedValue( - entry.values[0], - comparison.operator, - comparison.value as JsonPrimitive, + plan.comparisons.every((comparison) => + compareIndexedValue( + entry.values[0], + comparison.operator, + comparison.value as JsonPrimitive, + ), ), ) .flatMap((entry) => entry.ids); diff --git a/src/shared-utils.ts b/src/shared-utils.ts index 97389a2..36a6d07 100644 --- a/src/shared-utils.ts +++ b/src/shared-utils.ts @@ -73,3 +73,53 @@ export function ownValue( ? dictionary[key] : undefined; } + +export type AsyncOperationLimiter = ( + operation: () => Promise, +) => Promise; + +export function createAsyncOperationLimiter( + maximumConcurrency: number, +): AsyncOperationLimiter { + if ( + !Number.isInteger(maximumConcurrency) || + maximumConcurrency < 1 + ) { + throw new Error( + "Maximum concurrency must be a positive integer", + ); + } + + let active = 0; + const waiting: Array<() => void> = []; + + const acquire = async () => { + if (active < maximumConcurrency) { + active += 1; + return; + } + await new Promise((resolve) => { + waiting.push(resolve); + }); + }; + + const release = () => { + const next = waiting.shift(); + if (next) { + next(); + } else { + active -= 1; + } + }; + + return async ( + operation: () => Promise, + ): Promise => { + await acquire(); + try { + return await operation(); + } finally { + release(); + } + }; +} diff --git a/tests/parallel-write-pipeline.test.ts b/tests/parallel-write-pipeline.test.ts new file mode 100644 index 0000000..d104750 --- /dev/null +++ b/tests/parallel-write-pipeline.test.ts @@ -0,0 +1,307 @@ +import { describe, expect, it } from "vitest"; +import { + PreconditionFailedError, + type JsonDocument, + type ObjectStore, + type PutConditions, + type StoredObject, +} from "../src/core.js"; +import { + ContentAddressedTrieEngine, +} from "../src/engines/content-trie.js"; +import { + ImmutableSnapshotEngine, +} from "../src/engines/immutable-snapshot.js"; +import type { + CollectionIndexConfiguration, +} from "../src/secondary-index.js"; +import { + createAsyncOperationLimiter, +} from "../src/shared-utils.js"; +import { + snapshotHeadKey, +} from "../src/snapshot-protocol.js"; +import { + trieHeadKey, +} from "../src/trie-protocol.js"; + +const indexes: CollectionIndexConfiguration = { + notes: [ + { + name: "by-category", + fields: ["category"], + mode: "equality", + include: ["title", "lastModified"], + }, + { + name: "by-last-modified", + fields: ["lastModified"], + mode: "range", + include: ["title", "category"], + }, + ], +}; + +describe("bounded parallel write pipeline", () => { + it("limits an async operation group to three in flight", async () => { + const limit = createAsyncOperationLimiter(3); + let active = 0; + let maximumActive = 0; + + await Promise.all( + Array.from({ length: 8 }, () => + limit(async () => { + active += 1; + maximumActive = Math.max( + maximumActive, + active, + ); + await delay(5); + active -= 1; + }), + ), + ); + + expect(maximumActive).toBe(3); + }); + + it.each(["snapshot", "trie"] as const)( + "preserves %s protocol objects while overlapping immutable writes", + async (layout) => { + const serialStore = new DelayedStore(true); + const parallelStore = new DelayedStore(false); + const serial = engine(layout, serialStore); + const parallel = engine(layout, parallelStore); + const initial = documents(64); + await serial.putMany("notes", initial); + await parallel.putMany("notes", initial); + serialStore.resetMetrics(); + parallelStore.resetMetrics(); + + const replacement: JsonDocument = { + ...initial[17]!, + category: "updated-category", + body: "updated body", + lastModified: 999, + }; + await serial.put( + "notes", + replacement.id, + replacement, + ); + await parallel.put( + "notes", + replacement.id, + replacement, + ); + + expect( + await parallel.scan("notes"), + ).toEqual(await serial.scan("notes")); + expect( + decodedObjects(parallelStore.objects), + ).toEqual(decodedObjects(serialStore.objects)); + expect(parallelStore.writes).toBe( + serialStore.writes, + ); + expect(parallelStore.writtenBytes).toBe( + serialStore.writtenBytes, + ); + expect( + parallelStore.maximumConcurrentWrites, + ).toBeGreaterThan(1); + expect( + parallelStore.maximumConcurrentWrites, + ).toBeLessThanOrEqual(3); + expect( + serialStore.maximumConcurrentWrites, + ).toBe(1); + expect( + parallelStore.headStartedBeforeImmutableWritesSettled, + ).toBe(false); + expect( + parallelStore.completedWrites.at(-1), + ).toBe( + layout === "snapshot" + ? snapshotHeadKey("notes") + : trieHeadKey("notes"), + ); + }, + ); +}); + +function engine( + layout: "snapshot" | "trie", + store: ObjectStore, +) { + return layout === "snapshot" + ? new ImmutableSnapshotEngine( + store, + 40, + undefined, + false, + indexes, + ) + : new ContentAddressedTrieEngine( + store, + 40, + undefined, + false, + indexes, + ); +} + +function documents(count: number): JsonDocument[] { + return Array.from( + { length: count }, + (_, index) => ({ + id: `note-${String(index).padStart(6, "0")}`, + title: `Note ${index}`, + category: + `category-${String(index % 20).padStart(2, "0")}`, + body: `representative content ${index}`, + lastModified: index, + }), + ); +} + +function decodedObjects( + objects: Map, +) { + return [...objects.entries()] + .map(([key, object]) => ({ + key, + bytes: [...object.bytes], + value: JSON.parse( + new TextDecoder().decode(object.bytes), + ) as unknown, + })) + .sort((left, right) => + left.key.localeCompare(right.key), + ); +} + +class DelayedStore implements ObjectStore { + readonly objects = new Map(); + readonly completedWrites: string[] = []; + private tail = Promise.resolve(); + private etag = 0; + private activeWrites = 0; + private activeImmutableWrites = 0; + writes = 0; + writtenBytes = 0; + maximumConcurrentWrites = 0; + headStartedBeforeImmutableWritesSettled = false; + + constructor( + private readonly serializeWrites: boolean, + ) {} + + async get(key: string) { + await delay(1); + const object = this.objects.get(key); + return object + ? { + bytes: object.bytes.slice(), + etag: object.etag, + } + : null; + } + + put( + key: string, + bytes: Uint8Array, + conditions: PutConditions = {}, + ) { + const operation = () => + this.write(key, bytes, conditions); + if (!this.serializeWrites) { + return operation(); + } + const result = this.tail.then( + operation, + operation, + ); + this.tail = result.then( + () => undefined, + () => undefined, + ); + return result; + } + + delete(key: string) { + this.objects.delete(key); + return Promise.resolve(); + } + + list(prefix: string) { + return Promise.resolve( + [...this.objects.keys()] + .filter((key) => key.startsWith(prefix)) + .sort(), + ); + } + + resetMetrics() { + this.writes = 0; + this.writtenBytes = 0; + this.maximumConcurrentWrites = 0; + this.headStartedBeforeImmutableWritesSettled = + false; + this.completedWrites.length = 0; + } + + private async write( + key: string, + bytes: Uint8Array, + conditions: PutConditions, + ) { + const isHead = key.endsWith("/HEAD.json"); + if ( + isHead && + this.activeImmutableWrites > 0 + ) { + this.headStartedBeforeImmutableWritesSettled = + true; + } + this.activeWrites += 1; + if (!isHead) { + this.activeImmutableWrites += 1; + } + this.maximumConcurrentWrites = Math.max( + this.maximumConcurrentWrites, + this.activeWrites, + ); + try { + await delay(5); + const current = this.objects.get(key); + if ( + (conditions.ifNoneMatch && current) || + (conditions.ifMatch !== undefined && + current?.etag !== conditions.ifMatch) + ) { + throw new PreconditionFailedError(key); + } + const etag = String(++this.etag); + this.objects.set(key, { + bytes: bytes.slice(), + etag, + }); + this.writes += 1; + this.writtenBytes += bytes.byteLength; + this.completedWrites.push(key); + return { etag }; + } finally { + this.activeWrites -= 1; + if (!isHead) { + this.activeImmutableWrites -= 1; + } + } + } +} + +function delay(milliseconds: number) { + return new Promise((resolve) => + setTimeout(resolve, milliseconds), + ); +} diff --git a/tests/secondary-index.test.ts b/tests/secondary-index.test.ts index 16bf71f..c25398c 100644 --- a/tests/secondary-index.test.ts +++ b/tests/secondary-index.test.ts @@ -4,7 +4,9 @@ import path from "node:path"; import { describe, expect, it } from "vitest"; import { ContentAddressedTrieEngine, + buildSecondaryIndexPage, defineIndex, + idsFromSecondaryIndex, ImmutableSnapshotEngine, MemoryObjectCache, planSecondaryIndex, @@ -910,6 +912,50 @@ describe("secondary indexes", () => { ).toBeNull(); }); + it("applies every bounded range comparison before loading candidates", () => { + const definition = indexes.notes![1]!; + const plan = planSecondaryIndex( + [definition], + { + version: 1, + where: { + and: [ + { + field: "lastModified", + operator: "gte", + value: 2_500, + }, + { + field: "lastModified", + operator: "lt", + value: 2_525, + }, + ], + }, + }, + ); + expect(plan).not.toBeNull(); + if (!plan) { + throw new Error("Expected a range index plan"); + } + expect(plan.comparisons).toHaveLength(2); + + const page = buildSecondaryIndexPage( + definition, + Array.from({ length: 10_000 }, (_, index) => ({ + id: `note-${index}`, + title: `Note ${index}`, + lastModified: index, + })), + ); + expect(idsFromSecondaryIndex(page, plan)).toEqual( + Array.from( + { length: 25 }, + (_, index) => `note-${index + 2_500}`, + ), + ); + }); + it("falls back when equality values cannot be indexed", () => { expect( planSecondaryIndex(