fix(meeting-document): hold quota for in-flight generation claims
This commit is contained in:
parent
f4724ddf53
commit
b35676c75c
5 changed files with 994 additions and 179 deletions
|
|
@ -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<ProviderOutcome>
|
||||||
|
commit?: () => Promise<RpcResult>
|
||||||
|
findDocument?: RpcResult
|
||||||
|
} = {}): Harness & { run: () => ReturnType<typeof generateMeetingDocument> } {
|
||||||
|
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<ProviderOutcome>({ 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<typeof buildProviderRequest>[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)
|
||||||
|
})
|
||||||
|
|
@ -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<string, unknown>
|
||||||
|
}
|
||||||
|
|
||||||
|
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<RpcResult>
|
||||||
|
markFailed(actorId: string, idempotencyKey: string, code: GenerationFailureCode): Promise<DatabaseError | null>
|
||||||
|
commit(
|
||||||
|
actorId: string,
|
||||||
|
idempotencyKey: string,
|
||||||
|
result: ProviderDocumentResult,
|
||||||
|
latencyMs: number,
|
||||||
|
): Promise<RpcResult>
|
||||||
|
findDocument(actorId: string, documentId: string): Promise<RpcResult>
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Provider port: sends one document request and classifies the transport outcome. */
|
||||||
|
export interface MeetingDocumentProvider {
|
||||||
|
generate(request: ProviderRequest): Promise<ProviderOutcome>
|
||||||
|
}
|
||||||
|
|
||||||
|
export type GenerationLogger = (message: string, context?: Record<string, unknown>) => 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<string, unknown> {
|
||||||
|
return typeof value === 'object' && value !== null && !Array.isArray(value)
|
||||||
|
}
|
||||||
|
|
||||||
|
function result(status: number, body: Record<string, unknown>): 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<GenerationHttpResult> {
|
||||||
|
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<void>,
|
||||||
|
): Promise<GenerationHttpResult> {
|
||||||
|
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<GenerationHttpResult> {
|
||||||
|
const markFailed = async (code: GenerationFailureCode): Promise<void> => {
|
||||||
|
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
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -4,22 +4,19 @@ import { createServiceRoleClient } from '../_shared/quota.ts'
|
||||||
import {
|
import {
|
||||||
MeetingDocumentRequestError,
|
MeetingDocumentRequestError,
|
||||||
parseGenerateMeetingDocumentRequest,
|
parseGenerateMeetingDocumentRequest,
|
||||||
parseProviderDocumentResult,
|
|
||||||
} from '../_shared/meeting-document-contract.ts'
|
} from '../_shared/meeting-document-contract.ts'
|
||||||
import { buildMeetingDocumentSystemPrompt } from '../_shared/generative-ai-safety.ts'
|
|
||||||
import { readProviderKey } from '../_shared/provider-key.ts'
|
import { readProviderKey } from '../_shared/provider-key.ts'
|
||||||
|
import {
|
||||||
|
generateMeetingDocument,
|
||||||
|
type GenerationLogger,
|
||||||
|
type MeetingDocumentProvider,
|
||||||
|
type MeetingDocumentStore,
|
||||||
|
type ProviderOutcome,
|
||||||
|
} from './generation.ts'
|
||||||
|
|
||||||
interface ClaimResult {
|
type ServiceClient = ReturnType<typeof createServiceRoleClient>
|
||||||
claimed: boolean
|
|
||||||
status: 'processing' | 'succeeded' | 'failed'
|
const PROVIDER_TIMEOUT_MS = 60_000
|
||||||
documentId: string | null
|
|
||||||
meetingTitle: string
|
|
||||||
documentTitle: string
|
|
||||||
templateType: string
|
|
||||||
systemPrompt: string | null
|
|
||||||
transcript: string | null
|
|
||||||
model: string
|
|
||||||
}
|
|
||||||
|
|
||||||
function json(status: number, body: Record<string, unknown>): Response {
|
function json(status: number, body: Record<string, unknown>): Response {
|
||||||
return new Response(JSON.stringify(body), {
|
return new Response(JSON.stringify(body), {
|
||||||
|
|
@ -32,39 +29,84 @@ function json(status: number, body: Record<string, unknown>): Response {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
function isRecord(value: unknown): value is Record<string, unknown> {
|
const logError: GenerationLogger = (message, context) => {
|
||||||
return typeof value === 'object' && value !== null && !Array.isArray(value)
|
if (context === undefined) console.error(message)
|
||||||
|
else console.error(message, context)
|
||||||
}
|
}
|
||||||
|
|
||||||
function parseClaim(value: unknown): ClaimResult | null {
|
/** Supabase adapter for the persistence port (service-role RPCs). */
|
||||||
if (!isRecord(value)) return null
|
function createSupabaseStore(client: ServiceClient): MeetingDocumentStore {
|
||||||
if (
|
return {
|
||||||
typeof value.claimed !== 'boolean'
|
async claim(actorId, request) {
|
||||||
|| !['processing', 'succeeded', 'failed'].includes(String(value.status))
|
const { data, error } = await client.rpc('claim_meeting_document_generation_v1', {
|
||||||
|| (value.documentId !== null && typeof value.documentId !== 'string')
|
p_actor_id: actorId,
|
||||||
|| typeof value.meetingTitle !== 'string'
|
p_idempotency_key: request.idempotencyKey,
|
||||||
|| typeof value.documentTitle !== 'string'
|
p_meeting_id: request.meetingId,
|
||||||
|| typeof value.templateType !== 'string'
|
p_template_id: request.templateId,
|
||||||
|| (value.systemPrompt !== null && typeof value.systemPrompt !== 'string')
|
p_title: request.title,
|
||||||
|| (value.transcript !== null && typeof value.transcript !== 'string')
|
p_model: request.model,
|
||||||
|| typeof value.model !== 'string'
|
})
|
||||||
) {
|
return { data, error }
|
||||||
return null
|
},
|
||||||
|
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 {
|
/** Anthropic Messages API adapter for the provider port. */
|
||||||
const code = error?.code ?? ''
|
function createAnthropicProvider(apiKey: string): MeetingDocumentProvider {
|
||||||
const message = error?.message ?? ''
|
return {
|
||||||
if (code === '42501') return json(403, { error: 'forbidden' })
|
async generate(payload): Promise<ProviderOutcome> {
|
||||||
if (code === 'P0002') return json(404, { error: 'not_found' })
|
let response: Response
|
||||||
if (code === '22023') return json(400, { error: 'invalid_request' })
|
try {
|
||||||
if (message.includes('generation_quota_exceeded')) {
|
response = await fetch('https://api.anthropic.com/v1/messages', {
|
||||||
return json(429, { error: 'quota_exceeded' })
|
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) => {
|
Deno.serve(async (request: Request) => {
|
||||||
|
|
@ -72,30 +114,8 @@ Deno.serve(async (request: Request) => {
|
||||||
if (preflight) return preflight
|
if (preflight) return preflight
|
||||||
if (request.method !== 'POST') return json(405, { error: 'method_not_allowed' })
|
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<void> => {
|
|
||||||
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 {
|
try {
|
||||||
const actor = await requireUser(request)
|
const actor = await requireUser(request)
|
||||||
actorId = actor.id
|
|
||||||
|
|
||||||
let rawBody: unknown
|
let rawBody: unknown
|
||||||
try {
|
try {
|
||||||
|
|
@ -104,130 +124,19 @@ Deno.serve(async (request: Request) => {
|
||||||
throw new MeetingDocumentRequestError('Invalid JSON body')
|
throw new MeetingDocumentRequestError('Invalid JSON body')
|
||||||
}
|
}
|
||||||
const body = parseGenerateMeetingDocumentRequest(rawBody)
|
const body = parseGenerateMeetingDocumentRequest(rawBody)
|
||||||
idempotencyKey = body.idempotencyKey
|
|
||||||
|
|
||||||
const providerKey = readProviderKey('ANTHROPIC_API_KEY')
|
const providerKey = readProviderKey('ANTHROPIC_API_KEY')
|
||||||
if (providerKey.length === 0) {
|
if (providerKey.length === 0) {
|
||||||
return json(503, { error: 'provider_unavailable' })
|
return json(503, { error: 'provider_unavailable' })
|
||||||
}
|
}
|
||||||
|
|
||||||
const { data: claimData, error: claimError } = await serviceClient.rpc(
|
const outcome = await generateMeetingDocument(actor.id, body, {
|
||||||
'claim_meeting_document_generation_v1',
|
store: createSupabaseStore(createServiceRoleClient()),
|
||||||
{
|
provider: createAnthropicProvider(providerKey),
|
||||||
p_actor_id: actor.id,
|
now: () => Date.now(),
|
||||||
p_idempotency_key: body.idempotencyKey,
|
log: logError,
|
||||||
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,
|
|
||||||
})
|
})
|
||||||
|
return json(outcome.status, outcome.body)
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
if (error instanceof MeetingDocumentRequestError) {
|
if (error instanceof MeetingDocumentRequestError) {
|
||||||
return json(error.status, { error: 'invalid_request', message: error.message })
|
return json(error.status, { error: 'invalid_request', message: error.message })
|
||||||
|
|
|
||||||
|
|
@ -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;
|
||||||
|
|
@ -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;
|
||||||
Loading…
Add table
Add a link
Reference in a new issue