Skip to content
File

Blob: src/worker/objects/common/collection-object-domain.ts

typescript281 lines
1import { sha256Hex } from "@/worker/auth/crypto";
2import type { DavChangeType } from "@/worker/db/types";
3import type {
4 CollectionDomainError,
5 DomainChange,
6 DomainChangesSince,
7 DomainCollection,
8 DomainConditionalInput,
9 DomainDeleteObjectInput,
10 DomainGetObjectInput,
11 DomainMultigetInput,
12 DomainObject,
13 DomainObjectResources,
14 DomainPutObjectInput,
15 DomainSyncInput,
16 ParsedDomainObject,
17} from "@/worker/objects/common/collection-types";
18 
19export interface CollectionObjectSpec<
20 TDb,
21 TTx,
22 TCollection extends DomainCollection,
23 TObject extends DomainObject,
24 TPutInput extends DomainPutObjectInput,
25 TParsed extends ParsedDomainObject,
26 TError extends CollectionDomainError,
27> {
28 ensureCollection(db: TDb, subjectId: string, nowMs: number, collectionName: string): TCollection | null;
29 objectByName(db: TDb, collectionId: string, objectName: string): TObject | undefined;
30 objectByUid(db: TDb, collectionId: string, uid: string): TObject | undefined;
31 objectResourceFromHref(
32 db: TDb,
33 collection: TCollection,
34 href: string,
35 acceptedOrigins?: string[],
36 ): { href: string } | null;
37 objectResources(db: TDb, collection: TCollection, maxResults: number): DomainObjectResources<{ href: string }>;
38 changesSince(
39 db: TDb,
40 collection: TCollection,
41 token: string | null,
42 maxResults: number,
43 ): DomainChangesSince<DomainChange>;
44 latestChangeSeq(db: TDb): number;
45 syncTokenSeq(token: string | null): number | null;
46 syncToken(seq: number): string;
47 objectHref(collectionName: string, objectName: string): string;
48 parseObject(input: TPutInput, collection: TCollection): TParsed | Promise<TParsed>;
49 validateParsed?(parsed: TParsed, input: TPutInput, collection: TCollection): TError | null;
50 createObjectId(): string;
51 writePrecondition(
52 db: TDb,
53 input: DomainConditionalInput,
54 targetHref: string,
55 targetEtag: string | null,
56 ): TError | null;
57 transaction(db: TDb, callback: (tx: TTx) => void): void;
58 saveObject(
59 tx: TTx,
60 input: {
61 collection: TCollection;
62 objectName: string;
63 existing: TObject | undefined;
64 parsed: TParsed;
65 objectId: string;
66 etag: string;
67 version: number;
68 nowMs: number;
69 },
70 ): void;
71 deleteObject(tx: TTx, object: TObject): void;
72 recordChange(tx: TTx, collectionId: string, href: string, changeType: DavChangeType, nowMs: number): void;
73 collectionNotFound(): TError;
74 objectNotFound(): TError;
75 parseFailure(input: TPutInput, cause: unknown): TError;
76 uidAlreadyExists(): TError;
77 uidCannotChange(): TError;
78 multigetLimitExceeded(): TError;
79 invalidSyncToken(): TError;
80}
81 
82export function getCollectionObject<
83 TDb,
84 TTx,
85 TCollection extends DomainCollection,
86 TObject extends DomainObject,
87 TPutInput extends DomainPutObjectInput,
88 TParsed extends ParsedDomainObject,
89 TError extends CollectionDomainError,
90>(
91 db: TDb,
92 input: DomainGetObjectInput,
93 spec: CollectionObjectSpec<TDb, TTx, TCollection, TObject, TPutInput, TParsed, TError>,
94): { ok: true; object: TObject } | TError {
95 const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName);
96 if (!collection) return spec.collectionNotFound();
97 const object = spec.objectByName(db, collection.id, input.objectName);
98 if (!object) return spec.objectNotFound();
99 return { ok: true, object };
100}
101 
102export async function putCollectionObject<
103 TDb,
104 TTx,
105 TCollection extends DomainCollection,
106 TObject extends DomainObject,
107 TPutInput extends DomainPutObjectInput,
108 TParsed extends ParsedDomainObject,
109 TError extends CollectionDomainError,
110>(
111 db: TDb,
112 input: TPutInput,
113 spec: CollectionObjectSpec<TDb, TTx, TCollection, TObject, TPutInput, TParsed, TError>,
114): Promise<{ ok: true; created: boolean; object: TObject } | TError> {
115 const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName);
116 if (!collection) return spec.collectionNotFound();
117 
118 let parsed: TParsed;
119 try {
120 parsed = await spec.parseObject(input, collection);
121 } catch (cause) {
122 return spec.parseFailure(input, cause);
123 }
124 const parsedFailure = spec.validateParsed?.(parsed, input, collection);
125 if (parsedFailure) return parsedFailure;
126 
127 const existing = spec.objectByName(db, collection.id, input.objectName);
128 const uidOwner = spec.objectByUid(db, collection.id, parsed.uid);
129 if (uidOwner && uidOwner.name !== input.objectName) return spec.uidAlreadyExists();
130 if (existing && existing.uid !== parsed.uid) return spec.uidCannotChange();
131 
132 const href = spec.objectHref(collection.name, input.objectName);
133 const precondition = spec.writePrecondition(db, input, href, existing?.etag ?? null);
134 if (precondition) return precondition;
135 
136 const version = existing ? existing.version + 1 : 1;
137 const etag = `"${await sha256Hex(`${parsed.body}\n${version}`)}"`;
138 const objectId = existing?.id ?? spec.createObjectId();
139 const changeType = existing ? "updated" : "created";
140 
141 spec.transaction(db, (tx) => {
142 spec.saveObject(tx, {
143 collection,
144 objectName: input.objectName,
145 existing,
146 parsed,
147 objectId,
148 etag,
149 version,
150 nowMs: input.nowMs,
151 });
152 spec.recordChange(tx, collection.id, href, changeType, input.nowMs);
153 });
154 
155 return { ok: true, created: !existing, object: spec.objectByName(db, collection.id, input.objectName)! };
156}
157 
158export function deleteCollectionObject<
159 TDb,
160 TTx,
161 TCollection extends DomainCollection,
162 TObject extends DomainObject,
163 TPutInput extends DomainPutObjectInput,
164 TParsed extends ParsedDomainObject,
165 TError extends CollectionDomainError,
166>(
167 db: TDb,
168 input: DomainDeleteObjectInput,
169 spec: CollectionObjectSpec<TDb, TTx, TCollection, TObject, TPutInput, TParsed, TError>,
170): { ok: true } | TError {
171 const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName);
172 if (!collection) return spec.collectionNotFound();
173 const object = spec.objectByName(db, collection.id, input.objectName);
174 if (!object) return spec.objectNotFound();
175 
176 const href = spec.objectHref(collection.name, input.objectName);
177 const precondition = spec.writePrecondition(db, input, href, object.etag);
178 if (precondition) return precondition;
179 
180 spec.transaction(db, (tx) => {
181 spec.deleteObject(tx, object);
182 spec.recordChange(tx, collection.id, spec.objectHref(collection.name, object.name), "deleted", input.nowMs);
183 });
184 return { ok: true };
185}
186 
187export function collectionMultiget<
188 TDb,
189 TTx,
190 TCollection extends DomainCollection,
191 TObject extends DomainObject,
192 TPutInput extends DomainPutObjectInput,
193 TParsed extends ParsedDomainObject,
194 TError extends CollectionDomainError,
195 TResource extends { href: string },
196>(
197 db: TDb,
198 input: DomainMultigetInput,
199 spec: Omit<
200 CollectionObjectSpec<TDb, TTx, TCollection, TObject, TPutInput, TParsed, TError>,
201 "objectResourceFromHref" | "objectResources" | "changesSince"
202 > & {
203 objectResourceFromHref(
204 db: TDb,
205 collection: TCollection,
206 href: string,
207 acceptedOrigins?: string[],
208 ): TResource | null;
209 },
210): { ok: true; requested: { href: string; resource: TResource | null }[] } | TError {
211 const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName);
212 if (!collection) return spec.collectionNotFound();
213 if (input.hrefs.length > input.maxResults) return spec.multigetLimitExceeded();
214 return {
215 ok: true,
216 requested: input.hrefs.map((href) => ({
217 href,
218 resource: spec.objectResourceFromHref(db, collection, href, input.acceptedOrigins),
219 })),
220 };
221}
222 
223export function syncCollectionDomain<
224 TDb,
225 TTx,
226 TCollection extends DomainCollection,
227 TObject extends DomainObject,
228 TPutInput extends DomainPutObjectInput,
229 TParsed extends ParsedDomainObject,
230 TError extends CollectionDomainError,
231 TResource extends { href: string },
232 TChange extends DomainChange,
233>(
234 db: TDb,
235 input: DomainSyncInput,
236 spec: Omit<
237 CollectionObjectSpec<TDb, TTx, TCollection, TObject, TPutInput, TParsed, TError>,
238 "objectResourceFromHref" | "objectResources" | "changesSince"
239 > & {
240 objectResourceFromHref(db: TDb, collection: TCollection, href: string): TResource | null;
241 objectResources(db: TDb, collection: TCollection, maxResults: number): DomainObjectResources<TResource>;
242 changesSince(
243 db: TDb,
244 collection: TCollection,
245 token: string | null,
246 maxResults: number,
247 ): DomainChangesSince<TChange>;
248 },
249): { ok: true; changes: TChange[]; resources: TResource[]; syncToken: string; truncated: boolean } | TError {
250 const collection = spec.ensureCollection(db, input.subjectId, input.nowMs, input.collectionName);
251 if (!collection) return spec.collectionNotFound();
252 const since = spec.syncTokenSeq(input.token);
253 if (since === null) return spec.invalidSyncToken();
254 
255 if (!input.token) {
256 const initial = spec.objectResources(db, collection, input.maxResults);
257 const syncTokenValue = initial.truncated ? spec.syncToken(0) : spec.syncToken(spec.latestChangeSeq(db));
258 return {
259 ok: true,
260 changes: [],
261 resources: initial.resources,
262 syncToken: syncTokenValue,
263 truncated: initial.truncated,
264 };
265 }
266 
267 const { changes, truncated } = spec.changesSince(db, collection, input.token, input.maxResults);
268 const lastReturnedSeq = changes.length > 0 ? changes[changes.length - 1]!.seq : since;
269 const tokenSeq = truncated ? lastReturnedSeq : spec.latestChangeSeq(db);
270 return {
271 ok: true,
272 changes,
273 resources: changes
274 .filter((change) => change.changeType !== "deleted")
275 .map((change) => spec.objectResourceFromHref(db, collection, change.href))
276 .filter((resource): resource is TResource => resource !== null),
277 syncToken: spec.syncToken(tokenSeq),
278 truncated,
279 };
280}