panels-origin/api/rpc.ts
Cursor Agent 8717c00a83
Refactor payroll_http and excel to use core RPC functions
- Replace db.prepare() in payroll_http with callCoreFn for attendance,
  loans, loan receipts, and legacy payroll period routes
- Use respondRpc/respondApiError for standardized HTTP responses
- Add fn_loan_list/get, fn_loan_payment_get, and fn_payroll_period_* RPCs
  to 013-rpc-payroll.sql with grants
- Refactor excel storeDocument to fn_worker_document_store (S3 stays in TS)
- Refactor importExcel and upsertWorker to fn_workers_import_batch

Co-authored-by: alberto.martinez <alberto.martinez@mrdev.mx>
2026-09-03 23:57:56 +00:00

125 lines
3.4 KiB
TypeScript

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<string, unknown>;
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<T>(
db: Db,
fn: string,
payload: Record<string, unknown>,
): Promise<RpcEnvelope<T>> {
const row = await db.prepare(`SELECT ${fn}($1::jsonb) AS result`).get(
JSON.stringify(payload),
) as { result: unknown } | undefined;
const raw = row?.result;
if (typeof raw === "string") {
try {
const parsed = JSON.parse(raw) as RpcEnvelope<T>;
if (isEnvelope(parsed)) return parsed;
} catch {
/* fall through */
}
}
if (isEnvelope(raw)) return raw as RpcEnvelope<T>;
throw new Error(`La función ${fn} no devolvió un envelope RPC válido`);
}
export async function callCoreFn<T = unknown>(
db: Db,
fn: string,
payload: Record<string, unknown> = {},
opts: { route?: string; retries?: number } = {},
): Promise<RpcEnvelope<T>> {
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<T>(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<T>(
db: Db,
fn: string,
payload: Record<string, unknown> = {},
opts: { route?: string } = {},
): Promise<T> {
const envelope = await callCoreFn<T>(db, fn, payload, opts);
if (!envelope.ok) throw new RpcCallError(envelope);
return envelope.data as T;
}
export function unwrapRpc<T>(
envelope: RpcEnvelope<T>,
): { 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,
};
}