SPB Git forge

spb/polyllm

Public
15commits 1branches 0releases
2.2 MBsize
maindefault branch
13 days agolast push
TypeScript 97.4% SQL 1% JavaScript 0.9% CSS 0.6%
17.8 KB · 411 lines typescript
Raw Blame History
1import "server-only";2import { and, desc, eq, gte, lt, sql, type SQL } from "drizzle-orm";3import { getDb, usageRecords, conversations, projects } from "@/db";4import { isProviderId, type PolyModel, type ProviderId } from "@/lib/ai/core/types";5import { listRegistryModels } from "@/lib/ai/registry";6import { listCustomModels } from "@/lib/endpoints/service";7import { resolveRange, bucketKeys, previousPeriod, observedDays, type ResolvedRange, type UsageRangeKey, type Bucket } from "./time";8import { computeSavings, type SavingsOpportunity } from "./savings";910/** @deprecated kept for older imports — use `UsageRangeKey`. */11export type UsageRange = UsageRangeKey;1213export interface UsageQuery {14  range?: string | null;15  from?: string | null;16  to?: string | null;17  tz?: string | null;18  provider?: string | null;19  modelKey?: string | null;20  projectId?: string | null;21}2223export interface UsageFilters {24  provider: ProviderId | null;25  modelKey: string | null;26  projectId: string | null;27}2829export interface UsageKpis {30  requests: number;31  failures: number;32  stopped: number;33  errorRate: number;34  inputTokens: number;35  outputTokens: number;36  cachedTokens: number;37  reasoningTokens: number;38  totalTokens: number;39  costUsd: number;40  /** Requests whose cost is unknown (no pricing) — shown as a caveat. */41  unpricedRequests: number;42  avgLatencyMs: number;43  avgTtftMs: number | null;44  avgTokensPerSec: number | null;45  /** Average prompt size (input tokens per successful request). */46  avgContextTokens: number | null;47}4849export interface UsageSeriesPoint {50  bucket: string;51  requests: number;52  failures: number;53  inputTokens: number;54  outputTokens: number;55  costUsd: number;56}5758export interface ProviderAgg {59  provider: ProviderId;60  requests: number;61  failures: number;62  inputTokens: number;63  outputTokens: number;64  costUsd: number;65}6667export interface ModelAgg {68  modelKey: string;69  provider: ProviderId;70  requests: number;71  failures: number;72  inputTokens: number;73  outputTokens: number;74  cachedTokens: number;75  reasoningTokens: number;76  costUsd: number;77  avgLatencyMs: number;78  avgTtftMs: number | null;79  tokensPerSec: number | null;80  avgContextTokens: number | null;81}8283export interface UsageProjection {84  /** Cost per observed day in the range. */85  dailyAverageCost: number;86  /** dailyAverageCost × 30. */87  estimatedMonthlyCost: number;88  observedDays: number;89  previous: { costUsd: number; requests: number } | null;90  /** Relative change of cost vs the previous period of equal length; null when there is no previous data. */91  costTrendPct: number | null;92  requestsTrendPct: number | null;93}9495export interface RecentRecord {96  id: string;97  provider: ProviderId;98  modelKey: string;99  kind: string;100  status: string;101  errorCode: string | null;102  inputTokens: number;103  outputTokens: number;104  cachedTokens: number;105  reasoningTokens: number;106  costUsd: number | null;107  latencyMs: number | null;108  ttftMs: number | null;109  conversationId: string | null;110  conversationTitle: string | null;111  projectId: string | null;112  createdAt: string;113}114115export interface UsageSummary {116  range: { key: UsageRangeKey; from: string | null; to: string; bucket: Bucket; tz: string; days: number };117  filters: UsageFilters;118  kpis: UsageKpis;119  series: UsageSeriesPoint[];120  byProvider: ProviderAgg[];121  byModel: ModelAgg[];122  projection: UsageProjection;123  savings: SavingsOpportunity[];124  recent: RecentRecord[];125  /** Filter chips: providers/models seen in the range regardless of the active filter, and the user's projects. */126  facets: { providers: ProviderId[]; models: { modelKey: string; provider: ProviderId; requests: number }[]; projects: { id: string; name: string; icon: string | null; color: string | null }[] };127  generatedAt: string;128}129130/** Read the dashboard query from a request URL (`range`, `from`, `to`, `tz`, `provider`, `modelKey`, `projectId`). */131export function usageQueryFromUrl(url: string): UsageQuery {132  const sp = new URL(url).searchParams;133  return { range: sp.get("range"), from: sp.get("from"), to: sp.get("to"), tz: sp.get("tz"), provider: sp.get("provider"), modelKey: sp.get("modelKey"), projectId: sp.get("projectId") };134}135136const num = (v: unknown) => Number(v ?? 0);137const nul = (v: unknown) => (v === null || v === undefined ? null : Number(v));138139export function parseFilters(q: UsageQuery): UsageFilters {140  return {141    provider: q.provider && isProviderId(q.provider) ? q.provider : null,142    modelKey: q.modelKey && q.modelKey.length <= 160 ? q.modelKey : null,143    projectId: q.projectId && q.projectId.length <= 64 ? q.projectId : null,144  };145}146147function whereFor(userId: string, from: Date | null, to: Date, f: UsageFilters, opts: { ignoreProvider?: boolean; ignoreModel?: boolean } = {}): SQL {148  const conds: SQL[] = [eq(usageRecords.userId, userId), lt(usageRecords.createdAt, to)];149  if (from) conds.push(gte(usageRecords.createdAt, from));150  if (f.provider && !opts.ignoreProvider) conds.push(eq(usageRecords.provider, f.provider));151  if (f.modelKey && !opts.ignoreModel) conds.push(eq(usageRecords.modelKey, f.modelKey));152  if (f.projectId) {153    conds.push(sql`${usageRecords.conversationId} in (select ${conversations.id} from ${conversations} where ${conversations.userId} = ${userId} and ${conversations.projectId} = ${f.projectId})`);154  }155  return and(...conds)!;156}157158const TPS = sql<number>`avg(case when ${usageRecords.latencyMs} > 0 and ${usageRecords.outputTokens} > 0 and ${usageRecords.status} = 'ok' then ${usageRecords.outputTokens}::float8 / (greatest(${usageRecords.latencyMs} - coalesce(${usageRecords.ttftMs}, 0), 1)::float8 / 1000) end)`;159160async function totalsFor(where: SQL) {161  const [t] = await getDb()162    .select({163      requests: sql<number>`count(*)::int`,164      failures: sql<number>`sum(case when ${usageRecords.status} = 'error' then 1 else 0 end)::int`,165      stopped: sql<number>`sum(case when ${usageRecords.status} = 'stopped' then 1 else 0 end)::int`,166      inputTokens: sql<number>`coalesce(sum(${usageRecords.inputTokens}),0)::bigint`,167      outputTokens: sql<number>`coalesce(sum(${usageRecords.outputTokens}),0)::bigint`,168      cachedTokens: sql<number>`coalesce(sum(${usageRecords.cachedTokens}),0)::bigint`,169      reasoningTokens: sql<number>`coalesce(sum(${usageRecords.reasoningTokens}),0)::bigint`,170      costUsd: sql<number>`coalesce(sum(${usageRecords.costUsd}),0)::float8`,171      unpriced: sql<number>`sum(case when ${usageRecords.costUsd} is null and ${usageRecords.status} = 'ok' then 1 else 0 end)::int`,172      avgLatencyMs: sql<number>`coalesce(avg(${usageRecords.latencyMs}) filter (where ${usageRecords.status} = 'ok'),0)::float8`,173      avgTtftMs: sql<number | null>`avg(${usageRecords.ttftMs}) filter (where ${usageRecords.ttftMs} is not null and ${usageRecords.status} = 'ok')`,174      tps: sql<number | null>`${TPS}`,175      avgContext: sql<number | null>`avg(${usageRecords.inputTokens}) filter (where ${usageRecords.status} = 'ok')`,176      firstAt: sql<string | null>`min(${usageRecords.createdAt})`,177    })178    .from(usageRecords)179    .where(where);180  return t;181}182183export async function usageSummary(userId: string, query: UsageQuery, now = new Date()): Promise<UsageSummary> {184  const db = getDb();185  const range = resolveRange(query, now);186  const filters = parseFilters(query);187  const where = whereFor(userId, range.from, range.to, filters);188  const bucketExpr = range.bucket === "hour" ? sql<string>`to_char(${usageRecords.createdAt} at time zone ${range.tz}, 'YYYY-MM-DD"T"HH24:00')` : sql<string>`to_char(${usageRecords.createdAt} at time zone ${range.tz}, 'YYYY-MM-DD')`;189190  const [totals, seriesRows, byProviderRows, byModelRows, recentRows, projectRows, facetProviders, facetModels, registry, customModels] = await Promise.all([191    totalsFor(where),192    db193      .select({194        bucket: bucketExpr.as("bucket"),195        requests: sql<number>`count(*)::int`,196        failures: sql<number>`sum(case when ${usageRecords.status} = 'error' then 1 else 0 end)::int`,197        inputTokens: sql<number>`coalesce(sum(${usageRecords.inputTokens}),0)::bigint`,198        outputTokens: sql<number>`coalesce(sum(${usageRecords.outputTokens}),0)::bigint`,199        costUsd: sql<number>`coalesce(sum(${usageRecords.costUsd}),0)::float8`,200      })201      .from(usageRecords)202      .where(where)203      .groupBy(sql`1`)204      .orderBy(sql`1`),205    db206      .select({207        provider: usageRecords.provider,208        requests: sql<number>`count(*)::int`,209        failures: sql<number>`sum(case when ${usageRecords.status} = 'error' then 1 else 0 end)::int`,210        inputTokens: sql<number>`coalesce(sum(${usageRecords.inputTokens}),0)::bigint`,211        outputTokens: sql<number>`coalesce(sum(${usageRecords.outputTokens}),0)::bigint`,212        costUsd: sql<number>`coalesce(sum(${usageRecords.costUsd}),0)::float8`,213      })214      .from(usageRecords)215      .where(where)216      .groupBy(usageRecords.provider),217    db218      .select({219        modelKey: usageRecords.modelKey,220        provider: usageRecords.provider,221        requests: sql<number>`count(*)::int`,222        failures: sql<number>`sum(case when ${usageRecords.status} = 'error' then 1 else 0 end)::int`,223        inputTokens: sql<number>`coalesce(sum(${usageRecords.inputTokens}),0)::bigint`,224        outputTokens: sql<number>`coalesce(sum(${usageRecords.outputTokens}),0)::bigint`,225        cachedTokens: sql<number>`coalesce(sum(${usageRecords.cachedTokens}),0)::bigint`,226        reasoningTokens: sql<number>`coalesce(sum(${usageRecords.reasoningTokens}),0)::bigint`,227        costUsd: sql<number>`coalesce(sum(${usageRecords.costUsd}),0)::float8`,228        avgLatencyMs: sql<number>`coalesce(avg(${usageRecords.latencyMs}) filter (where ${usageRecords.status} = 'ok'),0)::float8`,229        avgTtftMs: sql<number | null>`avg(${usageRecords.ttftMs}) filter (where ${usageRecords.ttftMs} is not null and ${usageRecords.status} = 'ok')`,230        tokensPerSec: sql<number | null>`${TPS}`,231        avgContext: sql<number | null>`avg(${usageRecords.inputTokens}) filter (where ${usageRecords.status} = 'ok')`,232      })233      .from(usageRecords)234      .where(where)235      .groupBy(usageRecords.modelKey, usageRecords.provider)236      .orderBy(desc(sql`count(*)`)),237    db238      .select({239        id: usageRecords.id,240        provider: usageRecords.provider,241        modelKey: usageRecords.modelKey,242        kind: usageRecords.kind,243        status: usageRecords.status,244        errorCode: usageRecords.errorCode,245        inputTokens: usageRecords.inputTokens,246        outputTokens: usageRecords.outputTokens,247        cachedTokens: usageRecords.cachedTokens,248        reasoningTokens: usageRecords.reasoningTokens,249        costUsd: usageRecords.costUsd,250        latencyMs: usageRecords.latencyMs,251        ttftMs: usageRecords.ttftMs,252        conversationId: usageRecords.conversationId,253        conversationTitle: conversations.title,254        projectId: conversations.projectId,255        createdAt: usageRecords.createdAt,256      })257      .from(usageRecords)258      .leftJoin(conversations, eq(conversations.id, usageRecords.conversationId))259      .where(where)260      .orderBy(desc(usageRecords.createdAt))261      .limit(60),262    db.select({ id: projects.id, name: projects.name, icon: projects.icon, color: projects.color }).from(projects).where(and(eq(projects.userId, userId), eq(projects.archived, false))).orderBy(projects.sortOrder, projects.name),263    db264      .select({ provider: usageRecords.provider })265      .from(usageRecords)266      .where(whereFor(userId, range.from, range.to, filters, { ignoreProvider: true, ignoreModel: true }))267      .groupBy(usageRecords.provider),268    db269      .select({ modelKey: usageRecords.modelKey, provider: usageRecords.provider, requests: sql<number>`count(*)::int` })270      .from(usageRecords)271      .where(whereFor(userId, range.from, range.to, filters, { ignoreModel: true }))272      .groupBy(usageRecords.modelKey, usageRecords.provider)273      .orderBy(desc(sql`count(*)`))274      .limit(24),275    listRegistryModels({ includeDeprecated: true, includeHidden: true }),276    listCustomModels(userId).catch(() => [] as PolyModel[]),277  ]);278279  // Gap-filled series in the caller's time zone.280  const byBucket = new Map(seriesRows.map((s) => [s.bucket, s]));281  const keys = range.from ? bucketKeys(range.from, range.to, range.bucket, range.tz) : seriesRows.map((s) => s.bucket);282  const series: UsageSeriesPoint[] = keys.map((k) => {283    const s = byBucket.get(k);284    return { bucket: k, requests: num(s?.requests), failures: num(s?.failures), inputTokens: num(s?.inputTokens), outputTokens: num(s?.outputTokens), costUsd: num(s?.costUsd) };285  });286287  const requests = num(totals.requests);288  const failures = num(totals.failures);289  const kpis: UsageKpis = {290    requests,291    failures,292    stopped: num(totals.stopped),293    errorRate: requests ? failures / requests : 0,294    inputTokens: num(totals.inputTokens),295    outputTokens: num(totals.outputTokens),296    cachedTokens: num(totals.cachedTokens),297    reasoningTokens: num(totals.reasoningTokens),298    totalTokens: num(totals.inputTokens) + num(totals.outputTokens),299    costUsd: num(totals.costUsd),300    unpricedRequests: num(totals.unpriced),301    avgLatencyMs: num(totals.avgLatencyMs),302    avgTtftMs: nul(totals.avgTtftMs),303    avgTokensPerSec: nul(totals.tps),304    avgContextTokens: nul(totals.avgContext),305  };306307  // Projection + trend vs the previous period of equal length.308  const prev = previousPeriod(range);309  const prevTotals = prev ? await totalsFor(whereFor(userId, prev.from, prev.to, filters)) : null;310  const firstAt = totals.firstAt ? new Date(totals.firstAt) : null;311  const days = observedDays(range, now, firstAt);312  const dailyAverageCost = kpis.costUsd / days;313  const prevCost = prevTotals ? num(prevTotals.costUsd) : 0;314  const prevReq = prevTotals ? num(prevTotals.requests) : 0;315  const projection: UsageProjection = {316    dailyAverageCost,317    estimatedMonthlyCost: dailyAverageCost * 30,318    observedDays: Math.round(days * 100) / 100,319    previous: prevTotals ? { costUsd: prevCost, requests: prevReq } : null,320    costTrendPct: prevTotals && prevCost > 0 ? ((kpis.costUsd - prevCost) / prevCost) * 100 : null,321    requestsTrendPct: prevTotals && prevReq > 0 ? ((requests - prevReq) / prevReq) * 100 : null,322  };323324  const byModel: ModelAgg[] = byModelRows.map((m) => ({325    modelKey: m.modelKey,326    provider: m.provider,327    requests: num(m.requests),328    failures: num(m.failures),329    inputTokens: num(m.inputTokens),330    outputTokens: num(m.outputTokens),331    cachedTokens: num(m.cachedTokens),332    reasoningTokens: num(m.reasoningTokens),333    costUsd: num(m.costUsd),334    avgLatencyMs: num(m.avgLatencyMs),335    avgTtftMs: nul(m.avgTtftMs),336    tokensPerSec: nul(m.tokensPerSec),337    avgContextTokens: nul(m.avgContext),338  }));339340  const registryMap = new Map<string, PolyModel>([...registry, ...customModels].map((m) => [m.key, m]));341  const savings = computeSavings(byModel, registryMap);342343  return {344    range: { key: range.key, from: range.from?.toISOString() ?? null, to: range.to.toISOString(), bucket: range.bucket, tz: range.tz, days: range.days },345    filters,346    kpis,347    series,348    byProvider: byProviderRows.map((p) => ({ provider: p.provider, requests: num(p.requests), failures: num(p.failures), inputTokens: num(p.inputTokens), outputTokens: num(p.outputTokens), costUsd: num(p.costUsd) })),349    byModel,350    projection,351    savings,352    recent: recentRows.map((r) => ({ ...r, costUsd: r.costUsd ?? null, createdAt: r.createdAt.toISOString() })),353    facets: {354      providers: facetProviders.map((p) => p.provider),355      models: facetModels.map((m) => ({ modelKey: m.modelKey, provider: m.provider, requests: num(m.requests) })),356      projects: projectRows,357    },358    generatedAt: now.toISOString(),359  };360}361362/** Raw records for CSV export (same range/filters as the dashboard), newest first, capped. */363export async function usageRecordsForExport(userId: string, query: UsageQuery, limit = 50_000, now = new Date()) {364  const range = resolveRange(query, now);365  const filters = parseFilters(query);366  const rows = await getDb()367    .select({368      id: usageRecords.id,369      createdAt: usageRecords.createdAt,370      provider: usageRecords.provider,371      modelKey: usageRecords.modelKey,372      kind: usageRecords.kind,373      status: usageRecords.status,374      errorCode: usageRecords.errorCode,375      inputTokens: usageRecords.inputTokens,376      outputTokens: usageRecords.outputTokens,377      cachedTokens: usageRecords.cachedTokens,378      reasoningTokens: usageRecords.reasoningTokens,379      costUsd: usageRecords.costUsd,380      latencyMs: usageRecords.latencyMs,381      ttftMs: usageRecords.ttftMs,382      conversationId: usageRecords.conversationId,383      conversationTitle: conversations.title,384      projectId: conversations.projectId,385    })386    .from(usageRecords)387    .leftJoin(conversations, eq(conversations.id, usageRecords.conversationId))388    .where(whereFor(userId, range.from, range.to, filters))389    .orderBy(desc(usageRecords.createdAt))390    .limit(limit);391  return { range, filters, rows };392}393394export function usageCsv(rows: Awaited<ReturnType<typeof usageRecordsForExport>>["rows"]): string {395  const header = ["id", "created_at", "provider", "model_key", "kind", "status", "error_code", "input_tokens", "output_tokens", "cached_tokens", "reasoning_tokens", "cost_usd", "latency_ms", "ttft_ms", "conversation_id", "conversation_title", "project_id"];396  const esc = (v: unknown): string => {397    if (v === null || v === undefined) return "";398    const s = v instanceof Date ? v.toISOString() : String(v);399    // Neutralise spreadsheet formula injection and quote when needed.400    const safe = /^[=+\-@\t\r]/.test(s) ? `'${s}` : s;401    return /[",\n\r]/.test(safe) ? `"${safe.replace(/"/g, '""')}"` : safe;402  };403  const lines = [header.join(",")];404  for (const r of rows) {405    lines.push([r.id, r.createdAt, r.provider, r.modelKey, r.kind, r.status, r.errorCode, r.inputTokens, r.outputTokens, r.cachedTokens, r.reasoningTokens, r.costUsd, r.latencyMs, r.ttftMs, r.conversationId, r.conversationTitle, r.projectId].map(esc).join(","));406  }407  return `${lines.join("\r\n")}\r\n`;408}409410export type { ResolvedRange };411