(
+ "by-title",
+ ["title"],
+ "equality",
+ { include: ["lastModified"] },
+ ),
+ ],
+ });
+ const noteSummarySchema = {
+ parse(value: unknown) {
+ return value as Pick<
+ Note,
+ "id" | "title" | "lastModified"
+ >;
+ },
+ };
+
+ async function projected() {
+ const client = await createThimbleClient();
+ return client
+ .collection(notes)
+ .where((note) => note.title.eq("example"))
+ .select(
+ ["title", "lastModified"],
+ noteSummarySchema,
+ )
+ .get();
+ }
+
void [
ThimbleClient,
createThimbleClient,
@@ -196,6 +233,7 @@ try {
AuthService,
LocalObjectStore,
createArchiveManifest,
+ projected,
];
`,
);
diff --git a/site/src/data/docs.ts b/site/src/data/docs.ts
index 413312f..8e40168 100644
--- a/site/src/data/docs.ts
+++ b/site/src/data/docs.ts
@@ -187,6 +187,15 @@ export const docs: DocMeta[] = [
order: 10,
featured: true,
},
+ {
+ id: "authority-deployment",
+ title: "Authority deployment modes",
+ description:
+ "Choose an embedded or separately deployed authority and preserve the browser origin, secret, and operations boundaries.",
+ group: "Understand",
+ order: 15,
+ featured: true,
+ },
{
id: "diagrams",
title: "System diagrams",
@@ -258,7 +267,7 @@ export const docs: DocMeta[] = [
id: "queries-indexes",
title: "Queries and secondary indexes",
description:
- "Define typed collections, bounded predicates, deterministic query plans, and explicit immutable indexes.",
+ "Define typed collections, bounded predicates, deterministic query plans, explicit indexes, and covering projections.",
group: "Understand",
order: 75,
featured: true,
diff --git a/site/src/pages/index.astro b/site/src/pages/index.astro
index b110da0..b528f51 100644
--- a/site/src/pages/index.astro
+++ b/site/src/pages/index.astro
@@ -79,7 +79,7 @@ const websiteSchema = {
- Apache-2.0
- - 16.3 KB gzip browser build
+ - 18.7 KB gzip browser build
- 2-package base install
- Node.js 22+
@@ -220,11 +220,12 @@ const websiteSchema = {
.where((note) => note.title.eq("First note"))
.orderBy((note) => note.lastModified.desc())
.take(25)
- .get();
-
- Declared indexes narrow candidates. The client checks the complete
- predicate and falls back to a bounded scan when needed.
-
+ .select(["title", "lastModified"], noteCardSchema)
+ .get();
+
+ Declared indexes narrow candidates. Explicit covering fields can
+ satisfy a typed projection without loading complete documents.
+
diff --git a/site/src/pages/llms.txt.ts b/site/src/pages/llms.txt.ts
index 653fced..1f60a61 100644
--- a/site/src/pages/llms.txt.ts
+++ b/site/src/pages/llms.txt.ts
@@ -33,11 +33,12 @@ scope and one collection.
## Architecture and security
- [Architecture](${site.url}/docs/architecture/): Browser cache, authority, object storage, and scopes.
+- [Authority deployment modes](${site.url}/docs/authority-deployment/): Embedded and separate authority services, decision criteria, and same-origin requirements.
- [Security](${site.url}/security/): Threat model, encryption, key handling, and browser boundaries.
- [Authentication](${site.url}/docs/authentication/): External OIDC identities and revocable sessions.
- [Machine and service access](${site.url}/docs/service-access/): Entra roles, service principals, live viewers, and why there is no global admin key.
- [Object protocol](${site.url}/docs/protocol/): TDB1 envelopes, snapshots, tries, and conditional writes.
-- [Queries and indexes](${site.url}/docs/queries-indexes/): Typed predicates and developer-declared secondary indexes.
+- [Queries and indexes](${site.url}/docs/queries-indexes/): Typed predicates, developer-declared secondary indexes, and explicit covering projections.
- [Deletion and retention](${site.url}/docs/deletion-retention/): Tombstones, restore windows, and physical collection.
- [Logical migration](${site.url}/docs/migration/): Portable archives and database adapters.
diff --git a/site/tests/site.spec.ts b/site/tests/site.spec.ts
index 03eff5a..e5fd4a8 100644
--- a/site/tests/site.spec.ts
+++ b/site/tests/site.spec.ts
@@ -72,6 +72,27 @@ test("repository documentation renders with rewritten internal links", async ({
true,
);
await assertNoHorizontalOverflow(page);
+
+ await page.goto("/docs/authority-deployment/");
+ await expect(
+ page.getByRole("heading", {
+ level: 1,
+ name: "Authority deployment modes",
+ }),
+ ).toBeVisible();
+ await expect(
+ page.getByRole("heading", {
+ level: 2,
+ name: "Embedded authority",
+ }),
+ ).toBeVisible();
+ await expect(
+ page.getByRole("heading", {
+ level: 2,
+ name: "Separate authority service",
+ }),
+ ).toBeVisible();
+ await assertNoHorizontalOverflow(page);
});
test("documentation search returns relevant repository pages", async ({
diff --git a/src/browser/client.ts b/src/browser/client.ts
index 0a6aa63..2f4430b 100644
--- a/src/browser/client.ts
+++ b/src/browser/client.ts
@@ -35,15 +35,20 @@ import {
import {
evaluateThimbleQuery,
pointReadId,
+ queryFieldNames,
validateThimbleQuery,
type ThimbleQuery,
type ThimbleQueryResult,
} from "../query.js";
import {
idsFromSecondaryIndex,
+ documentsFromCoveringIndex,
+ encodeSecondaryIndexPage,
+ MAX_SECONDARY_INDEX_PAGE_BYTES,
planSecondaryIndex,
secondaryIndexDefinitionsEqual,
secondaryIndexPageFromJson,
+ validateProjectionFields,
type CollectionIndexConfiguration,
type SecondaryIndexPage,
type SecondaryIndexReference,
@@ -56,6 +61,7 @@ import {
} from "./cache.js";
import type {
JsonObjectReader,
+ PointReadBundleReader,
RemoteJsonObject,
} from "./remote-reader.js";
import { HttpObjectReadError } from "./remote-reader.js";
@@ -69,6 +75,9 @@ import {
export type ThimbleClientMetrics = {
remoteReads: number;
remoteBytes: number;
+ bundleReads: number;
+ bundleBytes: number;
+ bundleFallbacks: number;
notModified: number;
missing: number;
offlineFallbacks: number;
@@ -82,6 +91,9 @@ const MAX_BOUNDED_DECODED_BYTES = 16 * 1024 * 1024;
export class ThimbleClient {
private remoteReads = 0;
private remoteBytes = 0;
+ private bundleReads = 0;
+ private bundleBytes = 0;
+ private bundleFallbacks = 0;
private notModified = 0;
private missing = 0;
private offlineFallbacks = 0;
@@ -103,6 +115,7 @@ export class ThimbleClient {
constructor(
private readonly options: {
reader: JsonObjectReader;
+ bundleReader?: PointReadBundleReader;
cache: TieredObjectCache;
headTtlMs: number;
writeBaseUrl?: string;
@@ -120,10 +133,13 @@ export class ThimbleClient {
collectionIndexes?: CollectionIndexConfiguration;
layoutGeneration?: string;
configurationUrl?: string;
+ configurationCheckedAt?: number;
layoutCheckTtlMs?: number;
onLayoutChange?: () => void;
},
) {
+ this.layoutCheckedAt =
+ options.configurationCheckedAt ?? 0;
this.channel =
typeof BroadcastChannel === "undefined"
? null
@@ -203,17 +219,21 @@ export class ThimbleClient {
async queryDocuments(
collection: string,
query: ThimbleQuery,
+ projectionFields?: string[],
): Promise> {
validateThimbleQuery(query);
+ if (projectionFields) {
+ validateProjectionFields(projectionFields);
+ }
const pointId = pointReadId(query);
if (pointId) {
const document = await this.get(collection, pointId);
- return {
+ return projectQueryResult({
documents: document ? [document as unknown as T] : [],
plan: "point",
indexName: null,
scannedDocuments: document ? 1 : 0,
- };
+ }, projectionFields);
}
const definitions =
this.options.collectionIndexes?.[collection] ?? [];
@@ -228,6 +248,19 @@ export class ThimbleClient {
: await this.readHead(collection, generation);
const reference = head.indexes?.[indexPlan.definition.name];
if (reference) {
+ if (
+ typeof reference.decodedBytes !== "number" ||
+ reference.decodedBytes >
+ MAX_SECONDARY_INDEX_PAGE_BYTES
+ ) {
+ return projectQueryResult(evaluateThimbleQuery(
+ (await this.scanBounded(
+ collection,
+ query.maxScanDocuments ?? 1_000,
+ )) as unknown as T[],
+ query,
+ ), projectionFields);
+ }
const page = await this.readSecondaryIndexPage(
layout,
collection,
@@ -235,19 +268,27 @@ export class ThimbleClient {
reference.hash,
generation,
);
+ if (
+ encodeSecondaryIndexPage(page).byteLength !==
+ reference.decodedBytes
+ ) {
+ throw new Error(
+ `Secondary index ${indexPlan.definition.name} size does not match its collection head`,
+ );
+ }
if (
!secondaryIndexDefinitionsEqual(
page.definition,
indexPlan.definition,
)
) {
- return evaluateThimbleQuery(
+ return projectQueryResult(evaluateThimbleQuery(
(await this.scanBounded(
collection,
query.maxScanDocuments ?? 1_000,
)) as unknown as T[],
query,
- );
+ ), projectionFields);
}
if (reference.entries !== page.entries.length) {
throw new Error(
@@ -261,33 +302,45 @@ export class ThimbleClient {
`Secondary index ${indexPlan.definition.name} matched ${ids.length} documents, above the configured maximum of ${maximum}`,
);
}
- const documents = (
- await Promise.all(
- ids.map((id) => this.get(collection, id)),
- )
- ).filter(
- (document): document is JsonDocument =>
- document !== null,
- ) as unknown as T[];
+ const coveredDocuments = projectionFields
+ ? documentsFromCoveringIndex(
+ page,
+ indexPlan,
+ [
+ ...queryFieldNames(query),
+ ...projectionFields,
+ ],
+ )
+ : null;
+ const documents = coveredDocuments
+ ? (coveredDocuments as unknown as T[])
+ : ((
+ await Promise.all(
+ ids.map((id) => this.get(collection, id)),
+ )
+ ).filter(
+ (document): document is JsonDocument =>
+ document !== null,
+ ) as unknown as T[]);
const result = evaluateThimbleQuery(documents, {
...query,
maxScanDocuments: maximum,
});
- return {
+ return projectQueryResult({
...result,
plan: "index",
indexName: indexPlan.definition.name,
scannedDocuments: ids.length,
- };
+ }, projectionFields);
}
}
- return evaluateThimbleQuery(
+ return projectQueryResult(evaluateThimbleQuery(
(await this.scanBounded(
collection,
query.maxScanDocuments ?? 1_000,
)) as unknown as T[],
query,
- );
+ ), projectionFields);
}
explainQuery(
@@ -328,6 +381,14 @@ export class ThimbleClient {
): Promise {
await this.ensureLayoutCurrent(false);
const generation = this.currentGeneration();
+ const bundled = await this.readPointBundleIfCold(
+ collection,
+ id,
+ generation,
+ );
+ if (bundled.used) {
+ return bundled.document;
+ }
if (this.layoutFor(collection) === "snapshot") {
return this.getSnapshot(collection, id, generation);
}
@@ -674,11 +735,13 @@ export class ThimbleClient {
const head = bundle.objects.find((object) =>
object.key.endsWith("/HEAD.json"),
);
- const bundleLayout = bundle.objects.some((object) =>
- object.key.startsWith("content-snapshot/"),
- )
- ? "snapshot"
- : "trie";
+ const bundleLayout =
+ bundle.layout ??
+ (bundle.objects.some((object) =>
+ object.key.startsWith("content-snapshot/"),
+ )
+ ? "snapshot"
+ : "trie");
if (bundleLayout !== this.layoutFor(bundle.collection)) {
await this.handleLayoutChange();
throw new Error("Collection layout changed; reload required");
@@ -748,6 +811,9 @@ export class ThimbleClient {
resetMetrics(): void {
this.remoteReads = 0;
this.remoteBytes = 0;
+ this.bundleReads = 0;
+ this.bundleBytes = 0;
+ this.bundleFallbacks = 0;
this.notModified = 0;
this.missing = 0;
this.offlineFallbacks = 0;
@@ -758,6 +824,9 @@ export class ThimbleClient {
return {
remoteReads: this.remoteReads,
remoteBytes: this.remoteBytes,
+ bundleReads: this.bundleReads,
+ bundleBytes: this.bundleBytes,
+ bundleFallbacks: this.bundleFallbacks,
notModified: this.notModified,
missing: this.missing,
offlineFallbacks: this.offlineFallbacks,
@@ -1122,6 +1191,69 @@ export class ThimbleClient {
return result;
}
+ private async readPointBundleIfCold(
+ collection: string,
+ id: string,
+ generation: number,
+ ): Promise<{
+ used: boolean;
+ document: JsonDocument | null;
+ }> {
+ if (!this.options.bundleReader) {
+ return { used: false, document: null };
+ }
+ const layout = this.layoutFor(collection);
+ const headKey =
+ layout === "snapshot"
+ ? snapshotHeadKey(collection)
+ : trieHeadKey(collection);
+ const cachedHead = await this.options.cache.get(headKey);
+ this.assertGeneration(generation);
+ if (cachedHead) {
+ return { used: false, document: null };
+ }
+ this.remoteReads += 1;
+ this.bundleReads += 1;
+ let result;
+ try {
+ result = await this.options.bundleReader.get(
+ collection,
+ id,
+ );
+ } catch (error) {
+ if (
+ error instanceof HttpObjectReadError &&
+ (error.status === 401 || error.status === 403)
+ ) {
+ await this.handleAuthorizationFailure(error.status);
+ }
+ throw error;
+ }
+ this.assertGeneration(generation);
+ if (result.status === "fallback") {
+ this.bundleFallbacks += 1;
+ return { used: false, document: null };
+ }
+ this.remoteBytes += result.bytes;
+ this.bundleBytes += result.bytes;
+ if (
+ result.bundle.layout &&
+ result.bundle.layout !== layout
+ ) {
+ await this.handleLayoutChange();
+ throw new Error("Collection layout changed; reload required");
+ }
+ await this.applyBundle(
+ result.bundle,
+ false,
+ generation,
+ );
+ return {
+ used: true,
+ document: result.bundle.document,
+ };
+ }
+
private async handleAuthorizationFailure(
status: number,
): Promise {
@@ -1375,6 +1507,28 @@ function cacheEntryFromRemote(
};
}
+function projectQueryResult(
+ result: ThimbleQueryResult,
+ projectionFields?: string[],
+): ThimbleQueryResult {
+ if (!projectionFields) {
+ return result;
+ }
+ return {
+ ...result,
+ documents: result.documents.map((document) => {
+ const projection: JsonDocument = { id: document.id };
+ for (const field of projectionFields) {
+ const value = (document as Record)[field];
+ if (value !== undefined) {
+ projection[field] = structuredClone(value);
+ }
+ }
+ return projection as unknown as T;
+ }),
+ };
+}
+
async function hashId(id: string): Promise {
const digest = await crypto.subtle.digest(
"SHA-256",
diff --git a/src/browser/collection.ts b/src/browser/collection.ts
index edb0c20..967ec2c 100644
--- a/src/browser/collection.ts
+++ b/src/browser/collection.ts
@@ -15,6 +15,7 @@ import {
} from "../shared-utils.js";
import {
validateIndexConfiguration,
+ validateProjectionFields,
type CollectionIndexConfiguration,
type SecondaryIndexDefinition,
} from "../secondary-index.js";
@@ -29,6 +30,11 @@ export type CollectionDefinition = {
indexes?: SecondaryIndexDefinition[];
};
+export type ProjectedDocument<
+ T extends { id: string },
+ K extends Extract,
+> = Pick;
+
export type QueryFieldExpression<
T extends { id: string },
K extends Extract,
@@ -74,6 +80,7 @@ export interface CollectionClient {
queryDocuments?(
collection: string,
query: ThimbleQuery,
+ projectionFields?: string[],
): Promise>;
explainQuery?(
collection: string,
@@ -206,6 +213,7 @@ export class ThimbleCollection {
),
};
}
+
const pointId = pointReadId(query);
if (pointId) {
const document = await this.get(pointId);
@@ -219,6 +227,47 @@ export class ThimbleCollection {
return evaluateThimbleQuery(await this.scan(), query);
}
+ async queryProjection<
+ K extends Exclude, "id">,
+ >(
+ query: ThimbleQuery,
+ fields: K[],
+ schema: ThimbleSchema>,
+ ): Promise>> {
+ validateProjectionFields(fields);
+ if (this.client.queryDocuments) {
+ const result = await this.client.queryDocuments(
+ this.definition.name,
+ query,
+ fields,
+ );
+ return {
+ ...result,
+ documents: result.documents.map((document) =>
+ parseProjection(
+ this.definition.name,
+ schema,
+ document,
+ ),
+ ),
+ };
+ }
+ const result = evaluateThimbleQuery(
+ await this.scan(),
+ query,
+ );
+ return {
+ ...result,
+ documents: result.documents.map((document) =>
+ parseProjection(
+ this.definition.name,
+ schema,
+ projectDocument(document, fields),
+ ),
+ ),
+ };
+ }
+
where(
predicate: (fields: QueryFields) => QueryExpression,
): ThimbleQueryBuilder {
@@ -360,6 +409,21 @@ export class ThimbleQueryBuilder {
return this;
}
+ select<
+ K extends Exclude, "id">,
+ >(
+ fields: K[],
+ schema: ThimbleSchema>,
+ ): ThimbleProjectionQueryBuilder {
+ validateProjectionFields(fields);
+ return new ThimbleProjectionQueryBuilder(
+ this.collection,
+ this.queryValue,
+ fields,
+ schema,
+ );
+ }
+
get(): Promise> {
return this.collection.query(this.queryValue);
}
@@ -373,6 +437,44 @@ export class ThimbleQueryBuilder {
}
}
+export class ThimbleProjectionQueryBuilder<
+ T extends { id: string },
+ K extends Exclude, "id">,
+> {
+ constructor(
+ private readonly collection: ThimbleCollection,
+ private readonly query: ThimbleQuery,
+ private readonly fields: K[],
+ private readonly schema: ThimbleSchema<
+ ProjectedDocument
+ >,
+ ) {}
+
+ get(): Promise<
+ ThimbleQueryResult>
+ > {
+ return this.collection.queryProjection(
+ this.query,
+ this.fields,
+ this.schema,
+ );
+ }
+
+ explain(): QueryPlan {
+ return this.collection.explain(this.query);
+ }
+
+ toJSON(): {
+ query: ThimbleQuery;
+ select: K[];
+ } {
+ return {
+ query: structuredClone(this.query),
+ select: [...this.fields],
+ };
+ }
+}
+
function queryFields(): QueryFields {
return new Proxy(
{},
@@ -424,3 +526,38 @@ type QueryExpressionOperator =
? Operator
: never
: never;
+
+function parseProjection<
+ T extends { id: string },
+ K extends Extract,
+>(
+ collection: string,
+ schema: ThimbleSchema>,
+ value: unknown,
+): ProjectedDocument {
+ try {
+ return schema.parse(value);
+ } catch (error) {
+ throw new Error(
+ `Projection validation failed in collection ${collection}`,
+ { cause: error },
+ );
+ }
+}
+
+function projectDocument<
+ T extends { id: string },
+ K extends Exclude, "id">,
+>(
+ document: T,
+ fields: K[],
+): ProjectedDocument {
+ const projection = { id: document.id } as
+ ProjectedDocument;
+ for (const field of fields) {
+ if (document[field] !== undefined) {
+ projection[field] = structuredClone(document[field]);
+ }
+ }
+ return projection;
+}
diff --git a/src/browser/connect.ts b/src/browser/connect.ts
index b489a87..e8a5e1a 100644
--- a/src/browser/connect.ts
+++ b/src/browser/connect.ts
@@ -23,6 +23,7 @@ import { ThimbleClient } from "./client.js";
import {
EnvelopeJsonObjectReader,
HttpByteObjectReader,
+ HttpPointReadBundleReader,
ScopedJsonObjectReader,
} from "./remote-reader.js";
@@ -30,6 +31,7 @@ export type ThimbleAuthorityConfig = {
name: string;
provider: "local" | "azure" | "s3" | "r2";
readBaseUrl: string;
+ readBundleBaseUrl?: string;
headTtlMs: number;
cachePolicy: CachePolicy;
collectionLayouts: Record;
@@ -139,6 +141,12 @@ export async function createThimbleConnection(
config.readBaseUrl,
configurationUrl,
).toString();
+ const readBundleBaseUrl = config.readBundleBaseUrl
+ ? new URL(
+ config.readBundleBaseUrl,
+ configurationUrl,
+ ).toString()
+ : null;
const namespace = [
config.provider,
new URL(readBaseUrl).origin,
@@ -197,6 +205,16 @@ export async function createThimbleConnection(
);
const client = new ThimbleClient({
reader,
+ ...(readBundleBaseUrl
+ ? {
+ bundleReader: new HttpPointReadBundleReader(
+ readBundleBaseUrl,
+ config.scope.id,
+ fetchImplementation,
+ configurationUrl.toString(),
+ ),
+ }
+ : {}),
cache,
headTtlMs: config.headTtlMs,
csrfToken: config.csrfToken,
@@ -218,6 +236,7 @@ export async function createThimbleConnection(
collectionIndexes: config.collectionIndexes,
layoutGeneration: config.layoutGeneration,
configurationUrl: configurationUrl.toString(),
+ configurationCheckedAt: Date.now(),
layoutCheckTtlMs: options.layoutCheckTtlMs ?? 1_000,
onLayoutChange:
options.onLayoutChange ??
@@ -299,6 +318,9 @@ function validateConfig(value: unknown): ThimbleAuthorityConfig {
typeof value.name !== "string" ||
!isProvider(value.provider) ||
typeof value.readBaseUrl !== "string" ||
+ (value.readBundleBaseUrl !== undefined &&
+ (typeof value.readBundleBaseUrl !== "string" ||
+ value.readBundleBaseUrl.length === 0)) ||
typeof value.headTtlMs !== "number" ||
!Number.isFinite(value.headTtlMs) ||
value.headTtlMs < 0 ||
diff --git a/src/browser/main.ts b/src/browser/main.ts
index 2512fbd..e79c1f0 100644
--- a/src/browser/main.ts
+++ b/src/browser/main.ts
@@ -692,6 +692,9 @@ function renderMetrics(): void {
["Cache misses", metrics.cache.misses],
["Remote reads", metrics.remoteReads],
["Remote bytes", formatBytes(metrics.remoteBytes)],
+ ["Read bundles", metrics.bundleReads],
+ ["Bundle bytes", formatBytes(metrics.bundleBytes)],
+ ["Bundle fallbacks", metrics.bundleFallbacks],
["HEAD 304s", metrics.notModified],
["Offline fallbacks", metrics.offlineFallbacks],
["Memory entries", metrics.cache.memoryEntries],
diff --git a/src/browser/remote-reader.ts b/src/browser/remote-reader.ts
index bae1bff..9d8eb9b 100644
--- a/src/browser/remote-reader.ts
+++ b/src/browser/remote-reader.ts
@@ -4,6 +4,7 @@ import {
type EnvelopeKeyResolver,
} from "../envelope.js";
import { scopeStoragePrefix } from "../trie-protocol.js";
+import type { TrieReadBundle } from "../trie-protocol.js";
export type RemoteJsonObject =
| {
@@ -54,6 +55,23 @@ export interface ByteObjectReader {
): Promise;
}
+export type RemoteReadBundle =
+ | {
+ status: "found";
+ bundle: TrieReadBundle;
+ bytes: number;
+ }
+ | {
+ status: "fallback";
+ };
+
+export interface PointReadBundleReader {
+ get(
+ collection: string,
+ id: string,
+ ): Promise;
+}
+
export class HttpObjectReadError extends Error {
constructor(
readonly status: number,
@@ -125,6 +143,66 @@ export class HttpByteObjectReader implements ByteObjectReader {
}
}
+export class HttpPointReadBundleReader
+implements PointReadBundleReader {
+ constructor(
+ private readonly baseUrl: string,
+ private readonly scopeId: string,
+ private readonly fetchImplementation: typeof fetch = fetch,
+ private readonly origin = globalThis.location?.href ??
+ "http://127.0.0.1/",
+ ) {}
+
+ async get(
+ collection: string,
+ id: string,
+ ): Promise {
+ const url = new URL(this.baseUrl, this.origin);
+ url.pathname = [
+ url.pathname.replace(/\/+$/, ""),
+ encodeURIComponent(this.scopeId),
+ encodeURIComponent(collection),
+ encodeURIComponent(id),
+ ].join("/");
+ const response = await this.fetchImplementation.call(
+ globalThis,
+ url,
+ {
+ credentials: "same-origin",
+ cache: "no-store",
+ },
+ );
+ if (
+ response.status === 404 ||
+ response.status === 409 ||
+ response.status === 413
+ ) {
+ return { status: "fallback" };
+ }
+ if (!response.ok) {
+ throw new HttpObjectReadError(
+ response.status,
+ `bundle:${collection}/${id}`,
+ );
+ }
+ const body = await response.text();
+ const bundle = readBundleFromJson(JSON.parse(body) as unknown);
+ if (
+ bundle.collection !== collection ||
+ bundle.id !== id
+ ) {
+ throw new Error(
+ "Read bundle does not match the requested document",
+ );
+ }
+ return {
+ status: "found",
+ bundle,
+ bytes: new TextEncoder().encode(body).byteLength,
+ };
+ }
+}
+
export class EnvelopeJsonObjectReader implements JsonObjectReader {
constructor(
private readonly delegate: ByteObjectReader,
@@ -230,3 +308,42 @@ export function objectUrl(
url.search = query;
return url.toString();
}
+
+function readBundleFromJson(value: unknown): TrieReadBundle {
+ if (
+ typeof value !== "object" ||
+ value === null ||
+ Array.isArray(value) ||
+ !("collection" in value) ||
+ typeof value.collection !== "string" ||
+ !("id" in value) ||
+ typeof value.id !== "string" ||
+ !("revision" in value) ||
+ typeof value.revision !== "number" ||
+ !Number.isInteger(value.revision) ||
+ value.revision < 0 ||
+ !("objects" in value) ||
+ !Array.isArray(value.objects) ||
+ !value.objects.every(
+ (object) =>
+ typeof object === "object" &&
+ object !== null &&
+ !Array.isArray(object) &&
+ "key" in object &&
+ typeof object.key === "string" &&
+ "etag" in object &&
+ typeof object.etag === "string" &&
+ "value" in object,
+ ) ||
+ !("document" in value) ||
+ (value.document !== null &&
+ (typeof value.document !== "object" ||
+ Array.isArray(value.document))) ||
+ ("layout" in value &&
+ value.layout !== "trie" &&
+ value.layout !== "snapshot")
+ ) {
+ throw new Error("Read bundle response is malformed");
+ }
+ return value as TrieReadBundle;
+}
diff --git a/src/cloudflare-worker.ts b/src/cloudflare-worker.ts
index 981ccbe..dfb36da 100644
--- a/src/cloudflare-worker.ts
+++ b/src/cloudflare-worker.ts
@@ -34,6 +34,7 @@ import type {
JsonValue,
ObjectStore,
} from "./core.js";
+import { BoundedReadError } from "./core.js";
import { ContentAddressedTrieEngine } from "./engines/content-trie.js";
import { ImmutableSnapshotEngine } from "./engines/immutable-snapshot.js";
import { EnvelopeObjectStore } from "./envelope-store.js";
@@ -56,6 +57,7 @@ import {
import type { CollectionLayout } from "./snapshot-protocol.js";
import {
parseIndexConfiguration,
+ SecondaryIndexLimitError,
validateIndexConfiguration,
type CollectionIndexConfiguration,
} from "./secondary-index.js";
@@ -70,6 +72,7 @@ import {
studioDeletedDocuments,
studioScopes,
} from "./studio-api.js";
+import { readPointBundle } from "./read-bundle.js";
type RateLimitBinding = {
limit(options: { key: string }): Promise<{ success: boolean }>;
@@ -92,6 +95,7 @@ export type CloudflareAuthorityEnv = {
THIMBLE_DELETE_GRACE_DAYS?: string;
THIMBLE_MAINTENANCE_MODE?: string;
THIMBLE_STUDIO?: string;
+ THIMBLE_READ_BUNDLES?: string;
THIMBLE_STUDIO_ORIGIN?: string;
THIMBLE_COLLECTIONS?: string;
ENTRA_TENANT_ID?: string;
@@ -137,6 +141,7 @@ type Runtime = {
layoutGeneration: string;
maintenanceMode: boolean;
studioEnabled: boolean;
+ readBundlesEnabled: boolean;
studioOrigin: string | null;
oidcProviders: string[];
allowedOrigin: string;
@@ -148,6 +153,7 @@ export type CloudflareAuthorityOptions = {
collectionIndexes?: CollectionIndexConfiguration;
collections?: string[];
studio?: boolean;
+ readBundles?: boolean;
studioOrigin?: string;
};
@@ -166,26 +172,45 @@ export function createCloudflareAuthority(
);
return await route(await runtimePromise, request, env);
} catch (error) {
- if (!(error instanceof AuthError && error.status < 500)) {
- console.error(error);
+ const handledError =
+ error instanceof SecondaryIndexLimitError
+ ? new AuthError(
+ 413,
+ "secondary_index_too_large",
+ error.message,
+ )
+ : error;
+ if (
+ !(
+ handledError instanceof AuthError &&
+ handledError.status < 500
+ )
+ ) {
+ console.error(handledError);
}
- const status = error instanceof AuthError ? error.status : 500;
+ const status =
+ handledError instanceof AuthError
+ ? handledError.status
+ : 500;
const headers = new Headers();
- if (error instanceof AuthError && error.retryAfterSeconds) {
+ if (
+ handledError instanceof AuthError &&
+ handledError.retryAfterSeconds
+ ) {
headers.set(
"retry-after",
- String(error.retryAfterSeconds),
+ String(handledError.retryAfterSeconds),
);
}
return json(
{
error:
- error instanceof AuthError
- ? error.code
+ handledError instanceof AuthError
+ ? handledError.code
: "internal_error",
message:
- error instanceof AuthError
- ? error.message
+ handledError instanceof AuthError
+ ? handledError.message
: "Request failed",
},
status,
@@ -620,6 +645,9 @@ async function route(
name: "ThimbleDB",
provider: "r2",
readBaseUrl: "/api/objects",
+ ...(runtime.readBundlesEnabled
+ ? { readBundleBaseUrl: "/api/read-bundles" }
+ : {}),
headTtlMs: runtime.headTtlMs,
cachePolicy: "content",
collectionLayouts: runtime.collectionLayouts,
@@ -676,6 +704,50 @@ async function route(
});
}
+ const readBundleRoute =
+ /^\/api\/read-bundles\/([^/]+)\/([^/]+)\/([^/]+)$/.exec(
+ url.pathname,
+ );
+ if (
+ request.method === "GET" &&
+ runtime.readBundlesEnabled &&
+ readBundleRoute?.[1] &&
+ readBundleRoute[2] &&
+ readBundleRoute[3]
+ ) {
+ requireAuthenticated(authenticated);
+ const scopeId = decodePathSegment(readBundleRoute[1]);
+ requireGrant(authenticated.session.grants, scopeId, "read");
+ const collection = decodePathSegment(readBundleRoute[2]);
+ const id = decodePathSegment(readBundleRoute[3]);
+ const scope = await runtime.scope(scopeId);
+ try {
+ const bundle = await readPointBundle(
+ engineFor(runtime, scope, collection),
+ collection,
+ id,
+ );
+ return json(
+ bundle,
+ 200,
+ new Headers({
+ "x-thimble-bundle-objects": String(
+ bundle.objects.length,
+ ),
+ }),
+ );
+ } catch (error) {
+ if (error instanceof BoundedReadError) {
+ throw new AuthError(
+ 413,
+ "read_bundle_unavailable",
+ error.message,
+ );
+ }
+ throw error;
+ }
+ }
+
if (
request.method === "GET" &&
url.pathname.startsWith("/api/objects/")
@@ -1024,6 +1096,9 @@ async function createRuntime(
: configuredCollectionIndexes(env);
const studioEnabled =
options.studio ?? env.THIMBLE_STUDIO === "true";
+ const readBundlesEnabled =
+ options.readBundles ??
+ env.THIMBLE_READ_BUNDLES === "true";
const collections = studioEnabled
? studioCollectionCatalog({
collections:
@@ -1072,6 +1147,7 @@ async function createRuntime(
layoutGeneration,
maintenanceMode: env.THIMBLE_MAINTENANCE_MODE === "true",
studioEnabled,
+ readBundlesEnabled,
studioOrigin:
options.studioOrigin ??
env.THIMBLE_STUDIO_ORIGIN ??
diff --git a/src/engines/content-trie.ts b/src/engines/content-trie.ts
index c3b0671..25f795b 100644
--- a/src/engines/content-trie.ts
+++ b/src/engines/content-trie.ts
@@ -30,12 +30,14 @@ import {
type TrieLeafMetadata,
type TrieNode,
type TrieReadBundle,
+ type ReadBundleLimits,
type TrieRootNode,
type TrieStoredDocument,
type TrieTombstone,
} from "../trie-protocol.js";
import {
buildSecondaryIndexPage,
+ encodeSecondaryIndexPage,
secondaryIndexDefinitionsEqual,
secondaryIndexPageFromJson,
updateSecondaryIndexPage,
@@ -58,6 +60,13 @@ type TrieUpdate = {
second: string;
};
+type PreparedSecondaryIndex = {
+ key: string;
+ bytes: Uint8Array;
+ name: string;
+ reference: SecondaryIndexReference;
+};
+
export class ContentAddressedTrieEngine implements DatabaseEngine {
readonly name = "content-addressed-trie";
private casRetries = 0;
@@ -444,10 +453,11 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
);
}
const normalized = validateName(collection, "Collection");
- const indexes = await this.writeFreshIndexes(
+ const preparedIndexes = await this.prepareFreshIndexes(
normalized,
documents,
);
+ const indexes = await this.commitIndexes(preparedIndexes);
for (let attempt = 0; attempt < this.maxRetries; attempt += 1) {
const head = await this.loadHead(normalized);
const nextHead: TrieHead = {
@@ -532,6 +542,11 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
) {
return false;
}
+ const preparedIndexes = await this.prepareIndexes(
+ normalized,
+ head.state,
+ collapsedChanges,
+ );
const root =
head.state.rootHash === null
? this.emptyRoot()
@@ -631,10 +646,8 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
Object.keys(nextRoot.children).length === 0
? null
: await this.writeNode(normalized, nextRoot);
- const indexes = await this.writeIndexes(
- normalized,
- head.state,
- collapsedChanges,
+ const indexes = await this.commitIndexes(
+ preparedIndexes,
);
const nextHead: TrieHead = {
revision: head.state.revision + 1,
@@ -742,17 +755,25 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
async readBundle(
collection: string,
id: string,
+ limits?: ReadBundleLimits,
): Promise {
const normalized = validateName(collection, "Collection");
const head = await this.loadHead(normalized);
- const objects = [];
+ const objects: TrieReadBundle["objects"] = [];
+ let decodedBytes = 0;
if (head.object !== null) {
- objects.push({
+ decodedBytes = addBundleObject(
+ objects,
+ decodedBytes,
+ {
key: trieHeadKey(normalized),
etag: head.object.etag,
value: head.state as unknown as JsonValue,
- });
+ },
+ head.object.bytes.byteLength,
+ limits,
+ );
}
if (head.state.rootHash === null) {
@@ -762,6 +783,7 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
revision: head.state.revision,
document: null,
objects,
+ layout: "trie",
};
}
@@ -771,11 +793,17 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
head.state.rootHash,
"root",
);
- objects.push({
- key: trieNodeKey(normalized, head.state.rootHash),
- etag: root.object.etag,
- value: root.value as unknown as JsonValue,
- });
+ decodedBytes = addBundleObject(
+ objects,
+ decodedBytes,
+ {
+ key: trieNodeKey(normalized, head.state.rootHash),
+ etag: root.object.etag,
+ value: root.value as unknown as JsonValue,
+ },
+ root.object.bytes.byteLength,
+ limits,
+ );
const branchHash = root.value.children[first];
if (!branchHash) {
return {
@@ -784,6 +812,7 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
revision: head.state.revision,
document: null,
objects,
+ layout: "trie",
};
}
@@ -792,11 +821,17 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
branchHash,
"branch",
);
- objects.push({
- key: trieNodeKey(normalized, branchHash),
- etag: branch.object.etag,
- value: branch.value as unknown as JsonValue,
- });
+ decodedBytes = addBundleObject(
+ objects,
+ decodedBytes,
+ {
+ key: trieNodeKey(normalized, branchHash),
+ etag: branch.object.etag,
+ value: branch.value as unknown as JsonValue,
+ },
+ branch.object.bytes.byteLength,
+ limits,
+ );
const leafHash = branch.value.children[second];
if (!leafHash) {
return {
@@ -805,22 +840,42 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
revision: head.state.revision,
document: null,
objects,
+ layout: "trie",
};
}
+ if (limits) {
+ const metadata = branch.value.leafMetadata?.[second];
+ if (!metadata) {
+ throw new BoundedReadError(
+ "Trie leaf size metadata is unavailable for a bounded read bundle",
+ );
+ }
+ assertBundleCapacity(
+ objects.length + 1,
+ decodedBytes + metadata.decodedBytes,
+ limits,
+ );
+ }
const leaf = await this.loadNodeObject(
normalized,
leafHash,
"leaf",
);
- objects.push({
- key: trieNodeKey(
- normalized,
- leafHash,
- ),
- etag: leaf.object.etag,
- value: leaf.value as unknown as JsonValue,
- });
+ addBundleObject(
+ objects,
+ decodedBytes,
+ {
+ key: trieNodeKey(
+ normalized,
+ leafHash,
+ ),
+ etag: leaf.object.etag,
+ value: leaf.value as unknown as JsonValue,
+ },
+ leaf.object.bytes.byteLength,
+ limits,
+ );
return {
collection: normalized,
@@ -830,6 +885,7 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
ownValue(leaf.value.documents, id),
),
objects,
+ layout: "trie",
};
}
@@ -845,19 +901,18 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
return { object, state: decodeJson(object.bytes) };
}
- private async writeIndexes(
+ private async prepareIndexes(
collection: string,
head: TrieHead,
changes: SecondaryIndexChange[],
- ): Promise {
+ ): Promise {
const definitions = this.indexConfiguration[collection] ?? [];
if (definitions.length === 0) {
- return {};
+ return [];
}
let storedDocuments: TrieStoredDocument[] | undefined;
- const references =
- createDictionary();
+ const prepared: PreparedSecondaryIndex[] = [];
for (const definition of definitions) {
const currentReference = head.indexes?.[definition.name];
let currentPage: SecondaryIndexPage | null = null;
@@ -909,45 +964,66 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
definition,
changes,
);
- const bytes = encodeJson(page as unknown as JsonValue);
+ const bytes = encodeSecondaryIndexPage(page);
const hash = await this.addressNode(bytes);
- try {
- await this.store.put(
- trieIndexKey(collection, definition.name, hash),
- bytes,
- { ifNoneMatch: true },
- );
- } catch (error) {
- if (!isPreconditionFailure(error)) {
- throw error;
- }
- }
- references[definition.name] = {
- hash,
- entries: page.entries.length,
- decodedBytes: bytes.byteLength,
- };
+ prepared.push({
+ key: trieIndexKey(
+ collection,
+ definition.name,
+ hash,
+ ),
+ bytes,
+ name: definition.name,
+ reference: {
+ hash,
+ entries: page.entries.length,
+ decodedBytes: bytes.byteLength,
+ },
+ });
}
- return references;
+ return prepared;
}
- private async writeFreshIndexes(
+ private async prepareFreshIndexes(
collection: string,
documents: TrieStoredDocument[],
- ): Promise {
- const references =
- createDictionary();
+ ): Promise {
+ const prepared: PreparedSecondaryIndex[] = [];
for (const definition of this.indexConfiguration[collection] ?? []) {
const page = buildSecondaryIndexPage(
definition,
documents,
);
- const bytes = encodeJson(page as unknown as JsonValue);
+ const bytes = encodeSecondaryIndexPage(page);
const hash = await this.addressNode(bytes);
+ prepared.push({
+ key: trieIndexKey(
+ collection,
+ definition.name,
+ hash,
+ ),
+ bytes,
+ name: definition.name,
+ reference: {
+ hash,
+ entries: page.entries.length,
+ decodedBytes: bytes.byteLength,
+ },
+ });
+ }
+ return prepared;
+ }
+
+ private async commitIndexes(
+ prepared: PreparedSecondaryIndex[],
+ ): Promise {
+ const references =
+ createDictionary();
+ for (const index of prepared) {
try {
await this.store.put(
- trieIndexKey(collection, definition.name, hash),
- bytes,
+ index.key,
+ index.bytes,
{ ifNoneMatch: true },
);
} catch (error) {
@@ -955,11 +1031,7 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
throw error;
}
}
- references[definition.name] = {
- hash,
- entries: page.entries.length,
- decodedBytes: bytes.byteLength,
- };
+ references[index.name] = index.reference;
}
return references;
}
@@ -1281,6 +1353,40 @@ export class ContentAddressedTrieEngine implements DatabaseEngine {
}
}
+function addBundleObject(
+ objects: TrieReadBundle["objects"],
+ decodedBytes: number,
+ object: TrieReadBundle["objects"][number],
+ objectBytes: number,
+ limits?: ReadBundleLimits,
+): number {
+ const nextBytes = decodedBytes + objectBytes;
+ if (limits) {
+ assertBundleCapacity(
+ objects.length + 1,
+ nextBytes,
+ limits,
+ );
+ }
+ objects.push(object);
+ return nextBytes;
+}
+
+function assertBundleCapacity(
+ objects: number,
+ decodedBytes: number,
+ limits: ReadBundleLimits,
+): void {
+ if (
+ objects > limits.maxObjects ||
+ decodedBytes > limits.maxDecodedBytes
+ ) {
+ throw new BoundedReadError(
+ `Read bundle exceeds ${limits.maxObjects} objects or ${limits.maxDecodedBytes} decoded bytes`,
+ );
+ }
+}
+
async function hashBytes(bytes: Uint8Array): Promise {
const copy = new Uint8Array(new ArrayBuffer(bytes.byteLength));
copy.set(bytes);
diff --git a/src/engines/immutable-snapshot.ts b/src/engines/immutable-snapshot.ts
index 054692d..83de84b 100644
--- a/src/engines/immutable-snapshot.ts
+++ b/src/engines/immutable-snapshot.ts
@@ -28,11 +28,13 @@ import {
isTrieTombstone,
visibleTrieDocument,
type TrieReadBundle,
+ type ReadBundleLimits,
type TrieStoredDocument,
type TrieTombstone,
} from "../trie-protocol.js";
import {
buildSecondaryIndexPage,
+ encodeSecondaryIndexPage,
secondaryIndexDefinitionsEqual,
secondaryIndexPageFromJson,
type CollectionIndexConfiguration,
@@ -45,6 +47,13 @@ type LoadedHead = {
state: SnapshotHead;
};
+type PreparedSecondaryIndex = {
+ key: string;
+ bytes: Uint8Array;
+ name: string;
+ reference: SecondaryIndexReference;
+};
+
export class ImmutableSnapshotEngine implements DatabaseEngine {
readonly name = "immutable-snapshot";
private casRetries = 0;
@@ -369,35 +378,82 @@ export class ImmutableSnapshotEngine implements DatabaseEngine {
async readBundle(
collection: string,
id: string,
+ limits?: ReadBundleLimits,
): Promise {
const normalized = validateName(collection, "Collection");
- const loaded = await this.loadCurrent(normalized);
- const objects = [];
- if (loaded.head.object) {
- objects.push({
+ const head = await this.loadHead(normalized);
+ const objects: TrieReadBundle["objects"] = [];
+ let decodedBytes = 0;
+ if (head.object) {
+ decodedBytes = addBundleObject(
+ objects,
+ decodedBytes,
+ {
key: snapshotHeadKey(normalized),
- etag: loaded.head.object.etag,
- value: loaded.head.state as unknown as JsonValue,
- });
+ etag: head.object.etag,
+ value: head.state as unknown as JsonValue,
+ },
+ head.object.bytes.byteLength,
+ limits,
+ );
+ }
+ if (!head.state.snapshotHash) {
+ return {
+ collection: normalized,
+ id,
+ revision: head.state.revision,
+ document: null,
+ objects,
+ layout: "snapshot",
+ };
+ }
+ if (limits) {
+ if (typeof head.state.decodedBytes !== "number") {
+ throw new BoundedReadError(
+ "Snapshot size metadata is unavailable for a bounded read bundle",
+ );
+ }
+ assertBundleCapacity(
+ objects.length + 1,
+ decodedBytes + head.state.decodedBytes,
+ limits,
+ );
}
- if (loaded.pageObject && loaded.head.state.snapshotHash) {
- objects.push({
+ const pageObject = await this.store.get(
+ snapshotPageKey(
+ normalized,
+ head.state.snapshotHash,
+ ),
+ );
+ if (!pageObject) {
+ throw new Error(
+ `Snapshot ${head.state.snapshotHash} is missing`,
+ );
+ }
+ const page = decodeJson(pageObject.bytes);
+ addBundleObject(
+ objects,
+ decodedBytes,
+ {
key: snapshotPageKey(
normalized,
- loaded.head.state.snapshotHash,
+ head.state.snapshotHash,
),
- etag: loaded.pageObject.etag,
- value: loaded.page as unknown as JsonValue,
- });
- }
+ etag: pageObject.etag,
+ value: page as unknown as JsonValue,
+ },
+ pageObject.bytes.byteLength,
+ limits,
+ );
return {
collection: normalized,
id,
- revision: loaded.head.state.revision,
+ revision: head.state.revision,
document: visibleTrieDocument(
- ownValue(loaded.page.documents, id),
+ ownValue(page.documents, id),
),
objects,
+ layout: "snapshot",
};
}
@@ -421,8 +477,13 @@ export class ImmutableSnapshotEngine implements DatabaseEngine {
if (!updateResult.changed) {
return updateResult.result;
}
+
const page: SnapshotPage = { documents };
const pageBytes = encodeJson(page as unknown as JsonValue);
+ const preparedIndexes = await this.prepareIndexes(
+ normalized,
+ Object.values(documents),
+ );
const snapshotHash =
Object.keys(documents).length === 0
? null
@@ -430,9 +491,8 @@ export class ImmutableSnapshotEngine implements DatabaseEngine {
normalized,
pageBytes,
);
- const indexes = await this.writeIndexes(
- normalized,
- Object.values(documents),
+ const indexes = await this.commitIndexes(
+ preparedIndexes,
);
const nextHead: SnapshotHead = {
revision: loaded.head.state.revision + 1,
@@ -526,24 +586,47 @@ export class ImmutableSnapshotEngine implements DatabaseEngine {
return hash;
}
- private async writeIndexes(
+ private async prepareIndexes(
collection: string,
documents: TrieStoredDocument[],
- ): Promise {
+ ): Promise {
const definitions = this.indexConfiguration[collection] ?? [];
- const references =
- createDictionary();
+ const prepared: PreparedSecondaryIndex[] = [];
for (const definition of definitions) {
const page = buildSecondaryIndexPage(
definition,
documents,
);
- const bytes = encodeJson(page as unknown as JsonValue);
+ const bytes = encodeSecondaryIndexPage(page);
const hash = await this.addressSnapshot(bytes);
+ prepared.push({
+ key: snapshotIndexKey(
+ collection,
+ definition.name,
+ hash,
+ ),
+ bytes,
+ name: definition.name,
+ reference: {
+ hash,
+ entries: page.entries.length,
+ decodedBytes: bytes.byteLength,
+ },
+ });
+ }
+ return prepared;
+ }
+
+ private async commitIndexes(
+ prepared: PreparedSecondaryIndex[],
+ ): Promise {
+ const references =
+ createDictionary();
+ for (const index of prepared) {
try {
await this.store.put(
- snapshotIndexKey(collection, definition.name, hash),
- bytes,
+ index.key,
+ index.bytes,
{ ifNoneMatch: true },
);
} catch (error) {
@@ -551,11 +634,7 @@ export class ImmutableSnapshotEngine implements DatabaseEngine {
throw error;
}
}
- references[definition.name] = {
- hash,
- entries: page.entries.length,
- decodedBytes: bytes.byteLength,
- };
+ references[index.name] = index.reference;
}
return references;
}
@@ -617,6 +696,40 @@ export class ImmutableSnapshotEngine implements DatabaseEngine {
}
}
+function addBundleObject(
+ objects: TrieReadBundle["objects"],
+ decodedBytes: number,
+ object: TrieReadBundle["objects"][number],
+ objectBytes: number,
+ limits?: ReadBundleLimits,
+): number {
+ const nextBytes = decodedBytes + objectBytes;
+ if (limits) {
+ assertBundleCapacity(
+ objects.length + 1,
+ nextBytes,
+ limits,
+ );
+ }
+ objects.push(object);
+ return nextBytes;
+}
+
+function assertBundleCapacity(
+ objects: number,
+ decodedBytes: number,
+ limits: ReadBundleLimits,
+): void {
+ if (
+ objects > limits.maxObjects ||
+ decodedBytes > limits.maxDecodedBytes
+ ) {
+ throw new BoundedReadError(
+ `Read bundle exceeds ${limits.maxObjects} objects or ${limits.maxDecodedBytes} decoded bytes`,
+ );
+ }
+}
+
function assertUserDocument(document: JsonDocument): void {
if ("__thimbleTombstone" in document) {
throw new Error(
diff --git a/src/query.ts b/src/query.ts
index 91559e9..625f941 100644
--- a/src/query.ts
+++ b/src/query.ts
@@ -91,9 +91,21 @@ export function pointReadId(
) {
return where.value;
}
+
return null;
}
+export function queryFieldNames(
+ query: ThimbleQuery,
+): string[] {
+ const fields = new Set();
+ collectExpressionFields(query.where, fields);
+ for (const order of query.orderBy ?? []) {
+ fields.add(order.field);
+ }
+ return [...fields];
+}
+
export function validateThimbleQuery(
query: ThimbleQuery,
): void {
@@ -154,10 +166,12 @@ function validateExpression(
if (depth > 12) {
throw new Error("Query expression nesting exceeds 12 levels");
}
+
state.nodes += 1;
if (state.nodes > 500) {
throw new Error("Query expression exceeds 500 nodes");
}
+
if (!isRecord(expression)) {
throw new Error("Query expression is malformed");
}
@@ -227,6 +241,32 @@ function validateExpression(
validateExpression(expression.not, depth + 1, state);
}
+function collectExpressionFields(
+ expression: QueryExpression | undefined,
+ fields: Set,
+): void {
+ if (!expression) {
+ return;
+ }
+ if ("field" in expression) {
+ fields.add(expression.field);
+ return;
+ }
+ if ("and" in expression) {
+ expression.and.forEach((child) =>
+ collectExpressionFields(child, fields),
+ );
+ return;
+ }
+ if ("or" in expression) {
+ expression.or.forEach((child) =>
+ collectExpressionFields(child, fields),
+ );
+ return;
+ }
+ collectExpressionFields(expression.not, fields);
+}
+
function evaluateExpression(
document: T,
expression: QueryExpression,
diff --git a/src/read-bundle.ts b/src/read-bundle.ts
new file mode 100644
index 0000000..5a4ac99
--- /dev/null
+++ b/src/read-bundle.ts
@@ -0,0 +1,39 @@
+import type {
+ ContentAddressedTrieEngine,
+} from "./engines/content-trie.js";
+import type {
+ ImmutableSnapshotEngine,
+} from "./engines/immutable-snapshot.js";
+import type { JsonValue } from "./core.js";
+import { encodeJson } from "./shared-utils.js";
+import { BoundedReadError } from "./core.js";
+import type { TrieReadBundle } from "./trie-protocol.js";
+
+export const READ_BUNDLE_MAX_OBJECTS = 4;
+export const READ_BUNDLE_MAX_DECODED_BYTES =
+ 4 * 1024 * 1024;
+
+type ReadBundleEngine =
+ | ContentAddressedTrieEngine
+ | ImmutableSnapshotEngine;
+
+export function readPointBundle(
+ engine: ReadBundleEngine,
+ collection: string,
+ id: string,
+): Promise {
+ return engine.readBundle(collection, id, {
+ maxObjects: READ_BUNDLE_MAX_OBJECTS,
+ maxDecodedBytes: READ_BUNDLE_MAX_DECODED_BYTES,
+ }).then((bundle) => {
+ if (
+ encodeJson(bundle as unknown as JsonValue).byteLength >
+ READ_BUNDLE_MAX_DECODED_BYTES
+ ) {
+ throw new BoundedReadError(
+ `Read bundle response exceeds ${READ_BUNDLE_MAX_DECODED_BYTES} decoded bytes`,
+ );
+ }
+ return bundle;
+ });
+}
diff --git a/src/secondary-index.ts b/src/secondary-index.ts
index 48f6de5..0335e29 100644
--- a/src/secondary-index.ts
+++ b/src/secondary-index.ts
@@ -10,6 +10,7 @@ import type {
} from "./query.js";
import {
createDictionary,
+ encodeJson,
validateName,
} from "./shared-utils.js";
import {
@@ -18,17 +19,50 @@ import {
} from "./trie-protocol.js";
export type SecondaryIndexMode = "equality" | "range";
+export const MAX_COVERING_FIELDS = 8;
+export const MAX_COVERING_DOCUMENT_BYTES = 64 * 1024;
+export const MAX_SECONDARY_INDEX_PAGE_BYTES =
+ 4 * 1024 * 1024;
+
+export class SecondaryIndexLimitError extends Error {
+ constructor(message: string) {
+ super(message);
+ this.name = "SecondaryIndexLimitError";
+ }
+}
+
+export function validateProjectionFields(
+ fields: string[],
+): string[] {
+ if (
+ fields.length < 1 ||
+ fields.length > MAX_COVERING_FIELDS ||
+ new Set(fields).size !== fields.length ||
+ !fields.every(
+ (field) => isIndexFieldName(field) && field !== "id",
+ )
+ ) {
+ throw new Error(
+ "Query projection requires 1-8 unique safe non-ID fields",
+ );
+ }
+ return [...fields];
+}
export type SecondaryIndexDefinition = {
name: string;
fields: string[];
mode: SecondaryIndexMode;
+ include?: string[];
};
export function defineIndex(
name: string,
fields: Array>,
mode: SecondaryIndexMode = "equality",
+ options: {
+ include?: Array>;
+ } = {},
): SecondaryIndexDefinition {
const configuration = validateIndexConfiguration({
collection: [
@@ -36,6 +70,9 @@ export function defineIndex(
name,
fields,
mode,
+ ...(options.include
+ ? { include: options.include }
+ : {}),
},
],
});
@@ -67,6 +104,7 @@ export type SecondaryIndexPage = {
version: 1;
definition: SecondaryIndexDefinition;
entries: SecondaryIndexEntry[];
+ projections?: Record;
};
export type SecondaryIndexChange = {
@@ -104,15 +142,23 @@ export function validateIndexConfiguration(
definition.fields.length > 4 ||
new Set(definition.fields).size !==
definition.fields.length ||
- !definition.fields.every(
- (field) =>
- typeof field === "string" &&
- /^[A-Za-z0-9_-]{1,64}$/.test(field),
- ) ||
+ !definition.fields.every(isIndexFieldName) ||
(definition.mode !== "equality" &&
definition.mode !== "range") ||
(definition.mode === "range" &&
- definition.fields.length !== 1)
+ definition.fields.length !== 1) ||
+ (definition.include !== undefined &&
+ (!Array.isArray(definition.include) ||
+ definition.include.length < 1 ||
+ definition.include.length > MAX_COVERING_FIELDS ||
+ new Set(definition.include).size !==
+ definition.include.length ||
+ !definition.include.every(
+ (field) =>
+ isIndexFieldName(field) &&
+ field !== "id" &&
+ !definition.fields.includes(field),
+ )))
) {
throw new Error(
`Invalid secondary index configuration for ${collection}`,
@@ -123,6 +169,9 @@ export function validateIndexConfiguration(
name: definition.name,
fields: [...definition.fields],
mode: definition.mode,
+ ...(definition.include
+ ? { include: [...definition.include] }
+ : {}),
};
});
}
@@ -163,16 +212,20 @@ export function buildSecondaryIndexPage(
documents: Iterable,
): SecondaryIndexPage {
const entries = new Map();
+ const projections = definition.include
+ ? createDictionary()
+ : undefined;
for (const stored of documents) {
if (isTrieTombstone(stored)) {
continue;
}
- addDocument(entries, definition, stored);
+ addDocument(entries, projections, definition, stored);
}
return {
version: 1,
definition,
entries: sortEntries([...entries.values()]),
+ ...(projections ? { projections } : {}),
};
}
@@ -182,6 +235,11 @@ export function updateSecondaryIndexPage(
changes: SecondaryIndexChange[],
): SecondaryIndexPage {
const entries = new Map();
+ const projections = definition.include
+ ? createDictionary(
+ current?.projections,
+ )
+ : undefined;
for (const entry of current?.entries ?? []) {
entries.set(indexKey(entry.values), {
values: [...entry.values],
@@ -196,14 +254,23 @@ export function updateSecondaryIndexPage(
}
}
for (const change of changes) {
+ if (projections) {
+ delete projections[change.id];
+ }
if (change.document && !isTrieTombstone(change.document)) {
- addDocument(entries, definition, change.document);
+ addDocument(
+ entries,
+ projections,
+ definition,
+ change.document,
+ );
}
}
return {
version: 1,
definition,
entries: sortEntries([...entries.values()]),
+ ...(projections ? { projections } : {}),
};
}
@@ -316,6 +383,7 @@ export function secondaryIndexPageFromJson(
}).collection![0]!;
const keys = new Set();
const documentIds = new Set();
+ const indexedValues = new Map();
const entries = value.entries.map((entry) => {
if (
typeof entry !== "object" ||
@@ -346,16 +414,24 @@ export function secondaryIndexPageFromJson(
);
}
documentIds.add(id);
+ indexedValues.set(id, values);
}
return {
values,
ids,
};
});
+ const projections = parseProjections(
+ value.projections,
+ definition,
+ documentIds,
+ indexedValues,
+ );
return {
version: 1,
definition,
entries,
+ ...(projections ? { projections } : {}),
};
}
@@ -369,12 +445,18 @@ export function secondaryIndexDefinitionsEqual(
left.fields.length === right.fields.length &&
left.fields.every(
(field, index) => field === right.fields[index],
+ ) &&
+ (left.include?.length ?? 0) ===
+ (right.include?.length ?? 0) &&
+ (left.include ?? []).every(
+ (field, index) => field === right.include?.[index],
)
);
}
function addDocument(
entries: Map,
+ projections: Record | undefined,
definition: SecondaryIndexDefinition,
document: JsonDocument,
): void {
@@ -393,6 +475,158 @@ function addDocument(
entry.ids.sort();
}
entries.set(key, entry);
+ if (projections) {
+ projections[document.id] = projectIndexDocument(
+ definition,
+ document,
+ );
+ }
+}
+
+export function documentsFromCoveringIndex<
+ T extends { id: string },
+>(
+ page: SecondaryIndexPage,
+ plan: SecondaryIndexPlan,
+ requiredFields: Iterable,
+): JsonDocument[] | null {
+ if (!page.projections) {
+ return null;
+ }
+ const covered = new Set([
+ "id",
+ ...page.definition.fields,
+ ...(page.definition.include ?? []),
+ ]);
+ if (
+ [...requiredFields].some((field) => !covered.has(field))
+ ) {
+ return null;
+ }
+ const ids = idsFromSecondaryIndex(page, plan);
+ const documents: JsonDocument[] = [];
+ for (const id of ids) {
+ const projection = page.projections[id];
+ if (!projection) {
+ throw new Error(
+ `Secondary index projection is missing document ${id}`,
+ );
+ }
+ documents.push(structuredClone(projection));
+ }
+ return documents;
+}
+
+function projectIndexDocument(
+ definition: SecondaryIndexDefinition,
+ document: JsonDocument,
+): JsonDocument {
+ const projection: JsonDocument = { id: document.id };
+ for (const field of [
+ ...definition.fields,
+ ...(definition.include ?? []),
+ ]) {
+ if (document[field] !== undefined) {
+ projection[field] = structuredClone(document[field]);
+ }
+ }
+ if (
+ encodeJson(projection as unknown as JsonValue).byteLength >
+ MAX_COVERING_DOCUMENT_BYTES
+ ) {
+ throw new SecondaryIndexLimitError(
+ `Secondary index covering projection exceeds ${MAX_COVERING_DOCUMENT_BYTES} decoded bytes for ${document.id}`,
+ );
+ }
+
+ return projection;
+}
+
+export function encodeSecondaryIndexPage(
+ page: SecondaryIndexPage,
+): Uint8Array {
+ const bytes = encodeJson(page as unknown as JsonValue);
+ if (bytes.byteLength > MAX_SECONDARY_INDEX_PAGE_BYTES) {
+ throw new SecondaryIndexLimitError(
+ `Secondary index ${page.definition.name} exceeds ${MAX_SECONDARY_INDEX_PAGE_BYTES} decoded bytes`,
+ );
+ }
+ return bytes;
+}
+
+function isIndexFieldName(value: unknown): value is string {
+ return (
+ typeof value === "string" &&
+ /^[A-Za-z0-9_-]{1,64}$/.test(value) &&
+ value !== "__proto__" &&
+ value !== "prototype" &&
+ value !== "constructor"
+ );
+}
+
+function parseProjections(
+ value: unknown,
+ definition: SecondaryIndexDefinition,
+ documentIds: Set,
+ indexedValues: Map,
+): Record | undefined {
+ if (!definition.include) {
+ if (value !== undefined) {
+ throw new Error(
+ "Secondary index page has undeclared projections",
+ );
+ }
+ return undefined;
+ }
+ if (
+ typeof value !== "object" ||
+ value === null ||
+ Array.isArray(value)
+ ) {
+ throw new Error(
+ "Secondary index page is missing covering projections",
+ );
+ }
+ const allowedFields = new Set([
+ "id",
+ ...definition.fields,
+ ...definition.include,
+ ]);
+ const projections = createDictionary();
+ for (const [id, candidate] of Object.entries(value)) {
+ if (
+ !documentIds.has(id) ||
+ typeof candidate !== "object" ||
+ candidate === null ||
+ Array.isArray(candidate) ||
+ candidate.id !== id ||
+ Object.keys(candidate).some(
+ (field) => !allowedFields.has(field),
+ )
+ ) {
+ throw new Error(
+ "Secondary index covering projection is malformed",
+ );
+ }
+ const projection = candidate as JsonDocument;
+ const values = indexedValues.get(id)!;
+ for (let index = 0; index < definition.fields.length; index += 1) {
+ if (
+ projection[definition.fields[index]!] !== values[index]
+ ) {
+ throw new Error(
+ "Secondary index covering projection does not match its key",
+ );
+ }
+ }
+ projections[id] = structuredClone(projection);
+ }
+ if (Object.keys(projections).length !== documentIds.size) {
+ throw new Error(
+ "Secondary index page is missing covering projections",
+ );
+ }
+ return projections;
}
function flattenComparisons(
diff --git a/src/server.ts b/src/server.ts
index 46f4e6e..50eb0de 100644
--- a/src/server.ts
+++ b/src/server.ts
@@ -40,6 +40,7 @@ import type {
JsonValue,
ObjectStore,
} from "./core.js";
+import { BoundedReadError } from "./core.js";
import { ContentAddressedTrieEngine } from "./engines/content-trie.js";
import { ImmutableSnapshotEngine } from "./engines/immutable-snapshot.js";
import { EnvelopeObjectStore } from "./envelope-store.js";
@@ -69,6 +70,7 @@ import {
import type { CollectionLayout } from "./snapshot-protocol.js";
import {
parseIndexConfiguration,
+ SecondaryIndexLimitError,
validateIndexConfiguration,
type CollectionIndexConfiguration,
} from "./secondary-index.js";
@@ -83,6 +85,7 @@ import {
studioDeletedDocuments,
studioScopes,
} from "./studio-api.js";
+import { readPointBundle } from "./read-bundle.js";
type ScopeRuntime = {
material: ScopeMaterial;
@@ -105,6 +108,7 @@ type ServerContext = {
layoutGeneration: string;
maintenanceMode: boolean;
studioEnabled: boolean;
+ readBundlesEnabled: boolean;
studioOrigin: string | null;
developmentIdentity: boolean;
oidcProviders: string[];
@@ -123,6 +127,7 @@ export type NodeAuthorityOptions = {
collectionIndexes?: CollectionIndexConfiguration;
collections?: string[];
studio?: boolean;
+ readBundles?: boolean;
studioOrigin?: string;
};
@@ -172,6 +177,9 @@ async function createContext(
: requiredEnvironment("THIMBLE_ALLOWED_ORIGIN"));
const studioEnabled =
options.studio ?? process.env.THIMBLE_STUDIO === "true";
+ const readBundlesEnabled =
+ options.readBundles ??
+ process.env.THIMBLE_READ_BUNDLES === "true";
const studioOrigin =
options.studioOrigin ??
process.env.THIMBLE_STUDIO_ORIGIN ??
@@ -302,6 +310,7 @@ async function createContext(
maintenanceMode:
process.env.THIMBLE_MAINTENANCE_MODE === "true",
studioEnabled,
+ readBundlesEnabled,
studioOrigin,
developmentIdentity,
oidcProviders: [...identityAdapters.keys()].filter(
@@ -837,6 +846,9 @@ async function handleRequest(
name: "ThimbleDB",
provider: context.provider,
readBaseUrl: "/api/objects",
+ ...(context.readBundlesEnabled
+ ? { readBundleBaseUrl: "/api/read-bundles" }
+ : {}),
headTtlMs: context.headTtlMs,
cachePolicy: "content",
collectionLayouts: context.collectionLayouts,
@@ -891,6 +903,47 @@ async function handleRequest(
return;
}
+ const readBundleRoute =
+ /^\/api\/read-bundles\/([^/]+)\/([^/]+)\/([^/]+)$/.exec(
+ url.pathname,
+ );
+ if (
+ request.method === "GET" &&
+ context.readBundlesEnabled &&
+ readBundleRoute?.[1] &&
+ readBundleRoute[2] &&
+ readBundleRoute[3]
+ ) {
+ requireAuthenticated(authenticated);
+ const scopeId = decodePathSegment(readBundleRoute[1]);
+ requireGrant(authenticated.session.grants, scopeId, "read");
+ const collection = decodePathSegment(readBundleRoute[2]);
+ const id = decodePathSegment(readBundleRoute[3]);
+ const runtime = await context.scope(scopeId);
+ try {
+ const bundle = await readPointBundle(
+ engineFor(context, runtime, collection),
+ collection,
+ id,
+ );
+ sendJson(response, 200, bundle, {
+ "x-thimble-bundle-objects": String(
+ bundle.objects.length,
+ ),
+ });
+ } catch (error) {
+ if (error instanceof BoundedReadError) {
+ throw new AuthError(
+ 413,
+ "read_bundle_unavailable",
+ error.message,
+ );
+ }
+ throw error;
+ }
+ return;
+ }
+
if (
request.method === "GET" &&
url.pathname.startsWith("/api/objects/")
@@ -2107,26 +2160,51 @@ function handleServerError(
error: unknown,
response: ServerResponse,
): void {
- if (!(error instanceof AuthError && error.status < 500)) {
- console.error(error);
+ const handledError =
+ error instanceof SecondaryIndexLimitError
+ ? new AuthError(
+ 413,
+ "secondary_index_too_large",
+ error.message,
+ )
+ : error;
+ if (
+ !(
+ handledError instanceof AuthError &&
+ handledError.status < 500
+ )
+ ) {
+ console.error(handledError);
}
if (response.headersSent) {
response.end();
return;
}
- const status = error instanceof AuthError ? error.status : 500;
+ const status =
+ handledError instanceof AuthError
+ ? handledError.status
+ : 500;
const headers: Record = {};
- if (error instanceof AuthError && error.retryAfterSeconds) {
- headers["retry-after"] = String(error.retryAfterSeconds);
+ if (
+ handledError instanceof AuthError &&
+ handledError.retryAfterSeconds
+ ) {
+ headers["retry-after"] = String(
+ handledError.retryAfterSeconds,
+ );
}
sendJson(
response,
status,
{
error:
- error instanceof AuthError ? error.code : "internal_error",
+ handledError instanceof AuthError
+ ? handledError.code
+ : "internal_error",
message:
- error instanceof AuthError ? error.message : "Request failed",
+ handledError instanceof AuthError
+ ? handledError.message
+ : "Request failed",
},
headers,
);
diff --git a/src/trie-protocol.ts b/src/trie-protocol.ts
index 5ef215a..875b6f6 100644
--- a/src/trie-protocol.ts
+++ b/src/trie-protocol.ts
@@ -60,6 +60,12 @@ export type TrieReadBundle = {
revision: number;
document: JsonDocument | null;
objects: TrieBundleObject[];
+ layout?: "trie" | "snapshot";
+};
+
+export type ReadBundleLimits = {
+ maxObjects: number;
+ maxDecodedBytes: number;
};
export function trieCollectionPrefix(collection: string): string {
diff --git a/studio/src/main.ts b/studio/src/main.ts
index 2b0b7b1..9b1a80f 100644
--- a/studio/src/main.ts
+++ b/studio/src/main.ts
@@ -40,6 +40,7 @@ type StudioIndex = {
name: string;
fields: string[];
mode: "equality" | "range";
+ include?: string[];
};
active: boolean;
entries: number | null;
@@ -1073,7 +1074,7 @@ function renderIndexes(): void {
}
const table = document.createElement("table");
table.innerHTML =
- "| Index | Mode | Fields | Status | Entries |
";
+ "| Index | Mode | Fields | Covers | Status | Entries |
";
const body = document.createElement("tbody");
for (const index of collection.indexes) {
const row = document.createElement("tr");
@@ -1081,6 +1082,7 @@ function renderIndexes(): void {
index.definition.name,
index.definition.mode,
index.definition.fields.join(", "),
+ index.definition.include?.join(", ") ?? "None",
index.status,
index.entries === null ? "—" : String(index.entries),
]) {
diff --git a/templates/local-web/server.mjs b/templates/local-web/server.mjs
index c799c05..ee533b9 100644
--- a/templates/local-web/server.mjs
+++ b/templates/local-web/server.mjs
@@ -10,6 +10,7 @@ const { startNodeAuthority } = await import(
await startNodeAuthority({
studio: true,
+ readBundles: true,
collections: ["notes"],
collectionLayouts: {
notes: "snapshot",
@@ -20,6 +21,7 @@ await startNodeAuthority({
name: "by-title",
fields: ["title"],
mode: "equality",
+ include: ["lastModified"],
},
{
name: "by-last-modified",
diff --git a/templates/local-web/src/main.ts b/templates/local-web/src/main.ts
index 35078e2..ee95808 100644
--- a/templates/local-web/src/main.ts
+++ b/templates/local-web/src/main.ts
@@ -32,9 +32,37 @@ const noteSchema = {
},
};
+type NoteSummary = Pick<
+ Note,
+ "id" | "title" | "lastModified"
+>;
+
+const noteSummarySchema = {
+ parse(value: unknown): NoteSummary {
+ if (
+ typeof value !== "object" ||
+ value === null ||
+ !("id" in value) ||
+ typeof value.id !== "string" ||
+ !("title" in value) ||
+ typeof value.title !== "string" ||
+ !("lastModified" in value) ||
+ typeof value.lastModified !== "number"
+ ) {
+ throw new Error("Invalid note summary");
+ }
+ return value as NoteSummary;
+ },
+};
+
const noteDefinition = defineCollection("notes", noteSchema, {
indexes: [
- defineIndex("by-title", ["title"]),
+ defineIndex(
+ "by-title",
+ ["title"],
+ "equality",
+ { include: ["lastModified"] },
+ ),
defineIndex(
"by-last-modified",
["lastModified"],
@@ -90,6 +118,10 @@ element("filter").addEventListener(
.where((note) => note.title.eq(title))
.orderBy((note) => note.lastModified.desc())
.take(50)
+ .select(
+ ["title", "lastModified"],
+ noteSummarySchema,
+ )
.get();
status.textContent =
`${result.documents.length} notes via ${result.plan}` +
@@ -130,14 +162,21 @@ async function renderAll() {
renderNotes(result.documents);
}
-function renderNotes(items: Note[]) {
+function renderNotes(
+ items: Array<
+ Note | NoteSummary
+ >,
+) {
notesOutput.replaceChildren(
...items.map((note) => {
const article = document.createElement("article");
const title = document.createElement("h2");
title.textContent = note.title;
const body = document.createElement("p");
- body.textContent = note.body || "No body";
+ body.textContent =
+ "body" in note
+ ? note.body || "No body"
+ : "Use Show all to load the note body.";
const metadata = document.createElement("small");
metadata.textContent = new Date(
note.lastModified,
diff --git a/tests/browser-client.test.ts b/tests/browser-client.test.ts
index 8561a72..ca92502 100644
--- a/tests/browser-client.test.ts
+++ b/tests/browser-client.test.ts
@@ -11,6 +11,7 @@ import {
import { ThimbleClient } from "../src/browser/client.js";
import type {
JsonObjectReader,
+ PointReadBundleReader,
RemoteJsonObject,
} from "../src/browser/remote-reader.js";
import { HttpObjectReadError } from "../src/browser/remote-reader.js";
@@ -25,6 +26,80 @@ import {
} from "../src/trie-protocol.js";
describe("ThimbleDB browser client", () => {
+ it("uses one cold read bundle and reuses its cached objects", async () => {
+ const fixture = trieFixture("products", "product-bundle");
+ const bundleReader = new FakeBundleReader({
+ status: "found",
+ bytes: 512,
+ bundle: {
+ collection: "products",
+ id: fixture.id,
+ revision: 1,
+ document: fixture.document,
+ objects: [...fixture.objects].map(([key, object]) => ({
+ key,
+ etag: object.etag,
+ value: structuredClone(object.value),
+ })),
+ layout: "trie",
+ },
+ });
+ const objectReader = new FakeReader(fixture.objects);
+ const cache = cacheFor("content", uniqueName());
+ const client = new ThimbleClient({
+ reader: objectReader,
+ bundleReader,
+ cache,
+ headTtlMs: 10_000,
+ channelName: uniqueName(),
+ });
+
+ await expect(
+ client.get("products", fixture.id),
+ ).resolves.toEqual(fixture.document);
+ await expect(
+ client.get("products", fixture.id),
+ ).resolves.toEqual(fixture.document);
+
+ expect(bundleReader.calls).toBe(1);
+ expect(objectReader.calls).toBe(0);
+ expect(client.metrics()).toMatchObject({
+ remoteReads: 1,
+ remoteBytes: 512,
+ bundleReads: 1,
+ bundleBytes: 512,
+ bundleFallbacks: 0,
+ });
+ });
+
+ it("falls back to individual objects when a bundle is unavailable", async () => {
+ const fixture = trieFixture("products", "product-fallback");
+ const bundleReader = new FakeBundleReader({
+ status: "fallback",
+ });
+ const objectReader = new FakeReader(fixture.objects);
+ const cache = cacheFor("content", uniqueName());
+ const client = new ThimbleClient({
+ reader: objectReader,
+ bundleReader,
+ cache,
+ headTtlMs: 10_000,
+ channelName: uniqueName(),
+ });
+
+ await expect(
+ client.get("products", fixture.id),
+ ).resolves.toEqual(fixture.document);
+
+ expect(bundleReader.calls).toBe(1);
+ expect(objectReader.calls).toBe(4);
+ expect(client.metrics()).toMatchObject({
+ remoteReads: 5,
+ bundleReads: 1,
+ bundleFallbacks: 1,
+ });
+ });
+
it("serves warm content reads from memory and revalidates HEAD", async () => {
const fixture = trieFixture("products", "product-00001");
const reader = new FakeReader(fixture.objects);
@@ -956,6 +1031,7 @@ class FakeReader implements JsonObjectReader {
if (this.offline) {
throw new Error("offline");
}
+
if (this.error) {
throw this.error;
}
@@ -980,6 +1056,25 @@ class FakeReader implements JsonObjectReader {
}
}
+class FakeBundleReader implements PointReadBundleReader {
+ calls = 0;
+
+ constructor(
+ private readonly result:
+ | {
+ status: "found";
+ bundle: import("../src/trie-protocol.js").TrieReadBundle;
+ bytes: number;
+ }
+ | { status: "fallback" },
+ ) {}
+
+ get() {
+ this.calls += 1;
+ return Promise.resolve(structuredClone(this.result));
+ }
+}
+
function cacheFor(
policy: "content" | "locations",
databaseName: string,
diff --git a/tests/browser-connect.test.ts b/tests/browser-connect.test.ts
index 1207eb1..deb7704 100644
--- a/tests/browser-connect.test.ts
+++ b/tests/browser-connect.test.ts
@@ -110,6 +110,77 @@ describe("browser connection factory", () => {
client.close();
});
+ it("uses an advertised point-read bundle endpoint", async () => {
+ const requests: string[] = [];
+ const client = await createThimbleClient({
+ configurationUrl: "https://app.example.test/api/config",
+ persistentCache: false,
+ fetchImplementation: async (input) => {
+ const url = String(input);
+ requests.push(url);
+ if (url.endsWith("/api/config")) {
+ return Response.json({
+ ...browserConfig(false),
+ readBundleBaseUrl: "/api/read-bundles",
+ });
+ }
+ if (
+ url.endsWith(
+ "/api/read-bundles/user%3Auser-1/notes/note-1",
+ )
+ ) {
+ return Response.json({
+ collection: "notes",
+ id: "note-1",
+ revision: 1,
+ document: {
+ id: "note-1",
+ title: "Bundled",
+ },
+ objects: [
+ {
+ key: "content-snapshot/notes/HEAD.json",
+ etag: "head",
+ value: {
+ revision: 1,
+ snapshotHash: "snapshot-one",
+ },
+ },
+ {
+ key:
+ "content-snapshot/notes/snapshots/" +
+ "snapshot-one.json",
+ etag: "snapshot",
+ value: {
+ documents: {
+ "note-1": {
+ id: "note-1",
+ title: "Bundled",
+ },
+ },
+ },
+ },
+ ],
+ layout: "snapshot",
+ });
+ }
+ return new Response(null, { status: 404 });
+ },
+ });
+
+ await expect(
+ client.get("notes", "note-1"),
+ ).resolves.toEqual({
+ id: "note-1",
+ title: "Bundled",
+ });
+ expect(requests).toEqual([
+ "https://app.example.test/api/config",
+ "https://app.example.test/api/read-bundles/" +
+ "user%3Auser-1/notes/note-1",
+ ]);
+ });
+
it("rejects malformed authority configuration", async () => {
await expect(
createThimbleClient({
diff --git a/tests/cloudflare-worker.test.ts b/tests/cloudflare-worker.test.ts
index 038fcde..8d8111e 100644
--- a/tests/cloudflare-worker.test.ts
+++ b/tests/cloudflare-worker.test.ts
@@ -170,6 +170,24 @@ describe("Cloudflare Worker request parsing", () => {
await expect(missingOrigin.json()).resolves.toMatchObject({
error: "origin_rejected",
});
+
+ const disabledBundle = await createCloudflareAuthority().fetch(
+ new Request(
+ "https://db.example.test/api/read-bundles/public/notes/note-1",
+ ),
+ environment as never,
+ );
+ expect(disabledBundle.status).toBe(404);
+
+ const bundle = await createCloudflareAuthority({
+ readBundles: true,
+ }).fetch(
+ new Request(
+ "https://db.example.test/api/read-bundles/public/notes/note-1",
+ ),
+ environment as never,
+ );
+ expect(bundle.status).toBe(401);
});
it("does not apply Studio catalog limits when Studio is disabled", async () => {
diff --git a/tests/collection.test.ts b/tests/collection.test.ts
index ed95054..9f1af14 100644
--- a/tests/collection.test.ts
+++ b/tests/collection.test.ts
@@ -121,6 +121,150 @@ describe("typed collections", () => {
],
limit: 10,
});
+
+ });
+
+ it("requests explicit typed projections without full schema parsing", async () => {
+ type DetailedNote = Note & {
+ body: string;
+ lastModified: number;
+ };
+ const queryDocuments = vi.fn().mockResolvedValue({
+ documents: [
+ {
+ id: "note-1",
+ title: "Projected",
+ },
+ ],
+ plan: "index",
+ indexName: "by-title",
+ scannedDocuments: 1,
+ });
+ const collection = new ThimbleCollection(
+ client({ queryDocuments }),
+ defineCollection("notes"),
+ );
+ const projectedSchema = {
+ parse(value: unknown) {
+ if (
+ typeof value !== "object" ||
+ value === null ||
+ !("id" in value) ||
+ typeof value.id !== "string" ||
+ !("title" in value) ||
+ typeof value.title !== "string"
+ ) {
+ throw new Error("Invalid projected note");
+ }
+ return value as Pick;
+ },
+ };
+ const builder = collection
+ .where((note) => note.title.eq("Projected"))
+ .select(["title"], projectedSchema);
+
+ await expect(builder.get()).resolves.toMatchObject({
+ documents: [
+ {
+ id: "note-1",
+ title: "Projected",
+ },
+ ],
+ });
+ expect(queryDocuments).toHaveBeenCalledWith(
+ "notes",
+ {
+ version: 1,
+ where: {
+ field: "title",
+ operator: "eq",
+ value: "Projected",
+ },
+ },
+ ["title"],
+ );
+ expect(builder.toJSON()).toEqual({
+ query: {
+ version: 1,
+ where: {
+ field: "title",
+ operator: "eq",
+ value: "Projected",
+ },
+ },
+ select: ["title"],
+ });
+ });
+
+ it("rejects projected values that fail their projection schema", async () => {
+ type DetailedNote = Note & {
+ body: string;
+ };
+ const collection = new ThimbleCollection(
+ client({
+ queryDocuments: async () => ({
+ documents: [
+ {
+ id: "note-1",
+ title: 123,
+ } as never,
+ ],
+ plan: "index",
+ indexName: "by-title",
+ scannedDocuments: 1,
+ }),
+ }),
+ defineCollection("notes"),
+ );
+
+ await expect(
+ collection
+ .where((note) => note.title.eq("Projected"))
+ .select(["title"], {
+ parse(value: unknown) {
+ if (
+ typeof value !== "object" ||
+ value === null ||
+ !("id" in value) ||
+ typeof value.id !== "string" ||
+ !("title" in value) ||
+ typeof value.title !== "string"
+ ) {
+ throw new Error("Invalid projected note");
+ }
+ return value as Pick<
+ DetailedNote,
+ "id" | "title"
+ >;
+ },
+ })
+ .get(),
+ ).rejects.toThrow(
+ "Projection validation failed in collection notes",
+ );
+ });
+
+ it("rejects prototype-sensitive projection fields", () => {
+ type FlexibleNote = Note & {
+ __proto__?: string;
+ };
+ const collection = new ThimbleCollection(
+ client({}),
+ defineCollection("notes"),
+ );
+
+ expect(() =>
+ collection
+ .where((note) => note.title.eq("example"))
+ .select(["__proto__"], {
+ parse(value: unknown) {
+ return value as Pick<
+ FlexibleNote,
+ "id" | "__proto__"
+ >;
+ },
+ }),
+ ).toThrow("unique safe non-ID fields");
});
it("surfaces collection and document context on validation failure", async () => {
diff --git a/tests/e2e/authenticated-store.spec.ts b/tests/e2e/authenticated-store.spec.ts
index f604601..b4815c6 100644
--- a/tests/e2e/authenticated-store.spec.ts
+++ b/tests/e2e/authenticated-store.spec.ts
@@ -5,6 +5,12 @@ test("authenticates externally, reads, writes, persists cache, and logs out", as
browserName,
request,
}) => {
+ const bundleRequests: string[] = [];
+ page.on("request", (request) => {
+ if (request.url().includes("/api/read-bundles/")) {
+ bundleRequests.push(request.url());
+ }
+ });
const subject = `${browserName}-${crypto.randomUUID()}`;
const token = await request
.get(
@@ -50,10 +56,52 @@ test("authenticates externally, reads, writes, persists cache, and logs out", as
'"products": 128',
);
+ const oversizedProjection = await page.evaluate(async () => {
+ const config = await fetch("/api/config", {
+ credentials: "same-origin",
+ cache: "no-store",
+ }).then((response) => response.json()) as {
+ csrfToken: string;
+ layoutGeneration: string;
+ scope: { id: string };
+ };
+ const response = await fetch(
+ "/api/collections/products/documents/oversized-cover",
+ {
+ method: "POST",
+ credentials: "same-origin",
+ headers: {
+ "content-type": "application/json",
+ "x-thimble-csrf": config.csrfToken,
+ "x-thimble-scope": config.scope.id,
+ "x-thimble-layout-generation":
+ config.layoutGeneration,
+ },
+ body: JSON.stringify({
+ id: "oversized-cover",
+ sku: "OVERSIZED",
+ name: "x".repeat(70 * 1024),
+ priceCents: 1,
+ }),
+ },
+ );
+ return {
+ status: response.status,
+ body: await response.json(),
+ };
+ });
+ expect(oversizedProjection).toMatchObject({
+ status: 413,
+ body: {
+ error: "secondary_index_too_large",
+ },
+ });
+
await page.locator("#read-product").click();
await expect(page.locator("#product-output")).toContainText(
"product-00000",
);
+ expect(bundleRequests).toHaveLength(1);
await page.locator("#delete-product").click();
await expect(page.locator("#status")).toContainText(
diff --git a/tests/read-bundle.test.ts b/tests/read-bundle.test.ts
new file mode 100644
index 0000000..097dd0b
--- /dev/null
+++ b/tests/read-bundle.test.ts
@@ -0,0 +1,135 @@
+import { mkdtemp, rm } from "node:fs/promises";
+import os from "node:os";
+import path from "node:path";
+import { describe, expect, it } from "vitest";
+import { BoundedReadError } from "../src/core.js";
+import { ContentAddressedTrieEngine } from "../src/engines/content-trie.js";
+import { ImmutableSnapshotEngine } from "../src/engines/immutable-snapshot.js";
+import { LocalObjectStore } from "../src/providers/local.js";
+import { readPointBundle } from "../src/read-bundle.js";
+
+describe("bounded point-read bundles", () => {
+ it.each(["trie", "snapshot"] as const)(
+ "returns the current %s document and cache objects",
+ async (layout) => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), `thimble-bundle-${layout}-`),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const engine =
+ layout === "trie"
+ ? new ContentAddressedTrieEngine(store)
+ : new ImmutableSnapshotEngine(store);
+ await engine.put("notes", "note-1", {
+ id: "note-1",
+ title: "Bundled",
+ });
+
+ const bundle = await engine.readBundle(
+ "notes",
+ "note-1",
+ {
+ maxObjects: 4,
+ maxDecodedBytes: 1024 * 1024,
+ },
+ );
+
+ expect(bundle.layout).toBe(layout);
+ expect(bundle.document).toEqual({
+ id: "note-1",
+ title: "Bundled",
+ });
+ expect(bundle.objects).toHaveLength(
+ layout === "trie" ? 4 : 2,
+ );
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ },
+ );
+
+ it("rejects an oversized snapshot before loading its page", async () => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), "thimble-bundle-snapshot-limit-"),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const engine = new ImmutableSnapshotEngine(store);
+ await engine.put("notes", "note-1", {
+ id: "note-1",
+ body: "x".repeat(4_096),
+ });
+ let pageReads = 0;
+ const originalGet = store.get.bind(store);
+ store.get = async (key) => {
+ if (key.includes("/snapshots/")) {
+ pageReads += 1;
+ }
+ return originalGet(key);
+ };
+
+ await expect(
+ engine.readBundle("notes", "note-1", {
+ maxObjects: 4,
+ maxDecodedBytes: 100,
+ }),
+ ).rejects.toBeInstanceOf(BoundedReadError);
+ expect(pageReads).toBe(0);
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ });
+
+ it("rejects a trie bundle before loading a disallowed leaf", async () => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), "thimble-bundle-trie-limit-"),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const engine = new ContentAddressedTrieEngine(store);
+ await engine.put("notes", "note-1", {
+ id: "note-1",
+ title: "Bounded",
+ });
+ let nodeReads = 0;
+ const originalGet = store.get.bind(store);
+ store.get = async (key) => {
+ if (key.includes("/nodes/")) {
+ nodeReads += 1;
+ }
+ return originalGet(key);
+ };
+
+ await expect(
+ engine.readBundle("notes", "note-1", {
+ maxObjects: 3,
+ maxDecodedBytes: 1024 * 1024,
+ }),
+ ).rejects.toBeInstanceOf(BoundedReadError);
+ expect(nodeReads).toBe(2);
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ });
+
+ it("bounds the complete serialized response including the document copy", async () => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), "thimble-bundle-response-limit-"),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const engine = new ImmutableSnapshotEngine(store);
+ await engine.put("notes", "note-1", {
+ id: "note-1",
+ body: "x".repeat(2_200_000),
+ });
+
+ await expect(
+ readPointBundle(engine, "notes", "note-1"),
+ ).rejects.toBeInstanceOf(BoundedReadError);
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ });
+});
diff --git a/tests/secondary-index.test.ts b/tests/secondary-index.test.ts
index a545f88..16bf71f 100644
--- a/tests/secondary-index.test.ts
+++ b/tests/secondary-index.test.ts
@@ -23,6 +23,7 @@ import {
type RemoteJsonObject,
type SnapshotHead,
type TrieHead,
+ MAX_SECONDARY_INDEX_PAGE_BYTES,
} from "../src/index.js";
import { LocalObjectStore } from "../src/providers/local.js";
@@ -33,6 +34,49 @@ type Note = {
tags?: string[];
};
+const noteSummarySchema = {
+ parse(value: unknown): Pick<
+ Note,
+ "id" | "title" | "lastModified"
+ > {
+ if (
+ typeof value !== "object" ||
+ value === null ||
+ !("id" in value) ||
+ typeof value.id !== "string" ||
+ !("title" in value) ||
+ typeof value.title !== "string" ||
+ !("lastModified" in value) ||
+ typeof value.lastModified !== "number"
+ ) {
+ throw new Error("Invalid note summary");
+ }
+ return value as Pick<
+ Note,
+ "id" | "title" | "lastModified"
+ >;
+ },
+};
+
+const noteTagsSchema = {
+ parse(value: unknown): Pick {
+ if (
+ typeof value !== "object" ||
+ value === null ||
+ !("id" in value) ||
+ typeof value.id !== "string" ||
+ !("title" in value) ||
+ typeof value.title !== "string" ||
+ !("tags" in value) ||
+ !Array.isArray(value.tags) ||
+ !value.tags.every((tag) => typeof tag === "string")
+ ) {
+ throw new Error("Invalid note tags");
+ }
+ return value as Pick;
+ },
+};
+
const indexes: CollectionIndexConfiguration = {
notes: [
{
@@ -209,6 +253,7 @@ describe("secondary indexes", () => {
},
collectionIndexes: indexes,
});
+
const result = await client
.collection("notes")
.where((note) => note.title.eq("same"))
@@ -225,6 +270,475 @@ describe("secondary indexes", () => {
}
});
+ it.each(["snapshot", "trie"] as const)(
+ "maintains explicit covering projections for %s indexes",
+ async (layout) => {
+ const directory = await mkdtemp(
+ path.join(
+ os.tmpdir(),
+ `thimble-covering-${layout}-`,
+ ),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const covering: CollectionIndexConfiguration = {
+ notes: [
+ defineIndex(
+ "by-title",
+ ["title"],
+ "equality",
+ { include: ["lastModified"] },
+ ),
+ ],
+ };
+ const engine =
+ layout === "snapshot"
+ ? new ImmutableSnapshotEngine(
+ store,
+ 40,
+ address,
+ false,
+ covering,
+ )
+ : new ContentAddressedTrieEngine(
+ store,
+ 40,
+ address,
+ false,
+ covering,
+ );
+ await engine.putMany("notes", [
+ {
+ id: "note-1",
+ title: "same",
+ lastModified: 1,
+ tags: ["first"],
+ },
+ {
+ id: "note-2",
+ title: "same",
+ lastModified: 2,
+ tags: ["second"],
+ },
+ ]);
+
+ let page = await coveringPage(store, layout);
+ expect(page.projections).toEqual({
+ "note-1": {
+ id: "note-1",
+ title: "same",
+ lastModified: 1,
+ },
+ "note-2": {
+ id: "note-2",
+ title: "same",
+ lastModified: 2,
+ },
+ });
+
+ await engine.put("notes", "note-1", {
+ id: "note-1",
+ title: "same",
+ lastModified: 3,
+ tags: ["updated"],
+ });
+ page = await coveringPage(store, layout);
+ expect(page.projections?.["note-1"]).toEqual({
+ id: "note-1",
+ title: "same",
+ lastModified: 3,
+ });
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ },
+ );
+
+ it("serves an explicit projection without loading full documents", async () => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), "thimble-covering-query-"),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const covering: CollectionIndexConfiguration = {
+ notes: [
+ defineIndex(
+ "by-title",
+ ["title"],
+ "equality",
+ { include: ["lastModified"] },
+ ),
+ ],
+ };
+ const engine = new ImmutableSnapshotEngine(
+ store,
+ 40,
+ address,
+ false,
+ covering,
+ );
+ await engine.putMany("notes", [
+ {
+ id: "note-1",
+ title: "same",
+ lastModified: 2,
+ tags: ["private"],
+ },
+ {
+ id: "note-2",
+ title: "same",
+ lastModified: 1,
+ tags: ["private"],
+ },
+ ]);
+ const calls = new Map();
+ const client = new ThimbleClient({
+ reader: objectReader(store, calls),
+ cache: new TieredObjectCache(
+ new MemoryObjectCache(),
+ new NullPersistentCache(),
+ ),
+ headTtlMs: 10_000,
+ collectionLayouts: { notes: "snapshot" },
+ collectionIndexes: covering,
+ });
+
+ const result = await client
+ .collection("notes")
+ .where((note) => note.title.eq("same"))
+ .orderBy((note) => note.lastModified.asc())
+ .select(
+ ["title", "lastModified"],
+ noteSummarySchema,
+ )
+ .get();
+ const head = (await readHeadForIndex(
+ store,
+ "snapshot",
+ "by-title",
+ )) as SnapshotHead;
+
+ expect(result.plan).toBe("index");
+ expect(result.documents).toEqual([
+ {
+ id: "note-2",
+ title: "same",
+ lastModified: 1,
+ },
+ {
+ id: "note-1",
+ title: "same",
+ lastModified: 2,
+ },
+ ]);
+ expect(
+ calls.get(snapshotPageKey("notes", head.snapshotHash!)),
+ ).toBeUndefined();
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ });
+
+ it("loads full documents when a selected field is not covered", async () => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), "thimble-covering-fallback-"),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const covering: CollectionIndexConfiguration = {
+ notes: [
+ defineIndex(
+ "by-title",
+ ["title"],
+ "equality",
+ { include: ["lastModified"] },
+ ),
+ ],
+ };
+ const engine = new ImmutableSnapshotEngine(
+ store,
+ 40,
+ address,
+ false,
+ covering,
+ );
+ await engine.put("notes", "note-1", {
+ id: "note-1",
+ title: "same",
+ lastModified: 1,
+ tags: ["required"],
+ });
+
+ const calls = new Map();
+ const client = new ThimbleClient({
+ reader: objectReader(store, calls),
+ cache: new TieredObjectCache(
+ new MemoryObjectCache(),
+ new NullPersistentCache(),
+ ),
+ headTtlMs: 10_000,
+ collectionLayouts: { notes: "snapshot" },
+ collectionIndexes: covering,
+ });
+
+ const result = await client
+ .collection("notes")
+ .where((note) => note.title.eq("same"))
+ .select(
+ ["title", "tags"],
+ noteTagsSchema,
+ )
+ .get();
+ const head = (await readHeadForIndex(
+ store,
+ "snapshot",
+ "by-title",
+ )) as SnapshotHead;
+
+ expect(result.documents).toEqual([
+ {
+ id: "note-1",
+ title: "same",
+ tags: ["required"],
+ },
+ ]);
+ expect(
+ calls.get(snapshotPageKey("notes", head.snapshotHash!)),
+ ).toBe(1);
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ });
+
+ it.each(["snapshot", "trie"] as const)(
+ "rejects oversized %s projections before writing content objects",
+ async (layout) => {
+ const directory = await mkdtemp(
+ path.join(
+ os.tmpdir(),
+ `thimble-covering-size-${layout}-`,
+ ),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const covering: CollectionIndexConfiguration = {
+ notes: [
+ defineIndex(
+ "by-title",
+ ["title"],
+ "equality",
+ { include: ["tags"] },
+ ),
+ ],
+ };
+ const engine =
+ layout === "snapshot"
+ ? new ImmutableSnapshotEngine(
+ store,
+ 40,
+ address,
+ false,
+ covering,
+ )
+ : new ContentAddressedTrieEngine(
+ store,
+ 40,
+ address,
+ false,
+ covering,
+ );
+
+ await expect(
+ engine.put("notes", "note-1", {
+ id: "note-1",
+ title: "large",
+ lastModified: 1,
+ tags: ["x".repeat(70 * 1024)],
+ }),
+ ).rejects.toThrow(
+ "covering projection exceeds 65536 decoded bytes",
+ );
+ expect(
+ await store.list(
+ layout === "snapshot"
+ ? "content-snapshot/notes/"
+ : "content-trie/notes/",
+ ),
+ ).toEqual([]);
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ },
+ );
+
+ it("rejects aggregate covering index pages above four MiB", async () => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), "thimble-covering-page-size-"),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const covering: CollectionIndexConfiguration = {
+ notes: [
+ defineIndex(
+ "by-title",
+ ["title"],
+ "equality",
+ { include: ["tags"] },
+ ),
+ ],
+ };
+ const engine = new ImmutableSnapshotEngine(
+ store,
+ 40,
+ address,
+ false,
+ covering,
+ );
+
+ await expect(
+ engine.putMany(
+ "notes",
+ Array.from({ length: 70 }, (_, index) => ({
+ id: `note-${index}`,
+ title: `title-${index}`,
+ lastModified: index,
+ tags: ["x".repeat(60_000)],
+ })),
+ ),
+ ).rejects.toThrow(
+ `exceeds ${MAX_SECONDARY_INDEX_PAGE_BYTES} decoded bytes`,
+ );
+ expect(
+ await store.list("content-snapshot/notes/"),
+ ).toEqual([]);
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ });
+
+ it("projects stale covering definitions after bounded scan fallback", async () => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), "thimble-covering-stale-"),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const oldIndexes: CollectionIndexConfiguration = {
+ notes: [
+ defineIndex("by-title", ["title"]),
+ ],
+ };
+ await new ImmutableSnapshotEngine(
+ store,
+ 40,
+ address,
+ false,
+ oldIndexes,
+ ).put("notes", "note-1", {
+ id: "note-1",
+ title: "same",
+ lastModified: 1,
+ secret: "must-not-return",
+ });
+ const covering: CollectionIndexConfiguration = {
+ notes: [
+ defineIndex(
+ "by-title",
+ ["title"],
+ "equality",
+ { include: ["lastModified"] },
+ ),
+ ],
+ };
+ const result = await browserClient(
+ store,
+ "snapshot",
+ covering,
+ )
+ .collection("notes")
+ .where((note) => note.title.eq("same"))
+ .select(["title"], {
+ parse(value: unknown) {
+ return value as Pick;
+ },
+ })
+ .get();
+
+ expect(result.plan).toBe("scan");
+ expect(result.documents).toEqual([
+ {
+ id: "note-1",
+ title: "same",
+ },
+ ]);
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ });
+
+ it("bypasses oversized index pages from authenticated HEAD metadata", async () => {
+ const directory = await mkdtemp(
+ path.join(os.tmpdir(), "thimble-index-client-size-"),
+ );
+ try {
+ const store = new LocalObjectStore(directory);
+ const engine = new ImmutableSnapshotEngine(
+ store,
+ 40,
+ address,
+ false,
+ indexes,
+ );
+ await engine.put("notes", "note-1", {
+ id: "note-1",
+ title: "same",
+ lastModified: 1,
+ });
+ const headObject = await store.get(
+ snapshotHeadKey("notes"),
+ );
+ const head = JSON.parse(
+ Buffer.from(headObject!.bytes).toString("utf8"),
+ ) as SnapshotHead;
+ head.indexes!["by-title"]!.decodedBytes =
+ MAX_SECONDARY_INDEX_PAGE_BYTES + 1;
+ await store.put(
+ snapshotHeadKey("notes"),
+ Buffer.from(JSON.stringify(head)),
+ { ifMatch: headObject!.etag },
+ );
+ const calls = new Map();
+ const client = new ThimbleClient({
+ reader: objectReader(store, calls),
+ cache: new TieredObjectCache(
+ new MemoryObjectCache(),
+ new NullPersistentCache(),
+ ),
+ headTtlMs: 10_000,
+ collectionLayouts: { notes: "snapshot" },
+ collectionIndexes: indexes,
+ });
+
+ const result = await client
+ .collection("notes")
+ .where((note) => note.title.eq("same"))
+ .get();
+ const reference = head.indexes!["by-title"]!;
+
+ expect(result.plan).toBe("scan");
+ expect(
+ calls.get(
+ snapshotIndexKey(
+ "notes",
+ "by-title",
+ reference.hash,
+ ),
+ ),
+ ).toBeUndefined();
+ } finally {
+ await rm(directory, { recursive: true, force: true });
+ }
+ });
+
it("falls back for a stale definition and rebuilds it on the next trie write", async () => {
const directory = await mkdtemp(
path.join(os.tmpdir(), "thimble-index-definition-"),
@@ -327,6 +841,20 @@ describe("secondary indexes", () => {
expect(() =>
defineIndex("duplicate", ["title", "title"]),
).toThrow("Invalid secondary index configuration");
+ expect(() =>
+ defineIndex(
+ "invalid-cover",
+ ["title"],
+ "equality",
+ { include: ["title"] },
+ ),
+ ).toThrow("Invalid secondary index configuration");
+ expect(() =>
+ defineIndex(
+ "prototype-field",
+ ["__proto__" as keyof Note],
+ ),
+ ).toThrow("Invalid secondary index configuration");
expect(() =>
secondaryIndexPageFromJson({
@@ -348,6 +876,24 @@ describe("secondary indexes", () => {
],
}),
).toThrow("duplicate document");
+
+ expect(() =>
+ secondaryIndexPageFromJson({
+ version: 1,
+ definition: {
+ name: "by-title",
+ fields: ["title"],
+ mode: "equality",
+ include: ["lastModified"],
+ },
+ entries: [
+ {
+ values: ["First"],
+ ids: ["note-1"],
+ },
+ ],
+ }),
+ ).toThrow("missing covering projections");
});
it("does not use a sparse range index for ordering alone", () => {
@@ -450,6 +996,7 @@ describe("secondary indexes", () => {
} finally {
await rm(directory, { recursive: true, force: true });
}
+
});
it.each(["snapshot", "trie"] as const)(
@@ -633,6 +1180,53 @@ describe("secondary indexes", () => {
);
});
+async function readHeadForIndex(
+ store: LocalObjectStore,
+ layout: "snapshot" | "trie",
+ indexName: string,
+): Promise {
+ const key =
+ layout === "snapshot"
+ ? snapshotHeadKey("notes")
+ : trieHeadKey("notes");
+ const object = await store.get(key);
+ if (!object) {
+ throw new Error("Collection head is missing");
+ }
+ const head = JSON.parse(
+ Buffer.from(object.bytes).toString("utf8"),
+ ) as SnapshotHead | TrieHead;
+ if (!head.indexes?.[indexName]) {
+ throw new Error(`Index ${indexName} is missing`);
+ }
+ return head;
+}
+
+async function coveringPage(
+ store: LocalObjectStore,
+ layout: "snapshot" | "trie",
+) {
+ const head = await readHeadForIndex(
+ store,
+ layout,
+ "by-title",
+ );
+ const reference = head.indexes!["by-title"]!;
+ const key =
+ layout === "snapshot"
+ ? snapshotIndexKey("notes", "by-title", reference.hash)
+ : trieIndexKey("notes", "by-title", reference.hash);
+ const object = await store.get(key);
+ if (!object) {
+ throw new Error("Covering index page is missing");
+ }
+ return secondaryIndexPageFromJson(
+ JSON.parse(
+ Buffer.from(object.bytes).toString("utf8"),
+ ) as JsonValue,
+ );
+}
+
async function readHead(
store: LocalObjectStore,
layout: "snapshot" | "trie",