import type { Db } from "./db.ts"; import type { RpcEnvelope } from "./http_errors.ts"; import { enrichInfraError } from "./http_errors.ts"; const DEADLOCK = "40P01"; export class RpcCallError extends Error { constructor( readonly envelope: RpcEnvelope, ) { super(envelope.message); this.name = "RpcCallError"; } } function isEnvelope(value: unknown): value is RpcEnvelope { if (!value || typeof value !== "object") return false; const o = value as Record; return typeof o.ok === "boolean" && typeof o.code === "string" && typeof o.message === "string" && typeof o.layer === "string"; } function postgresCode(e: unknown): string | undefined { if (!e || typeof e !== "object") return undefined; return (e as { code?: string }).code; } function postgresMessage(e: unknown): string { if (e instanceof Error) return e.message; return String(e); } async function invokeOnce( db: Db, fn: string, payload: Record, ): Promise> { const row = await db.prepare(`SELECT ${fn}($1::jsonb) AS result`).get( payload, ) as { result: unknown } | undefined; const raw = row?.result; if (typeof raw === "string") { try { const parsed = JSON.parse(raw) as RpcEnvelope; if (isEnvelope(parsed)) return parsed; } catch { /* fall through */ } } if (isEnvelope(raw)) return raw as RpcEnvelope; throw new Error(`La función ${fn} no devolvió un envelope RPC válido`); } export async function callCoreFn( db: Db, fn: string, payload: Record = {}, opts: { route?: string; retries?: number } = {}, ): Promise> { const route = opts.route ?? fn; const retries = opts.retries ?? 1; let lastErr: unknown; for (let attempt = 0; attempt <= retries; attempt++) { try { const envelope = await invokeOnce(db, fn, payload); if (!envelope.ok) return envelope; return envelope; } catch (e) { lastErr = e; const code = postgresCode(e); if (code === DEADLOCK && attempt < retries) continue; const infra = enrichInfraError( fn, route, code ?? "exception", postgresMessage(e), ); return { ok: false, code: code === DEADLOCK ? "CONFLICT" : "INTERNAL", layer: "api", message: infra.message, context: { ...infra.context, sqlstate: code ?? null }, data: null, errors: null, }; } } const infra = enrichInfraError(fn, route, "exception", postgresMessage(lastErr)); return { ok: false, code: "INTERNAL", layer: "api", message: infra.message, context: infra.context, data: null, errors: null, }; } export async function callCoreFnData( db: Db, fn: string, payload: Record = {}, opts: { route?: string } = {}, ): Promise { const envelope = await callCoreFn(db, fn, payload, opts); if (!envelope.ok) throw new RpcCallError(envelope); return envelope.data as T; } export function unwrapRpc( envelope: RpcEnvelope, ): { data?: T; error?: string; code?: string; status?: number } { if (envelope.ok) return { data: envelope.data as T }; return { error: envelope.message, code: String(envelope.code), status: envelope.code === "NOT_FOUND" ? 404 : envelope.code === "CONFLICT" ? 409 : 400, }; }