/
githubmirror
/
tldraw
Обзор
Документация
Войти
/
githubmirror
/
tldraw
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
apps/dotcom/sync-worker/src/adminRoutes.ts
971 строка
34 KB
Mime Čuvalo
feat(dotcom): add board stats to the admin page (#9914)
11 авг 2026, 17:28
Не верифицирован
11 авг 2026, 17:28
c47a5bd
Код
Авторство
О чём код?
import { AdminFileAssetsResponseBody, AdminFileStatsResponseBody, FILE_PREFIX, FeatureFlagKey, FriendsAndFamilyEntry, LOCAL_FILE_PREFIX, PUBLISH_PREFIX, ROOM_PREFIX, TlaFile, WELCOME_CREATE_SOURCE, } from '@tldraw/dotcom-shared' import { assert, retry, sleep, uniqueId } from '@tldraw/utils' import { createRouter } from '@tldraw/worker-shared' import { StatusError, json } from 'itty-router' import { sql } from 'kysely' import PQueue from 'p-queue' import { getUploadObjectName } from './assetAssociation' import { summarizeSnapshotDocuments } from './fileStats' import { MAX_ATTEMPTS } from './outboxDrain' import { createPostgresConnectionPool } from './postgres' import { getR2KeyForRoom } from './r2' import { getFileSnapshot, returnFileSnapshot } from './routes/tla/getFileSnapshot' import { type Environment } from './types' import { undeleteFile } from './undeleteFile' import { getFileEffectProcessor, getRoomDurableObject, getUserDurableObject, } from './utils/durableObjects' import { FEATURE_FLAG_KEYS, getFeatureFlagsAdmin, setFeatureFlag } from './utils/featureFlags' import { getFriendsAndFamilyList, parseFriendsAndFamilyEmails, setFriendsAndFamilyList, } from './utils/mcpFriendsAndFamily' import { getClerkClient, requireAdminAccess, requireAuth } from './utils/tla/getAuth' /** * Resolves the admin's emails to user ids once, at save time, so the request path only ever has to * compare the userId `getAuth` already gives it. One query for the whole list rather than a lookup * per address, and an email with no tldraw account fails the save instead of being stored as an * entry that can never match. */ async function resolveFriendsAndFamilyUsers( env: Environment, emails: string[] ): Promise<FriendsAndFamilyEntry[]> { if (!emails.length) return [] const db = createPostgresConnectionPool(env, '/app/admin/mcp-friends-and-family') try { const rows = await db .selectFrom('user') .select(['id', 'email']) .where(sql<string>`lower(email)`, 'in', emails) .execute() const byEmail = new Map(rows.map((row) => [row.email.toLowerCase(), row])) const unknown = emails.filter((email) => !byEmail.has(email)) if (unknown.length) { throw new StatusError(400, `No tldraw account for: ${unknown.join(', ')}`) } // Store the address as the database has it, not as the admin typed it, so the panel shows the // account's real email. return emails.map((email) => { const row = byEmail.get(email)! return { userId: row.id, email: row.email } }) } finally { await db.destroy() } } async function requireUser(env: Environment, q: string) { const db = createPostgresConnectionPool(env, '/app/admin/user') const userRow = await db .selectFrom('user') .where((eb) => eb.or([eb('email', '=', q), eb('id', '=', q)])) .selectAll() .executeTakeFirst() if (!userRow) { throw new StatusError(404, 'User not found ' + q) } return userRow } export const adminRoutes = createRouter<Environment>() .all('/app/admin/*', async (req, env) => { const auth = await requireAuth(req, env) await requireAdminAccess(env, auth) }) .get('/app/admin/user', async (res, env) => { const q = res.query['q'] if (typeof q !== 'string') { return new Response('Missing query param', { status: 400 }) } const userRow = await requireUser(env, q) const db = createPostgresConnectionPool(env, '/app/admin/user') try { const memberships = await db .selectFrom('group_user') .innerJoin('group', 'group.id', 'group_user.groupId') .where('group_user.userId', '=', userRow.id) .select(['group.id', 'group.name', 'group.isDeleted', 'group_user.role']) .execute() const files = await db .selectFrom('file') .leftJoin('group_file', 'group_file.fileId', 'file.id') .where('file.isDeleted', '=', false) .where((eb) => eb.or([ eb('file.ownerId', '=', userRow.id), eb( 'group_file.groupId', 'in', memberships.length ? memberships.map((m) => m.id) : [''] ), ]) ) // distinctOn dedupes before the limit is applied, so a file matching both the owner // clause and multiple group memberships (join fanout) still only counts once against // the 500 cap. .distinctOn('file.id') .orderBy('file.id') .selectAll('file') .limit(500) .execute() return json({ user: userRow, memberships, files }) } finally { await db.destroy() } }) .get('/app/admin/outbox', async (res, env) => { const db = createPostgresConnectionPool(env, '/app/admin/outbox') try { const stats = await db .selectFrom('effect_outbox') .select((eb) => [ eb.fn.countAll<number>().filterWhere('attempts', '<', MAX_ATTEMPTS).as('pending'), eb.fn.countAll<number>().filterWhere('attempts', '>=', MAX_ATTEMPTS).as('parked'), eb.fn.min('createdAt').filterWhere('attempts', '<', MAX_ATTEMPTS).as('oldestPending'), ]) .executeTakeFirstOrThrow() return json({ outbox: { pending: Number(stats.pending), parked: Number(stats.parked), oldestPendingAgeSeconds: stats.oldestPending ? Math.round((Date.now() - new Date(stats.oldestPending).getTime()) / 1000) : null, }, }) } finally { await db.destroy() } }) .get('/app/admin/feature-flags', getFeatureFlagsAdmin) .post('/app/admin/feature-flags', async (req, env) => { const body: any = await req.json() const { flag, enabled, percentage } = body if (typeof flag !== 'string') { throw new StatusError(400, 'flag (string) is required') } if (enabled !== undefined && typeof enabled !== 'boolean') { throw new StatusError(400, 'enabled must be a boolean') } if ( percentage !== undefined && (typeof percentage !== 'number' || percentage < 0 || percentage > 100) ) { throw new StatusError(400, 'percentage must be a number between 0 and 100') } if (!FEATURE_FLAG_KEYS.includes(flag as FeatureFlagKey)) { throw new StatusError(400, `Invalid flag. Must be one of: ${FEATURE_FLAG_KEYS.join(', ')}`) } const update: { enabled?: boolean; percentage?: number } = {} if (enabled !== undefined) update.enabled = enabled if (percentage !== undefined) update.percentage = percentage await setFeatureFlag(env, flag as FeatureFlagKey, update) return json({ success: true, flag, ...update }) }) .get('/app/admin/mcp-friends-and-family', async (_req, env) => { return json({ entries: await getFriendsAndFamilyList(env) }) }) .post('/app/admin/mcp-friends-and-family', async (req, env) => { const body: any = await req.json() // Parsing before resolving means a typo is rejected at the point someone can still fix it, // rather than sitting in the list looking like it grants access while matching nothing. let emails: string[] try { emails = parseFriendsAndFamilyEmails(body?.entries) } catch (e) { throw new StatusError(400, e instanceof Error ? e.message : String(e)) } const entries = await resolveFriendsAndFamilyUsers(env, emails) await setFriendsAndFamilyList(env, entries) return json({ success: true, entries }) }) .post('/app/admin/create_legacy_file', async (_res, env) => { const slug = uniqueId() await getRoomDurableObject(env, slug).__admin__createLegacyRoom(slug) return json({ slug }) }) .post('/app/admin/hard_delete_file/:fileId', async (res, env) => { const fileId = res.params.fileId assert(typeof fileId === 'string', 'fileId is required') const pg = createPostgresConnectionPool(env, '/app/admin/hard_delete_file') const file = await pg.selectFrom('file').where('id', '=', fileId).selectAll().executeTakeFirst() if (!file) { if (await maybeHardDeleteLegacyFile({ id: fileId, env })) { return new Response('deleted') } else { return new Response('File not found', { status: 404 }) } } return await hardDeleteAppFile({ pg, file, env }) }) .post('/app/admin/undelete_file/:fileId', async (res, env) => { const fileId = res.params.fileId assert(typeof fileId === 'string', 'fileId is required') const pg = createPostgresConnectionPool(env, '/app/admin/undelete_file') const outcome = await undeleteFile(pg, fileId) if (outcome.result === 'not_found') { return new Response('File not found', { status: 404 }) } if (outcome.result === 'not_deleted') { return new Response('File is not deleted', { status: 400 }) } if (outcome.result === 'group_deleted') { return new Response('Owning workspace is deleted — restore the workspace first', { status: 409, }) } // Best-effort nudge so the restore's effects land promptly; the sweep backstops it, so a // poke failure must not fail the request (the restore already committed; a retry would // just 400 with 'File is not deleted'). await getFileEffectProcessor(env) .poke() .catch(() => {}) // Hard-reboot every affected user so the restored rows replicate to their session (the // isDeleted false-flip case flagged in dotcom-shared mutators.ts). Best-effort: the // writes already committed, so a reboot failure must not surface as a 500 (a retry would // just 400 with 'File is not deleted' without rebooting anyone). Bounded concurrency keeps // us inside the worker's connection budget for workspaces with many members. const rebootQueue = new PQueue({ concurrency: 5 }) await rebootQueue.addAll( outcome.rebootUserIds.map((userId) => async () => { try { await getUserDurableObject(env, userId).admin_forceHardReboot(userId) } catch (e) { console.error(`Failed to reboot user ${userId} after undeleting file ${fileId}`, e) } }) ) return json({ success: true }) }) // Deleted files the user OWNS: legacy direct owner, their home workspace (group id = user // id), or a workspace where they hold the owner role. Mere memberships and guest files are // excluded — the per-row Undelete button restores files, so the list must only contain files // the user legitimately owns. Queried from Postgres because the user's replicated store // filters out their own deleted files (see fetchEverythingSql). .get('/app/admin/user/deleted_files', async (res, env) => { const q = res.query['q'] if (typeof q !== 'string') { return new Response('Missing query param', { status: 400 }) } const userRow = await requireUser(env, q) const pg = createPostgresConnectionPool(env, '/app/admin/user/deleted_files') const files = await pg .selectFrom('file') .leftJoin('group', 'group.id', 'file.owningGroupId') .leftJoin('group_user', (join) => join .onRef('group_user.groupId', '=', 'file.owningGroupId') .on('group_user.userId', '=', userRow.id) ) .where('file.isDeleted', '=', true) // A deleted file whose owning workspace is also deleted can't be restored (undeleteFile // blocks it, see undeleteFile.ts) until the workspace is restored first, so hide it here. .where((eb) => eb.or([eb('group.isDeleted', '=', false), eb('file.owningGroupId', 'is', null)]) ) .where((eb) => eb.or([ eb('file.ownerId', '=', userRow.id), eb('file.owningGroupId', '=', userRow.id), eb( 'file.owningGroupId', 'in', eb .selectFrom('group_user') .select('group_user.groupId') .where('group_user.userId', '=', userRow.id) .where('group_user.role', '=', 'owner') ), ]) ) .selectAll('file') .select(['group.name as workspaceName', 'group_user.role as workspaceRole']) .orderBy('file.updatedAt', 'desc') .execute() return json(files) }) .post('/app/admin/delete_user', async (res, env) => { const q = res.query['q'] if (typeof q !== 'string') { return new Response('Missing query param', { status: 400 }) } const userRow = await requireUser(env, q) await performUserDeletion(userRow, env) return new Response('User deleted', { status: 200 }) }) .get('/app/admin/delete_user_sse', async (res, env) => { const q = res.query['q'] if (typeof q !== 'string') { return new Response('Missing query param', { status: 400 }) } const userRow = await requireUser(env, q) return new Response( new ReadableStream({ async start(controller) { try { // Helper function to send progress events const sendProgress = (step: string, message: string, details?: any) => { const event = { type: 'progress', step, message, timestamp: Date.now(), details, } controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify(event)}\n\n`)) } sendProgress('starting', 'Beginning user deletion process...', { userId: userRow.id }) await performUserDeletion(userRow, env, sendProgress) // Send completion event const completionEvent = { type: 'complete', step: 'finished', message: 'User deletion completed successfully', timestamp: Date.now(), details: { userId: userRow.id }, } controller.enqueue( new TextEncoder().encode(`data: ${JSON.stringify(completionEvent)}\n\n`) ) } catch (error) { // Send error event const errorEvent = { type: 'error', step: 'error', message: error instanceof Error ? error.message : 'Unknown error occurred', timestamp: Date.now(), details: { error: error instanceof Error ? error.stack : String(error) }, } controller.enqueue(new TextEncoder().encode(`data: ${JSON.stringify(errorEvent)}\n\n`)) } finally { controller.close() } }, }), { headers: { 'Content-Type': 'text/event-stream', 'Cache-Control': 'no-cache', Connection: 'keep-alive', 'Access-Control-Allow-Origin': '*', 'Access-Control-Allow-Headers': 'Cache-Control', }, } ) }) // Read-only asset health report for a file: is each asset's object still in the uploads // bucket, and is it associated with the file? Explains files stuck in a zero-progress // association loop. .get('/app/admin/file-assets/:slug', async (res, env) => { const slug = res.params.slug assert(typeof slug === 'string', 'slug is required') const pg = createPostgresConnectionPool(env, '/app/admin/file-assets') const file = await pg .selectFrom('file') .where('id', '=', slug) .select(['id', 'name', 'ownerId', 'owningGroupId', 'isDeleted', 'createSource']) .executeTakeFirst() const snapshot = await getFileSnapshot(env, slug, true) if (!snapshot) { throw new StatusError(404, `No persisted snapshot for ${slug}`) } // Mirrors how the association pass parses asset records (see associatePendingAssets) const userContentUrl = env.USER_CONTENT_URL const assets: Array<{ assetId: string objectName: string src: string fileIdMeta: string | null associated: boolean oldFormatUrl: boolean inBucket: boolean | null sizeBytes: number | null }> = [] let totalShapes = 0 const shapesByType: Record<string, number> = {} // Assets the association pass can never act on (bookmarks, non-http srcs, R2-invalid // object names). Counted instead of reported as missing uploads. let external = 0 for (const { state } of snapshot.documents) { const record = state as any if (record.typeName === 'shape') { totalShapes++ shapesByType[record.type] = (shapesByType[record.type] ?? 0) + 1 continue } if (record.typeName !== 'asset') continue const src = record.props?.src if (!src) continue const objectName = getUploadObjectName(record) if (!objectName) { external++ continue } const fileIdMeta = record.meta?.fileId ?? null const associated = fileIdMeta === slug assets.push({ assetId: record.id, objectName, src, fileIdMeta, associated, oldFormatUrl: associated && src.startsWith('http') && !!userContentUrl && !src.startsWith(userContentUrl), // null until the head check settles; a failed check stays null so R2 flakiness // doesn't read as a confirmed-missing object inBucket: null, sizeBytes: null, }) } // Bounded concurrency keeps us inside the worker's connection budget; persistent head // failures become warnings rather than failing the report const warnings: string[] = [] const headQueue = new PQueue({ concurrency: 5 }) await headQueue.addAll( assets.map((asset) => async () => { try { const head = await retry(() => env.UPLOADS.head(asset.objectName), { attempts: 2, waitDuration: 500, }) asset.inBucket = !!head asset.sizeBytes = head?.size ?? null } catch (e) { warnings.push(`head failed for ${asset.objectName}: ${e}`) } }) ) // Cross-check the asset table both ways: which fileId the DB thinks owns each referenced // object, and rows claimed by this file whose objects the snapshot no longer references const referencedSet = new Set(assets.map((a) => a.objectName)) const [dbRowsForReferenced, rowsForThisFile] = await Promise.all([ referencedSet.size > 0 ? pg .selectFrom('asset') .where('objectName', 'in', [...referencedSet]) .select(['objectName', 'fileId']) .execute() : [], pg.selectFrom('asset').where('fileId', '=', slug).select(['objectName']).execute(), ]) const dbFileIdByObjectName = new Map(dbRowsForReferenced.map((r) => [r.objectName, r.fileId])) const orphaned = rowsForThisFile.filter((row) => !referencedSet.has(row.objectName)).length // Mirrors loadCreateSourceData: exists means seeding from this source would find content. // Readonly and snapshot prefixes need slug translation to check, so they report null (not // checked), as does a failed check. let source: { raw: string; exists: boolean | null } | null = null if (file?.createSource) { const raw = file.createSource const [prefix, id] = raw.split('/') let exists: boolean | null = null try { if (raw === WELCOME_CREATE_SOURCE || prefix === LOCAL_FILE_PREFIX) { exists = true } else if (prefix === FILE_PREFIX && id) { exists = !!(await env.ROOMS.head(getR2KeyForRoom({ slug: id, isApp: true }))) } else if (prefix === PUBLISH_PREFIX && id) { exists = !!(await pg .selectFrom('file') .where('publishedSlug', '=', id) .where('published', '=', true) .select('id') .executeTakeFirst()) } else if (prefix === ROOM_PREFIX && id) { exists = !!(await env.ROOMS.head(getR2KeyForRoom({ slug: id, isApp: false }))) } } catch (e) { warnings.push(`createSource check failed for ${raw}: ${e}`) exists = null } source = { raw, exists } } let associated = 0 let oldFormatUrls = 0 let missingInBucket = 0 let headFailures = 0 let totalSizeBytes = 0 let largestSizeBytes = 0 for (const a of assets) { if (a.associated) associated++ if (a.oldFormatUrl) oldFormatUrls++ if (a.inBucket === false) missingInBucket++ if (a.inBucket === null) headFailures++ if (a.sizeBytes !== null) { totalSizeBytes += a.sizeBytes largestSizeBytes = Math.max(largestSizeBytes, a.sizeBytes) } } const report: AdminFileAssetsResponseBody = { file: file ?? null, source, shapes: { total: totalShapes, byType: shapesByType }, assets: { // Every asset record in the snapshot; the upload-oriented counts below exclude // the `external` ones total: assets.length + external, associated, pending: assets.length - associated, external, oldFormatUrls, missingInBucket, headFailures, totalSizeBytes, largestSizeBytes, problems: assets .filter((a) => !a.associated || a.inBucket !== true) .map((a) => ({ assetId: a.assetId, objectName: a.objectName, src: a.src, fileIdMeta: a.fileIdMeta, inBucket: a.inBucket, dbRow: dbFileIdByObjectName.has(a.objectName) ? { fileId: dbFileIdByObjectName.get(a.objectName)! } : null, })), }, dbRows: { forThisFile: rowsForThisFile.length, orphaned }, warnings, } return json(report) }) // A board's shape without its contents. Answers "how big and how unusual is this board" for // perf reports, migration bugs, and support threads without anyone having to open it — and // without putting anything a user typed into the report. Read AdminFileStatsResponseBody // before adding a field: staying content-free is the point of this endpoint. .get('/app/admin/file-stats/:slug', async (res, env) => { const slug = res.params.slug assert(typeof slug === 'string', 'slug is required') const warnings: string[] = [] const pg = createPostgresConnectionPool(env, '/app/admin/file-stats') const [fileRow, snapshot, head] = await Promise.all([ pg .selectFrom('file') .where('id', '=', slug) .select([ 'ownerId', 'owningGroupId', 'createdAt', 'updatedAt', 'isDeleted', 'isEmpty', 'published', 'shared', 'sharedLinkType', 'createSource', ]) .executeTakeFirst(), getFileSnapshot(env, slug, true), env.ROOMS.head(getR2KeyForRoom({ slug, isApp: true })).catch((e) => { // Label only: an R2 error stringifies to the object key, which names the board console.error('file-stats snapshot head failed', e) warnings.push('snapshot head failed') return null }), ]) if (!snapshot) { throw new StatusError(404, `No persisted snapshot for ${slug}`) } const summary = summarizeSnapshotDocuments(snapshot.documents) // file_visitor, comment_thread, and comment all have a fileId index. file_state is // deliberately not counted here: its primary key is (userId, fileId), so counting by fileId // would sequentially scan the whole table. const countRows = async ( label: string, query: Promise<{ count: number | string | bigint } | undefined> ) => { try { return Number((await query)?.count ?? 0) } catch (e) { // Label only: a query error can carry table, column, and parameter detail console.error(`file-stats ${label} count failed`, e) warnings.push(`${label} count failed`) return 0 } } const [visitors, commentThreads, comments] = await Promise.all([ countRows( 'file_visitor', pg .selectFrom('file_visitor') .where('fileId', '=', slug) .select((eb) => eb.fn.countAll<number>().as('count')) .executeTakeFirst() ), countRows( 'comment_thread', pg .selectFrom('comment_thread') .where('fileId', '=', slug) .where('isDeleted', '=', false) .select((eb) => eb.fn.countAll<number>().as('count')) .executeTakeFirst() ), countRows( 'comment', pg .selectFrom('comment') .where('fileId', '=', slug) .where('isDeleted', '=', false) .select((eb) => eb.fn.countAll<number>().as('count')) .executeTakeFirst() ), ]) const schema = snapshot.schema as | { schemaVersion?: number; sequences?: Record<string, number> } | undefined const createSourceKind = fileRow?.createSource?.split('/')[0] ?? null const { recordsByTypeName, ...snapshotStats } = summary const stats: AdminFileStatsResponseBody = { file: fileRow ? { ownerType: fileRow.ownerId ? 'user' : fileRow.owningGroupId ? 'group' : 'none', createdAt: fileRow.createdAt, updatedAt: fileRow.updatedAt, isDeleted: fileRow.isDeleted, isEmpty: fileRow.isEmpty, published: fileRow.published, shared: fileRow.shared, sharedLinkType: fileRow.sharedLinkType, createSourceKind, } : null, snapshot: { sizeBytes: head?.size ?? null, clock: snapshot.clock ?? null, documentClock: snapshot.documentClock ?? null, tombstones: Object.keys(snapshot.tombstones ?? {}).length, records: snapshot.documents.length, recordsByTypeName, schemaVersion: schema?.schemaVersion ?? null, sequences: schema?.sequences ?? null, }, ...snapshotStats, collaboration: { visitors, commentThreads, comments }, warnings, } return json(stats) }) .get('/app/admin/download-tldr/:fileSlug', async (res, env) => { const fileSlug = res.params.fileSlug assert(typeof fileSlug === 'string', 'fileSlug is required') return await returnFileSnapshot(env, fileSlug, true) }) .get('/app/admin/download-legacy-tldr/:fileSlug', async (res, env) => { const fileSlug = res.params.fileSlug assert(typeof fileSlug === 'string', 'fileSlug is required') return await returnFileSnapshot(env, fileSlug, false) }) // The current welcome template (the file new workspaces fork their first file from), or // null when none is set and the committed default is used. Also reports whether the marked // file is still live and published: the resolver silently falls back to the default if it // isn't, so the admin needs to see a stale pointer rather than assume it's working. See // resolveWelcomeSnapshot. .get('/app/admin/welcome-template', async (_res, env) => { const pg = createPostgresConnectionPool(env, '/app/admin/welcome-template') const row = await pg.selectFrom('welcome_template').selectAll().executeTakeFirst() if (!row) return json(null) const file = await pg .selectFrom('file') .where('id', '=', row.fileId) .select(['published', 'isDeleted']) .executeTakeFirst() const live = !!file && !file.isDeleted && file.published return json({ ...row, live }) }) // Mark a published file as the welcome template. We store its publishedSlug, so the file // must be published first; new workspaces then fork its published snapshot. .post('/app/admin/welcome-template', async (req, env) => { const { fileId } = (await req.json()) as { fileId?: unknown } assert(typeof fileId === 'string' && fileId.length > 0, 'fileId (string) is required') const pg = createPostgresConnectionPool(env, '/app/admin/welcome-template') const file = await pg .selectFrom('file') .where('id', '=', fileId) .select(['id', 'published', 'publishedSlug', 'isDeleted']) .executeTakeFirst() if (!file) throw new StatusError(404, `File not found: ${fileId}`) if (!file.published) { throw new StatusError(400, 'File must be published before it can be the welcome template') } const updatedAt = Date.now() await pg .insertInto('welcome_template') .values({ id: true, fileId: file.id, publishedSlug: file.publishedSlug, updatedAt }) .onConflict((oc) => oc .column('id') .doUpdateSet({ fileId: file.id, publishedSlug: file.publishedSlug, updatedAt }) ) .execute() // Return the same shape as GET, including `live`, so the admin UI doesn't flash the // "not published" warning right after a successful set. const live = !file.isDeleted && file.published return json({ fileId: file.id, publishedSlug: file.publishedSlug, updatedAt, live }) }) // Clear the welcome template, reverting new workspaces to the committed default snapshot. .post('/app/admin/welcome-template/clear', async (_res, env) => { const pg = createPostgresConnectionPool(env, '/app/admin/welcome-template') await pg.deleteFrom('welcome_template').execute() return json({ cleared: true }) }) async function maybeHardDeleteLegacyFile({ id, env }: { id: string; env: Environment }) { return await getRoomDurableObject(env, id).__admin__hardDeleteIfLegacy() } async function hardDeleteAppFile({ pg, file, env, }: { env: Environment pg: ReturnType<typeof createPostgresConnectionPool> file: TlaFile }) { if (!file.isDeleted) { // do soft delete first if not done already; the outbox trigger records it await pg.updateTable('file').set('isDeleted', true).where('id', '=', file.id).execute() } // Session kicks and R2/room cleanup ride the terminal delete-row effect written by the // DELETE FROM file below, delivered via the post-delete poke() (sweep backstop ~30s); the // soft-delete row's effect is staleness-guarded, so it skips harmlessly if it runs after // the row is gone. // clean up assets eagerly const assets = await pg.selectFrom('asset').where('fileId', '=', file.id).selectAll().execute() for (const asset of assets) { await env.UPLOADS.delete(asset.objectName) // TODO: bust caches // it's tricky though. calling caches.default.delete() will only delete the cache entry // in the local datacenter so we'd need to do a global cache bust with the REST API // either that or maintain a KV store of deleted assets and check that before serving // could maybe use a bloom filter if that hurts perf too much. // although how would the bloom filter sync across workers 🤔 // since cache entries last a year we could store a timestamp in the KV and clean it periodically // or just let it grow forever, it's not that big. // const cacheUrl = new URL(`${appOrigin}/app/uploads/${asset.objectName}`) // console.log('Busting our cache entry', asset.objectName) // await caches.default.delete(cacheUrl) // console.log('Busting resize worker cache entry') // await env.IMAGE_RESIZE_WORKER.bustCache(cacheUrl.toString()) } // hard delete file (this will trigger a cascade delete of all remaining related records & R2 objects) await pg.deleteFrom('file').where('id', '=', file.id).execute() // Nudge the outbox so the delete's effects land promptly instead of waiting for the 30s // alarm sweep. poke() is cheap: it just schedules an alarm. Best-effort nudge: the sweep // backstops it, so a poke failure must not fail the request after the delete committed. await getFileEffectProcessor(env) .poke() .catch(() => {}) return new Response('Deleted', { status: 200 }) } async function deleteUserFromAnalytics( userId: string, env: Environment, sendProgress?: (step: string, message: string, details?: any) => void ) { if (!env.ANALYTICS_API_URL || !env.ANALYTICS_API_TOKEN) { sendProgress?.( 'analytics', 'Skipping analytics deletion - missing configuration (ANALYTICS_API_URL or ANALYTICS_API_TOKEN)' ) return } try { const response = await fetch(`${env.ANALYTICS_API_URL}/api/user-deletion`, { method: 'POST', headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${env.ANALYTICS_API_TOKEN}`, }, body: JSON.stringify({ clerk_id: userId, }), signal: AbortSignal.timeout(30000), }) if (!response.ok) { const errorText = await response.text().catch(() => 'Unknown error') throw new Error(`Analytics API returned ${response.status}: ${errorText}`) } const result = (await response.json()) as { success: boolean } sendProgress?.('analytics', 'Successfully deleted user data from analytics', { success: result.success, }) } catch (error) { const errorMessage = error instanceof Error ? error.message : String(error) console.error('Failed to delete user from analytics:', errorMessage) sendProgress?.('analytics', `Warning: Analytics deletion failed - ${errorMessage}`) } } async function performUserDeletion( userRow: any, env: any, sendProgress?: (step: string, message: string, details?: any) => void ) { const pg = createPostgresConnectionPool(env, '/app/admin/delete_user') // Step 1: Find all groups the user is the only owner of // This includes their home group (group.id = user.id) and any other groups they solely own sendProgress?.('groups', 'Finding groups to delete...') // Get all groups where this user is an owner const userOwnedGroupMemberships = await pg .selectFrom('group_user') .where('userId', '=', userRow.id) .where('role', '=', 'owner') .select('groupId') .execute() const groupsToDelete: string[] = [] for (const membership of userOwnedGroupMemberships) { // Check if this user is the only owner of this group const ownerCount = await pg .selectFrom('group_user') .where('groupId', '=', membership.groupId) .where('role', '=', 'owner') .select((eb) => eb.fn.countAll().as('count')) .executeTakeFirst() if (ownerCount && Number(ownerCount.count) === 1) { groupsToDelete.push(membership.groupId) } } sendProgress?.('groups', `Found ${groupsToDelete.length} groups to delete`, { groupCount: groupsToDelete.length, groupIds: groupsToDelete, }) // Step 2: Soft delete groups (the cleanup_deleted_group_trigger will soft delete their files) if (groupsToDelete.length > 0) { sendProgress?.('groups', 'Soft deleting groups...') await pg.updateTable('group').set('isDeleted', true).where('id', 'in', groupsToDelete).execute() } // Step 3: Get all files to hard delete const filesToDelete = new Map<string, TlaFile>() if (groupsToDelete.length > 0) { const groupFiles = await pg .selectFrom('file') .where('owningGroupId', 'in', groupsToDelete) .selectAll() .execute() for (const file of groupFiles) { filesToDelete.set(file.id, file) } } sendProgress?.('files', `Found ${filesToDelete.size} files to delete`, { fileCount: filesToDelete.size, }) // Allow time for soft deletes to propagate if (groupsToDelete.length > 0 || filesToDelete.size > 0) { await sleep(3000) } // Now hard delete all files for (const file of filesToDelete.values()) { sendProgress?.('files', `Hard deleting file '${file.name}' (${file.id})`) await hardDeleteAppFile({ pg, file, env }) } sendProgress?.('database', 'Cleaning up database records...') // Step 5: Hard delete groups and user in a transaction await pg.transaction().execute(async (tx) => { // Clean up tables that don't have CASCADE delete constraints await tx.deleteFrom('user_mutation_number').where('userId', '=', userRow.id).execute() // Clean up assets that reference this user (nullable foreign key) await tx.deleteFrom('asset').where('userId', '=', userRow.id).execute() // Remove user from all groups they're a member of (including ones they don't solely own) await tx.deleteFrom('group_user').where('userId', '=', userRow.id).execute() // Hard delete the groups (this will cascade delete group_user and group_file entries) if (groupsToDelete.length > 0) { await tx.deleteFrom('group').where('id', 'in', groupsToDelete).execute() } // Delete the user row (this will cascade delete any remaining related records) await tx.deleteFrom('user').where('id', '=', userRow.id).execute() }) sendProgress?.('clerk', 'Deleting user from Clerk...') // Delete user from Clerk const clerk = getClerkClient(env) await clerk.users.deleteUser(userRow.id) // Delete user from analytics service sendProgress?.('analytics', 'Deleting user from analytics...') await deleteUserFromAnalytics(userRow.id, env, sendProgress) sendProgress?.('durable_object', 'Cleaning up user durable object state...') // Clean up user durable object state and R2 data const user = getUserDurableObject(env, userRow.id) await user.admin_delete(userRow.id) }