From b35676c75cf91a2e9b19b2940ac40703c536d7ba Mon Sep 17 00:00:00 2001 From: Yun Chan Date: Mon, 28 Sep 2026 00:53:47 +0900 Subject: [PATCH] fix(meeting-document): hold quota for in-flight generation claims --- .../generation.test.ts | 219 ++++++++++++++ .../generate-meeting-document/generation.ts | 273 ++++++++++++++++++ .../generate-meeting-document/index.ts | 267 ++++++----------- ...ng_document_generation_in_flight_quota.sql | 252 ++++++++++++++++ ...-document-generation-quota.integration.sql | 162 +++++++++++ 5 files changed, 994 insertions(+), 179 deletions(-) create mode 100644 server/supabase/functions/generate-meeting-document/generation.test.ts create mode 100644 server/supabase/functions/generate-meeting-document/generation.ts create mode 100644 server/supabase/migrations/20260928000037_meeting_document_generation_in_flight_quota.sql create mode 100644 server/supabase/tests/meeting-document-generation-quota.integration.sql diff --git a/server/supabase/functions/generate-meeting-document/generation.test.ts b/server/supabase/functions/generate-meeting-document/generation.test.ts new file mode 100644 index 0000000..ac6a904 --- /dev/null +++ b/server/supabase/functions/generate-meeting-document/generation.test.ts @@ -0,0 +1,219 @@ +import { + assertEquals, + assertRejects, +} from 'https://deno.land/std@0.224.0/assert/mod.ts' +import type { GenerateMeetingDocumentRequest } from '../_shared/meeting-document-contract.ts' +import { + buildProviderRequest, + type GenerationFailureCode, + generateMeetingDocument, + mapDatabaseError, + type MeetingDocumentStore, + parseClaim, + PROVIDER_MAX_TOKENS, + type ProviderOutcome, + type ProviderRequest, + type RpcResult, +} from './generation.ts' + +const actorId = '44444444-4444-4444-8444-444444444444' +const request: GenerateMeetingDocumentRequest = { + meetingId: '11111111-1111-4111-8111-111111111111', + templateId: '22222222-2222-4222-8222-222222222222', + idempotencyKey: '33333333-3333-4333-8333-333333333333', + title: 'Weekly summary', + model: 'claude-haiku-4-5-20251001', +} + +const claimedPayload = { + claimed: true, + status: 'processing', + documentId: null, + meetingTitle: 'Weekly sync', + documentTitle: 'Weekly summary', + templateType: 'custom', + systemPrompt: 'Summarize.', + transcript: 'A: hello', + model: 'claude-haiku-4-5-20251001', +} + +const providerBody = { + content: [{ type: 'text', text: 'Generated body' }], + usage: { input_tokens: 10, output_tokens: 20 }, +} + +interface Harness { + store: MeetingDocumentStore + providerCalls: ProviderRequest[] + failures: GenerationFailureCode[] + commits: number + logs: string[] +} + +function harness(options: { + claim?: RpcResult + provider?: () => Promise + commit?: () => Promise + findDocument?: RpcResult +} = {}): Harness & { run: () => ReturnType } { + const state: Harness = { + store: undefined as unknown as MeetingDocumentStore, + providerCalls: [], + failures: [], + commits: 0, + logs: [], + } + state.store = { + claim: () => Promise.resolve(options.claim ?? { data: claimedPayload, error: null }), + markFailed: (_actor, _key, code) => { + state.failures.push(code) + return Promise.resolve(null) + }, + commit: () => { + state.commits += 1 + return options.commit?.() ?? Promise.resolve({ + data: { idempotent: false, document: { id: 'doc-1' } }, + error: null, + }) + }, + findDocument: () => Promise.resolve(options.findDocument ?? { data: { id: 'doc-1' }, error: null }), + } + let clock = 1_000 + return Object.assign(state, { + run: () => generateMeetingDocument(actorId, request, { + store: state.store, + provider: { + generate: (payload: ProviderRequest) => { + state.providerCalls.push(payload) + return options.provider?.() ?? Promise.resolve({ kind: 'ok', body: providerBody }) + }, + }, + now: () => (clock += 250), + log: (message: string) => { state.logs.push(message) }, + }), + }) +} + +Deno.test('a claim rejected for quota never reaches the provider', async () => { + const h = harness({ + claim: { data: null, error: { code: 'P0001', message: 'generation_quota_exceeded' } }, + }) + assertEquals(await h.run(), { status: 429, body: { error: 'quota_exceeded' } }) + assertEquals(h.providerCalls.length, 0) + assertEquals(h.commits, 0) +}) + +Deno.test('regression r1-17: an unexpected error after claim releases the in-flight quota unit', async () => { + const h = harness({ commit: () => Promise.reject(new Error('socket hang up')) }) + await assertRejects(() => h.run(), Error, 'socket hang up') + assertEquals(h.failures, ['commit_failed']) +}) + +Deno.test('an unexpected provider adapter error also releases the claim', async () => { + const h = harness({ provider: () => Promise.reject(new TypeError('boom')) }) + await assertRejects(() => h.run(), TypeError) + assertEquals(h.failures, ['commit_failed']) + assertEquals(h.commits, 0) +}) + +Deno.test('successful generation commits once and returns the document', async () => { + const h = harness() + assertEquals(await h.run(), { + status: 200, + body: { document: { id: 'doc-1' }, idempotent: false }, + }) + assertEquals(h.providerCalls.length, 1) + assertEquals(h.failures, []) +}) + +Deno.test('provider failures mark the claim failed with the matching code', async () => { + const cases: Array<[ProviderOutcome, number, string, GenerationFailureCode]> = [ + [{ kind: 'timeout' }, 504, 'provider_timeout', 'provider_timeout'], + [{ kind: 'network_error' }, 502, 'provider_request_failed', 'provider_request_failed'], + [{ kind: 'http_error', status: 529 }, 502, 'provider_request_failed', 'provider_request_failed'], + [{ kind: 'ok', body: null }, 502, 'provider_invalid_response', 'provider_invalid_response'], + ] + for (const [outcome, status, error, failure] of cases) { + const h = harness({ provider: () => Promise.resolve(outcome) }) + assertEquals(await h.run(), { status, body: { error } }) + assertEquals(h.failures, [failure]) + assertEquals(h.commits, 0) + } +}) + +Deno.test('commit quota rejection marks quota_exceeded and returns 429', async () => { + const h = harness({ + commit: () => Promise.resolve({ data: null, error: { code: 'P0001', message: 'generation_quota_exceeded' } }), + }) + assertEquals(await h.run(), { status: 429, body: { error: 'quota_exceeded' } }) + assertEquals(h.failures, ['quota_exceeded']) +}) + +Deno.test('other commit errors mark commit_failed and return 500', async () => { + const h = harness({ + commit: () => Promise.resolve({ data: null, error: { code: '55000', message: 'not committable' } }), + }) + assertEquals(await h.run(), { status: 500, body: { error: 'commit_failed' } }) + assertEquals(h.failures, ['commit_failed']) +}) + +Deno.test('claim without prompt or transcript is released', async () => { + const h = harness({ claim: { data: { ...claimedPayload, transcript: null }, error: null } }) + assertEquals(await h.run(), { status: 500, body: { error: 'internal_error' } }) + assertEquals(h.failures, ['commit_failed']) + assertEquals(h.providerCalls.length, 0) +}) + +Deno.test('replays return the stored document or a conflict without provider work', async () => { + const succeeded = harness({ + claim: { + data: { ...claimedPayload, claimed: false, status: 'succeeded', documentId: 'doc-1', systemPrompt: null, transcript: null }, + error: null, + }, + }) + assertEquals(await succeeded.run(), { + status: 200, + body: { document: { id: 'doc-1' }, idempotent: true }, + }) + + const inProgress = harness({ + claim: { data: { ...claimedPayload, claimed: false, systemPrompt: null, transcript: null }, error: null }, + }) + assertEquals(await inProgress.run(), { status: 409, body: { error: 'generation_in_progress' } }) + + const failed = harness({ + claim: { data: { ...claimedPayload, claimed: false, status: 'failed', systemPrompt: null, transcript: null }, error: null }, + }) + assertEquals(await failed.run(), { status: 409, body: { error: 'generation_failed' } }) + + for (const h of [succeeded, inProgress, failed]) { + assertEquals(h.providerCalls.length, 0) + assertEquals(h.failures, []) + } +}) + +Deno.test('parseClaim rejects a null claimed flag', () => { + assertEquals(parseClaim({ ...claimedPayload, claimed: null }), null) + assertEquals(parseClaim(claimedPayload)?.claimed, true) +}) + +Deno.test('mapDatabaseError keeps the HTTP contract', () => { + const log = () => {} + assertEquals(mapDatabaseError({ code: '42501' }, log).status, 403) + assertEquals(mapDatabaseError({ code: 'P0002' }, log).status, 404) + assertEquals(mapDatabaseError({ code: '22023' }, log).status, 400) + assertEquals(mapDatabaseError({ code: 'P0001', message: 'generation_quota_exceeded' }, log), { + status: 429, + body: { error: 'quota_exceeded' }, + }) + assertEquals(mapDatabaseError(null, log), { status: 500, body: { error: 'internal_error' } }) +}) + +Deno.test('buildProviderRequest bounds output and wraps the template prompt', () => { + const payload = buildProviderRequest({ ...claimedPayload, status: 'processing' } as Parameters[0]) + assertEquals(payload.max_tokens, PROVIDER_MAX_TOKENS) + assertEquals(payload.stream, false) + assertEquals(payload.model, 'claude-haiku-4-5-20251001') + assertEquals(payload.messages.length, 1) + assertEquals(payload.system.includes('Summarize.'), true) +}) diff --git a/server/supabase/functions/generate-meeting-document/generation.ts b/server/supabase/functions/generate-meeting-document/generation.ts new file mode 100644 index 0000000..18eb3a0 --- /dev/null +++ b/server/supabase/functions/generate-meeting-document/generation.ts @@ -0,0 +1,273 @@ +// Application core of generate-meeting-document. +// +// The flow (claim -> provider -> commit, with a failure marker on every exit +// after a successful claim) depends only on the ports below, so it can be +// exercised without Supabase, Anthropic or an HTTP runtime. index.ts is the +// adapter that wires real clients into these ports and turns the result into a +// Response. +// +// Quota: the claim RPC holds one unit of the user's allowance while the request +// is 'processing' (migration 20260928000037). Every path that leaves a claimed +// request unfinished must therefore call markFailed so the unit is released +// immediately instead of waiting for the lease to expire. + +import { buildMeetingDocumentSystemPrompt } from '../_shared/generative-ai-safety.ts' +import { + type GenerateMeetingDocumentRequest, + parseProviderDocumentResult, + type ProviderDocumentResult, +} from '../_shared/meeting-document-contract.ts' + +export type GenerationFailureCode = + | 'provider_timeout' + | 'provider_request_failed' + | 'provider_invalid_response' + | 'quota_exceeded' + | 'commit_failed' + +export interface DatabaseError { + code?: string + message?: string +} + +export interface RpcResult { + data: unknown + error: DatabaseError | null +} + +export interface GenerationHttpResult { + status: number + body: Record +} + +export interface ClaimResult { + claimed: boolean + status: 'processing' | 'succeeded' | 'failed' + documentId: string | null + meetingTitle: string + documentTitle: string + templateType: string + systemPrompt: string | null + transcript: string | null + model: string +} + +export interface ProviderRequest { + model: string + max_tokens: number + system: string + messages: Array<{ role: 'user'; content: string }> + stream: false +} + +/** Outcome of one provider call, already classified by the adapter. */ +export type ProviderOutcome = + | { kind: 'ok'; body: unknown } + | { kind: 'timeout' } + | { kind: 'network_error' } + | { kind: 'http_error'; status: number } + +/** Persistence port: the four RPC/table operations the flow needs. */ +export interface MeetingDocumentStore { + claim(actorId: string, request: GenerateMeetingDocumentRequest): Promise + markFailed(actorId: string, idempotencyKey: string, code: GenerationFailureCode): Promise + commit( + actorId: string, + idempotencyKey: string, + result: ProviderDocumentResult, + latencyMs: number, + ): Promise + findDocument(actorId: string, documentId: string): Promise +} + +/** Provider port: sends one document request and classifies the transport outcome. */ +export interface MeetingDocumentProvider { + generate(request: ProviderRequest): Promise +} + +export type GenerationLogger = (message: string, context?: Record) => void + +export interface GenerationDeps { + store: MeetingDocumentStore + provider: MeetingDocumentProvider + now: () => number + log: GenerationLogger +} + +export const PROVIDER_MAX_TOKENS = 4096 + +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value) +} + +function result(status: number, body: Record): GenerationHttpResult { + return { status, body } +} + +export function parseClaim(value: unknown): ClaimResult | null { + if (!isRecord(value)) return null + if ( + typeof value.claimed !== 'boolean' + || !['processing', 'succeeded', 'failed'].includes(String(value.status)) + || (value.documentId !== null && typeof value.documentId !== 'string') + || typeof value.meetingTitle !== 'string' + || typeof value.documentTitle !== 'string' + || typeof value.templateType !== 'string' + || (value.systemPrompt !== null && typeof value.systemPrompt !== 'string') + || (value.transcript !== null && typeof value.transcript !== 'string') + || typeof value.model !== 'string' + ) { + return null + } + return value as unknown as ClaimResult +} + +export function isQuotaExceeded(error: DatabaseError | null): boolean { + return error?.message?.includes('generation_quota_exceeded') === true +} + +export function mapDatabaseError(error: DatabaseError | null, log: GenerationLogger): GenerationHttpResult { + const code = error?.code ?? '' + if (code === '42501') return result(403, { error: 'forbidden' }) + if (code === 'P0002') return result(404, { error: 'not_found' }) + if (code === '22023') return result(400, { error: 'invalid_request' }) + if (isQuotaExceeded(error)) return result(429, { error: 'quota_exceeded' }) + log('meeting document database operation failed', { code: code || 'unknown' }) + return result(500, { error: 'internal_error' }) +} + +export function buildProviderRequest( + claim: ClaimResult & { systemPrompt: string; transcript: string }, +): ProviderRequest { + return { + model: claim.model, + max_tokens: PROVIDER_MAX_TOKENS, + system: buildMeetingDocumentSystemPrompt(claim.systemPrompt), + messages: [{ + role: 'user', + content: `회의 제목: ${claim.meetingTitle}\n\n전사록:\n${claim.transcript}`, + }], + stream: false, + } +} + +async function replayClaim( + actorId: string, + claim: ClaimResult, + deps: GenerationDeps, +): Promise { + if (claim.status === 'succeeded' && claim.documentId !== null) { + const { data: document, error } = await deps.store.findDocument(actorId, claim.documentId) + if (error !== null || document === null || document === undefined) { + deps.log('idempotent meeting document lookup failed', { code: error?.code ?? 'missing' }) + return result(500, { error: 'internal_error' }) + } + return result(200, { document, idempotent: true }) + } + return result(409, { + error: claim.status === 'processing' ? 'generation_in_progress' : 'generation_failed', + }) +} + +async function runClaimedGeneration( + actorId: string, + request: GenerateMeetingDocumentRequest, + claim: ClaimResult & { systemPrompt: string; transcript: string }, + deps: GenerationDeps, + markFailed: (code: GenerationFailureCode) => Promise, +): Promise { + const startedAt = deps.now() + const outcome = await deps.provider.generate(buildProviderRequest(claim)) + + if (outcome.kind === 'timeout') { + await markFailed('provider_timeout') + return result(504, { error: 'provider_timeout' }) + } + if (outcome.kind === 'network_error') { + await markFailed('provider_request_failed') + return result(502, { error: 'provider_request_failed' }) + } + if (outcome.kind === 'http_error') { + deps.log('meeting document provider request failed', { status: outcome.status }) + await markFailed('provider_request_failed') + return result(502, { error: 'provider_request_failed' }) + } + + const providerResult = parseProviderDocumentResult(outcome.body) + if (providerResult === null) { + deps.log('meeting document provider returned an invalid response shape') + await markFailed('provider_invalid_response') + return result(502, { error: 'provider_invalid_response' }) + } + + const { data: commitData, error: commitError } = await deps.store.commit( + actorId, + request.idempotencyKey, + providerResult, + Math.max(0, deps.now() - startedAt), + ) + if (commitError !== null) { + const quotaExceeded = isQuotaExceeded(commitError) + await markFailed(quotaExceeded ? 'quota_exceeded' : 'commit_failed') + if (quotaExceeded) return result(429, { error: 'quota_exceeded' }) + deps.log('meeting document atomic commit failed', { code: commitError.code ?? 'unknown' }) + return result(500, { error: 'commit_failed' }) + } + if (!isRecord(commitData) || !isRecord(commitData.document)) { + deps.log('meeting document commit returned an invalid shape') + return result(500, { error: 'internal_error' }) + } + + return result(200, { + document: commitData.document, + idempotent: commitData.idempotent === true, + }) +} + +/** + * Claims the request, calls the provider and commits the document. + * Throws only for unexpected adapter failures; a claimed request is marked + * failed before such an error propagates, so its quota unit is released. + */ +export async function generateMeetingDocument( + actorId: string, + request: GenerateMeetingDocumentRequest, + deps: GenerationDeps, +): Promise { + const markFailed = async (code: GenerationFailureCode): Promise => { + try { + const error = await deps.store.markFailed(actorId, request.idempotencyKey, code) + if (error !== null) { + deps.log('meeting document failure marker failed', { code: error.code ?? 'unknown' }) + } + } catch { + deps.log('meeting document failure marker failed', { code: 'exception' }) + } + } + + const { data: claimData, error: claimError } = await deps.store.claim(actorId, request) + if (claimError !== null) return mapDatabaseError(claimError, deps.log) + + const claim = parseClaim(claimData) + if (claim === null) { + deps.log('meeting document claim returned an invalid shape') + return result(500, { error: 'internal_error' }) + } + if (!claim.claimed) return await replayClaim(actorId, claim, deps) + + if (claim.systemPrompt === null || claim.transcript === null) { + await markFailed('commit_failed') + return result(500, { error: 'internal_error' }) + } + const readyClaim = { ...claim, systemPrompt: claim.systemPrompt, transcript: claim.transcript } + + try { + return await runClaimedGeneration(actorId, request, readyClaim, deps, markFailed) + } catch (error) { + // Release the in-flight quota unit held by this claim; the request is not + // going to be committed. A no-op if commit already moved it out of + // 'processing'. + await markFailed('commit_failed') + throw error + } +} diff --git a/server/supabase/functions/generate-meeting-document/index.ts b/server/supabase/functions/generate-meeting-document/index.ts index 07ec2ab..a4d6256 100644 --- a/server/supabase/functions/generate-meeting-document/index.ts +++ b/server/supabase/functions/generate-meeting-document/index.ts @@ -4,22 +4,19 @@ import { createServiceRoleClient } from '../_shared/quota.ts' import { MeetingDocumentRequestError, parseGenerateMeetingDocumentRequest, - parseProviderDocumentResult, } from '../_shared/meeting-document-contract.ts' -import { buildMeetingDocumentSystemPrompt } from '../_shared/generative-ai-safety.ts' import { readProviderKey } from '../_shared/provider-key.ts' +import { + generateMeetingDocument, + type GenerationLogger, + type MeetingDocumentProvider, + type MeetingDocumentStore, + type ProviderOutcome, +} from './generation.ts' -interface ClaimResult { - claimed: boolean - status: 'processing' | 'succeeded' | 'failed' - documentId: string | null - meetingTitle: string - documentTitle: string - templateType: string - systemPrompt: string | null - transcript: string | null - model: string -} +type ServiceClient = ReturnType + +const PROVIDER_TIMEOUT_MS = 60_000 function json(status: number, body: Record): Response { return new Response(JSON.stringify(body), { @@ -32,39 +29,84 @@ function json(status: number, body: Record): Response { }) } -function isRecord(value: unknown): value is Record { - return typeof value === 'object' && value !== null && !Array.isArray(value) +const logError: GenerationLogger = (message, context) => { + if (context === undefined) console.error(message) + else console.error(message, context) } -function parseClaim(value: unknown): ClaimResult | null { - if (!isRecord(value)) return null - if ( - typeof value.claimed !== 'boolean' - || !['processing', 'succeeded', 'failed'].includes(String(value.status)) - || (value.documentId !== null && typeof value.documentId !== 'string') - || typeof value.meetingTitle !== 'string' - || typeof value.documentTitle !== 'string' - || typeof value.templateType !== 'string' - || (value.systemPrompt !== null && typeof value.systemPrompt !== 'string') - || (value.transcript !== null && typeof value.transcript !== 'string') - || typeof value.model !== 'string' - ) { - return null +/** Supabase adapter for the persistence port (service-role RPCs). */ +function createSupabaseStore(client: ServiceClient): MeetingDocumentStore { + return { + async claim(actorId, request) { + const { data, error } = await client.rpc('claim_meeting_document_generation_v1', { + p_actor_id: actorId, + p_idempotency_key: request.idempotencyKey, + p_meeting_id: request.meetingId, + p_template_id: request.templateId, + p_title: request.title, + p_model: request.model, + }) + return { data, error } + }, + async markFailed(actorId, idempotencyKey, code) { + const { error } = await client.rpc('fail_meeting_document_generation_v1', { + p_actor_id: actorId, + p_idempotency_key: idempotencyKey, + p_error_code: code, + }) + return error + }, + async commit(actorId, idempotencyKey, providerResult, latencyMs) { + const { data, error } = await client.rpc('commit_meeting_document_generation_v1', { + p_actor_id: actorId, + p_idempotency_key: idempotencyKey, + p_content: providerResult.content, + p_latency_ms: latencyMs, + p_input_tokens: providerResult.inputTokens, + p_output_tokens: providerResult.outputTokens, + }) + return { data, error } + }, + async findDocument(actorId, documentId) { + const { data, error } = await client + .from('meeting_documents') + .select('*') + .eq('id', documentId) + .eq('user_id', actorId) + .maybeSingle() + return { data, error } + }, } - return value as unknown as ClaimResult } -function databaseErrorResponse(error: { code?: string; message?: string } | null): Response { - const code = error?.code ?? '' - const message = error?.message ?? '' - if (code === '42501') return json(403, { error: 'forbidden' }) - if (code === 'P0002') return json(404, { error: 'not_found' }) - if (code === '22023') return json(400, { error: 'invalid_request' }) - if (message.includes('generation_quota_exceeded')) { - return json(429, { error: 'quota_exceeded' }) +/** Anthropic Messages API adapter for the provider port. */ +function createAnthropicProvider(apiKey: string): MeetingDocumentProvider { + return { + async generate(payload): Promise { + let response: Response + try { + response = await fetch('https://api.anthropic.com/v1/messages', { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + 'x-api-key': apiKey, + 'anthropic-version': '2023-06-01', + }, + body: JSON.stringify(payload), + signal: AbortSignal.timeout(PROVIDER_TIMEOUT_MS), + }) + } catch (error) { + const timedOut = error instanceof DOMException && error.name === 'TimeoutError' + return timedOut ? { kind: 'timeout' } : { kind: 'network_error' } + } + if (!response.ok) return { kind: 'http_error', status: response.status } + try { + return { kind: 'ok', body: await response.json() } + } catch { + return { kind: 'ok', body: null } + } + }, } - console.error('meeting document database operation failed', { code: code || 'unknown' }) - return json(500, { error: 'internal_error' }) } Deno.serve(async (request: Request) => { @@ -72,30 +114,8 @@ Deno.serve(async (request: Request) => { if (preflight) return preflight if (request.method !== 'POST') return json(405, { error: 'method_not_allowed' }) - let actorId: string | null = null - let idempotencyKey: string | null = null - const serviceClient = createServiceRoleClient() - - const markFailed = async (errorCode: - | 'provider_timeout' - | 'provider_request_failed' - | 'provider_invalid_response' - | 'quota_exceeded' - | 'commit_failed'): Promise => { - if (actorId === null || idempotencyKey === null) return - const { error } = await serviceClient.rpc('fail_meeting_document_generation_v1', { - p_actor_id: actorId, - p_idempotency_key: idempotencyKey, - p_error_code: errorCode, - }) - if (error !== null) { - console.error('meeting document failure marker failed', { code: error.code ?? 'unknown' }) - } - } - try { const actor = await requireUser(request) - actorId = actor.id let rawBody: unknown try { @@ -104,130 +124,19 @@ Deno.serve(async (request: Request) => { throw new MeetingDocumentRequestError('Invalid JSON body') } const body = parseGenerateMeetingDocumentRequest(rawBody) - idempotencyKey = body.idempotencyKey const providerKey = readProviderKey('ANTHROPIC_API_KEY') if (providerKey.length === 0) { return json(503, { error: 'provider_unavailable' }) } - const { data: claimData, error: claimError } = await serviceClient.rpc( - 'claim_meeting_document_generation_v1', - { - p_actor_id: actor.id, - p_idempotency_key: body.idempotencyKey, - p_meeting_id: body.meetingId, - p_template_id: body.templateId, - p_title: body.title, - p_model: body.model, - }, - ) - if (claimError !== null) return databaseErrorResponse(claimError) - - const claim = parseClaim(claimData) - if (claim === null) { - console.error('meeting document claim returned an invalid shape') - return json(500, { error: 'internal_error' }) - } - if (!claim.claimed) { - if (claim.status === 'succeeded' && claim.documentId !== null) { - const { data: document, error } = await serviceClient - .from('meeting_documents') - .select('*') - .eq('id', claim.documentId) - .eq('user_id', actor.id) - .maybeSingle() - if (error !== null || document === null) { - console.error('idempotent meeting document lookup failed', { code: error?.code ?? 'missing' }) - return json(500, { error: 'internal_error' }) - } - return json(200, { document, idempotent: true }) - } - return json(409, { - error: claim.status === 'processing' ? 'generation_in_progress' : 'generation_failed', - }) - } - if (claim.systemPrompt === null || claim.transcript === null) { - await markFailed('commit_failed') - return json(500, { error: 'internal_error' }) - } - - const startedAt = Date.now() - let providerResponse: Response - try { - providerResponse = await fetch('https://api.anthropic.com/v1/messages', { - method: 'POST', - headers: { - 'Content-Type': 'application/json', - 'x-api-key': providerKey, - 'anthropic-version': '2023-06-01', - }, - body: JSON.stringify({ - model: claim.model, - max_tokens: 4096, - system: buildMeetingDocumentSystemPrompt(claim.systemPrompt), - messages: [{ - role: 'user', - content: `회의 제목: ${claim.meetingTitle}\n\n전사록:\n${claim.transcript}`, - }], - stream: false, - }), - signal: AbortSignal.timeout(60_000), - }) - } catch (error) { - const timedOut = error instanceof DOMException && error.name === 'TimeoutError' - await markFailed(timedOut ? 'provider_timeout' : 'provider_request_failed') - return json(timedOut ? 504 : 502, { - error: timedOut ? 'provider_timeout' : 'provider_request_failed', - }) - } - - if (!providerResponse.ok) { - console.error('meeting document provider request failed', { status: providerResponse.status }) - await markFailed('provider_request_failed') - return json(502, { error: 'provider_request_failed' }) - } - - let providerBody: unknown - try { - providerBody = await providerResponse.json() - } catch { - providerBody = null - } - const providerResult = parseProviderDocumentResult(providerBody) - if (providerResult === null) { - console.error('meeting document provider returned an invalid response shape') - await markFailed('provider_invalid_response') - return json(502, { error: 'provider_invalid_response' }) - } - - const { data: commitData, error: commitError } = await serviceClient.rpc( - 'commit_meeting_document_generation_v1', - { - p_actor_id: actor.id, - p_idempotency_key: body.idempotencyKey, - p_content: providerResult.content, - p_latency_ms: Math.max(0, Date.now() - startedAt), - p_input_tokens: providerResult.inputTokens, - p_output_tokens: providerResult.outputTokens, - }, - ) - if (commitError !== null) { - const quotaExceeded = commitError.message?.includes('generation_quota_exceeded') === true - await markFailed(quotaExceeded ? 'quota_exceeded' : 'commit_failed') - if (quotaExceeded) return json(429, { error: 'quota_exceeded' }) - console.error('meeting document atomic commit failed', { code: commitError.code ?? 'unknown' }) - return json(500, { error: 'commit_failed' }) - } - if (!isRecord(commitData) || !isRecord(commitData.document)) { - console.error('meeting document commit returned an invalid shape') - return json(500, { error: 'internal_error' }) - } - - return json(200, { - document: commitData.document, - idempotent: commitData.idempotent === true, + const outcome = await generateMeetingDocument(actor.id, body, { + store: createSupabaseStore(createServiceRoleClient()), + provider: createAnthropicProvider(providerKey), + now: () => Date.now(), + log: logError, }) + return json(outcome.status, outcome.body) } catch (error) { if (error instanceof MeetingDocumentRequestError) { return json(error.status, { error: 'invalid_request', message: error.message }) diff --git a/server/supabase/migrations/20260928000037_meeting_document_generation_in_flight_quota.sql b/server/supabase/migrations/20260928000037_meeting_document_generation_in_flight_quota.sql new file mode 100644 index 0000000..6919eba --- /dev/null +++ b/server/supabase/migrations/20260928000037_meeting_document_generation_in_flight_quota.sql @@ -0,0 +1,252 @@ +-- Meeting document generation: count in-flight claims against the allowance. +-- +-- claim_meeting_document_generation_v1 used to read daily_usage without any +-- lock or reservation, while usage was charged only by +-- commit_meeting_document_generation_v1 after the provider call. N parallel +-- requests with fresh idempotency keys therefore all passed the claim and each +-- paid for a 4096-token Anthropic call (Opus on paid tiers) before commit +-- rejected N-1 of them with generation_quota_exceeded. +-- +-- A request row in status 'processing' now holds one unit of the allowance: +-- * claim takes the same per-user/per-feature advisory lock as commit and +-- rejects when recorded usage plus in-flight claims would exceed the base +-- limit plus remaining overage credits; +-- * fail_meeting_document_generation_v1 releases the unit (status -> failed); +-- * commit converts it into recorded usage (status -> succeeded). +-- A crashed worker cannot hold a unit forever: in-flight rows stop counting +-- after a 10 minute lease, well beyond the 60 s provider timeout. +-- Replaying an existing idempotency key skips the quota check, since it never +-- starts new provider work. +-- +-- Also: `INSERT ... ON CONFLICT DO NOTHING RETURNING true INTO inserted` leaves +-- `inserted` NULL on a replay, so the old payload carried "claimed": null and +-- the edge function rejected it as an invalid shape (500) instead of returning +-- the idempotent result. The payload now always carries a boolean. + +CREATE INDEX IF NOT EXISTS meeting_document_generation_in_flight_idx + ON public.meeting_document_generation_requests(user_id, quota_feature, created_at) + WHERE status = 'processing'; + +CREATE OR REPLACE FUNCTION public.claim_meeting_document_generation_v1( + p_actor_id uuid, + p_idempotency_key uuid, + p_meeting_id uuid, + p_template_id uuid, + p_title text, + p_model text +) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY DEFINER +SET search_path = pg_catalog, public, auth, extensions +AS $$ +DECLARE + meeting_row public.meetings; + template_row public.user_templates; + transcript_value text; + transcript_digest text; + request_digest text; + request_row public.meeting_document_generation_requests; + inserted boolean := false; + tier_value text := 'free'; + overage_value integer := 0; + quota_feature_value text; + quota_limit_value integer; + quota_period_value text; + current_usage bigint := 0; + in_flight bigint := 0; + is_replay boolean := false; + safe_title text := trim(p_title); +BEGIN + IF p_actor_id IS NULL OR p_idempotency_key IS NULL OR p_meeting_id IS NULL OR p_template_id IS NULL THEN + RAISE EXCEPTION 'generation_identifiers_required' USING ERRCODE = '22023'; + END IF; + IF char_length(safe_title) NOT BETWEEN 1 AND 160 THEN + RAISE EXCEPTION 'invalid_document_title' USING ERRCODE = '22023'; + END IF; + IF p_model NOT IN ('claude-haiku-4-5-20251001', 'claude-sonnet-4-6', 'claude-opus-4-6') THEN + RAISE EXCEPTION 'invalid_generation_model' USING ERRCODE = '22023'; + END IF; + + SELECT * INTO meeting_row FROM public.meetings WHERE id = p_meeting_id; + IF meeting_row.id IS NULL THEN + RAISE EXCEPTION 'meeting_not_found' USING ERRCODE = 'P0002'; + END IF; + IF meeting_row.user_id <> p_actor_id AND NOT ( + meeting_row.team_id IS NOT NULL + AND ( + EXISTS ( + SELECT 1 FROM public.teams + WHERE id = meeting_row.team_id AND owner_id = p_actor_id + ) + OR EXISTS ( + SELECT 1 FROM public.team_members + WHERE team_id = meeting_row.team_id + AND user_id = p_actor_id + AND role IN ('owner', 'admin') + ) + ) + ) THEN + RAISE EXCEPTION 'meeting_generation_forbidden' USING ERRCODE = '42501'; + END IF; + + SELECT * INTO template_row + FROM public.user_templates + WHERE id = p_template_id + AND user_id = p_actor_id + AND template_kind = 'meeting_document'; + IF template_row.id IS NULL THEN + RAISE EXCEPTION 'meeting_template_not_found' USING ERRCODE = 'P0002'; + END IF; + + SELECT nullif(string_agg( + CASE WHEN nullif(trim(transcript.speaker), '') IS NULL + THEN transcript.text + ELSE trim(transcript.speaker) || ': ' || transcript.text + END, + E'\n' ORDER BY transcript.segment_index + ), '') + INTO transcript_value + FROM public.transcripts AS transcript + WHERE transcript.meeting_id = meeting_row.id; + + transcript_value := coalesce( + transcript_value, + nullif(trim(meeting_row.edited_transcript), ''), + nullif(trim(meeting_row.raw_transcript), '') + ); + IF transcript_value IS NULL THEN + RAISE EXCEPTION 'meeting_transcript_required' USING ERRCODE = '22023'; + END IF; + IF char_length(transcript_value) > 48000 THEN + RAISE EXCEPTION 'meeting_transcript_too_large' USING ERRCODE = '22023'; + END IF; + + SELECT coalesce(subscription.tier, 'free'), coalesce(subscription.overage_credits, 0) + INTO tier_value, overage_value + FROM public.subscriptions AS subscription + WHERE subscription.user_id = p_actor_id; + tier_value := coalesce(tier_value, 'free'); + overage_value := coalesce(overage_value, 0); + + IF tier_value = 'free' AND p_model <> 'claude-haiku-4-5-20251001' THEN + RAISE EXCEPTION 'generation_model_not_allowed' USING ERRCODE = '42501'; + END IF; + + quota_feature_value := CASE + WHEN p_model LIKE '%sonnet%' THEN 'llm_sonnet' + WHEN p_model LIKE '%opus%' THEN 'llm_opus' + ELSE 'llm_haiku' + END; + quota_limit_value := CASE + WHEN tier_value = 'free' AND quota_feature_value = 'llm_haiku' THEN 250 + WHEN tier_value = 'free' THEN 0 + WHEN tier_value = 'pro' AND quota_feature_value = 'llm_haiku' THEN 1500 + WHEN tier_value = 'pro' AND quota_feature_value = 'llm_sonnet' THEN 300 + WHEN tier_value = 'pro' AND quota_feature_value = 'llm_opus' THEN 50 + WHEN tier_value = 'pro_plus' AND quota_feature_value = 'llm_haiku' THEN -1 + WHEN tier_value = 'pro_plus' AND quota_feature_value = 'llm_sonnet' THEN 1500 + WHEN tier_value = 'pro_plus' AND quota_feature_value = 'llm_opus' THEN 300 + WHEN tier_value = 'team' AND quota_feature_value = 'llm_haiku' THEN -1 + WHEN tier_value = 'team' AND quota_feature_value = 'llm_sonnet' THEN 3000 + WHEN tier_value = 'team' AND quota_feature_value = 'llm_opus' THEN 600 + WHEN tier_value = 'enterprise' THEN -1 + ELSE 0 + END; + quota_period_value := CASE WHEN tier_value = 'free' THEN 'weekly' ELSE 'daily' END; + IF quota_limit_value = 0 THEN + RAISE EXCEPTION 'generation_quota_exceeded' USING ERRCODE = 'P0001'; + END IF; + + IF quota_limit_value > 0 THEN + -- Same lock key as commit_meeting_document_generation_v1, so claims and + -- commits for one user and feature see each other's in-flight rows. + PERFORM pg_advisory_xact_lock(hashtextextended(p_actor_id::text || ':' || quota_feature_value, 0)); + + SELECT EXISTS ( + SELECT 1 FROM public.meeting_document_generation_requests + WHERE user_id = p_actor_id AND idempotency_key = p_idempotency_key + ) INTO is_replay; + + IF NOT is_replay THEN + SELECT coalesce(sum(usage.count), 0) + INTO current_usage + FROM public.daily_usage AS usage + WHERE usage.user_id = p_actor_id + AND usage.feature = quota_feature_value + AND usage.date >= CASE quota_period_value + WHEN 'weekly' THEN current_date - 6 + ELSE current_date + END; + + SELECT count(*) + INTO in_flight + FROM public.meeting_document_generation_requests AS pending + WHERE pending.user_id = p_actor_id + AND pending.quota_feature = quota_feature_value + AND pending.status = 'processing' + AND pending.created_at > now() - interval '10 minutes'; + + -- Remaining allowance = unused base units + overage credits. With no + -- in-flight claims this is exactly the previous rule + -- (usage >= limit AND overage <= 0). + IF in_flight >= greatest(quota_limit_value - current_usage, 0) + greatest(overage_value, 0) THEN + RAISE EXCEPTION 'generation_quota_exceeded' USING ERRCODE = 'P0001'; + END IF; + END IF; + END IF; + + transcript_digest := encode(extensions.digest(transcript_value, 'sha256'), 'hex'); + request_digest := encode(extensions.digest( + jsonb_build_object( + 'meeting_id', meeting_row.id, + 'template_id', template_row.id, + 'template_revision', template_row.revision, + 'transcript_hash', transcript_digest, + 'title', safe_title, + 'model', p_model + )::text, + 'sha256' + ), 'hex'); + + INSERT INTO public.meeting_document_generation_requests( + user_id, idempotency_key, meeting_id, template_id, request_hash, + document_title, model, quota_feature, quota_limit, quota_period, + template_revision, template_type, transcript_hash, status + ) + VALUES ( + p_actor_id, p_idempotency_key, meeting_row.id, template_row.id, request_digest, + safe_title, p_model, quota_feature_value, quota_limit_value, quota_period_value, + template_row.revision, template_row.template_type, transcript_digest, 'processing' + ) + ON CONFLICT DO NOTHING + RETURNING true INTO inserted; + + SELECT * INTO request_row + FROM public.meeting_document_generation_requests + WHERE user_id = p_actor_id AND idempotency_key = p_idempotency_key; + + IF NOT coalesce(inserted, false) THEN + IF request_row.request_hash <> request_digest THEN + RAISE EXCEPTION 'generation_idempotency_conflict' USING ERRCODE = '22023'; + END IF; + END IF; + + RETURN jsonb_build_object( + 'claimed', coalesce(inserted, false), + 'status', request_row.status, + 'documentId', request_row.document_id, + 'meetingTitle', coalesce(meeting_row.title, 'Meeting'), + 'documentTitle', request_row.document_title, + 'templateType', request_row.template_type, + 'systemPrompt', CASE WHEN coalesce(inserted, false) THEN template_row.system_prompt ELSE NULL END, + 'transcript', CASE WHEN coalesce(inserted, false) THEN transcript_value ELSE NULL END, + 'model', request_row.model + ); +END; +$$; + +REVOKE ALL ON FUNCTION public.claim_meeting_document_generation_v1(uuid, uuid, uuid, uuid, text, text) + FROM PUBLIC, anon, authenticated; +GRANT EXECUTE ON FUNCTION public.claim_meeting_document_generation_v1(uuid, uuid, uuid, uuid, text, text) + TO service_role; diff --git a/server/supabase/tests/meeting-document-generation-quota.integration.sql b/server/supabase/tests/meeting-document-generation-quota.integration.sql new file mode 100644 index 0000000..faf646d --- /dev/null +++ b/server/supabase/tests/meeting-document-generation-quota.integration.sql @@ -0,0 +1,162 @@ +\set ON_ERROR_STOP on + +-- Regression (red-team r1-17): claim_meeting_document_generation_v1 used to +-- read daily_usage without counting in-flight claims, so N parallel requests +-- with fresh idempotency keys all passed the quota check and each paid for a +-- provider call before commit rejected N-1 of them. A 'processing' request now +-- holds one unit of the allowance until it is failed, committed or its lease +-- expires. + +BEGIN; + +CREATE OR REPLACE FUNCTION pg_temp.assert_true(condition boolean, message text) +RETURNS void +LANGUAGE plpgsql +AS $$ +BEGIN + IF condition IS NOT TRUE THEN + RAISE EXCEPTION 'assertion_failed: %', message; + END IF; +END; +$$; + +-- Returns the claim payload, or {"error": SQLERRM} when the claim raises. +CREATE OR REPLACE FUNCTION pg_temp.try_claim(p_actor uuid, p_key uuid, p_meeting uuid, p_template uuid, p_model text) +RETURNS jsonb +LANGUAGE plpgsql +AS $$ +BEGIN + RETURN public.claim_meeting_document_generation_v1( + p_actor, p_key, p_meeting, p_template, 'Quota fixture document', p_model + ); +EXCEPTION WHEN OTHERS THEN + RETURN jsonb_build_object('error', SQLERRM); +END; +$$; + +INSERT INTO auth.users ( + id, aud, role, email, encrypted_password, email_confirmed_at, + raw_app_meta_data, raw_user_meta_data, created_at, updated_at +) VALUES + ( + '37000000-0000-4000-8000-000000000001', 'authenticated', 'authenticated', + 'meeting-doc-quota-free@example.invalid', crypt('fixture-password', gen_salt('bf')), now(), + '{"provider":"email","providers":["email"]}'::jsonb, '{}'::jsonb, now(), now() + ), + ( + '37000000-0000-4000-8000-000000000002', 'authenticated', 'authenticated', + 'meeting-doc-quota-unlimited@example.invalid', crypt('fixture-password', gen_salt('bf')), now(), + '{"provider":"email","providers":["email"]}'::jsonb, '{}'::jsonb, now(), now() + ); + +UPDATE public.subscriptions SET tier = 'free', status = 'active', overage_credits = 0 +WHERE user_id = '37000000-0000-4000-8000-000000000001'; +UPDATE public.subscriptions SET tier = 'pro_plus', status = 'active', overage_credits = 0 +WHERE user_id = '37000000-0000-4000-8000-000000000002'; + +INSERT INTO public.meetings (id, user_id, title, status, raw_transcript) +VALUES + ('37100000-0000-4000-8000-000000000001', '37000000-0000-4000-8000-000000000001', + 'Quota fixture meeting', 'completed', 'Speaker one talked about the roadmap.'), + ('37100000-0000-4000-8000-000000000002', '37000000-0000-4000-8000-000000000002', + 'Unlimited fixture meeting', 'completed', 'Speaker two talked about hiring.'); + +INSERT INTO public.user_templates ( + id, user_id, template_kind, name, template_type, system_prompt, is_builtin +) VALUES + ('37200000-0000-4000-8000-000000000001', '37000000-0000-4000-8000-000000000001', + 'meeting_document', 'Quota fixture template', 'custom', 'Summarize the meeting.', false), + ('37200000-0000-4000-8000-000000000002', '37000000-0000-4000-8000-000000000002', + 'meeting_document', 'Unlimited fixture template', 'custom', 'Summarize the meeting.', false); + +-- One weekly Haiku unit left for the free user. +INSERT INTO public.daily_usage (user_id, date, feature, count) +VALUES ('37000000-0000-4000-8000-000000000001', CURRENT_DATE, 'llm_haiku', 249); + +DO $$ +DECLARE + actor constant uuid := '37000000-0000-4000-8000-000000000001'; + meeting constant uuid := '37100000-0000-4000-8000-000000000001'; + template constant uuid := '37200000-0000-4000-8000-000000000001'; + haiku constant text := 'claude-haiku-4-5-20251001'; + first jsonb; + parallel jsonb; + replay jsonb; + after_release jsonb; + after_commit jsonb; + overage_claim jsonb; + overage_parallel jsonb; + after_lease jsonb; + usage_count integer; +BEGIN + first := pg_temp.try_claim(actor, '37300000-0000-4000-8000-000000000001', meeting, template, haiku); + PERFORM pg_temp.assert_true((first->>'claimed')::boolean, 'last weekly unit can be claimed'); + + -- The bug: a second fresh key used to pass because only daily_usage was read. + parallel := pg_temp.try_claim(actor, '37300000-0000-4000-8000-000000000002', meeting, template, haiku); + PERFORM pg_temp.assert_true( + parallel->>'error' = 'generation_quota_exceeded', + 'an in-flight claim holds the last unit, so a parallel claim is rejected before provider work: ' || parallel::text + ); + + -- Replaying the in-flight key is idempotent, not a quota error. + replay := pg_temp.try_claim(actor, '37300000-0000-4000-8000-000000000001', meeting, template, haiku); + PERFORM pg_temp.assert_true( + replay->>'error' IS NULL AND NOT (replay->>'claimed')::boolean AND replay->>'status' = 'processing', + 'replaying the in-flight key reports processing: ' || replay::text + ); + + -- A failed request releases the unit it held. + PERFORM public.fail_meeting_document_generation_v1(actor, '37300000-0000-4000-8000-000000000001', 'provider_timeout'); + after_release := pg_temp.try_claim(actor, '37300000-0000-4000-8000-000000000003', meeting, template, haiku); + PERFORM pg_temp.assert_true((after_release->>'claimed')::boolean, + 'failure releases the in-flight unit: ' || after_release::text); + + -- Commit converts the held unit into recorded usage; the allowance is spent. + PERFORM public.commit_meeting_document_generation_v1( + actor, '37300000-0000-4000-8000-000000000003', 'Generated body', 10, 1, 1 + ); + SELECT count INTO usage_count FROM public.daily_usage + WHERE user_id = actor AND date = CURRENT_DATE AND feature = 'llm_haiku'; + PERFORM pg_temp.assert_true(usage_count = 250, 'commit records exactly one unit'); + after_commit := pg_temp.try_claim(actor, '37300000-0000-4000-8000-000000000004', meeting, template, haiku); + PERFORM pg_temp.assert_true(after_commit->>'error' = 'generation_quota_exceeded', + 'exhausted allowance rejects new claims: ' || after_commit::text); + + -- One overage credit covers exactly one in-flight request. + UPDATE public.subscriptions SET overage_credits = 1 WHERE user_id = actor; + overage_claim := pg_temp.try_claim(actor, '37300000-0000-4000-8000-000000000005', meeting, template, haiku); + PERFORM pg_temp.assert_true((overage_claim->>'claimed')::boolean, + 'an overage credit allows one claim past the base limit: ' || overage_claim::text); + overage_parallel := pg_temp.try_claim(actor, '37300000-0000-4000-8000-000000000006', meeting, template, haiku); + PERFORM pg_temp.assert_true(overage_parallel->>'error' = 'generation_quota_exceeded', + 'a single overage credit is not spent twice by parallel claims: ' || overage_parallel::text); + + -- A crashed request stops holding its unit once its lease expires. + UPDATE public.meeting_document_generation_requests + SET created_at = now() - interval '11 minutes' + WHERE user_id = actor AND idempotency_key = '37300000-0000-4000-8000-000000000005'; + after_lease := pg_temp.try_claim(actor, '37300000-0000-4000-8000-000000000007', meeting, template, haiku); + PERFORM pg_temp.assert_true((after_lease->>'claimed')::boolean, + 'an expired in-flight lease no longer holds a unit: ' || after_lease::text); +END; +$$; + +DO $$ +DECLARE + actor constant uuid := '37000000-0000-4000-8000-000000000002'; + meeting constant uuid := '37100000-0000-4000-8000-000000000002'; + template constant uuid := '37200000-0000-4000-8000-000000000002'; + first jsonb; + second jsonb; +BEGIN + first := pg_temp.try_claim(actor, '37400000-0000-4000-8000-000000000001', meeting, template, 'claude-haiku-4-5-20251001'); + second := pg_temp.try_claim(actor, '37400000-0000-4000-8000-000000000002', meeting, template, 'claude-haiku-4-5-20251001'); + PERFORM pg_temp.assert_true( + (first->>'claimed')::boolean AND (second->>'claimed')::boolean, + 'unlimited allowances are not throttled by in-flight claims' + ); +END; +$$; + +ROLLBACK;