import { resolveVoiceInputProviderRegistry, type Settings } from "@opengeni/config"; import { VOICE_INPUT_ACCEPTED_MIME_TYPES } from "@opengeni/contracts"; import { filenameForMimeType, isAcceptedMimeType, normalizeMimeType, TRANSCRIPTION_PROVIDER_REQUEST_TIMEOUT_MILLISECONDS, type TranscriptionAvailabilityContext, type TranscriptionProvider, type TranscriptionService, TranscriptionServiceError, } from "@opengeni/core"; import type { Database } from "@opengeni/db"; import { createAzureOpenAiTranscriptionProvider } from "./providers/azure-openai"; import { createCodexSubscriptionTranscriptionProvider } from "./providers/codex-subscription"; import { createOpenAiTranscriptionProvider } from "./providers/openai"; import { createXaiSubscriptionTranscriptionProvider } from "./providers/xai-subscription"; export function createTranscriptionService(input: { settings: Settings; db: Database; fetch?: typeof fetch; codexFetch?: typeof fetch; probeCodex?: (context?: TranscriptionAvailabilityContext) => boolean | Promise; /** Test seam for exercising timeout and late-completion behavior quickly. */ providerRequestTimeoutMilliseconds?: number; /** Test seam for evaluating persisted absolute deadlines after delayed setup. */ now?: () => Date; }): TranscriptionService { const providers: TranscriptionProvider[] = resolveVoiceInputProviderRegistry(input.settings).map( (config) => { switch (config.kind) { case "openai": return createOpenAiTranscriptionProvider({ ...config, ...(input.fetch ? { fetch: input.fetch } : {}), }); case "azure-openai": return createAzureOpenAiTranscriptionProvider({ ...config, ...(input.fetch ? { fetch: input.fetch } : {}), }); case "codex-subscription": return createCodexSubscriptionTranscriptionProvider({ settings: input.settings, db: input.db, ...(input.codexFetch ? { fetch: input.codexFetch } : {}), ...(input.probeCodex ? { probe: input.probeCodex } : {}), }); case "supergrok-subscription": return createXaiSubscriptionTranscriptionProvider({ settings: input.settings, db: input.db, ...(input.fetch ? { fetch: input.fetch } : {}), }); default: throw new Error("Unsupported voice-input provider."); } }, ); const limits = { maxDurationSeconds: input.settings.voiceInputMaxDurationSeconds, maxSizeBytes: input.settings.voiceInputMaxSizeBytes, acceptedMimeTypes: [...VOICE_INPUT_ACCEPTED_MIME_TYPES], }; const providerRequestTimeoutMilliseconds = input.providerRequestTimeoutMilliseconds ?? TRANSCRIPTION_PROVIDER_REQUEST_TIMEOUT_MILLISECONDS; const now = input.now ?? (() => new Date()); return { limits: () => limits, async available(context) { return (await Promise.all(providers.map((provider) => provider.available(context)))).some( Boolean, ); }, async selectProvider(context) { return (await firstAvailable(orderedProviders(providers, context), context))?.id ?? null; }, async transcribe(request) { const mimeType = normalizeMimeType(request.mimeType); if (!isAcceptedMimeType(mimeType, limits.acceptedMimeTypes)) { throw new TranscriptionServiceError({ code: "not_supported", message: "Unsupported audio format.", }); } if (request.audio.byteLength > limits.maxSizeBytes) { throw new TranscriptionServiceError({ code: "too_large", message: "Audio is too large.", }); } if ( request.durationSeconds !== undefined && (!Number.isFinite(request.durationSeconds) || request.durationSeconds < 0 || request.durationSeconds > limits.maxDurationSeconds) ) { throw new TranscriptionServiceError({ code: "invalid_audio", message: "Invalid audio duration.", }); } const provider = request.providerId ? await exactAvailable(providers, request.providerId, { workspaceId: request.workspaceId, subjectId: request.subjectId, }) : await firstAvailable(orderedProviders(providers, request), request); if (!provider) { throw new TranscriptionServiceError({ fallbackSafe: true, code: "unavailable", message: "Transcription is unavailable.", }); } if (provider.supportsServerDeadline !== true) { throw new TranscriptionServiceError({ code: "unavailable", message: "Transcription provider does not support bounded requests.", }); } const startedAt = performance.now(); const remainingMilliseconds = request.providerDeadlineAt ? remainingTranscriptionProviderRequestMilliseconds(request.providerDeadlineAt, now()) : providerRequestTimeoutMilliseconds; if ( request.providerDeadlineAt && (!Number.isFinite(remainingMilliseconds) || remainingMilliseconds <= 0) ) { throw new TranscriptionServiceError({ code: "timeout", message: "Transcription provider deadline expired.", retryable: true, }); } const deadline = createProviderRequestDeadline(request.signal, remainingMilliseconds); let result: { text: string; languages: string[] }; try { result = await provider.transcribe({ audio: request.audio, mimeType, filename: filenameForMimeType(mimeType), workspaceId: request.workspaceId, accountId: request.accountId, subjectId: request.subjectId, requestId: request.requestId, signal: deadline.signal, }); if (deadline.timedOut && !request.signal?.aborted) { throw new TranscriptionServiceError({ code: "timeout", message: "Transcription provider timed out.", retryable: true, }); } } catch (error) { if (deadline.timedOut && !request.signal?.aborted) { throw new TranscriptionServiceError({ code: "timeout", message: "Transcription provider timed out.", retryable: true, }); } if ( error instanceof TranscriptionServiceError && error.fallbackSafe && !request.providerId && request.fallbackEnabled !== false && !request.signal?.aborted ) { const excludedProviders = [...(request.excludedProviders ?? []), provider.id]; if ( await firstAvailable( orderedProviders(providers, { ...request, excludedProviders }), request, ) ) { return await this.transcribe({ ...request, excludedProviders }); } } throw error; } finally { deadline.dispose(); } return { ...result, providerId: provider.id, audioSeconds: request.durationSeconds ?? 0, latencyMs: Math.round(performance.now() - startedAt), }; }, }; } export function remainingTranscriptionProviderRequestMilliseconds( providerDeadlineAt: Date, now: Date, ): number { return providerDeadlineAt.getTime() - now.getTime(); } function createProviderRequestDeadline( parentSignal: AbortSignal | undefined, timeoutMilliseconds: number, ): { signal: AbortSignal; readonly timedOut: boolean; dispose: () => void } { const controller = new AbortController(); let timedOut = false; const timeout = setTimeout( () => { timedOut = true; controller.abort(new DOMException("Transcription provider timed out", "TimeoutError")); }, Math.max(1, Math.ceil(timeoutMilliseconds)), ); const abortFromParent = () => { controller.abort(parentSignal?.reason); }; if (parentSignal) { if (parentSignal.aborted) abortFromParent(); else parentSignal.addEventListener("abort", abortFromParent, { once: true }); } return { signal: controller.signal, get timedOut() { return timedOut; }, dispose: () => { clearTimeout(timeout); parentSignal?.removeEventListener("abort", abortFromParent); }, }; } export function orderedProviders( providers: readonly TranscriptionProvider[], context: TranscriptionAvailabilityContext, ): TranscriptionProvider[] { const preferred = providers.filter((provider) => provider.id === context.preferredProvider); const ordered = context.preferredProvider ? context.fallbackEnabled === false ? preferred : [...preferred, ...providers.filter((provider) => provider.id !== context.preferredProvider)] : context.fallbackEnabled === false ? providers.slice(0, 1) : [...providers]; const remaining = context.afterProvider ? ordered.slice(ordered.findIndex((provider) => provider.id === context.afterProvider) + 1) : ordered; return remaining.filter((provider) => !context.excludedProviders?.includes(provider.id)); } async function firstAvailable( providers: readonly TranscriptionProvider[], context: TranscriptionAvailabilityContext, ) { for (const provider of providers) { try { if (await provider.available(context)) return provider; } catch { /* No audio sent; another configured provider may be ready. */ } } return null; } async function exactAvailable( providers: readonly TranscriptionProvider[], providerId: string, context: TranscriptionAvailabilityContext, ) { const provider = providers.find((candidate) => candidate.id === providerId); return provider && (await provider.available(context)) ? provider : null; }