diff --git a/app/pages/pm.vue b/app/pages/pm.vue
index a445cb8..587a6ce 100644
--- a/app/pages/pm.vue
+++ b/app/pages/pm.vue
@@ -128,6 +128,18 @@
+
+ {{ error }}
+
+
@@ -2793,7 +2805,17 @@ const handlePilotTransmission = async (message: string, source: 'text' | 'ptt' =
completedScenario.value = activeScenario.value
currentScreen.value = 'complete'
}
- } catch (e) {
+ } catch (e: any) {
+ const status = e?.status ?? e?.response?.status
+ if (status === 404) {
+ // Session expired (idle timeout) or backend restarted without the
+ // persisted session — reset to scenario selection for a clean restart.
+ pmLog.error('SESSION GONE (404) — resetting to scenario selection', { session: backendSessionId.value })
+ backendSessionId.value = null
+ error.value = 'Session expired — please start the scenario again.'
+ currentScreen.value = 'scenario'
+ return
+ }
pmLog.error('TRANSMIT FAILED', { transcript, session: backendSessionId.value, error: e })
console.error('Backend transmission failed', e)
setLastTransmission(`${prefix}: ${transcript} (backend failed)`)
@@ -2885,6 +2907,7 @@ const startMonitoring = async (flightPlan: any, scenario: Scenario) => {
// 4. Create a backend session with the flight-plan-derived variables
try {
+ error.value = ''
const session = await radioBackend.createSession(
scenario.startFlow,
backendVariables,
diff --git a/server/api/atc/ptt.post.ts b/server/api/atc/ptt.post.ts
index 6760663..0cd49cb 100644
--- a/server/api/atc/ptt.post.ts
+++ b/server/api/atc/ptt.post.ts
@@ -1,6 +1,6 @@
// server/api/atc/ptt.post.ts
import { createError, readBody } from "h3";
-import { writeFile, rm } from "node:fs/promises";
+import { writeFile, rm, readFile } from "node:fs/promises";
import { join } from "node:path";
import { tmpdir } from "node:os";
import { randomUUID } from "node:crypto";
@@ -9,6 +9,8 @@ import { getOpenAIClient } from "../../utils/openai";
import { createReadStream } from "node:fs";
import { TransmissionLog } from "../../models/TransmissionLog";
import { getUserFromEvent } from "../../utils/auth";
+import { enforceRateLimit, getClientIp } from "../../utils/rateLimit";
+import { recordUsage } from "../../utils/usage";
type AudioFormat = 'wav' | 'mp3' | 'ogg' | 'webm'
@@ -66,6 +68,13 @@ function decodeAudioPayload(encoded: string): Buffer {
return buffer;
}
+function wavDurationSeconds(buffer: Buffer): number | undefined {
+ if (buffer.length < 44 || buffer.toString('ascii', 0, 4) !== 'RIFF') return undefined;
+ const byteRate = buffer.readUInt32LE(28);
+ if (!byteRate) return undefined;
+ return Math.round(((buffer.length - 44) / byteRate) * 100) / 100;
+}
+
async function convertToWav(inputPath: string, outputPath: string) {
await sh("ffmpeg", [
"-y", "-i", inputPath,
@@ -76,7 +85,12 @@ async function convertToWav(inputPath: string, outputPath: string) {
]);
}
+const PTT_RATE_LIMIT_PER_MINUTE = 20;
+
export default defineEventHandler(async (event) => {
+ const user = await getUserFromEvent(event);
+ enforceRateLimit(event, 'atc-ptt', user ? String(user._id) : getClientIp(event), PTT_RATE_LIMIT_PER_MINUTE);
+
const body = await readBody(event);
if (!body.audio || !body.moduleId || !body.lessonId) {
@@ -115,6 +129,34 @@ export default defineEventHandler(async (event) => {
const transcribedText = transcription.text.trim();
+ // Audio length for usage accounting. The Whisper input is 16kHz mono WAV
+ // in the normal path; fall back to a byte-rate estimate if header parsing fails.
+ let audioSeconds: number | undefined;
+ try {
+ const wavBuffer = audioFileForWhisper === tmpAudioInput && format === 'wav'
+ ? audioBuffer
+ : await readFile(audioFileForWhisper);
+ audioSeconds = wavDurationSeconds(wavBuffer) ?? Math.round((wavBuffer.length / 32000) * 100) / 100;
+ } catch {
+ audioSeconds = Math.round((audioBuffer.length / 32000) * 100) / 100;
+ }
+
+ // Prefer the explicit top-level sessionId (Python backend session).
+ // Fall back to the legacy context.flags.session_id for older clients.
+ const sessionId = body.sessionId
+ ?? (typeof body.context?.flags?.session_id === 'string' ? body.context.flags.session_id : undefined);
+
+ // Whisper bills the audio even when nothing was recognized.
+ await recordUsage({
+ user: user?._id ? String(user._id) : undefined,
+ sessionId,
+ kind: 'stt',
+ provider: 'openai',
+ model: 'whisper-1',
+ endpoint: '/api/atc/ptt',
+ audioSeconds,
+ });
+
if (!transcribedText) {
throw createError({ statusCode: 400, statusMessage: "No speech detected in audio" });
}
@@ -125,12 +167,6 @@ export default defineEventHandler(async (event) => {
}
try {
- const user = await getUserFromEvent(event);
- // Prefer the explicit top-level sessionId (Python backend session).
- // Fall back to the legacy context.flags.session_id for older clients.
- const sessionId = body.sessionId
- ?? (typeof body.context?.flags?.session_id === 'string' ? body.context.flags.session_id : undefined);
-
await TransmissionLog.create({
user: user?._id,
role: "pilot",
diff --git a/server/api/atc/say.post.ts b/server/api/atc/say.post.ts
index 256183c..772e466 100644
--- a/server/api/atc/say.post.ts
+++ b/server/api/atc/say.post.ts
@@ -7,6 +7,9 @@ import {normalize, TTS_MODEL, normalizeATC} from "../../utils/normalize";
import { getServerRuntimeConfig } from "../../utils/runtimeConfig";
import {request} from "node:http";
import { TransmissionLog } from "../../models/TransmissionLog";
+import { requireUserSession } from "../../utils/auth";
+import { enforceRateLimit } from "../../utils/rateLimit";
+import { recordUsage } from "../../utils/usage";
function outDir() {
@@ -157,7 +160,12 @@ async function speachesTTS(
return Buffer.from(arr);
}
+const SAY_RATE_LIMIT_PER_MINUTE = 60;
+
export default defineEventHandler(async (event) => {
+ const user = await requireUserSession(event);
+ enforceRateLimit(event, 'atc-say', String(user._id), SAY_RATE_LIMIT_PER_MINUTE);
+
const runtimeConfig = getServerRuntimeConfig();
const body = await readBody<{
text?: string;
@@ -169,10 +177,14 @@ export default defineEventHandler(async (event) => {
tag?: string;
format?: AudioFmt | "smallest";
sessionId?: string;
+ /**
+ * Client already ran the radiotelephony normalizer on the text.
+ * Skip server-side normalizeATC — double-normalizing corrupts it
+ * (e.g. expandAirports spells the city name "MAIN" letter-by-letter).
+ */
+ preNormalized?: boolean;
}>(event);
- // const user = await requireUserSession(event);
-
const rawSessionId = typeof body?.sessionId === "string"
? body.sessionId.trim()
: "";
@@ -185,7 +197,7 @@ export default defineEventHandler(async (event) => {
const voice = (body?.voice || runtimeConfig.voiceId).trim();
const speed = Math.max(0.5, Math.min(2.0, body?.speed || 1.0));
- const normalized = normalizeATC(raw);
+ const normalized = body?.preNormalized ? raw : normalizeATC(raw);
if (!normalized) throw createError({ statusCode: 400, statusMessage: "normalized text empty" });
// Routing
@@ -342,9 +354,20 @@ export default defineEventHandler(async (event) => {
}
};
+ await recordUsage({
+ user: String(user._id),
+ sessionId,
+ kind: 'tts',
+ // Cache hits cost nothing regardless of which provider filled the cache.
+ provider: cacheHit ? 'cache' : ttsProvider,
+ model: modelUsed,
+ endpoint: '/api/atc/say',
+ characters: normalized.length,
+ });
+
try {
await TransmissionLog.create({
- // user: user._id,
+ user: user._id,
role: "atc",
channel: "say",
direction: "outgoing",
diff --git a/server/api/service/cron/waitlist-drip.get.ts b/server/api/service/cron/waitlist-drip.get.ts
index e8dba73..71cd8d0 100644
--- a/server/api/service/cron/waitlist-drip.get.ts
+++ b/server/api/service/cron/waitlist-drip.get.ts
@@ -12,6 +12,8 @@ import {
renderInvitationText,
} from '../../../utils/invitations'
import { buildWeeklyKpiReport, renderWeeklyKpiEmail, renderWeeklyKpiText } from '../../../utils/kpiReport'
+import { requireCronSecret } from '../../../utils/cron'
+import { maybeSendUsageQuotaAlert } from '../../../utils/usageAlert'
const DAY_MS = 1000 * 60 * 60 * 24
const INVITATION_DELAY_DAYS = 5
@@ -93,7 +95,9 @@ async function sendWeeklyKpiReportIfDue(now: Date) {
}
}
-export default defineEventHandler(async () => {
+export default defineEventHandler(async (event) => {
+ requireCronSecret(event)
+
const now = new Date()
const invitationCutoff = new Date(now.getTime() - INVITATION_DELAY_DAYS * DAY_MS)
const feedbackCutoff = new Date(now.getTime() - FEEDBACK_DELAY_DAYS * DAY_MS)
@@ -188,9 +192,15 @@ export default defineEventHandler(async () => {
const kpiReport = await sendWeeklyKpiReportIfDue(now)
console.log(`[waitlist-drip] weekly KPI report: ${kpiReport.sent ? 'sent' : `skipped (${kpiReport.skipped})`}`)
+ const usageAlert = await maybeSendUsageQuotaAlert(now)
+ if (usageAlert.sent) {
+ console.log(`[waitlist-drip] usage quota alert sent ($${usageAlert.costUsd?.toFixed(4)})`)
+ }
+
return {
invitationsSent,
feedbackRequests,
kpiReport,
+ usageAlert,
}
})
diff --git a/server/api/service/cron/weekly-kpi-report.get.ts b/server/api/service/cron/weekly-kpi-report.get.ts
index bd1e1f0..b6cf588 100644
--- a/server/api/service/cron/weekly-kpi-report.get.ts
+++ b/server/api/service/cron/weekly-kpi-report.get.ts
@@ -1,18 +1,12 @@
-import { createError, getQuery } from 'h3'
import { sendMail } from '../../../utils/notifications'
import { buildWeeklyKpiReport, renderWeeklyKpiEmail, renderWeeklyKpiText } from '../../../utils/kpiReport'
+import { requireCronSecret } from '../../../utils/cron'
+import { maybeSendUsageQuotaAlert } from '../../../utils/usageAlert'
const DEFAULT_KPI_RECIPIENT = 'opensquawk-kpi@faktorxmensch.com'
export default defineEventHandler(async (event) => {
- const secret = process.env.KPI_CRON_SECRET?.trim()
- if (secret) {
- const query = getQuery(event)
- const provided = typeof query.secret === 'string' ? query.secret : ''
- if (provided !== secret) {
- throw createError({ statusCode: 401, statusMessage: 'Invalid KPI cron secret.' })
- }
- }
+ requireCronSecret(event)
const report = await buildWeeklyKpiReport()
const to = process.env.KPI_EMAIL_TO || DEFAULT_KPI_RECIPIENT
@@ -25,6 +19,8 @@ export default defineEventHandler(async (event) => {
html: renderWeeklyKpiEmail(report),
})
+ const usageAlert = await maybeSendUsageQuotaAlert()
+
return {
success: true,
sent,
@@ -34,5 +30,7 @@ export default defineEventHandler(async (event) => {
totals: report.totals,
products: report.products,
smartGoals: report.smartGoals,
+ aiUsage: report.aiUsage,
+ usageAlert,
}
})
diff --git a/server/api/service/tools/latency.get.ts b/server/api/service/tools/latency.get.ts
index f779f47..02356bc 100644
--- a/server/api/service/tools/latency.get.ts
+++ b/server/api/service/tools/latency.get.ts
@@ -2,6 +2,9 @@
import { createError } from 'h3'
import { getOpenAIClient } from '../../../utils/openai'
import { getServerRuntimeConfig } from '../../../utils/runtimeConfig'
+import { requireUserSession } from '../../../utils/auth'
+import { enforceRateLimit } from '../../../utils/rateLimit'
+import { recordUsage } from '../../../utils/usage'
const SYSTEM_PROMPT =
'Check if the pilot readback contains ALL of: Frankfurt or EDDF, FL320, and 120.8 MHz. ' +
@@ -10,7 +13,10 @@ const SYSTEM_PROMPT =
const READBACK =
'Lufthanser four seven eight cleared fra via NORDA1A, climb 5000 feet, expect flight level tree too zero, dep 120 decimal 8, squawk 4213.';
-export default defineEventHandler(async () => {
+export default defineEventHandler(async (event) => {
+ const user = await requireUserSession(event)
+ enforceRateLimit(event, 'tools-latency', String(user._id), 5)
+
const client = getOpenAIClient()
const { llmModel } = getServerRuntimeConfig()
const model = llmModel || 'chatgpt-5-nano'
@@ -32,6 +38,16 @@ export default defineEventHandler(async () => {
const validResult = Number.isInteger(parsed) && parsed >= 0 && parsed <= 2 ? parsed : null
const latencyMs = Date.now() - started
+ await recordUsage({
+ user: String(user._id),
+ kind: 'llm',
+ provider: 'openai',
+ model,
+ endpoint: '/api/service/tools/latency',
+ inputTokens: response.usage?.prompt_tokens,
+ outputTokens: response.usage?.completion_tokens,
+ })
+
return {
result: validResult,
raw,
diff --git a/server/middleware/auth.global.ts b/server/middleware/auth.global.ts
index 9f030c6..67a1386 100644
--- a/server/middleware/auth.global.ts
+++ b/server/middleware/auth.global.ts
@@ -6,9 +6,6 @@ export default defineEventHandler(async (event) => {
if (!url.pathname.startsWith('/api/')) {
return
}
- if (url.pathname.startsWith('/api/atc/say')) {
- return
- }
if (url.pathname.startsWith('/api/service/')) {
return
}
diff --git a/server/models/UsageAlertDelivery.ts b/server/models/UsageAlertDelivery.ts
new file mode 100644
index 0000000..344cfe7
--- /dev/null
+++ b/server/models/UsageAlertDelivery.ts
@@ -0,0 +1,21 @@
+import mongoose from 'mongoose'
+
+export interface UsageAlertDeliveryDocument extends mongoose.Document {
+ monthKey: string
+ thresholdUsd: number
+ costUsd: number
+ recipient: string
+ sentAt: Date
+}
+
+const usageAlertDeliverySchema = new mongoose.Schema({
+ monthKey: { type: String, required: true, unique: true },
+ thresholdUsd: { type: Number, required: true },
+ costUsd: { type: Number, required: true },
+ recipient: { type: String, required: true },
+ sentAt: { type: Date, default: () => new Date() },
+})
+
+export const UsageAlertDelivery =
+ (mongoose.models.UsageAlertDelivery as mongoose.Model | undefined) ||
+ mongoose.model('UsageAlertDelivery', usageAlertDeliverySchema)
diff --git a/server/models/UsageEvent.ts b/server/models/UsageEvent.ts
new file mode 100644
index 0000000..ebbc448
--- /dev/null
+++ b/server/models/UsageEvent.ts
@@ -0,0 +1,45 @@
+import mongoose from 'mongoose'
+
+const { Schema } = mongoose
+
+export type UsageKind = 'stt' | 'tts' | 'llm'
+export type UsageProvider = 'openai' | 'speaches' | 'piper' | 'cache'
+
+export interface UsageEventDocument extends mongoose.Document {
+ user?: mongoose.Types.ObjectId
+ sessionId?: string
+ kind: UsageKind
+ provider: UsageProvider
+ model: string
+ endpoint: string
+ /** Audio length for STT in seconds */
+ audioSeconds?: number
+ /** Synthesized text length for TTS */
+ characters?: number
+ inputTokens?: number
+ outputTokens?: number
+ /** Estimated cost in USD (0 for self-hosted providers and cache hits) */
+ costUsd: number
+ createdAt: Date
+}
+
+const usageEventSchema = new mongoose.Schema({
+ user: { type: Schema.Types.ObjectId, ref: 'User', index: true },
+ sessionId: { type: String },
+ kind: { type: String, enum: ['stt', 'tts', 'llm'], required: true },
+ provider: { type: String, enum: ['openai', 'speaches', 'piper', 'cache'], required: true },
+ model: { type: String, required: true },
+ endpoint: { type: String, required: true },
+ audioSeconds: { type: Number },
+ characters: { type: Number },
+ inputTokens: { type: Number },
+ outputTokens: { type: Number },
+ costUsd: { type: Number, required: true, default: 0 },
+ createdAt: { type: Date, default: () => new Date(), index: true },
+})
+
+usageEventSchema.index({ user: 1, createdAt: -1 })
+
+export const UsageEvent =
+ (mongoose.models.UsageEvent as mongoose.Model | undefined) ||
+ mongoose.model('UsageEvent', usageEventSchema)
diff --git a/server/utils/cron.ts b/server/utils/cron.ts
new file mode 100644
index 0000000..0e9151f
--- /dev/null
+++ b/server/utils/cron.ts
@@ -0,0 +1,37 @@
+import type { H3Event } from 'h3'
+import { createError, getHeader, getQuery } from 'h3'
+
+let warnedMissingSecret = false
+
+/**
+ * Guards cron endpoints. Accepts the secret via `?secret=` query param
+ * (matches the existing KPI_CRON_SECRET convention, easy to use in a
+ * Coolify scheduled-task URL) or via the `x-cron-secret` header.
+ *
+ * When neither CRON_SECRET nor KPI_CRON_SECRET is configured the request
+ * is allowed with a loud warning so existing deployments keep working
+ * until the env var is set.
+ */
+export function requireCronSecret(event: H3Event) {
+ const secret = (process.env.CRON_SECRET || process.env.KPI_CRON_SECRET || '').trim()
+
+ if (!secret) {
+ if (!warnedMissingSecret) {
+ console.warn(
+ '[cron] CRON_SECRET is not set — cron endpoints are publicly callable. ' +
+ 'Set CRON_SECRET and append ?secret= to your scheduled-task URLs.',
+ )
+ warnedMissingSecret = true
+ }
+ return
+ }
+
+ const query = getQuery(event)
+ const provided =
+ (typeof query.secret === 'string' ? query.secret : '') ||
+ (getHeader(event, 'x-cron-secret') || '')
+
+ if (provided !== secret) {
+ throw createError({ statusCode: 401, statusMessage: 'Invalid cron secret.' })
+ }
+}
diff --git a/server/utils/kpiReport.ts b/server/utils/kpiReport.ts
index 624850e..ce84b9f 100644
--- a/server/utils/kpiReport.ts
+++ b/server/utils/kpiReport.ts
@@ -3,6 +3,7 @@ import { InvitationCode } from '../models/InvitationCode'
import { LandingAnalyticsEvent } from '../models/LandingAnalyticsEvent'
import { ProductUsageSession } from '../models/ProductUsageSession'
import { WaitlistEntry } from '../models/WaitlistEntry'
+import { summarizeUsage, getRollingCostUsd, type UsageSummary } from './usage'
export type KpiProduct = 'classroom' | 'liveatc'
@@ -53,6 +54,7 @@ export type WeeklyKpiReport = {
smartGoals: SmartGoal[]
landingTrend: { date: string; views: number; scrolled: number; waitlistEntries: number }[]
feedback: FeedbackItem[]
+ aiUsage: UsageSummary & { rolling30dCostUsd: number }
}
function startOfUtcDay(date: Date) {
@@ -189,6 +191,11 @@ export async function buildWeeklyKpiReport(now = new Date()): Promise= 1_000_000) return `${Math.round((value / 1_000_000) * 10) / 10}M`
+ if (value >= 1_000) return `${Math.round((value / 1_000) * 10) / 10}k`
+ return String(value)
+}
+
function escapeHtml(value: unknown) {
return String(value ?? '')
.replace(/&/g, '&')
@@ -334,6 +352,16 @@ export function renderWeeklyKpiText(report: WeeklyKpiReport) {
'Produkte:',
...report.products.map((product) => `- ${product.product}: Waitlist ${product.waitlistEntries}, User ${product.activatedUsers}, Nutzung ${formatSeconds(product.averageUsageSeconds)} avg, Feedback ${product.feedbackCount}`),
'',
+ 'AI-Nutzung (diese Woche):',
+ `- Kosten: ${formatUsd(report.aiUsage.costUsd)} (letzte 30 Tage: ${formatUsd(report.aiUsage.rolling30dCostUsd)})`,
+ `- STT: ${formatSeconds(report.aiUsage.sttSeconds)} Audio (${report.aiUsage.byKind.stt.events} Requests, ${formatUsd(report.aiUsage.byKind.stt.costUsd)})`,
+ `- TTS: ${formatTokens(report.aiUsage.ttsCharacters)} Zeichen (${report.aiUsage.byKind.tts.events} Requests, ${formatUsd(report.aiUsage.byKind.tts.costUsd)})`,
+ `- LLM: ${formatTokens(report.aiUsage.llmInputTokens)} in / ${formatTokens(report.aiUsage.llmOutputTokens)} out Tokens (${report.aiUsage.byKind.llm.events} Requests, ${formatUsd(report.aiUsage.byKind.llm.costUsd)})`,
+ 'Top User nach Kosten:',
+ ...(report.aiUsage.topUsers.length
+ ? report.aiUsage.topUsers.map((u) => `- ${u.email}: ${formatUsd(u.costUsd)} (${u.events} Requests)`)
+ : ['- keine Nutzungsdaten in diesem Zeitraum']),
+ '',
'Feedback:',
...report.feedback.map((item) => `- ${item.product} ${item.excitement}/5 ${item.from}: ${item.summary}`),
]
@@ -373,6 +401,33 @@ export function renderWeeklyKpiEmail(report: WeeklyKpiReport) {
`)
.join('')
+ const usageTopUserRows = report.aiUsage.topUsers.length
+ ? report.aiUsage.topUsers.map((u) => `
+
+ | ${escapeHtml(u.email)} |
+ ${u.events} |
+ ${escapeHtml(formatUsd(u.costUsd))} |
+
`).join('')
+ : '| Keine Nutzungsdaten in diesem Zeitraum. |
'
+
+ const usageKindRows = (['stt', 'tts', 'llm'] as const)
+ .map((kind) => {
+ const row = report.aiUsage.byKind[kind]
+ const detail = kind === 'stt'
+ ? `${formatSeconds(report.aiUsage.sttSeconds)} Audio`
+ : kind === 'tts'
+ ? `${formatTokens(report.aiUsage.ttsCharacters)} Zeichen`
+ : `${formatTokens(report.aiUsage.llmInputTokens)} in / ${formatTokens(report.aiUsage.llmOutputTokens)} out`
+ return `
+
+ | ${kind.toUpperCase()} |
+ ${row.events} |
+ ${escapeHtml(detail)} |
+ ${escapeHtml(formatUsd(row.costUsd))} |
+
`
+ })
+ .join('')
+
const feedbackRows = report.feedback.length
? report.feedback.map((item) => `
@@ -432,6 +487,11 @@ export function renderWeeklyKpiEmail(report: WeeklyKpiReport) {
| KPI | Ziel | Ist | Status |
${goalRows}
Produkte
| Produkt | Waitlist | Neue User | Nutzung Ø | Feedback |
${productRows}
+ AI-Nutzung & Kosten
+ Diese Woche: ${escapeHtml(formatUsd(report.aiUsage.costUsd))} · Letzte 30 Tage: ${escapeHtml(formatUsd(report.aiUsage.rolling30dCostUsd))}
+ | Art | Requests | Volumen | Kosten |
${usageKindRows}
+ Top User nach Kosten
+ | User | Requests | Kosten |
${usageTopUserRows}
Landing Verlauf
| Datum | Views | Scrolled | Waitlist |
${trendRows}
Feedback
diff --git a/server/utils/rateLimit.ts b/server/utils/rateLimit.ts
new file mode 100644
index 0000000..5eb57a4
--- /dev/null
+++ b/server/utils/rateLimit.ts
@@ -0,0 +1,65 @@
+import type { H3Event } from 'h3'
+import { createError } from 'h3'
+
+interface WindowEntry {
+ windowStart: number
+ count: number
+}
+
+// In-memory fixed-window limiter. Per process — good enough for a single
+// Coolify instance; swap for a shared store if the app ever scales out.
+const windows = new Map()
+
+const CLEANUP_INTERVAL_MS = 10 * 60 * 1000
+let lastCleanup = Date.now()
+
+function cleanup(now: number, windowMs: number) {
+ if (now - lastCleanup < CLEANUP_INTERVAL_MS) return
+ lastCleanup = now
+ for (const [key, entry] of windows) {
+ if (now - entry.windowStart > windowMs * 2) {
+ windows.delete(key)
+ }
+ }
+}
+
+export function getClientIp(event: H3Event): string {
+ const forwarded = event.node.req.headers['x-forwarded-for']
+ if (typeof forwarded === 'string' && forwarded.trim()) {
+ return forwarded.split(',')[0]!.trim()
+ }
+ return event.node.req.socket?.remoteAddress || 'unknown'
+}
+
+/**
+ * Throws 429 when `key` exceeds `limit` requests per `windowMs`.
+ * Use a user id as key where available, client IP otherwise.
+ */
+export function enforceRateLimit(
+ event: H3Event,
+ bucket: string,
+ key: string,
+ limit: number,
+ windowMs = 60_000,
+) {
+ const now = Date.now()
+ cleanup(now, windowMs)
+
+ const mapKey = `${bucket}:${key}`
+ const entry = windows.get(mapKey)
+
+ if (!entry || now - entry.windowStart >= windowMs) {
+ windows.set(mapKey, { windowStart: now, count: 1 })
+ return
+ }
+
+ entry.count += 1
+ if (entry.count > limit) {
+ const retryAfter = Math.ceil((entry.windowStart + windowMs - now) / 1000)
+ event.node.res.setHeader('Retry-After', String(Math.max(retryAfter, 1)))
+ throw createError({
+ statusCode: 429,
+ statusMessage: 'Too many requests — slow down and try again shortly.',
+ })
+ }
+}
diff --git a/server/utils/usage.ts b/server/utils/usage.ts
new file mode 100644
index 0000000..00a0c81
--- /dev/null
+++ b/server/utils/usage.ts
@@ -0,0 +1,165 @@
+import mongoose from 'mongoose'
+import { UsageEvent, type UsageKind, type UsageProvider } from '../models/UsageEvent'
+
+// OpenAI list prices (USD). Self-hosted providers (speaches/piper) cost 0.
+// Update here when models or prices change.
+const LLM_PRICES_PER_1M_TOKENS: Record = {
+ 'gpt-5-nano': { input: 0.05, output: 0.4 },
+ 'gpt-5-mini': { input: 0.25, output: 2.0 },
+ 'gpt-5': { input: 1.25, output: 10.0 },
+ 'gpt-4o-mini': { input: 0.15, output: 0.6 },
+ 'gpt-4o': { input: 2.5, output: 10.0 },
+}
+
+const STT_PRICES_PER_MINUTE: Record = {
+ 'whisper-1': 0.006,
+ 'gpt-4o-transcribe': 0.006,
+ 'gpt-4o-mini-transcribe': 0.003,
+}
+
+const TTS_PRICES_PER_1M_CHARS: Record = {
+ 'tts-1': 15,
+ 'tts-1-hd': 30,
+ 'gpt-4o-mini-tts': 12,
+}
+
+export interface UsageInput {
+ user?: mongoose.Types.ObjectId | string | null
+ sessionId?: string
+ kind: UsageKind
+ provider: UsageProvider
+ model: string
+ endpoint: string
+ audioSeconds?: number
+ characters?: number
+ inputTokens?: number
+ outputTokens?: number
+}
+
+export function estimateCostUsd(input: UsageInput): number {
+ if (input.provider !== 'openai') {
+ return 0
+ }
+
+ if (input.kind === 'llm') {
+ const price = LLM_PRICES_PER_1M_TOKENS[input.model]
+ if (!price) return 0
+ const inputCost = ((input.inputTokens || 0) / 1_000_000) * price.input
+ const outputCost = ((input.outputTokens || 0) / 1_000_000) * price.output
+ return inputCost + outputCost
+ }
+
+ if (input.kind === 'stt') {
+ const perMinute = STT_PRICES_PER_MINUTE[input.model]
+ if (!perMinute) return 0
+ return ((input.audioSeconds || 0) / 60) * perMinute
+ }
+
+ if (input.kind === 'tts') {
+ const per1M = TTS_PRICES_PER_1M_CHARS[input.model]
+ if (!per1M) return 0
+ return ((input.characters || 0) / 1_000_000) * per1M
+ }
+
+ return 0
+}
+
+/** Fire-and-forget usage recording — must never break the calling endpoint. */
+export async function recordUsage(input: UsageInput): Promise {
+ try {
+ await UsageEvent.create({
+ user: input.user || undefined,
+ sessionId: input.sessionId,
+ kind: input.kind,
+ provider: input.provider,
+ model: input.model,
+ endpoint: input.endpoint,
+ audioSeconds: input.audioSeconds,
+ characters: input.characters,
+ inputTokens: input.inputTokens,
+ outputTokens: input.outputTokens,
+ costUsd: estimateCostUsd(input),
+ })
+ } catch (error) {
+ console.warn('[usage] Recording usage event failed', error)
+ }
+}
+
+export interface UsageSummary {
+ events: number
+ costUsd: number
+ sttSeconds: number
+ ttsCharacters: number
+ llmInputTokens: number
+ llmOutputTokens: number
+ byKind: Record
+ topUsers: Array<{ email: string; events: number; costUsd: number }>
+}
+
+export async function summarizeUsage(periodStart: Date, periodEnd: Date): Promise {
+ const match = { createdAt: { $gte: periodStart, $lt: periodEnd } }
+
+ const [totals, kinds, topUsers] = await Promise.all([
+ UsageEvent.aggregate([
+ { $match: match },
+ {
+ $group: {
+ _id: null,
+ events: { $sum: 1 },
+ costUsd: { $sum: '$costUsd' },
+ sttSeconds: { $sum: { $ifNull: ['$audioSeconds', 0] } },
+ ttsCharacters: { $sum: { $ifNull: ['$characters', 0] } },
+ llmInputTokens: { $sum: { $ifNull: ['$inputTokens', 0] } },
+ llmOutputTokens: { $sum: { $ifNull: ['$outputTokens', 0] } },
+ },
+ },
+ ]),
+ UsageEvent.aggregate([
+ { $match: match },
+ { $group: { _id: '$kind', events: { $sum: 1 }, costUsd: { $sum: '$costUsd' } } },
+ ]),
+ UsageEvent.aggregate([
+ { $match: { ...match, user: { $ne: null } } },
+ { $group: { _id: '$user', events: { $sum: 1 }, costUsd: { $sum: '$costUsd' } } },
+ { $sort: { costUsd: -1, events: -1 } },
+ { $limit: 5 },
+ { $lookup: { from: 'users', localField: '_id', foreignField: '_id', as: 'userDoc' } },
+ ]),
+ ])
+
+ const t = totals[0] || {}
+ const byKind: UsageSummary['byKind'] = {
+ stt: { events: 0, costUsd: 0 },
+ tts: { events: 0, costUsd: 0 },
+ llm: { events: 0, costUsd: 0 },
+ }
+ for (const row of kinds) {
+ if (row._id === 'stt' || row._id === 'tts' || row._id === 'llm') {
+ byKind[row._id as UsageKind] = { events: row.events || 0, costUsd: row.costUsd || 0 }
+ }
+ }
+
+ return {
+ events: t.events || 0,
+ costUsd: t.costUsd || 0,
+ sttSeconds: t.sttSeconds || 0,
+ ttsCharacters: t.ttsCharacters || 0,
+ llmInputTokens: t.llmInputTokens || 0,
+ llmOutputTokens: t.llmOutputTokens || 0,
+ byKind,
+ topUsers: topUsers.map((row: any) => ({
+ email: row.userDoc?.[0]?.email || String(row._id),
+ events: row.events || 0,
+ costUsd: row.costUsd || 0,
+ })),
+ }
+}
+
+export async function getRollingCostUsd(days: number, now = new Date()): Promise {
+ const start = new Date(now.getTime() - days * 24 * 60 * 60 * 1000)
+ const rows = await UsageEvent.aggregate([
+ { $match: { createdAt: { $gte: start, $lt: now } } },
+ { $group: { _id: null, costUsd: { $sum: '$costUsd' } } },
+ ])
+ return rows[0]?.costUsd || 0
+}
diff --git a/server/utils/usageAlert.ts b/server/utils/usageAlert.ts
new file mode 100644
index 0000000..7858890
--- /dev/null
+++ b/server/utils/usageAlert.ts
@@ -0,0 +1,73 @@
+import { UsageAlertDelivery } from '../models/UsageAlertDelivery'
+import { sendMail } from './notifications'
+import { getRollingCostUsd, summarizeUsage } from './usage'
+
+const DEFAULT_THRESHOLD_USD = 5
+const DEFAULT_RECIPIENT = 'opensquawk-kpi@faktorxmensch.com'
+
+function getThresholdUsd() {
+ const parsed = Number.parseFloat(process.env.USAGE_ALERT_USD || '')
+ return Number.isFinite(parsed) && parsed > 0 ? parsed : DEFAULT_THRESHOLD_USD
+}
+
+/**
+ * Sends an alert mail when the rolling 30-day AI cost crosses the threshold
+ * (USAGE_ALERT_USD, default $5). At most one alert per calendar month.
+ * Designed to be called from cron endpoints — never throws.
+ */
+export async function maybeSendUsageQuotaAlert(now = new Date()) {
+ try {
+ const thresholdUsd = getThresholdUsd()
+ const costUsd = await getRollingCostUsd(30, now)
+
+ if (costUsd < thresholdUsd) {
+ return { sent: false, costUsd, thresholdUsd }
+ }
+
+ const monthKey = now.toISOString().slice(0, 7)
+ const alreadySent = await UsageAlertDelivery.exists({ monthKey })
+ if (alreadySent) {
+ return { sent: false, skipped: 'already-sent', costUsd, thresholdUsd }
+ }
+
+ const recipient = process.env.KPI_EMAIL_TO || DEFAULT_RECIPIENT
+ const periodStart = new Date(now.getTime() - 30 * 24 * 60 * 60 * 1000)
+ const summary = await summarizeUsage(periodStart, now)
+
+ const topUsers = summary.topUsers.length
+ ? summary.topUsers.map((u) => `- ${u.email}: $${u.costUsd.toFixed(4)} (${u.events} Requests)`).join('\n')
+ : '- keine User zugeordnet'
+
+ const text = [
+ `Die geschätzten AI-Kosten der letzten 30 Tage liegen bei $${costUsd.toFixed(4)} und haben die Schwelle von $${thresholdUsd.toFixed(2)} überschritten.`,
+ '',
+ `STT: ${Math.round(summary.sttSeconds / 60)} Minuten Audio ($${summary.byKind.stt.costUsd.toFixed(4)})`,
+ `TTS: ${summary.ttsCharacters} Zeichen ($${summary.byKind.tts.costUsd.toFixed(4)})`,
+ `LLM: ${summary.llmInputTokens} in / ${summary.llmOutputTokens} out Tokens ($${summary.byKind.llm.costUsd.toFixed(4)})`,
+ '',
+ 'Top User nach Kosten:',
+ topUsers,
+ '',
+ 'Es wird maximal eine Warnung pro Kalendermonat verschickt.',
+ ].join('\n')
+
+ const mailAccepted = await sendMail({
+ to: recipient,
+ subject: `OpenSquawk Kosten-Alarm: $${costUsd.toFixed(2)} in 30 Tagen (Schwelle $${thresholdUsd.toFixed(2)})`,
+ text,
+ })
+
+ await UsageAlertDelivery.create({
+ monthKey,
+ thresholdUsd,
+ costUsd,
+ recipient,
+ sentAt: now,
+ })
+
+ return { sent: true, mailAccepted, costUsd, thresholdUsd }
+ } catch (error) {
+ console.warn('[usage-alert] Quota alert check failed', error)
+ return { sent: false, error: String(error) }
+ }
+}
diff --git a/shared/composables/flightlab/useFlightLabAudio.ts b/shared/composables/flightlab/useFlightLabAudio.ts
index c67eade..7c475f0 100644
--- a/shared/composables/flightlab/useFlightLabAudio.ts
+++ b/shared/composables/flightlab/useFlightLabAudio.ts
@@ -1,5 +1,6 @@
// shared/composables/flightlab/useFlightLabAudio.ts
import { ref, computed } from 'vue'
+import { useAuthStore } from '~/stores/auth'
import type { FlightLabSound } from '../../data/flightlab/types'
import { getReadabilityProfile, createNoiseGenerators } from '../../utils/radioEffects'
import type { PizzicatoLite as PizzicatoLiteType } from '../../utils/pizzicatoLite'
@@ -134,9 +135,11 @@ export function useFlightLabAudio() {
}
isSpeaking.value = true
try {
- // Call the existing TTS API
+ // Call the existing TTS API (now requires an authenticated session)
+ const auth = useAuthStore()
const res = await $fetch('/api/atc/say', {
method: 'POST',
+ headers: auth.accessToken ? { Authorization: `Bearer ${auth.accessToken}` } : undefined,
body: {
text,
level: options?.readability ?? 5,