diff --git a/cloudflare/src/bindings.ts b/cloudflare/src/bindings.ts index 815346a..238720f 100644 --- a/cloudflare/src/bindings.ts +++ b/cloudflare/src/bindings.ts @@ -15,6 +15,7 @@ export interface ElyR2PutOptions { export interface ElyR2Bucket { get(key: string): Promise; put(key: string, value: ArrayBuffer, options?: ElyR2PutOptions): Promise; + delete(key: string): Promise; } export interface ElyD1PreparedStatement { diff --git a/cloudflare/src/storage.ts b/cloudflare/src/storage.ts index dfeb1d5..6c885bb 100644 --- a/cloudflare/src/storage.ts +++ b/cloudflare/src/storage.ts @@ -138,6 +138,11 @@ export async function getVerifiedObject( return payload; } +export async function deleteKnownObject(bucket: ElyR2Bucket, key: string): Promise { + assertKnownObjectKey(key); + await bucket.delete(key); +} + function assertKnownObjectKey(key: string): void { const matches = [ /^sync-payloads\/[a-z0-9][a-z0-9-]{1,31}\/[a-f0-9]{64}\/[a-z0-9][a-z0-9._-]{0,127}\/[a-z0-9][a-z0-9._-]{0,127}\/[a-f0-9]{64}\.bin$/, diff --git a/cloudflare/src/sync_reset.ts b/cloudflare/src/sync_reset.ts new file mode 100644 index 0000000..bea48a6 --- /dev/null +++ b/cloudflare/src/sync_reset.ts @@ -0,0 +1,280 @@ +import type { AuthContext } from "./auth.js"; +import type { ElyD1PreparedStatement, Env } from "./bindings.js"; +import { StorageObjectError, deleteKnownObject } from "./storage.js"; + +const RESET_CONFIRMATION = "delete-cloud-sync-data"; +const IDEMPOTENCY_KEY_PATTERN = /^[a-zA-Z0-9._:-]{16,128}$/; +const SYNC_RESET_EVENT_QUERY = ` + SELECT actor_device_id, outcome, created_at + FROM audit_events + WHERE user_id = ? AND event_id = ? AND event_type = 'sync.reset' +`; +const SYNC_RESET_COUNTS_QUERY = ` + SELECT + (SELECT COUNT(*) FROM sync_objects WHERE user_id = ?) AS objects, + (SELECT COUNT(*) FROM sync_change_log WHERE user_id = ?) AS changes, + (SELECT COUNT(*) FROM sync_snapshots WHERE user_id = ?) AS snapshots, + (SELECT COUNT(*) FROM sync_tombstones WHERE user_id = ?) AS tombstones +`; +const SYNC_RESET_R2_KEYS_QUERY = ` + SELECT payload_r2_key AS r2_key + FROM sync_objects + WHERE user_id = ? AND payload_r2_key IS NOT NULL + UNION + SELECT r2_key + FROM sync_snapshots + WHERE user_id = ? + ORDER BY r2_key ASC +`; +const SYNC_RESET_AUDIT_INSERT_QUERY = ` + INSERT INTO audit_events ( + event_id, + user_id, + actor_device_id, + event_type, + subject_type, + subject_id, + outcome, + metadata_hash, + created_at + ) VALUES (?, ?, ?, 'sync.reset', 'sync', ?, 'success', NULL, ?) + ON CONFLICT(event_id) DO NOTHING +`; +const SYNC_OBJECTS_DELETE_QUERY = "DELETE FROM sync_objects WHERE user_id = ?"; +const SYNC_CHANGE_LOG_DELETE_QUERY = "DELETE FROM sync_change_log WHERE user_id = ?"; +const SYNC_SNAPSHOTS_DELETE_QUERY = "DELETE FROM sync_snapshots WHERE user_id = ?"; +const SYNC_TOMBSTONES_DELETE_QUERY = "DELETE FROM sync_tombstones WHERE user_id = ?"; + +export interface SyncResetDocument { + version: 1; + user_id: string; + device_id: string; + idempotency_key: string; + reset_at: number; + deleted: SyncResetDeletedDocument; +} + +export interface SyncResetDeletedDocument { + objects: number; + changes: number; + snapshots: number; + tombstones: number; + r2_objects: number; +} + +interface SyncResetRequest { + idempotencyKey: string; +} + +interface SyncResetEventRow { + actor_device_id: unknown; + outcome: unknown; + created_at: unknown; +} + +interface SyncResetCountsRow { + objects: unknown; + changes: unknown; + snapshots: unknown; + tombstones: unknown; +} + +interface SyncResetR2KeyRow { + r2_key: unknown; +} + +type RequestBody = Record; + +export class SyncResetRequestError extends Error { + constructor(message: string) { + super(message); + this.name = "SyncResetRequestError"; + } +} + +export class SyncResetPersistenceError extends Error { + constructor(message: string) { + super(message); + this.name = "SyncResetPersistenceError"; + } +} + +export async function syncResetDocument( + request: Request, + env: Env, + context: AuthContext, + nowSeconds = Math.floor(Date.now() / 1000), +): Promise { + const deviceId = currentDeviceId(context); + const reset = await syncResetRequest(request); + const eventId = syncResetEventId(context.userId, reset.idempotencyKey); + const existingEvent = await env.ELY_DB.prepare(SYNC_RESET_EVENT_QUERY) + .bind(context.userId, eventId) + .first(); + if (existingEvent !== null) { + return existingResetDocument(context, deviceId, reset, existingEvent); + } + + const counts = await syncResetCounts(env, context.userId); + const r2Keys = await syncResetR2Keys(env, context.userId); + for (const key of r2Keys) { + await deleteResetObject(env, key); + } + await env.ELY_DB.batch(syncResetStatements(env, context.userId, deviceId, eventId, nowSeconds)); + + return { + version: 1, + user_id: context.userId, + device_id: deviceId, + idempotency_key: reset.idempotencyKey, + reset_at: nowSeconds, + deleted: { ...counts, r2_objects: r2Keys.length }, + }; +} + +function existingResetDocument( + context: AuthContext, + deviceId: string, + reset: SyncResetRequest, + row: SyncResetEventRow, +): SyncResetDocument { + if (row.actor_device_id !== deviceId || row.outcome !== "success") { + throw new SyncResetRequestError("sync_reset_replay_mismatch"); + } + return { + version: 1, + user_id: context.userId, + device_id: deviceId, + idempotency_key: reset.idempotencyKey, + reset_at: integer(row.created_at, "created_at"), + deleted: emptyDeletedDocument(), + }; +} + +async function syncResetCounts( + env: Env, + userId: string, +): Promise> { + const row = await env.ELY_DB.prepare(SYNC_RESET_COUNTS_QUERY) + .bind(userId, userId, userId, userId) + .first(); + if (row === null) { + throw new SyncResetPersistenceError("sync_reset_counts_missing"); + } + return { + objects: integer(row.objects, "objects"), + changes: integer(row.changes, "changes"), + snapshots: integer(row.snapshots, "snapshots"), + tombstones: integer(row.tombstones, "tombstones"), + }; +} + +async function syncResetR2Keys(env: Env, userId: string): Promise { + const result = await env.ELY_DB.prepare(SYNC_RESET_R2_KEYS_QUERY) + .bind(userId, userId) + .all(); + return result.results.map(r2Key); +} + +function syncResetStatements( + env: Env, + userId: string, + deviceId: string, + eventId: string, + nowSeconds: number, +): ElyD1PreparedStatement[] { + return [ + env.ELY_DB.prepare(SYNC_CHANGE_LOG_DELETE_QUERY).bind(userId), + env.ELY_DB.prepare(SYNC_TOMBSTONES_DELETE_QUERY).bind(userId), + env.ELY_DB.prepare(SYNC_SNAPSHOTS_DELETE_QUERY).bind(userId), + env.ELY_DB.prepare(SYNC_OBJECTS_DELETE_QUERY).bind(userId), + env.ELY_DB.prepare(SYNC_RESET_AUDIT_INSERT_QUERY).bind( + eventId, + userId, + deviceId, + userId, + nowSeconds, + ), + ]; +} + +async function deleteResetObject(env: Env, key: string): Promise { + try { + await deleteKnownObject(env.ELY_STORAGE, key); + } catch (error) { + if (error instanceof StorageObjectError) { + throw new SyncResetPersistenceError(error.message); + } + throw error; + } +} + +async function syncResetRequest(request: Request): Promise { + const body = await requestBody(request); + assertOnlyFields(body, ["version", "confirmation", "idempotency_key"]); + if (body.version !== 1) { + throw new SyncResetRequestError("version_invalid"); + } + if (body.confirmation !== RESET_CONFIRMATION) { + throw new SyncResetRequestError("confirmation_invalid"); + } + return { idempotencyKey: idempotencyKey(body.idempotency_key) }; +} + +async function requestBody(request: Request): Promise { + let value: unknown; + try { + value = await request.json(); + } catch { + throw new SyncResetRequestError("json_invalid"); + } + if (typeof value !== "object" || value === null || Array.isArray(value)) { + throw new SyncResetRequestError("body_invalid"); + } + return value as RequestBody; +} + +function assertOnlyFields(value: RequestBody, fields: string[]): void { + const allowed = new Set(fields); + for (const field of Object.keys(value)) { + if (!allowed.has(field)) { + throw new SyncResetRequestError(`unexpected_field:${field}`); + } + } +} + +function currentDeviceId(context: AuthContext): string { + if (context.deviceId === undefined) { + throw new SyncResetRequestError("device_context_required"); + } + return context.deviceId; +} + +function idempotencyKey(value: unknown): string { + if (typeof value !== "string" || !IDEMPOTENCY_KEY_PATTERN.test(value)) { + throw new SyncResetRequestError("idempotency_key_invalid"); + } + return value; +} + +function r2Key(row: SyncResetR2KeyRow): string { + if (typeof row.r2_key !== "string") { + throw new SyncResetPersistenceError("r2_key_invalid"); + } + return row.r2_key; +} + +function integer(value: unknown, label: string): number { + if (typeof value !== "number" || !Number.isSafeInteger(value) || value < 0) { + throw new SyncResetPersistenceError(`${label}_invalid`); + } + return value; +} + +function emptyDeletedDocument(): SyncResetDeletedDocument { + return { objects: 0, changes: 0, snapshots: 0, tombstones: 0, r2_objects: 0 }; +} + +function syncResetEventId(userId: string, idempotencyKey: string): string { + return `sync-reset:${userId}:${idempotencyKey}`; +} diff --git a/cloudflare/src/sync_routes.ts b/cloudflare/src/sync_routes.ts index cd2da43..6eea2cd 100644 --- a/cloudflare/src/sync_routes.ts +++ b/cloudflare/src/sync_routes.ts @@ -8,6 +8,11 @@ import { SyncPushRequestError, syncPushDocument, } from "./sync_push.js"; +import { + SyncResetPersistenceError, + SyncResetRequestError, + syncResetDocument, +} from "./sync_reset.js"; import { SyncSnapshotConflictError, SyncSnapshotNotFoundError, @@ -35,6 +40,9 @@ export async function handleSyncRoute( if (url.pathname === "/api/sync/status") { return handleSyncStatus(request, env); } + if (url.pathname === "/api/sync/reset") { + return handleSyncReset(request, env); + } return null; } @@ -180,3 +188,35 @@ function handleSyncStatus(request: Request, env: Env): Promise { }, ); } + +function handleSyncReset(request: Request, env: Env): Promise { + return withApprovedDeviceApiControls( + request, + env, + "sync.reset", + ["POST"], + async (context) => { + try { + return jsonResponse(await syncResetDocument(request, env, context), 200, { + "Cache-Control": "no-store", + }); + } catch (error) { + if (error instanceof SyncResetRequestError) { + return jsonResponse( + { error: "invalid_sync_reset" }, + 400, + { "Cache-Control": "no-store" }, + ); + } + if (error instanceof SyncResetPersistenceError) { + return jsonResponse( + { error: "sync_reset_failed" }, + 500, + { "Cache-Control": "no-store" }, + ); + } + throw error; + } + }, + ); +} diff --git a/cloudflare/tests/api_controls.test.ts b/cloudflare/tests/api_controls.test.ts index d7a78ca..8ef3adf 100644 --- a/cloudflare/tests/api_controls.test.ts +++ b/cloudflare/tests/api_controls.test.ts @@ -290,6 +290,9 @@ function testR2Bucket(): Env["ELY_STORAGE"] { }, }); }, + delete() { + return Promise.resolve(); + }, }; } diff --git a/cloudflare/tests/devices_test_support.ts b/cloudflare/tests/devices_test_support.ts index a72b3ce..2f91946 100644 --- a/cloudflare/tests/devices_test_support.ts +++ b/cloudflare/tests/devices_test_support.ts @@ -14,6 +14,7 @@ export interface TestEnvOptions { d1?: RecordedD1Database; kvEntries?: [string, string][]; kvReads?: string[]; + r2Deletes?: string[]; r2Gets?: string[]; r2Objects?: [string, ArrayBuffer][]; r2Puts?: RecordedR2Put[]; @@ -47,7 +48,12 @@ export function testEnv(options: TestEnvOptions): Env { return Promise.resolve(values.get(key) ?? null); }, }, - ELY_STORAGE: testR2Bucket(options.r2Puts, options.r2Objects, options.r2Gets), + ELY_STORAGE: testR2Bucket( + options.r2Puts, + options.r2Objects, + options.r2Gets, + options.r2Deletes, + ), ELY_RATE_LIMITER: { limit(): Promise<{ success: boolean }> { return Promise.resolve({ success: true }); @@ -125,6 +131,7 @@ function testR2Bucket( puts: RecordedR2Put[] = [], objects: [string, ArrayBuffer][] = [], gets: string[] = [], + deletes: string[] = [], ): Env["ELY_STORAGE"] { const values = new Map(objects); return { @@ -149,5 +156,10 @@ function testR2Bucket( }, }); }, + delete(key: string) { + deletes.push(key); + values.delete(key); + return Promise.resolve(); + }, }; } diff --git a/cloudflare/tests/index.test.ts b/cloudflare/tests/index.test.ts index c39ab67..1aaba4e 100644 --- a/cloudflare/tests/index.test.ts +++ b/cloudflare/tests/index.test.ts @@ -371,6 +371,9 @@ function testR2Bucket(): Env["ELY_STORAGE"] { }, }); }, + delete() { + return Promise.resolve(); + }, }; } diff --git a/cloudflare/tests/storage.test.ts b/cloudflare/tests/storage.test.ts index 88cd587..1db4335 100644 --- a/cloudflare/tests/storage.test.ts +++ b/cloudflare/tests/storage.test.ts @@ -6,6 +6,7 @@ import type { ElyR2Bucket, ElyR2PutOptions } from "../src/bindings.js"; import { StorageObjectError, crashAttachmentKey, + deleteKnownObject, exportObjectKey, getVerifiedObject, pluginAssetKey, @@ -168,6 +169,24 @@ describe("R2 storage contracts", () => { assert.deepEqual(new Uint8Array(downloaded ?? new ArrayBuffer(0)), new Uint8Array(payload)); }); + + it("deletes only known object key shapes", async () => { + const bucket = recordedR2Bucket(); + const key = syncSnapshotKey({ + region: "us-east", + userHash: USER_HASH, + snapshotId: "snapshot-01", + }); + + await deleteKnownObject(bucket, key); + + assert.deepEqual(bucket.deletes, [key]); + await assert.rejects( + () => deleteKnownObject(bucket, "sync-snapshots/../bad.bin"), + StorageObjectError, + ); + assert.deepEqual(bucket.deletes, [key]); + }); }); interface RecordedPut { @@ -177,12 +196,15 @@ interface RecordedPut { } interface RecordedR2Bucket extends ElyR2Bucket { + deletes: string[]; puts: RecordedPut[]; } function recordedR2Bucket(payload?: ArrayBuffer): RecordedR2Bucket { + const deletes: string[] = []; const puts: RecordedPut[] = []; return { + deletes, puts, get() { if (payload === undefined) { @@ -202,6 +224,10 @@ function recordedR2Bucket(payload?: ArrayBuffer): RecordedR2Bucket { }, }); }, + delete(key: string) { + deletes.push(key); + return Promise.resolve(); + }, }; } diff --git a/cloudflare/tests/sync_reset_routes.test.ts b/cloudflare/tests/sync_reset_routes.test.ts new file mode 100644 index 0000000..5a6fb7f --- /dev/null +++ b/cloudflare/tests/sync_reset_routes.test.ts @@ -0,0 +1,226 @@ +import assert from "node:assert/strict"; +import { createHash } from "node:crypto"; +import { describe, it } from "node:test"; + +import { authSessionCacheKvKey, authTokenHash } from "../src/auth.js"; +import { handleRequest } from "../src/index.js"; +import { ACCESS_TOKEN, sessionDocument, testD1Database, testEnv } from "./devices_test_support.js"; + +const USER_ID = "user-01"; +const DEVICE_ID = "device-01"; +const IDEMPOTENCY_KEY = "sync-reset-000001"; +const PAYLOAD_HASH = "a".repeat(64); +const USER_HASH = sha256(bytes(USER_ID)); +const PAYLOAD_KEY = `sync-payloads/us-east/${USER_HASH}/tabs/tab-01/${PAYLOAD_HASH}.bin`; +const SNAPSHOT_KEY = `sync-snapshots/us-east/${USER_HASH}/snapshot-01.bin`; + +describe("sync reset routes", () => { + it("deletes cloud sync data and records an audit event", async () => { + const r2Deletes: string[] = []; + const tokenHash = await authTokenHash(ACCESS_TOKEN); + const d1 = testD1Database({ + firstRows: [{ device_id: DEVICE_ID }, null, resetCountsRow()], + allRows: [{ r2_key: PAYLOAD_KEY }, { r2_key: SNAPSHOT_KEY }], + }); + + const response = await handleRequest( + syncResetRequest(syncResetBody()), + testEnv({ + d1, + r2Deletes, + kvEntries: [[authSessionCacheKvKey("local", tokenHash), sessionDocument(DEVICE_ID)]], + }), + ); + + const body = (await response.json()) as { + version: number; + user_id: string; + device_id: string; + idempotency_key: string; + reset_at: number; + deleted: Record; + }; + assert.equal(response.status, 200); + assert.equal(response.headers.get("cache-control"), "no-store"); + assert.equal(body.version, 1); + assert.equal(body.user_id, USER_ID); + assert.equal(body.device_id, DEVICE_ID); + assert.equal(body.idempotency_key, IDEMPOTENCY_KEY); + assert.ok(Number.isSafeInteger(body.reset_at)); + assert.deepEqual(body.deleted, { + objects: 4, + changes: 9, + snapshots: 2, + tombstones: 1, + r2_objects: 2, + }); + assert.deepEqual(r2Deletes, [PAYLOAD_KEY, SNAPSHOT_KEY]); + assert.equal(d1.batches[0], 5); + assert.ok(d1.queries[1]?.includes("FROM audit_events")); + assert.ok(d1.queries[2]?.includes("FROM sync_objects")); + assert.ok(d1.queries[3]?.includes("UNION")); + assert.ok(d1.queries[4]?.includes("DELETE FROM sync_change_log")); + assert.ok(d1.queries[8]?.includes("INSERT INTO audit_events")); + assert.deepEqual(d1.binds[1], [USER_ID, syncResetEventId()]); + assert.deepEqual(d1.binds[2], [USER_ID, USER_ID, USER_ID, USER_ID]); + assert.deepEqual(d1.binds[3], [USER_ID, USER_ID]); + assert.deepEqual(d1.binds[8]?.slice(0, 4), [syncResetEventId(), USER_ID, DEVICE_ID, USER_ID]); + }); + + it("returns an idempotent reset document for existing audit events", async () => { + const r2Deletes: string[] = []; + const tokenHash = await authTokenHash(ACCESS_TOKEN); + const d1 = testD1Database({ + firstRows: [ + { device_id: DEVICE_ID }, + { actor_device_id: DEVICE_ID, outcome: "success", created_at: 1_780_001_000 }, + ], + allRows: [{ r2_key: PAYLOAD_KEY }], + }); + + const response = await handleRequest( + syncResetRequest(syncResetBody()), + testEnv({ + d1, + r2Deletes, + kvEntries: [[authSessionCacheKvKey("local", tokenHash), sessionDocument(DEVICE_ID)]], + }), + ); + + assert.equal(response.status, 200); + assert.deepEqual(await response.json(), { + version: 1, + user_id: USER_ID, + device_id: DEVICE_ID, + idempotency_key: IDEMPOTENCY_KEY, + reset_at: 1_780_001_000, + deleted: { objects: 0, changes: 0, snapshots: 0, tombstones: 0, r2_objects: 0 }, + }); + assert.deepEqual(r2Deletes, []); + assert.equal(d1.queries.length, 2); + assert.deepEqual(d1.batches, []); + }); + + it("rejects replay mismatches before deleting sync data", async () => { + const r2Deletes: string[] = []; + const tokenHash = await authTokenHash(ACCESS_TOKEN); + const d1 = testD1Database({ + firstRows: [ + { device_id: DEVICE_ID }, + { actor_device_id: "device-02", outcome: "success", created_at: 1_780_001_000 }, + ], + }); + + const response = await handleRequest( + syncResetRequest(syncResetBody()), + testEnv({ + d1, + r2Deletes, + kvEntries: [[authSessionCacheKvKey("local", tokenHash), sessionDocument(DEVICE_ID)]], + }), + ); + + assert.equal(response.status, 400); + assert.deepEqual(await response.json(), { error: "invalid_sync_reset" }); + assert.deepEqual(r2Deletes, []); + assert.deepEqual(d1.batches, []); + }); + + it("rejects missing confirmation before reset reads", async () => { + const r2Deletes: string[] = []; + const tokenHash = await authTokenHash(ACCESS_TOKEN); + const d1 = testD1Database({ firstRows: [{ device_id: DEVICE_ID }] }); + + const response = await handleRequest( + syncResetRequest(syncResetBody({ confirmation: "delete" })), + testEnv({ + d1, + r2Deletes, + kvEntries: [[authSessionCacheKvKey("local", tokenHash), sessionDocument(DEVICE_ID)]], + }), + ); + + assert.equal(response.status, 400); + assert.deepEqual(await response.json(), { error: "invalid_sync_reset" }); + assert.deepEqual(r2Deletes, []); + assert.equal(d1.queries.length, 1); + assert.deepEqual(d1.batches, []); + }); + + it("rejects revoked devices before reading reset bodies", async () => { + const tokenHash = await authTokenHash(ACCESS_TOKEN); + const d1 = testD1Database({ firstRows: [null] }); + + const response = await handleRequest( + syncResetRequest(syncResetBody()), + testEnv({ + d1, + kvEntries: [[authSessionCacheKvKey("local", tokenHash), sessionDocument(DEVICE_ID)]], + }), + ); + + assert.equal(response.status, 403); + assert.deepEqual(await response.json(), { error: "device_not_approved" }); + assert.equal(d1.queries.length, 1); + assert.deepEqual(d1.batches, []); + }); + + it("fails closed when stored R2 keys are malformed", async () => { + const r2Deletes: string[] = []; + const tokenHash = await authTokenHash(ACCESS_TOKEN); + const d1 = testD1Database({ + firstRows: [{ device_id: DEVICE_ID }, null, resetCountsRow()], + allRows: [{ r2_key: "sync-snapshots/../bad.bin" }], + }); + + const response = await handleRequest( + syncResetRequest(syncResetBody()), + testEnv({ + d1, + r2Deletes, + kvEntries: [[authSessionCacheKvKey("local", tokenHash), sessionDocument(DEVICE_ID)]], + }), + ); + + assert.equal(response.status, 500); + assert.deepEqual(await response.json(), { error: "sync_reset_failed" }); + assert.deepEqual(r2Deletes, []); + assert.deepEqual(d1.batches, []); + }); +}); + +function syncResetRequest(body: Record): Request { + return new Request("https://elydora.test/api/sync/reset", { + method: "POST", + headers: { + authorization: `Bearer ${ACCESS_TOKEN}`, + "content-type": "application/json", + }, + body: JSON.stringify(body), + }); +} + +function syncResetBody(overrides: Record = {}): Record { + return { + version: 1, + confirmation: "delete-cloud-sync-data", + idempotency_key: IDEMPOTENCY_KEY, + ...overrides, + }; +} + +function resetCountsRow(overrides: Record = {}): Record { + return { objects: 4, changes: 9, snapshots: 2, tombstones: 1, ...overrides }; +} + +function syncResetEventId(): string { + return `sync-reset:${USER_ID}:${IDEMPOTENCY_KEY}`; +} + +function bytes(value: string): ArrayBuffer { + return new TextEncoder().encode(value).buffer; +} + +function sha256(payload: ArrayBuffer): string { + return createHash("sha256").update(new Uint8Array(payload)).digest("hex"); +}