import { newId, redactHeaders, type ConcreteNetwork, type FetchRequest, type Plan } from "@fetcha/core"; import { db, domainProfiles, fetchRequests, requestAttempts, routingMetrics, usageEvents, sql } from "@fetcha/db"; import { foldRouteStat, preferredRoute, routeKey, type AttemptRecord, type ExecutionResult } from "@fetcha/routing"; import { priceRequest } from "./pricing"; export interface RequestContext { requestId: string; organizationId: string; projectId: string; apiKeyId: string | null; source: "api" | "playground" | "sdk" | "crawl"; plan: Plan; clientIp: string | null; userAgent: string | null; logLevel: "none" | "metadata" | "headers" | "full"; } export async function createRequestRow(ctx: RequestContext, req: FetchRequest, domain: string): Promise { await db.insert(fetchRequests).values({ id: ctx.requestId, organizationId: ctx.organizationId, projectId: ctx.projectId, apiKeyId: ctx.apiKeyId, source: ctx.source, url: req.url, domain, method: req.method, requestedNetwork: req.network, country: req.country ?? null, region: req.region ?? null, city: req.city ?? null, sessionId: req.session ?? null, browser: req.browser, format: req.format, status: "pending", requestHeaders: ctx.logLevel === "headers" || ctx.logLevel === "full" ? redactHeaders(req.headers) : null, clientIp: ctx.clientIp, userAgent: ctx.userAgent, }); } export async function recordAttempt(requestId: string, a: AttemptRecord): Promise { await db.insert(requestAttempts).values({ id: a.attemptId, requestId, attemptNo: a.attemptNo, provider: a.provider, network: a.network, mode: a.mode, country: a.country, sessionKey: a.sessionKey, outcome: a.outcome, httpStatus: a.httpStatus, errorCode: a.errorCode, errorDetail: a.errorDetail ?? (a.blockVendor ? `vendor=${a.blockVendor} profile=${a.profileId ?? "-"}` : a.profileId ? `profile=${a.profileId}` : null), blockReason: a.blockReason, durationMs: a.durationMs, bytesIn: a.bytesIn, bytesOut: a.bytesOut, unitPricePerGb: a.unitPricePerGb, costUsd: a.costUsd, routingScore: a.routingScore, timing: a.timing ? { ...a.timing } : null, }); } export interface CompletionSummary { priceUsd: number; costUsd: number; } /** Persist a successful (or blocked-but-answered) execution: request row, usage ledger, domain intelligence, routing metrics. */ export async function completeRequest(ctx: RequestContext, req: FetchRequest, result: ExecutionResult, latencyMs: number): Promise { const network = result.network; const price = priceRequest({ plan: ctx.plan, network, bytes: result.bytesIn + result.bytesOut, upstreamCostUsd: result.costUsd, attempts: result.attempts.length, success: result.body.success }); const finalStatus = result.body.success ? "success" : "failed"; await db .update(fetchRequests) .set({ status: finalStatus, httpStatus: result.body.status, errorCode: result.body.success ? null : "TARGET_BLOCKED", errorMessage: result.body.success ? null : "The target blocked every route we tried.", finalUrl: result.finalUrl, network, mode: result.mode, browser: result.mode === "browser", attempts: result.attempts.length, latencyMs, bytesIn: result.bytesIn, bytesOut: result.bytesOut, costUsd: result.costUsd, priceUsd: price.totalUsd, responseHeaders: ctx.logLevel === "headers" || ctx.logLevel === "full" ? redactHeaders(result.body.headers) : null, timing: result.body.metadata.timing ? { ...result.body.metadata.timing } : null, completedAt: new Date(), }) .where(sql`${fetchRequests.id} = ${ctx.requestId}`); await writeUsage(ctx, result, price.totalUsd, network); await Promise.all([updateDomainProfile(result.domain, result.attempts, result.body.success, result.browserRequired), updateRoutingMetrics(result.attempts)]); return { priceUsd: price.totalUsd, costUsd: result.costUsd }; } export async function failRequest(ctx: RequestContext, code: string, message: string, attempts: AttemptRecord[], latencyMs: number, domain: string): Promise { const costUsd = attempts.reduce((s, a) => s + a.costUsd, 0); await db .update(fetchRequests) .set({ status: "failed", errorCode: code, errorMessage: message, attempts: attempts.length, latencyMs, costUsd, completedAt: new Date() }) .where(sql`${fetchRequests.id} = ${ctx.requestId}`); // Failed requests still count against the request quota (abuse control) but are not priced. await db.insert(usageEvents).values({ id: newId("usage"), organizationId: ctx.organizationId, projectId: ctx.projectId, requestId: ctx.requestId, metric: "request", quantity: 1, unit: "count", costUsd: 0, upstreamCostUsd: costUsd, }); if (attempts.length) await Promise.all([updateDomainProfile(domain, attempts, false, false), updateRoutingMetrics(attempts)]); } async function writeUsage(ctx: RequestContext, result: ExecutionResult, priceUsd: number, network: ConcreteNetwork | null) { const bytes = result.bytesIn + result.bytesOut; const rows = [ { metric: "request", quantity: 1, unit: "count", costUsd: priceUsd, upstreamCostUsd: result.costUsd }, { metric: "bandwidth", quantity: bytes, unit: "bytes", costUsd: 0, upstreamCostUsd: 0 }, ]; if (network === "residential" || network === "isp") rows.push({ metric: "residential_bandwidth", quantity: bytes, unit: "bytes", costUsd: 0, upstreamCostUsd: 0 }); if (network === "mobile") rows.push({ metric: "mobile_bandwidth", quantity: bytes, unit: "bytes", costUsd: 0, upstreamCostUsd: 0 }); await db.insert(usageEvents).values( rows.map((r) => ({ id: newId("usage"), organizationId: ctx.organizationId, projectId: ctx.projectId, requestId: ctx.requestId, ...r })), ); } async function updateDomainProfile(domain: string, attempts: AttemptRecord[], success: boolean, browserRequired: boolean) { if (!domain) return; const [existing] = await db.select().from(domainProfiles).where(sql`${domainProfiles.domain} = ${domain}`).limit(1); const stats = { ...(existing?.routeStats ?? {}) }; let blocks = 0; let captchas = 0; for (const a of attempts) { foldRouteStat(stats, routeKey(a.provider, a.network), { ok: a.outcome === "success", blocked: a.outcome === "blocked", latencyMs: a.durationMs, costUsd: a.costUsd }); if (a.outcome === "blocked") blocks++; if (a.blockReason === "captcha" || a.blockReason === "cloudflare_challenge") captchas++; } const browserInc = browserRequired ? 1 : 0; const pref = preferredRoute(stats); const totalLatency = attempts.reduce((s, a) => s + a.durationMs, 0); const n = (existing?.requests ?? 0) + 1; const avgLatency = existing ? existing.avgLatencyMs + (totalLatency - existing.avgLatencyMs) / n : totalLatency; await db .insert(domainProfiles) .values({ domain, preferredNetwork: pref?.network ?? null, preferredProvider: pref?.provider ?? null, requests: 1, successes: success ? 1 : 0, blocks, captchas, browserRequired: browserInc, avgLatencyMs: totalLatency, routeStats: stats, lastSeenAt: new Date(), }) .onConflictDoUpdate({ target: domainProfiles.domain, set: { preferredNetwork: pref?.network ?? null, preferredProvider: pref?.provider ?? null, requests: sql`${domainProfiles.requests} + 1`, successes: sql`${domainProfiles.successes} + ${success ? 1 : 0}`, blocks: sql`${domainProfiles.blocks} + ${blocks}`, captchas: sql`${domainProfiles.captchas} + ${captchas}`, browserRequired: sql`${domainProfiles.browserRequired} + ${browserInc}`, avgLatencyMs: avgLatency, routeStats: stats, lastSeenAt: new Date(), updatedAt: new Date(), }, }); } async function updateRoutingMetrics(attempts: AttemptRecord[]) { const bucket = new Date(); bucket.setMinutes(0, 0, 0); for (const a of attempts) { const ok = a.outcome === "success" ? 1 : 0; const blocked = a.outcome === "blocked" ? 1 : 0; const errors = ok || blocked ? 0 : 1; await db .insert(routingMetrics) .values({ id: newId("evt"), bucket, provider: a.provider, network: a.network, country: a.country ?? "", requests: 1, successes: ok, blocked, errors, latencySumMs: a.durationMs, bytes: a.bytesIn + a.bytesOut, costUsd: a.costUsd, }) .onConflictDoUpdate({ target: [routingMetrics.bucket, routingMetrics.provider, routingMetrics.network, routingMetrics.country], set: { requests: sql`${routingMetrics.requests} + 1`, successes: sql`${routingMetrics.successes} + ${ok}`, blocked: sql`${routingMetrics.blocked} + ${blocked}`, errors: sql`${routingMetrics.errors} + ${errors}`, latencySumMs: sql`${routingMetrics.latencySumMs} + ${a.durationMs}`, bytes: sql`${routingMetrics.bytes} + ${a.bytesIn + a.bytesOut}`, costUsd: sql`${routingMetrics.costUsd} + ${a.costUsd}`, }, }); } }