import "server-only"; import { getDb, fetchRequests, requestAttempts, apiKeys, projects, proxySessions, usageEvents, eq, and, desc, sql, gte, lt, isNull, isNotNull, count, ilike, gt } from "@fetcha/db"; import type { SQL } from "drizzle-orm"; import { REQUEST_NETWORKS, filterWindow, type RequestFilters } from "@/lib/requests-filters"; /** * Drizzle queries backing the customer dashboard (Overview, Requests, Sessions, Usage, Analytics). * Every query is scoped to the caller's organization AND current project. Upstream data * (`costUsd`, `upstreamCostUsd`, attempt `provider`/`errorDetail`, session `provider`/`stickyKey`) * is never selected unless explicitly allowed by `organization.providerVisibility`. */ export interface Scope { organizationId: string; projectId: string; } const num = (v: unknown): number => (v === null || v === undefined ? 0 : Number(v)); const numOrNull = (v: unknown): number | null => (v === null || v === undefined ? null : Number(v)); function scoped(scope: Scope): SQL { return and(eq(fetchRequests.organizationId, scope.organizationId), eq(fetchRequests.projectId, scope.projectId))!; } function monthStart(d = new Date()): Date { return new Date(Date.UTC(d.getUTCFullYear(), d.getUTCMonth(), 1)); } // --------------------------------------------------------------------------- // Aggregates // --------------------------------------------------------------------------- export interface RequestStats { total: number; successful: number; failed: number; pending: number; successRate: number | null; bytes: number; spendUsd: number; avgLatencyMs: number | null; p50: number | null; p95: number | null; p99: number | null; avgAttempts: number | null; } /** Aggregate stats for a time window (inclusive from, exclusive to). */ export async function getRequestStats(scope: Scope, from: Date, to: Date | null = null): Promise { const db = getDb(); const [row] = await db .select({ total: count(), successful: sql`count(*) filter (where ${fetchRequests.status} = 'success')`.mapWith(num), failed: sql`count(*) filter (where ${fetchRequests.status} = 'failed')`.mapWith(num), pending: sql`count(*) filter (where ${fetchRequests.status} = 'pending')`.mapWith(num), bytes: sql`coalesce(sum(${fetchRequests.bytesIn} + ${fetchRequests.bytesOut}), 0)`.mapWith(num), spendUsd: sql`coalesce(sum(${fetchRequests.priceUsd}), 0)`.mapWith(num), avgLatencyMs: sql`avg(${fetchRequests.latencyMs})`.mapWith(numOrNull), p50: sql`percentile_cont(0.5) within group (order by ${fetchRequests.latencyMs})`.mapWith(numOrNull), p95: sql`percentile_cont(0.95) within group (order by ${fetchRequests.latencyMs})`.mapWith(numOrNull), p99: sql`percentile_cont(0.99) within group (order by ${fetchRequests.latencyMs})`.mapWith(numOrNull), avgAttempts: sql`avg(${fetchRequests.attempts}) filter (where ${fetchRequests.attempts} > 0)`.mapWith(numOrNull), }) .from(fetchRequests) .where(and(scoped(scope), gte(fetchRequests.createdAt, from), to ? lt(fetchRequests.createdAt, to) : undefined)); const r = row!; const decided = r.successful + r.failed; return { ...r, successRate: decided > 0 ? (r.successful / decided) * 100 : null }; } export async function getMonthStats(scope: Scope, month = monthStart()): Promise { const next = new Date(Date.UTC(month.getUTCFullYear(), month.getUTCMonth() + 1, 1)); return getRequestStats(scope, month, next); } export interface SeriesPoint { /** ISO timestamp of the bucket start (UTC). */ t: string; success: number; failed: number; p50: number | null; p95: number | null; } /** Requests + latency percentiles bucketed by hour or day, with empty buckets filled in. */ export async function getRequestSeries(scope: Scope, from: Date, bucket: "hour" | "day", to: Date = new Date()): Promise { const db = getDb(); // `bucket` is a closed union, inlined (not bound) so SELECT and GROUP BY use the identical expression. const trunc = sql`date_trunc(${sql.raw(bucket === "hour" ? "'hour'" : "'day'")}, ${fetchRequests.createdAt} at time zone 'UTC')`; const rows = await db .select({ t: sql`to_char(${trunc}, 'YYYY-MM-DD"T"HH24:MI:SS"Z"')`, success: sql`count(*) filter (where ${fetchRequests.status} = 'success')`.mapWith(num), failed: sql`count(*) filter (where ${fetchRequests.status} = 'failed')`.mapWith(num), p50: sql`percentile_cont(0.5) within group (order by ${fetchRequests.latencyMs})`.mapWith(numOrNull), p95: sql`percentile_cont(0.95) within group (order by ${fetchRequests.latencyMs})`.mapWith(numOrNull), }) .from(fetchRequests) .where(and(scoped(scope), gte(fetchRequests.createdAt, from), lt(fetchRequests.createdAt, to))) .groupBy(trunc) .orderBy(trunc); const byKey = new Map(rows.map((r) => [r.t, r])); const step = bucket === "hour" ? 3_600_000 : 86_400_000; const start = bucket === "hour" ? Date.UTC(from.getUTCFullYear(), from.getUTCMonth(), from.getUTCDate(), from.getUTCHours()) : Date.UTC(from.getUTCFullYear(), from.getUTCMonth(), from.getUTCDate()); const out: SeriesPoint[] = []; for (let ts = start; ts < to.getTime(); ts += step) { const key = new Date(ts).toISOString().replace(/\.\d{3}Z$/, "Z"); const r = byKey.get(key); out.push({ t: key, success: r?.success ?? 0, failed: r?.failed ?? 0, p50: r?.p50 ?? null, p95: r?.p95 ?? null }); } return out; } export interface NetworkShare { network: (typeof REQUEST_NETWORKS)[number]; requests: number; percent: number; } /** Distribution of resolved network classes. Always returns all four classes (0 % rows included). */ export async function getNetworkDistribution(scope: Scope, from: Date): Promise<{ shares: NetworkShare[]; resolved: number; unresolved: number }> { const db = getDb(); const rows = await db .select({ network: fetchRequests.network, requests: count() }) .from(fetchRequests) .where(and(scoped(scope), gte(fetchRequests.createdAt, from))) .groupBy(fetchRequests.network); const byNet = new Map(rows.map((r) => [r.network, r.requests])); const resolved = rows.filter((r) => r.network && (REQUEST_NETWORKS as readonly string[]).includes(r.network)).reduce((a, r) => a + r.requests, 0); const unresolved = rows.filter((r) => !r.network || !(REQUEST_NETWORKS as readonly string[]).includes(r.network)).reduce((a, r) => a + r.requests, 0); const shares = REQUEST_NETWORKS.map((network) => { const requests = byNet.get(network) ?? 0; return { network, requests, percent: resolved > 0 ? (requests / resolved) * 100 : 0 }; }); return { shares, resolved, unresolved }; } export async function getActiveSessionCount(scope: Scope): Promise { const db = getDb(); const [row] = await db .select({ n: count() }) .from(proxySessions) .where(and(eq(proxySessions.organizationId, scope.organizationId), eq(proxySessions.projectId, scope.projectId), eq(proxySessions.status, "active"), gt(proxySessions.expiresAt, new Date()))); return row?.n ?? 0; } export interface OnboardingState { hasApiKey: boolean; hasRequest: boolean; } export async function getOnboardingState(scope: Scope): Promise { const db = getDb(); const [key] = await db .select({ id: apiKeys.id }) .from(apiKeys) .where(and(eq(apiKeys.organizationId, scope.organizationId), eq(apiKeys.projectId, scope.projectId), isNull(apiKeys.revokedAt))) .limit(1); const [req] = await db.select({ id: fetchRequests.id }).from(fetchRequests).where(scoped(scope)).limit(1); return { hasApiKey: Boolean(key), hasRequest: Boolean(req) }; } /** True when the project has at least one request ever (used to tell "no data yet" from "no match"). */ export async function projectHasRequests(scope: Scope): Promise { const db = getDb(); const [req] = await db.select({ id: fetchRequests.id }).from(fetchRequests).where(scoped(scope)).limit(1); return Boolean(req); } // --------------------------------------------------------------------------- // Request log // --------------------------------------------------------------------------- const requestListColumns = { id: fetchRequests.id, createdAt: fetchRequests.createdAt, completedAt: fetchRequests.completedAt, url: fetchRequests.url, finalUrl: fetchRequests.finalUrl, domain: fetchRequests.domain, method: fetchRequests.method, source: fetchRequests.source, status: fetchRequests.status, httpStatus: fetchRequests.httpStatus, errorCode: fetchRequests.errorCode, requestedNetwork: fetchRequests.requestedNetwork, network: fetchRequests.network, country: fetchRequests.country, region: fetchRequests.region, city: fetchRequests.city, sessionId: fetchRequests.sessionId, latencyMs: fetchRequests.latencyMs, attempts: fetchRequests.attempts, bytesIn: fetchRequests.bytesIn, bytesOut: fetchRequests.bytesOut, priceUsd: fetchRequests.priceUsd, cached: fetchRequests.cached, }; export interface RequestListRow { id: string; createdAt: Date; completedAt: Date | null; url: string; finalUrl: string | null; domain: string; method: string; source: string; status: string; httpStatus: number | null; errorCode: string | null; requestedNetwork: string; network: string | null; country: string | null; region: string | null; city: string | null; sessionId: string | null; latencyMs: number | null; attempts: number; bytesIn: number; bytesOut: number; priceUsd: number; cached: boolean; } function escapeLike(s: string): string { return s.replace(/[\\%_]/g, (c) => `\\${c}`); } function requestFilterWhere(scope: Scope, f: RequestFilters): SQL { const { from, to } = filterWindow(f); const parts: Array = [scoped(scope), gte(fetchRequests.createdAt, from), to ? lt(fetchRequests.createdAt, to) : undefined]; if (f.status) parts.push(eq(fetchRequests.status, f.status)); if (f.domain) parts.push(ilike(fetchRequests.domain, `%${escapeLike(f.domain)}%`)); if (f.network) parts.push(eq(fetchRequests.network, f.network)); if (f.country) parts.push(eq(fetchRequests.country, f.country)); if (f.httpStatus) parts.push(eq(fetchRequests.httpStatus, f.httpStatus)); if (f.requestId) parts.push(eq(fetchRequests.id, f.requestId)); if (f.source) parts.push(eq(fetchRequests.source, f.source)); return and(...parts)!; } export async function listRequests(scope: Scope, f: RequestFilters, pageSize: number): Promise<{ rows: RequestListRow[]; total: number; page: number; pageCount: number }> { const db = getDb(); const where = requestFilterWhere(scope, f); const [{ total }] = await db.select({ total: count() }).from(fetchRequests).where(where); const pageCount = Math.max(1, Math.ceil(total / pageSize)); const page = Math.min(f.page, pageCount); const rows = await db .select(requestListColumns) .from(fetchRequests) .where(where) .orderBy(desc(fetchRequests.createdAt), desc(fetchRequests.id)) .limit(pageSize) .offset((page - 1) * pageSize); return { rows: rows as RequestListRow[], total, page, pageCount }; } /** Batched reader for CSV export (offset-based, deterministic ordering). */ export async function listRequestsBatch(scope: Scope, f: RequestFilters, offset: number, limit: number): Promise { const db = getDb(); const rows = await db .select(requestListColumns) .from(fetchRequests) .where(requestFilterWhere(scope, f)) .orderBy(desc(fetchRequests.createdAt), desc(fetchRequests.id)) .limit(limit) .offset(offset); return rows as RequestListRow[]; } export async function getRecentRequests(scope: Scope, limit = 8): Promise { const db = getDb(); const rows = await db.select(requestListColumns).from(fetchRequests).where(scoped(scope)).orderBy(desc(fetchRequests.createdAt)).limit(limit); return rows as RequestListRow[]; } // --------------------------------------------------------------------------- // Request detail // --------------------------------------------------------------------------- export interface RequestAttemptView { id: string; attemptNo: number; network: string; country: string | null; outcome: string; httpStatus: number | null; errorCode: string | null; blockReason: string | null; durationMs: number; bytesIn: number; bytesOut: number; routingScore: number | null; timing: Record | null; createdAt: Date; /** Only populated when the organization has provider visibility enabled. */ provider: string | null; } export interface RequestDetail extends RequestListRow { organizationId: string; projectId: string; projectName: string; apiKey: { id: string; name: string; prefix: string; last4: string; revoked: boolean } | null; browser: boolean; format: string; errorMessage: string | null; requestHeaders: Record | null; responseHeaders: Record | null; timing: Record | null; attemptRows: RequestAttemptView[]; } /** Loads one request for the organization (any project of the org), or null. */ export async function getRequestDetail(organizationId: string, id: string, opts: { providerVisibility: boolean }): Promise { const db = getDb(); const [row] = await db .select({ ...requestListColumns, organizationId: fetchRequests.organizationId, projectId: fetchRequests.projectId, projectName: projects.name, browser: fetchRequests.browser, format: fetchRequests.format, errorMessage: fetchRequests.errorMessage, requestHeaders: fetchRequests.requestHeaders, responseHeaders: fetchRequests.responseHeaders, timing: fetchRequests.timing, keyId: apiKeys.id, keyName: apiKeys.name, keyPrefix: apiKeys.keyPrefix, keyLast4: apiKeys.last4, keyRevokedAt: apiKeys.revokedAt, }) .from(fetchRequests) .innerJoin(projects, eq(projects.id, fetchRequests.projectId)) .leftJoin(apiKeys, eq(apiKeys.id, fetchRequests.apiKeyId)) .where(and(eq(fetchRequests.id, id), eq(fetchRequests.organizationId, organizationId))) .limit(1); if (!row) return null; const attempts = await db .select({ id: requestAttempts.id, attemptNo: requestAttempts.attemptNo, network: requestAttempts.network, country: requestAttempts.country, outcome: requestAttempts.outcome, httpStatus: requestAttempts.httpStatus, errorCode: requestAttempts.errorCode, blockReason: requestAttempts.blockReason, durationMs: requestAttempts.durationMs, bytesIn: requestAttempts.bytesIn, bytesOut: requestAttempts.bytesOut, routingScore: requestAttempts.routingScore, timing: requestAttempts.timing, createdAt: requestAttempts.createdAt, provider: opts.providerVisibility ? requestAttempts.provider : sql`null`, }) .from(requestAttempts) .where(eq(requestAttempts.requestId, id)) .orderBy(requestAttempts.attemptNo); const { keyId, keyName, keyPrefix, keyLast4, keyRevokedAt, ...rest } = row; return { ...(rest as Omit), apiKey: keyId && keyName && keyPrefix && keyLast4 ? { id: keyId, name: keyName, prefix: keyPrefix, last4: keyLast4, revoked: Boolean(keyRevokedAt) } : null, attemptRows: attempts.map((a) => ({ ...a, provider: a.provider ?? null })), }; } // --------------------------------------------------------------------------- // Sessions // --------------------------------------------------------------------------- export interface SessionRow { id: string; label: string | null; projectId: string; projectName: string; network: string; country: string | null; region: string | null; city: string | null; status: "active" | "expired" | "closed"; requestCount: number; lastUsedAt: Date | null; expiresAt: Date; createdAt: Date; } export async function listProjectSessions(scope: Scope, limit = 200): Promise { const db = getDb(); const rows = await db .select({ id: proxySessions.id, label: proxySessions.label, projectId: proxySessions.projectId, projectName: projects.name, network: proxySessions.network, country: proxySessions.country, region: proxySessions.region, city: proxySessions.city, status: proxySessions.status, requestCount: proxySessions.requestCount, lastUsedAt: proxySessions.lastUsedAt, expiresAt: proxySessions.expiresAt, createdAt: proxySessions.createdAt, }) .from(proxySessions) .innerJoin(projects, eq(projects.id, proxySessions.projectId)) .where(and(eq(proxySessions.organizationId, scope.organizationId), eq(proxySessions.projectId, scope.projectId))) .orderBy(desc(proxySessions.createdAt)) .limit(limit); const now = Date.now(); return rows.map((r) => ({ ...r, status: r.status === "closed" ? "closed" : r.status === "expired" || r.expiresAt.getTime() <= now ? "expired" : "active", })); } // --------------------------------------------------------------------------- // Usage // --------------------------------------------------------------------------- export interface UsageMonth { month: string; // YYYY-MM from: Date; to: Date; requests: number; successful: number; bytes: number; bandwidthBytes: number; residentialBytes: number; mobileBytes: number; browserSeconds: number; spendUsd: number; daily: Array<{ day: string; requests: number; bandwidthBytes: number; spendUsd: number }>; ledger: Array<{ id: string; metric: string; quantity: number; unit: string; priceUsd: number; requestId: string | null; createdAt: Date }>; } export function parseMonth(input: string | undefined, now = new Date()): { key: string; from: Date; to: Date } { let y = now.getUTCFullYear(); let m = now.getUTCMonth(); if (input && /^\d{4}-\d{2}$/.test(input)) { const [yy, mm] = input.split("-").map(Number); if (mm! >= 1 && mm! <= 12) { y = yy!; m = mm! - 1; } } const from = new Date(Date.UTC(y, m, 1)); const to = new Date(Date.UTC(y, m + 1, 1)); return { key: `${y}-${String(m + 1).padStart(2, "0")}`, from, to }; } export function lastMonths(n: number, now = new Date()): Array<{ key: string; label: string }> { const out: Array<{ key: string; label: string }> = []; for (let i = 0; i < n; i++) { const d = new Date(Date.UTC(now.getUTCFullYear(), now.getUTCMonth() - i, 1)); out.push({ key: `${d.getUTCFullYear()}-${String(d.getUTCMonth() + 1).padStart(2, "0")}`, label: new Intl.DateTimeFormat("en-US", { month: "short", year: "numeric", timeZone: "UTC" }).format(d) }); } return out; } export async function getUsageMonth(scope: Scope, monthInput?: string): Promise { const db = getDb(); const { key, from, to } = parseMonth(monthInput); const usageScope = and(eq(usageEvents.organizationId, scope.organizationId), eq(usageEvents.projectId, scope.projectId), gte(usageEvents.createdAt, from), lt(usageEvents.createdAt, to)); const [reqStats, [totals], daily, ledger] = await Promise.all([ getRequestStats(scope, from, to), db .select({ bandwidthBytes: sql`coalesce(sum(${usageEvents.quantity}) filter (where ${usageEvents.metric} = 'bandwidth'), 0)`.mapWith(num), residentialBytes: sql`coalesce(sum(${usageEvents.quantity}) filter (where ${usageEvents.metric} = 'residential_bandwidth'), 0)`.mapWith(num), mobileBytes: sql`coalesce(sum(${usageEvents.quantity}) filter (where ${usageEvents.metric} = 'mobile_bandwidth'), 0)`.mapWith(num), browserSeconds: sql`coalesce(sum(${usageEvents.quantity}) filter (where ${usageEvents.metric} = 'browser_seconds'), 0)`.mapWith(num), spendUsd: sql`coalesce(sum(${usageEvents.costUsd}), 0)`.mapWith(num), }) .from(usageEvents) .where(usageScope), db .select({ day: sql`to_char(date_trunc('day', ${usageEvents.createdAt} at time zone 'UTC'), 'YYYY-MM-DD')`, requests: sql`coalesce(sum(${usageEvents.quantity}) filter (where ${usageEvents.metric} = 'request'), 0)`.mapWith(num), // `bandwidth` is the total; `residential_bandwidth`/`mobile_bandwidth` are class breakdowns of the same bytes. bandwidthBytes: sql`coalesce(sum(${usageEvents.quantity}) filter (where ${usageEvents.metric} = 'bandwidth'), 0)`.mapWith(num), spendUsd: sql`coalesce(sum(${usageEvents.costUsd}), 0)`.mapWith(num), }) .from(usageEvents) .where(usageScope) .groupBy(sql`1`) .orderBy(sql`1 desc`), db .select({ id: usageEvents.id, metric: usageEvents.metric, quantity: usageEvents.quantity, unit: usageEvents.unit, priceUsd: usageEvents.costUsd, requestId: usageEvents.requestId, createdAt: usageEvents.createdAt, }) .from(usageEvents) .where(usageScope) .orderBy(desc(usageEvents.createdAt)) .limit(100), ]); return { month: key, from, to, requests: reqStats.total, successful: reqStats.successful, bytes: reqStats.bytes, bandwidthBytes: totals?.bandwidthBytes ?? 0, residentialBytes: totals?.residentialBytes ?? 0, mobileBytes: totals?.mobileBytes ?? 0, browserSeconds: totals?.browserSeconds ?? 0, // Prefer the immutable ledger; fall back to billed request prices when no usage events were recorded. spendUsd: totals && totals.spendUsd > 0 ? totals.spendUsd : reqStats.spendUsd, daily, ledger, }; } // --------------------------------------------------------------------------- // Analytics // --------------------------------------------------------------------------- export interface NetworkSuccess { network: string; requests: number; successful: number; successRate: number | null; } export async function getSuccessByNetwork(scope: Scope, from: Date): Promise { const db = getDb(); const rows = await db .select({ network: fetchRequests.network, requests: count(), successful: sql`count(*) filter (where ${fetchRequests.status} = 'success')`.mapWith(num), failed: sql`count(*) filter (where ${fetchRequests.status} = 'failed')`.mapWith(num), }) .from(fetchRequests) .where(and(scoped(scope), gte(fetchRequests.createdAt, from), isNotNull(fetchRequests.network))) .groupBy(fetchRequests.network); const byNet = new Map(rows.map((r) => [r.network, r])); return REQUEST_NETWORKS.map((network) => { const r = byNet.get(network); const decided = (r?.successful ?? 0) + (r?.failed ?? 0); return { network, requests: r?.requests ?? 0, successful: r?.successful ?? 0, successRate: decided > 0 ? ((r?.successful ?? 0) / decided) * 100 : null }; }); } export interface DomainStat { domain: string; requests: number; successful: number; successRate: number | null; avgLatencyMs: number | null; blocked: number; } export async function getTopDomains(scope: Scope, from: Date, limit = 10): Promise { const db = getDb(); const rows = await db .select({ domain: fetchRequests.domain, requests: count(), successful: sql`count(*) filter (where ${fetchRequests.status} = 'success')`.mapWith(num), failed: sql`count(*) filter (where ${fetchRequests.status} = 'failed')`.mapWith(num), avgLatencyMs: sql`avg(${fetchRequests.latencyMs})`.mapWith(numOrNull), blocked: sql`count(*) filter (where ${fetchRequests.errorCode} = 'TARGET_BLOCKED')`.mapWith(num), }) .from(fetchRequests) .where(and(scoped(scope), gte(fetchRequests.createdAt, from))) .groupBy(fetchRequests.domain) .orderBy(desc(count())) .limit(limit); return rows.map((r) => { const decided = r.successful + r.failed; return { domain: r.domain, requests: r.requests, successful: r.successful, successRate: decided > 0 ? (r.successful / decided) * 100 : null, avgLatencyMs: r.avgLatencyMs, blocked: r.blocked }; }); } export interface ErrorCodeStat { code: string; requests: number; percent: number; } export async function getErrorBreakdown(scope: Scope, from: Date, limit = 8): Promise { const db = getDb(); const rows = await db .select({ code: fetchRequests.errorCode, requests: count() }) .from(fetchRequests) .where(and(scoped(scope), gte(fetchRequests.createdAt, from), eq(fetchRequests.status, "failed"), isNotNull(fetchRequests.errorCode))) .groupBy(fetchRequests.errorCode) .orderBy(desc(count())) .limit(limit); const total = rows.reduce((a, r) => a + r.requests, 0); return rows.map((r) => ({ code: r.code ?? "UNKNOWN", requests: r.requests, percent: total > 0 ? (r.requests / total) * 100 : 0 })); } export interface CountryStat { country: string | null; requests: number; successful: number; successRate: number | null; avgLatencyMs: number | null; } export async function getCountries(scope: Scope, from: Date, limit = 10): Promise { const db = getDb(); const rows = await db .select({ country: fetchRequests.country, requests: count(), successful: sql`count(*) filter (where ${fetchRequests.status} = 'success')`.mapWith(num), failed: sql`count(*) filter (where ${fetchRequests.status} = 'failed')`.mapWith(num), avgLatencyMs: sql`avg(${fetchRequests.latencyMs})`.mapWith(numOrNull), }) .from(fetchRequests) .where(and(scoped(scope), gte(fetchRequests.createdAt, from))) .groupBy(fetchRequests.country) .orderBy(desc(count())) .limit(limit); return rows.map((r) => { const decided = r.successful + r.failed; return { country: r.country, requests: r.requests, successful: r.successful, successRate: decided > 0 ? (r.successful / decided) * 100 : null, avgLatencyMs: r.avgLatencyMs }; }); }