File
Blob: src/worker/db/d1/repository.ts
| 1 | import { desc, eq } from "drizzle-orm"; |
| 2 | |
| 3 | import { randomHex } from "@/worker/auth/bytes"; |
| 4 | import { generateHostLabel } from "@/worker/auth/host-labels"; |
| 5 | import type { DavScope, JsonObject } from "@/worker/db/types"; |
| 6 | |
| 7 | import type { ControlPlaneDb } from "./client"; |
| 8 | import { |
| 9 | auditEvents, |
| 10 | patProjections, |
| 11 | subjects, |
| 12 | type NewAuditEventRow, |
| 13 | type NewPatProjectionRow, |
| 14 | type PatProjectionRow, |
| 15 | type SubjectRow, |
| 16 | } from "./schema"; |
| 17 | |
| 18 | export interface TesseraProfile { |
| 19 | sub: string; |
| 20 | email?: string | null; |
| 21 | displayName?: string | null; |
| 22 | } |
| 23 | |
| 24 | export interface BootstrappedSubject { |
| 25 | subject: SubjectRow; |
| 26 | created: boolean; |
| 27 | } |
| 28 | |
| 29 | function storageId(): string { |
| 30 | return `stg_${randomHex(16)}`; |
| 31 | } |
| 32 | |
| 33 | function auditId(): string { |
| 34 | return `aud_${randomHex(16)}`; |
| 35 | } |
| 36 | |
| 37 | function auditRow(input: { |
| 38 | subjectId?: string | null; |
| 39 | actorSubjectId?: string | null; |
| 40 | eventType: string; |
| 41 | createdAtMs: number; |
| 42 | data?: JsonObject; |
| 43 | }): NewAuditEventRow { |
| 44 | return { |
| 45 | id: auditId(), |
| 46 | subjectId: input.subjectId ?? null, |
| 47 | actorSubjectId: input.actorSubjectId ?? null, |
| 48 | eventType: input.eventType, |
| 49 | createdAtMs: input.createdAtMs, |
| 50 | data: input.data ?? {}, |
| 51 | }; |
| 52 | } |
| 53 | |
| 54 | export async function createAuditEvent( |
| 55 | db: ControlPlaneDb, |
| 56 | input: { |
| 57 | subjectId?: string | null; |
| 58 | actorSubjectId?: string | null; |
| 59 | eventType: string; |
| 60 | createdAtMs: number; |
| 61 | data?: JsonObject; |
| 62 | }, |
| 63 | ): Promise<void> { |
| 64 | await db.insert(auditEvents).values(auditRow(input)); |
| 65 | } |
| 66 | |
| 67 | export async function getSubjectById(db: ControlPlaneDb, subjectId: string): Promise<SubjectRow | undefined> { |
| 68 | return await db.query.subjects.findFirst({ where: eq(subjects.id, subjectId) }); |
| 69 | } |
| 70 | |
| 71 | export async function getSubjectByHostLabel(db: ControlPlaneDb, hostLabel: string): Promise<SubjectRow | undefined> { |
| 72 | return await db.query.subjects.findFirst({ where: eq(subjects.hostLabel, hostLabel) }); |
| 73 | } |
| 74 | |
| 75 | export async function bootstrapSubject( |
| 76 | db: ControlPlaneDb, |
| 77 | profile: TesseraProfile, |
| 78 | nowMs = Date.now(), |
| 79 | ): Promise<BootstrappedSubject> { |
| 80 | const existing = await getSubjectById(db, profile.sub); |
| 81 | if (existing) { |
| 82 | await db |
| 83 | .update(subjects) |
| 84 | .set({ |
| 85 | email: profile.email ?? null, |
| 86 | displayName: profile.displayName ?? null, |
| 87 | lastLoginAtMs: nowMs, |
| 88 | }) |
| 89 | .where(eq(subjects.id, profile.sub)); |
| 90 | |
| 91 | return { subject: (await getSubjectById(db, profile.sub)) ?? existing, created: false }; |
| 92 | } |
| 93 | |
| 94 | for (let attempt = 0; attempt < 8; attempt += 1) { |
| 95 | const subject = { |
| 96 | id: profile.sub, |
| 97 | storageId: storageId(), |
| 98 | hostLabel: generateHostLabel(), |
| 99 | email: profile.email ?? null, |
| 100 | displayName: profile.displayName ?? null, |
| 101 | createdAtMs: nowMs, |
| 102 | lastLoginAtMs: nowMs, |
| 103 | hostRotatedAtMs: null, |
| 104 | disabledAtMs: null, |
| 105 | }; |
| 106 | |
| 107 | try { |
| 108 | await db.insert(subjects).values(subject); |
| 109 | return { subject, created: true }; |
| 110 | } catch (error) { |
| 111 | const racedSubject = await getSubjectById(db, profile.sub); |
| 112 | if (racedSubject) return { subject: racedSubject, created: false }; |
| 113 | if (attempt === 7) throw error; |
| 114 | } |
| 115 | } |
| 116 | |
| 117 | throw new Error("Unable to bootstrap subject"); |
| 118 | } |
| 119 | |
| 120 | export async function listPatProjections(db: ControlPlaneDb, subjectId: string): Promise<PatProjectionRow[]> { |
| 121 | return await db |
| 122 | .select() |
| 123 | .from(patProjections) |
| 124 | .where(eq(patProjections.subjectId, subjectId)) |
| 125 | .orderBy(desc(patProjections.createdAtMs)); |
| 126 | } |
| 127 | |
| 128 | export async function getPatProjection( |
| 129 | db: ControlPlaneDb, |
| 130 | subjectId: string, |
| 131 | patId: string, |
| 132 | ): Promise<PatProjectionRow | undefined> { |
| 133 | const row = await db.query.patProjections.findFirst({ where: eq(patProjections.id, patId) }); |
| 134 | return row?.subjectId === subjectId ? row : undefined; |
| 135 | } |
| 136 | |
| 137 | export async function createPatProjection(db: ControlPlaneDb, row: NewPatProjectionRow): Promise<void> { |
| 138 | await db.insert(patProjections).values(row); |
| 139 | } |
| 140 | |
| 141 | export async function updatePatProjection( |
| 142 | db: ControlPlaneDb, |
| 143 | subjectId: string, |
| 144 | patId: string, |
| 145 | patch: { |
| 146 | name?: string; |
| 147 | scopes?: DavScope[]; |
| 148 | expiresAtMs?: number | null; |
| 149 | revokedAtMs?: number | null; |
| 150 | lastUsedAtMs?: number | null; |
| 151 | }, |
| 152 | ): Promise<void> { |
| 153 | const existing = await getPatProjection(db, subjectId, patId); |
| 154 | if (!existing) return; |
| 155 | await db |
| 156 | .update(patProjections) |
| 157 | .set({ |
| 158 | name: patch.name ?? existing.name, |
| 159 | scopes: patch.scopes ?? existing.scopes, |
| 160 | expiresAtMs: patch.expiresAtMs === undefined ? existing.expiresAtMs : patch.expiresAtMs, |
| 161 | revokedAtMs: patch.revokedAtMs === undefined ? existing.revokedAtMs : patch.revokedAtMs, |
| 162 | lastUsedAtMs: patch.lastUsedAtMs === undefined ? existing.lastUsedAtMs : patch.lastUsedAtMs, |
| 163 | }) |
| 164 | .where(eq(patProjections.id, patId)); |
| 165 | } |