import { sha256Hex } from "@/worker/auth/crypto"; import type { DavChangeType } from "@/worker/db/types"; import type { CollectionDomainError, DomainChange, DomainChangesSince, DomainCollection, DomainConditionalInput, DomainDeleteObjectInput, DomainGetObjectInput, DomainMultigetInput, DomainObject, DomainObjectResources, DomainPutObjectInput, DomainSyncInput, ParsedDomainObject, } from "@/worker/objects/common/collection-types"; export interface CollectionObjectSpec< TDb, TTx, TCollection extends DomainCollection, TObject extends DomainObject, TPutInput extends DomainPutObjectInput, TParsed extends ParsedDomainObject, TError extends CollectionDomainError, > { ensureCollection(db: TDb, subjectId: string, nowMs: number, collectionName: string): TCollection | null; objectByName(db: TDb, collectionId: string, objectName: string): TObject | undefined; objectByUid(db: TDb, collectionId: string, uid: string): TObject | undefined; objectResourceFromHref( db: TDb, collection: TCollection, href: string, acceptedOrigins?: string[], ): { href: string } | null; objectResources(db: TDb, collection: TCollection, maxResults: number): DomainObjectResources<{ href: string }>; changesSince( db: TDb, collection: TCollection, token: string | null, maxResults: number, ): DomainChangesSince; latestChangeSeq(db: TDb): number; syncTokenSeq(token: string | null): number | null; syncToken(seq: number): string; objectHref(collectionName: string, objectName: string): string; parseObject(input: TPutInput, collection: TCollection): TParsed | Promise; validateParsed?(parsed: TParsed, input: TPutInput, collection: TCollection): TError | null; createObjectId(): string; writePrecondition( db: TDb, input: DomainConditionalInput, targetHref: string, targetEtag: string | null, ): TError | null; transaction(db: TDb, callback: (tx: TTx) => void): void; saveObject( tx: TTx, input: { collection: TCollection; objectName: string; existing: TObject | undefined; parsed: TParsed; objectId: string; etag: string; version: number; nowMs: number; }, ): void; deleteObject(tx: TTx, object: TObject): void; recordChange(tx: TTx, collectionId: string, href: string, changeType: DavChangeType, nowMs: number): void; collectionNotFound(): TError; objectNotFound(): TError; parseFailure(input: TPutInput, cause: unknown): TError; uidAlreadyExists(): TError; uidCannotChange(): TError; multigetLimitExceeded(): TError; invalidSyncToken(): TError; } export function getCollectionObject< TDb, TTx, TCollection extends DomainCollection, TObject extends DomainObject, TPutInput extends DomainPutObjectInput, TParsed extends ParsedDomainObject, TError extends CollectionDomainError, >( db: TDb, input: DomainGetObjectInput, spec: CollectionObjectSpec, ): { ok: true; object: TObject } | TError { const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName); if (!collection) return spec.collectionNotFound(); const object = spec.objectByName(db, collection.id, input.objectName); if (!object) return spec.objectNotFound(); return { ok: true, object }; } export async function putCollectionObject< TDb, TTx, TCollection extends DomainCollection, TObject extends DomainObject, TPutInput extends DomainPutObjectInput, TParsed extends ParsedDomainObject, TError extends CollectionDomainError, >( db: TDb, input: TPutInput, spec: CollectionObjectSpec, ): Promise<{ ok: true; created: boolean; object: TObject } | TError> { const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName); if (!collection) return spec.collectionNotFound(); let parsed: TParsed; try { parsed = await spec.parseObject(input, collection); } catch (cause) { return spec.parseFailure(input, cause); } const parsedFailure = spec.validateParsed?.(parsed, input, collection); if (parsedFailure) return parsedFailure; const existing = spec.objectByName(db, collection.id, input.objectName); const uidOwner = spec.objectByUid(db, collection.id, parsed.uid); if (uidOwner && uidOwner.name !== input.objectName) return spec.uidAlreadyExists(); if (existing && existing.uid !== parsed.uid) return spec.uidCannotChange(); const href = spec.objectHref(collection.name, input.objectName); const precondition = spec.writePrecondition(db, input, href, existing?.etag ?? null); if (precondition) return precondition; const version = existing ? existing.version + 1 : 1; const etag = `"${await sha256Hex(`${parsed.body}\n${version}`)}"`; const objectId = existing?.id ?? spec.createObjectId(); const changeType = existing ? "updated" : "created"; spec.transaction(db, (tx) => { spec.saveObject(tx, { collection, objectName: input.objectName, existing, parsed, objectId, etag, version, nowMs: input.nowMs, }); spec.recordChange(tx, collection.id, href, changeType, input.nowMs); }); return { ok: true, created: !existing, object: spec.objectByName(db, collection.id, input.objectName)! }; } export function deleteCollectionObject< TDb, TTx, TCollection extends DomainCollection, TObject extends DomainObject, TPutInput extends DomainPutObjectInput, TParsed extends ParsedDomainObject, TError extends CollectionDomainError, >( db: TDb, input: DomainDeleteObjectInput, spec: CollectionObjectSpec, ): { ok: true } | TError { const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName); if (!collection) return spec.collectionNotFound(); const object = spec.objectByName(db, collection.id, input.objectName); if (!object) return spec.objectNotFound(); const href = spec.objectHref(collection.name, input.objectName); const precondition = spec.writePrecondition(db, input, href, object.etag); if (precondition) return precondition; spec.transaction(db, (tx) => { spec.deleteObject(tx, object); spec.recordChange(tx, collection.id, spec.objectHref(collection.name, object.name), "deleted", input.nowMs); }); return { ok: true }; } export function collectionMultiget< TDb, TTx, TCollection extends DomainCollection, TObject extends DomainObject, TPutInput extends DomainPutObjectInput, TParsed extends ParsedDomainObject, TError extends CollectionDomainError, TResource extends { href: string }, >( db: TDb, input: DomainMultigetInput, spec: Omit< CollectionObjectSpec, "objectResourceFromHref" | "objectResources" | "changesSince" > & { objectResourceFromHref( db: TDb, collection: TCollection, href: string, acceptedOrigins?: string[], ): TResource | null; }, ): { ok: true; requested: { href: string; resource: TResource | null }[] } | TError { const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName); if (!collection) return spec.collectionNotFound(); if (input.hrefs.length > input.maxResults) return spec.multigetLimitExceeded(); return { ok: true, requested: input.hrefs.map((href) => ({ href, resource: spec.objectResourceFromHref(db, collection, href, input.acceptedOrigins), })), }; } export function syncCollectionDomain< TDb, TTx, TCollection extends DomainCollection, TObject extends DomainObject, TPutInput extends DomainPutObjectInput, TParsed extends ParsedDomainObject, TError extends CollectionDomainError, TResource extends { href: string }, TChange extends DomainChange, >( db: TDb, input: DomainSyncInput, spec: Omit< CollectionObjectSpec, "objectResourceFromHref" | "objectResources" | "changesSince" > & { objectResourceFromHref(db: TDb, collection: TCollection, href: string): TResource | null; objectResources(db: TDb, collection: TCollection, maxResults: number): DomainObjectResources; changesSince( db: TDb, collection: TCollection, token: string | null, maxResults: number, ): DomainChangesSince; }, ): { ok: true; changes: TChange[]; resources: TResource[]; syncToken: string; truncated: boolean } | TError { const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName); if (!collection) return spec.collectionNotFound(); const since = spec.syncTokenSeq(input.token); if (since === null) return spec.invalidSyncToken(); if (!input.token) { const initial = spec.objectResources(db, collection, input.maxResults); const syncTokenValue = initial.truncated ? spec.syncToken(0) : spec.syncToken(spec.latestChangeSeq(db)); return { ok: true, changes: [], resources: initial.resources, syncToken: syncTokenValue, truncated: initial.truncated, }; } const { changes, truncated } = spec.changesSince(db, collection, input.token, input.maxResults); const lastReturnedSeq = changes.length > 0 ? changes[changes.length - 1]!.seq : since; const tokenSeq = truncated ? lastReturnedSeq : spec.latestChangeSeq(db); return { ok: true, changes, resources: changes .filter((change) => change.changeType !== "deleted") .map((change) => spec.objectResourceFromHref(db, collection, change.href)) .filter((resource): resource is TResource => resource !== null), syncToken: spec.syncToken(tokenSeq), truncated, }; }