From 8717c00a83b6dc01d13d057cbf7d37452c462dca Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 3 Sep 2026 23:57:36 +0000 Subject: [PATCH 1/2] 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 --- api/excel.ts | 273 ++++++++--------- api/http_errors.ts | 130 ++++++++ api/payroll_http.ts | 166 ++++++----- api/rpc.ts | 125 ++++++++ db/core/changesets/006-rpc-infra.sql | 147 ++++++++++ db/core/changesets/007-rpc-catalogs.sql | 123 ++++++++ db/core/changesets/013-rpc-payroll.sql | 374 ++++++++++++++++++++++++ 7 files changed, 1106 insertions(+), 232 deletions(-) create mode 100644 api/http_errors.ts create mode 100644 api/rpc.ts create mode 100644 db/core/changesets/006-rpc-infra.sql create mode 100644 db/core/changesets/007-rpc-catalogs.sql diff --git a/api/excel.ts b/api/excel.ts index a2639ed..e31b617 100644 --- a/api/excel.ts +++ b/api/excel.ts @@ -3,9 +3,10 @@ import type { Db } from "./db.ts"; import { encryptBytes } from "./docs_crypto.ts"; import { sha256Hex } from "./crypto.ts"; import { canonicalRiskCode, normalizeWorker, validateWorkerFields, formatNss, normUpper, type WorkerInput } from "./mx.ts"; -import { refreshPipeline, lastInsertId } from "./db.ts"; +import { refreshPipeline } from "./db.ts"; import { resolveCompany } from "./companies.ts"; import { companyDocKey, projectDocKey, putObject, workerDocKey } from "./storage.ts"; +import { callCoreFn, RpcCallError } from "./rpc.ts"; export const IMPORT_COLUMNS = [ "NOMBRE", @@ -214,82 +215,34 @@ export async function storeDocument( imss_baja_at?: string | null; } = {}, ) { - const allowed = await db.prepare("SELECT code FROM document_types WHERE code = ?").get(type); - if (!allowed) throw new Error("Tipo de documento no válido"); - const { iv, cipher } = await encryptBytes(bytes); const hash = await sha256Hex(bytes); const storage = `${crypto.randomUUID()}.enc`; await putObject(workerDocKey(workerId, storage), cipher); - await db.prepare("UPDATE documents SET is_current = false WHERE worker_id = ? AND type_code = ?").run( - workerId, - type, - ); const issuedAt = meta.issued_at || null; const expiresAt = meta.expires_at || null; const imssCompanyId = meta.imss_company_id || null; const today = new Date().toISOString().slice(0, 10); const imssAltaAt = meta.imss_alta_at || (type === "alta_imss" ? today : null); const imssBajaAt = meta.imss_baja_at || (type === "baja_imss" ? today : null); - const movementDate = type === "baja_imss" ? imssBajaAt : imssAltaAt; - await db.prepare( - `INSERT INTO documents - (worker_id, type_code, original_name, mime, size_bytes, sha256, iv, storage_name, is_current, - parse_status, issued_at, expires_at, imss_company_id, imss_alta_at, uploaded_by_id, uploaded_by_name) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, true, 'manual', ?, ?, ?, ?, ?, '')`, - ).run( - workerId, - type, - filename, + const env = await callCoreFn(db, "core.fn_worker_document_store", { + worker_id: workerId, + type_code: type, + original_name: filename, mime, - bytes.byteLength, - hash, + size_bytes: bytes.byteLength, + sha256: hash, iv, - storage, - issuedAt || movementDate, - expiresAt, - imssCompanyId, - movementDate, - userId, - ); - - if (type === "alta_imss") { - const companyId = imssCompanyId || - (await db.prepare("SELECT company_id FROM workers WHERE id = ?").get(workerId) as { company_id: number } | undefined) - ?.company_id || - null; - const company = companyId - ? await db.prepare("SELECT id, code FROM companies WHERE id = ?").get(companyId) as - | { id: number; code: string } - | undefined - : undefined; - if (!company) throw new Error("Empresa patrón no válida"); - await db.prepare( - `UPDATE workers SET - imss_status = 'alta', - imss_company_id = ?, - imss_alta_at = ?, - imss_baja_at = NULL, - company_id = ?, - hire_type = ?, - updated_at = now() - WHERE id = ?`, - ).run(company.id, imssAltaAt, company.id, company.code, workerId); - } - - if (type === "baja_imss") { - await db.prepare( - `UPDATE workers SET - imss_status = 'baja_imss', - imss_baja_at = ?, - imss_company_id = NULL, - company_id = NULL, - hire_type = '', - updated_at = now() - WHERE id = ?`, - ).run(imssBajaAt, workerId); - } + storage_name: storage, + uploaded_by_id: userId, + issued_at: issuedAt, + expires_at: expiresAt, + imss_company_id: imssCompanyId, + imss_alta_at: imssAltaAt, + imss_baja_at: imssBajaAt, + }); + if (!env.ok) throw new RpcCallError(env); await refreshPipeline(db, workerId); } @@ -357,79 +310,78 @@ async function fetchPhoto(url: string): Promise { } } +export type ImportReport = { + inserted: number; + existed: number; + errors: number; + photos: number; + existed_rows: { sheet: string; row: number; nombre: string; curp: string; matched: string }[]; + error_rows: { sheet: string; row: number; nombre: string; curp: string; messages: string[] }[]; +}; + +type ImportBatchResult = { + inserted: number; + existed: number; + errors: number; + existed_rows: ImportReport["existed_rows"]; + error_rows: ImportReport["error_rows"]; +}; + +function workerToImportRow( + input: ReturnType & { company_id: number; tenant_id?: number | null }, + sheet: string, + row: number, + projectId: number | null, +) { + return { + sheet, + row, + first_name: input.first_name, + middle_name: input.middle_name, + last_name_p: input.last_name_p, + last_name_m: input.last_name_m, + curp: input.curp, + rfc: input.rfc, + nss: input.nss, + phone: input.phone, + email: input.email, + address: input.address, + blood_type: input.blood_type, + hire_type: input.hire_type, + company_id: input.company_id, + position: input.position, + risk_code: input.risk_code, + work_type: input.work_type, + daily_wage: input.daily_wage, + needs_badge: input.needs_badge, + status: input.status, + project_id: input.status === "activo" ? projectId : null, + tenant_id: input.tenant_id ?? 1, + }; +} + async function upsertWorker( db: Db, input: ReturnType & { company_id: number; tenant_id?: number | null }, projectId: number | null, ) { + const env = await callCoreFn(db, "core.fn_workers_import_batch", { + tenant_id: input.tenant_id ?? 1, + project_id: projectId, + rows: [workerToImportRow(input, "IMPORT", 1, projectId)], + }); + if (!env.ok) throw new RpcCallError(env); + const data = env.data!; + const errorRow = data.error_rows?.[0]; + if (errorRow?.messages?.length) throw new Error(errorRow.messages[0]); const existing = await findExisting(db, input.curp, input.rfc, input.nss); - if (existing) { - await db.prepare( - `UPDATE workers SET - first_name=?, middle_name=?, last_name_p=?, last_name_m=?, - curp=?, rfc=?, nss=?, phone=?, email=?, address=?, blood_type=?, - hire_type=?, company_id=?, tenant_id=COALESCE(?, tenant_id), position=?, risk_code=?, work_type=?, daily_wage=?, - needs_badge=?, status=?, updated_at=now() - WHERE id=?`, - ).run( - input.first_name, - input.middle_name, - input.last_name_p, - input.last_name_m, - input.curp, - input.rfc, - input.nss, - input.phone, - input.email, - input.address, - input.blood_type, - input.hire_type, - input.company_id, - input.tenant_id ?? null, - input.position, - input.risk_code, - input.work_type, - input.daily_wage, - input.needs_badge, - input.status, - existing.id, - ); - if (projectId) await assign(db, existing.id, projectId); - await refreshPipeline(db, existing.id, input.tenant_id ?? null); - const matched = existing.curp === input.curp ? "CURP" : existing.rfc === input.rfc ? "RFC" : "NSS"; - return { id: existing.id, action: "existed" as const, matched }; - } - await db.prepare( - `INSERT INTO workers - (first_name, middle_name, last_name_p, last_name_m, curp, rfc, nss, phone, email, address, - blood_type, hire_type, company_id, position, risk_code, work_type, daily_wage, needs_badge, status, tenant_id) - VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, - ).run( - input.first_name, - input.middle_name, - input.last_name_p, - input.last_name_m, - input.curp, - input.rfc, - input.nss, - input.phone, - input.email, - input.address, - input.blood_type, - input.hire_type, - input.company_id, - input.position, - input.risk_code, - input.work_type, - input.daily_wage, - input.needs_badge, - input.status, - input.tenant_id ?? null, - ); - const id = await lastInsertId(db); - if (projectId) await assign(db, id, projectId); - await refreshPipeline(db, id, input.tenant_id ?? null); - return { id, action: "inserted" as const, matched: null as string | null }; + if (!existing) throw new Error("No se pudo guardar el trabajador"); + await refreshPipeline(db, existing.id, input.tenant_id ?? null); + return { + id: existing.id, + action: (data.inserted ?? 0) > 0 ? "inserted" as const : "existed" as const, + matched: data.existed_rows?.[0]?.matched ?? null, + }; } async function assign(db: Db, workerId: number, projectId: number) { @@ -443,15 +395,6 @@ async function assign(db: Db, workerId: number, projectId: number) { ).run(workerId, projectId); } -export type ImportReport = { - inserted: number; - existed: number; - errors: number; - photos: number; - existed_rows: { sheet: string; row: number; nombre: string; curp: string; matched: string }[]; - error_rows: { sheet: string; row: number; nombre: string; curp: string; messages: string[] }[]; -}; - export async function importExcel( db: Db, bytes: Uint8Array, @@ -496,6 +439,9 @@ export async function importExcel( return report; } + const batchRows: ReturnType[] = []; + const photoJobs: { curp: string; rfc: string; nss: string; url: string; tenant_id: number }[] = []; + for (const job of jobs) { for (const { row, data } of job.rows) { if (isEmptyRow(data)) continue; @@ -539,28 +485,41 @@ export async function importExcel( company_id: company!.id, tenant_id: company!.tenant_id ?? 1, }; - const res = await upsertWorker(db, n, n.status === "activo" ? projectId : null); - if (res.action === "inserted") report.inserted++; - else { - report.existed++; - report.existed_rows.push({ - sheet: job.sheet, - row, - nombre, - curp: n.curp, - matched: res.matched || "CURP", - }); - } + batchRows.push(workerToImportRow(n, job.sheet, row, n.status === "activo" ? projectId : null)); const url = data["URL FOTO"]; if (url) { - const photo = await fetchPhoto(url); - if (photo) { - await storeDocument(db, res.id, "foto", "foto-import.jpg", "image/jpeg", photo, userId); - report.photos++; - } + photoJobs.push({ curp: n.curp, rfc: n.rfc, nss: n.nss, url, tenant_id: n.tenant_id }); } } } + + if (batchRows.length > 0) { + const env = await callCoreFn(db, "core.fn_workers_import_batch", { + project_id: projectId, + rows: batchRows, + }); + if (!env.ok) throw new RpcCallError(env); + const data = env.data!; + report.inserted = data.inserted ?? 0; + report.existed = data.existed ?? 0; + report.errors += data.errors ?? 0; + report.existed_rows.push(...(data.existed_rows ?? [])); + report.error_rows.push(...(data.error_rows ?? [])); + for (const row of batchRows) { + const worker = await findExisting(db, row.curp, row.rfc, row.nss); + if (worker) await refreshPipeline(db, worker.id, row.tenant_id ?? null); + } + } + + for (const photo of photoJobs) { + const worker = await findExisting(db, photo.curp, photo.rfc, photo.nss); + if (!worker) continue; + const photoBytes = await fetchPhoto(photo.url); + if (photoBytes) { + await storeDocument(db, worker.id, "foto", "foto-import.jpg", "image/jpeg", photoBytes, userId); + report.photos++; + } + } return report; } diff --git a/api/http_errors.ts b/api/http_errors.ts new file mode 100644 index 0000000..183977a --- /dev/null +++ b/api/http_errors.ts @@ -0,0 +1,130 @@ +import type { Context } from "hono"; + +export type RpcCode = + | "OK" + | "CREATED" + | "VALIDATION" + | "UNAUTHORIZED" + | "FORBIDDEN" + | "NOT_FOUND" + | "CONFLICT" + | "INTERNAL" + | "NETWORK"; + +export type ResponseLayer = "db" | "api" | "front"; + +export type RpcEnvelope = { + ok: boolean; + code: RpcCode | string; + layer: ResponseLayer; + message: string; + context?: Record | null; + data?: T | null; + errors?: Record | string[] | null; +}; + +export type ApiResponse = RpcEnvelope & { status: number }; + +const GENERIC_MESSAGES = new Set([ + "ok", "error", "operación exitosa", "operacion exitosa", "algo salió mal", + "algo salio mal", "error interno", "error interno del servidor", +]); + +export function isGenericMessage(message: string): boolean { + return GENERIC_MESSAGES.has(message.trim().toLowerCase()); +} + +export function mapRpcToStatus(code: string): number { + switch (code) { + case "OK": + return 200; + case "CREATED": + return 201; + case "VALIDATION": + return 400; + case "UNAUTHORIZED": + return 401; + case "FORBIDDEN": + return 403; + case "NOT_FOUND": + return 404; + case "CONFLICT": + return 409; + case "NETWORK": + return 503; + case "INTERNAL": + default: + return 500; + } +} + +export function respondRpc(c: Context, envelope: RpcEnvelope) { + const status = mapRpcToStatus(String(envelope.code)); + const body: ApiResponse = { ...envelope, status }; + return c.json(body, status as 200 | 201 | 400 | 401 | 403 | 404 | 409 | 500 | 503); +} + +export function respondApiError( + c: Context, + code: RpcCode, + message: string, + context?: Record, + errors?: Record | string[] | null, +) { + if (isGenericMessage(message)) { + throw new Error(`API error message is too generic: ${message}`); + } + const status = mapRpcToStatus(code); + const body: ApiResponse = { + ok: false, + code, + status, + layer: "api", + message, + context: context ?? {}, + data: null, + errors: errors ?? null, + }; + return c.json(body, status as 400 | 401 | 403 | 404 | 409 | 500 | 503); +} + +export function enrichInfraError( + fn: string, + route: string, + cause: string, + detail?: string, +): { message: string; context: Record } { + const message = detail + ? `Falló la llamada a ${fn} en ${route}: ${cause} — ${detail}` + : `Falló la llamada a ${fn} en ${route}: ${cause}`; + return { + message, + context: { fn, route, cause }, + }; +} + +export function routeLabel(c: Context): string { + return `${c.req.method} ${c.req.path}`; +} + +export function onAppError(err: unknown, c: Context) { + console.error(err); + if (err && typeof err === "object" && (err as RpcEnvelope).layer) { + const env = err as RpcEnvelope; + return respondRpc(c, env); + } + const message = err instanceof Error + ? `Error no controlado en ${routeLabel(c)}: ${err.message}` + : `Error no controlado en ${routeLabel(c)}`; + const status = 500; + return c.json({ + ok: false, + code: "INTERNAL", + status, + layer: "api", + message, + context: { route: routeLabel(c) }, + data: null, + errors: null, + }, status); +} diff --git a/api/payroll_http.ts b/api/payroll_http.ts index df72e83..5280453 100644 --- a/api/payroll_http.ts +++ b/api/payroll_http.ts @@ -2,9 +2,11 @@ import type { Context, Hono } from "hono"; import type { AuthUser } from "./auth.ts"; import { tenantScope } from "./auth.ts"; import { requireCoreAuth } from "./scope.ts"; -import { lastInsertId, projectById, projectMustBe, type Db } from "./db.ts"; +import { projectById, projectMustBe, type Db } from "./db.ts"; import { generateLoanReceiptPdf } from "./pdf.ts"; import { fullName } from "./mx.ts"; +import { callCoreFn } from "./rpc.ts"; +import { respondApiError, respondRpc, routeLabel } from "./http_errors.ts"; import { addAdminLine, addJornalWorker, @@ -47,13 +49,13 @@ export function registerPayrollRoutes(app: App) { const from = c.req.query("from") ?? ""; const to = c.req.query("to") ?? ""; const db = c.get("db"); - const rows = await db.prepare( - `SELECT a.*, w.first_name, w.last_name_p FROM attendance a - JOIN workers w ON w.id=a.worker_id - WHERE a.project_id=? AND a.work_date BETWEEN ? AND ? - ORDER BY a.work_date, w.last_name_p`, - ).all(projectId, from, to); - return c.json({ attendance: rows }); + const env = await callCoreFn<{ attendance: unknown[] }>( + db, + "core.fn_attendance_list", + { project_id: projectId, from, to }, + { route: routeLabel(c) }, + ); + return respondRpc(c, env); }); app.put("/v1/attendance", ...requireCoreAuth, async (c) => { @@ -69,7 +71,14 @@ export function registerPayrollRoutes(app: App) { ["activo", "pausado"], "No se registra asistencia en un proyecto concluido o cancelado", ); - if (blocked) return c.json({ error: blocked.error }, blocked.status); + if (blocked) { + return respondApiError( + c, + blocked.status === 404 ? "NOT_FOUND" : "VALIDATION", + blocked.error, + { route: routeLabel(c), project_id: body.project_id }, + ); + } const result = await setAttendance(db, tid(c), body); if ("error" in result) { return c.json({ error: result.error, other: result.other }, result.status === 409 ? 409 : 400); @@ -114,11 +123,13 @@ export function registerPayrollRoutes(app: App) { app.get("/v1/loans", ...requireCoreAuth, async (c) => { const db = c.get("db"); const workerId = c.req.query("worker_id"); - const sql = workerId - ? "SELECT l.*, w.first_name, w.last_name_p FROM loans l JOIN workers w ON w.id=l.worker_id WHERE l.worker_id=? ORDER BY l.id DESC" - : "SELECT l.*, w.first_name, w.last_name_p FROM loans l JOIN workers w ON w.id=l.worker_id ORDER BY l.id DESC"; - const loans = workerId ? await db.prepare(sql).all(Number(workerId)) : await db.prepare(sql).all(); - return c.json({ loans }); + const env = await callCoreFn<{ loans: unknown[] }>( + db, + "core.fn_loan_list", + workerId ? { worker_id: Number(workerId) } : {}, + { route: routeLabel(c) }, + ); + return respondRpc(c, env); }); app.post("/v1/loans", ...requireCoreAuth, async (c) => { @@ -154,10 +165,7 @@ export function registerPayrollRoutes(app: App) { app.get("/v1/loans/:id/recibo", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); - const loan = await db.prepare( - `SELECT l.*, w.first_name, w.middle_name, w.last_name_p, w.last_name_m - FROM loans l JOIN workers w ON w.id=l.worker_id WHERE l.id=?`, - ).get(id) as { + const env = await callCoreFn<{ loan: { delivered: number; commission_pct: number; commission_amount: number; @@ -172,8 +180,14 @@ export function registerPayrollRoutes(app: App) { last_name_p: string; last_name_m: string; note: string | null; - } | undefined; - if (!loan) return c.json({ error: "Préstamo no encontrado" }, 404); + } }>( + db, + "core.fn_loan_get", + { id }, + { route: routeLabel(c) }, + ); + if (!env.ok) return respondRpc(c, env); + const loan = env.data!.loan; const bytes = await generateLoanReceiptPdf({ title: "Recibo de aceptación de préstamo", subtitle: `Registrado ${String(loan.created_at).slice(0, 10)}`, @@ -202,14 +216,7 @@ export function registerPayrollRoutes(app: App) { app.get("/v1/loan-payments/:id/recibo", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); - const pay = await db.prepare( - `SELECT p.*, l.delivered, l.commission_pct, w.first_name, w.middle_name, w.last_name_p, w.last_name_m, wk.week_end - FROM loan_payments p - JOIN loans l ON l.id=p.loan_id - JOIN workers w ON w.id=l.worker_id - JOIN payroll_weeks wk ON wk.id=p.week_id - WHERE p.id=?`, - ).get(id) as { + const env = await callCoreFn<{ payment: { amount: number; installment_n: number; label: string; @@ -220,8 +227,14 @@ export function registerPayrollRoutes(app: App) { last_name_p: string; last_name_m: string; week_end: string; - } | undefined; - if (!pay) return c.json({ error: "Recibo no encontrado" }, 404); + } }>( + db, + "core.fn_loan_payment_get", + { id }, + { route: routeLabel(c) }, + ); + if (!env.ok) return respondRpc(c, env); + const pay = env.data!.payment; const bytes = await generateLoanReceiptPdf({ title: "Recibo de pago de préstamo", subtitle: `Descuento del sábado ${pay.week_end}`, @@ -380,65 +393,63 @@ export function registerPayrollRoutes(app: App) { ["activo", "pausado"], "No se genera nómina de un proyecto concluido o cancelado", ); - if (blocked) return c.json({ error: blocked.error }, blocked.status); - await db.prepare( - "INSERT INTO payroll_periods (project_id, week_start, week_end, status) VALUES (?, ?, ?, 'draft')", - ).run(body.project_id, body.week_start, body.week_end); - const periodId = await lastInsertId(db); - const workers = await db.prepare( - `SELECT w.id, w.daily_wage FROM workers w - JOIN assignments a ON a.worker_id=w.id AND a.project_id=? AND a.active=true - WHERE w.status='activo'`, - ).all(body.project_id) as { id: number; daily_wage: number }[]; - const ins = db.prepare( - `INSERT INTO payroll_lines (period_id, worker_id, days, daily_wage, gross, discounts, loan_payment, net) - VALUES (?, ?, ?, ?, ?, ?, ?, ?)`, - ); - for (const w of workers) { - const att = await db.prepare( - `SELECT COUNT(*) AS n FROM attendance - WHERE worker_id=? AND project_id=? AND work_date BETWEEN ? AND ? AND present=true`, - ).get(w.id, body.project_id, body.week_start, body.week_end) as { n: number }; - const days = Number(att.n); - const gross = days * w.daily_wage; - const extra = Number(body.extra_discounts?.[w.id] ?? 0); - const net = gross - extra; - await ins.run(periodId, w.id, days, w.daily_wage, gross, extra, 0, net); + if (blocked) { + return respondApiError( + c, + blocked.status === 404 ? "NOT_FOUND" : "VALIDATION", + blocked.error, + { route: routeLabel(c), project_id: body.project_id }, + ); } - return c.json({ id: periodId }); + const env = await callCoreFn<{ id: number }>( + db, + "core.fn_payroll_period_create", + { + project_id: body.project_id, + week_start: body.week_start, + week_end: body.week_end, + extra_discounts: body.extra_discounts ?? {}, + }, + { route: routeLabel(c) }, + ); + return respondRpc(c, env); }); app.get("/v1/payroll/periods", ...requireCoreAuth, async (c) => { const db = c.get("db"); const projectId = c.req.query("project_id"); - const sql = projectId - ? `SELECT pe.*, p.name AS project_name FROM payroll_periods pe JOIN projects p ON p.id=pe.project_id WHERE pe.project_id=? ORDER BY pe.id DESC` - : `SELECT pe.*, p.name AS project_name FROM payroll_periods pe JOIN projects p ON p.id=pe.project_id ORDER BY pe.id DESC`; - const periods = projectId ? await db.prepare(sql).all(Number(projectId)) : await db.prepare(sql).all(); - return c.json({ periods }); + const env = await callCoreFn<{ periods: unknown[] }>( + db, + "core.fn_payroll_period_list", + projectId ? { project_id: Number(projectId) } : {}, + { route: routeLabel(c) }, + ); + return respondRpc(c, env); }); app.get("/v1/payroll/periods/:id", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); - const period = await db.prepare( - `SELECT pe.*, p.name AS project_name FROM payroll_periods pe JOIN projects p ON p.id=pe.project_id WHERE pe.id=?`, - ).get(id); - if (!period) return c.json({ error: "Periodo no encontrado" }, 404); - const lines = await db.prepare( - `SELECT l.*, w.first_name, w.last_name_p, w.position FROM payroll_lines l - JOIN workers w ON w.id=l.worker_id WHERE l.period_id=? ORDER BY w.last_name_p`, - ).all(id); - return c.json({ period, lines }); + const env = await callCoreFn<{ period: unknown; lines: unknown[] }>( + db, + "core.fn_payroll_period_get", + { id }, + { route: routeLabel(c) }, + ); + return respondRpc(c, env); }); app.get("/v1/payroll/periods/:id/csv", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); - const lines = await db.prepare( - `SELECT w.first_name, w.last_name_p, l.days, l.daily_wage, l.gross, l.discounts, l.loan_payment, l.net - FROM payroll_lines l JOIN workers w ON w.id=l.worker_id WHERE l.period_id=?`, - ).all(id) as Record[]; + const env = await callCoreFn<{ period: unknown; lines: Array> }>( + db, + "core.fn_payroll_period_get", + { id }, + { route: routeLabel(c) }, + ); + if (!env.ok) return respondRpc(c, env); + const lines = env.data?.lines ?? []; const header = "Nombre,Apellido,Dias,Jornal,Bruto,Descuentos,Prestamo,Neto"; const rows = lines.map((l) => [l.first_name, l.last_name_p, l.days, l.daily_wage, l.gross, l.discounts, l.loan_payment, l.net].join(",") @@ -453,7 +464,12 @@ export function registerPayrollRoutes(app: App) { const id = Number(c.req.param("id")); const { status } = await c.req.json<{ status: string }>(); const db = c.get("db"); - await db.prepare("UPDATE payroll_periods SET status=? WHERE id=?").run(status, id); - return c.json({ ok: true }); + const env = await callCoreFn( + db, + "core.fn_payroll_period_patch", + { id, status }, + { route: routeLabel(c) }, + ); + return respondRpc(c, env); }); } diff --git a/api/rpc.ts b/api/rpc.ts new file mode 100644 index 0000000..9855fce --- /dev/null +++ b/api/rpc.ts @@ -0,0 +1,125 @@ +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( + JSON.stringify(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, + }; +} diff --git a/db/core/changesets/006-rpc-infra.sql b/db/core/changesets/006-rpc-infra.sql new file mode 100644 index 0000000..c669017 --- /dev/null +++ b/db/core/changesets/006-rpc-infra.sql @@ -0,0 +1,147 @@ +--liquibase formatted sql +-- PANELS · core · RPC infrastructure (envelope helpers) + +--changeset panel:core-006a-rpc-generic-messages endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.rpc_generic_messages() +RETURNS text[] +LANGUAGE sql +IMMUTABLE +AS $$ + SELECT ARRAY[ + 'OK', 'Error', 'ERROR', 'Operación exitosa', 'Operacion exitosa', + 'Algo salió mal', 'Algo salio mal', 'Error interno', 'Error interno del servidor' + ]::text[]; +$$; + +--changeset panel:core-006b-rpc-assert-message endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.rpc_assert_message(p_message text) +RETURNS void +LANGUAGE plpgsql +IMMUTABLE +AS $$ +BEGIN + IF p_message IS NULL OR btrim(p_message) = '' THEN + RAISE EXCEPTION 'RPC message must not be empty'; + END IF; + IF p_message = ANY (core.rpc_generic_messages()) THEN + RAISE EXCEPTION 'RPC message is too generic: %', p_message; + END IF; +END; +$$; + +--changeset panel:core-006c-rpc-ok endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.rpc_ok( + p_data jsonb, + p_message text, + p_ctx jsonb DEFAULT '{}'::jsonb +) +RETURNS jsonb +LANGUAGE plpgsql +IMMUTABLE +AS $$ +BEGIN + PERFORM core.rpc_assert_message(p_message); + RETURN jsonb_build_object( + 'ok', true, + 'code', 'OK', + 'layer', 'db', + 'message', p_message, + 'context', COALESCE(p_ctx, '{}'::jsonb), + 'data', p_data, + 'errors', NULL + ); +END; +$$; + +--changeset panel:core-006d-rpc-created endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.rpc_created( + p_data jsonb, + p_message text, + p_ctx jsonb DEFAULT '{}'::jsonb +) +RETURNS jsonb +LANGUAGE plpgsql +IMMUTABLE +AS $$ +BEGIN + PERFORM core.rpc_assert_message(p_message); + RETURN jsonb_build_object( + 'ok', true, + 'code', 'CREATED', + 'layer', 'db', + 'message', p_message, + 'context', COALESCE(p_ctx, '{}'::jsonb), + 'data', p_data, + 'errors', NULL + ); +END; +$$; + +--changeset panel:core-006e-rpc-err endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.rpc_err( + p_code text, + p_message text, + p_ctx jsonb DEFAULT '{}'::jsonb, + p_errors jsonb DEFAULT NULL +) +RETURNS jsonb +LANGUAGE plpgsql +IMMUTABLE +AS $$ +BEGIN + PERFORM core.rpc_assert_message(p_message); + RETURN jsonb_build_object( + 'ok', false, + 'code', p_code, + 'layer', 'db', + 'message', p_message, + 'context', COALESCE(p_ctx, '{}'::jsonb), + 'data', NULL, + 'errors', p_errors + ); +END; +$$; + +--changeset panel:core-006f-rpc-map-exception endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.rpc_from_exception( + p_fn text, + p_sqlstate text, + p_message text, + p_detail text DEFAULT NULL +) +RETURNS jsonb +LANGUAGE plpgsql +IMMUTABLE +AS $$ +DECLARE + v_code text := 'INTERNAL'; + v_msg text; + v_ctx jsonb := jsonb_build_object('fn', p_fn, 'sqlstate', p_sqlstate); +BEGIN + IF p_sqlstate = '23505' THEN + v_code := 'CONFLICT'; + v_msg := format('CONFLICT en %s: %s', p_fn, COALESCE(p_detail, p_message)); + ELSIF p_sqlstate = '23503' THEN + v_code := 'VALIDATION'; + v_msg := format('Referencia inválida en %s: %s', p_fn, COALESCE(p_detail, p_message)); + ELSIF p_sqlstate = '23514' THEN + v_code := 'VALIDATION'; + v_msg := format('Restricción de valor en %s: %s', p_fn, COALESCE(p_detail, p_message)); + ELSIF p_sqlstate = '40P01' THEN + v_code := 'CONFLICT'; + v_msg := format('Deadlock en %s: otra operación modificó los mismos registros, intente de nuevo', p_fn); + ELSE + v_code := 'INTERNAL'; + v_msg := format('Error inesperado en %s (sqlstate=%s): %s', p_fn, p_sqlstate, p_message); + END IF; + RETURN core.rpc_err(v_code, v_msg, v_ctx); +END; +$$; + +--changeset panel:core-006g-rpc-grants endDelimiter:; splitStatements:true +GRANT EXECUTE ON FUNCTION core.rpc_generic_messages() TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.rpc_assert_message(text) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.rpc_ok(jsonb, text, jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.rpc_created(jsonb, text, jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.rpc_err(text, text, jsonb, jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.rpc_from_exception(text, text, text, text) TO panels_core_app; diff --git a/db/core/changesets/007-rpc-catalogs.sql b/db/core/changesets/007-rpc-catalogs.sql new file mode 100644 index 0000000..49eade4 --- /dev/null +++ b/db/core/changesets/007-rpc-catalogs.sql @@ -0,0 +1,123 @@ +--liquibase formatted sql +-- PANELS · core · RPC catálogos y configuración tenant + +--changeset panel:core-007a-fn-catalogs endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.fn_catalogs(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_risks jsonb; + v_themes jsonb; + v_doc_types jsonb; + v_proj_doc_types jsonb; + v_co_doc_types jsonb; + v_companies jsonb; +BEGIN + SELECT COALESCE(jsonb_agg(to_jsonb(r) ORDER BY r.code), '[]'::jsonb) + INTO v_risks FROM risk_levels r; + SELECT COALESCE(jsonb_agg(to_jsonb(t) ORDER BY t.id), '[]'::jsonb) + INTO v_themes FROM badge_themes t; + SELECT COALESCE(jsonb_agg(to_jsonb(d) ORDER BY d.code), '[]'::jsonb) + INTO v_doc_types FROM document_types d; + SELECT COALESCE(jsonb_agg(to_jsonb(p) ORDER BY p.code), '[]'::jsonb) + INTO v_proj_doc_types FROM project_document_types p; + SELECT COALESCE(jsonb_agg(to_jsonb(c) ORDER BY c.code), '[]'::jsonb) + INTO v_co_doc_types FROM company_document_types c; + SELECT COALESCE(jsonb_agg(row_to_json(x)::jsonb ORDER BY x.kind, lower(x.name)), '[]'::jsonb) + INTO v_companies + FROM ( + SELECT c.*, p.code AS parent_code, p.name AS parent_name, + (SELECT COUNT(*) FROM workers w WHERE w.company_id = c.id) AS worker_count + FROM companies c + LEFT JOIN companies p ON p.id = c.parent_id + ORDER BY CASE c.kind WHEN 'principal' THEN 0 ELSE 1 END, lower(c.name) + ) x; + RETURN core.rpc_ok( + jsonb_build_object( + 'risks', v_risks, + 'themes', v_themes, + 'document_types', v_doc_types, + 'project_document_types', v_proj_doc_types, + 'company_document_types', v_co_doc_types, + 'companies', v_companies + ), + format('Catálogos cargados: %s empresas, %s tipos de documento de personal', + jsonb_array_length(v_companies), jsonb_array_length(v_doc_types)), + jsonb_build_object('fn', 'fn_catalogs', 'company_count', jsonb_array_length(v_companies)) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_catalogs', SQLSTATE, SQLERRM); +END; +$$; + +--changeset panel:core-007b-fn-tenant-timezone-get endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.fn_tenant_timezone_get(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_tid integer := (payload->>'tenant_id')::integer; + v_tz text; +BEGIN + IF v_tid IS NULL THEN + RETURN core.rpc_ok( + jsonb_build_object('timezone', 'America/Mexico_City'), + 'Zona horaria por defecto America/Mexico_City (sin tenant_id en la solicitud)', + jsonb_build_object('fn', 'fn_tenant_timezone_get', 'tenant_id', NULL) + ); + END IF; + SELECT timezone INTO v_tz FROM tenant_settings WHERE tenant_id = v_tid; + v_tz := COALESCE(v_tz, 'America/Mexico_City'); + RETURN core.rpc_ok( + jsonb_build_object('timezone', v_tz), + format('Zona horaria del tenant %s: %s', v_tid, v_tz), + jsonb_build_object('fn', 'fn_tenant_timezone_get', 'tenant_id', v_tid) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_tenant_timezone_get', SQLSTATE, SQLERRM); +END; +$$; + +--changeset panel:core-007c-fn-tenant-settings-upsert endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.fn_tenant_settings_upsert(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_tid integer := (payload->>'tenant_id')::integer; + v_tz text := btrim(payload->>'timezone'); +BEGIN + IF v_tid IS NULL THEN + RETURN core.rpc_err('VALIDATION', + 'fn_tenant_settings_upsert: tenant_id es obligatorio para guardar configuración', + jsonb_build_object('fn', 'fn_tenant_settings_upsert', 'field', 'tenant_id')); + END IF; + IF v_tz IS NULL OR v_tz = '' THEN + RETURN core.rpc_err('VALIDATION', + format('fn_tenant_settings_upsert: timezone vacío para tenant_id=%s', v_tid), + jsonb_build_object('fn', 'fn_tenant_settings_upsert', 'tenant_id', v_tid, 'field', 'timezone')); + END IF; + INSERT INTO tenant_settings (tenant_id, timezone, updated_at) + VALUES (v_tid, v_tz, now()) + ON CONFLICT (tenant_id) DO UPDATE SET timezone = EXCLUDED.timezone, updated_at = now(); + RETURN core.rpc_ok( + jsonb_build_object('timezone', v_tz), + format('Zona horaria del tenant %s actualizada a %s', v_tid, v_tz), + jsonb_build_object('fn', 'fn_tenant_settings_upsert', 'tenant_id', v_tid) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_tenant_settings_upsert', SQLSTATE, SQLERRM); +END; +$$; + +--changeset panel:core-007d-fn-catalogs-grants endDelimiter:; splitStatements:true +GRANT EXECUTE ON FUNCTION core.fn_catalogs(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_tenant_timezone_get(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_tenant_settings_upsert(jsonb) TO panels_core_app; diff --git a/db/core/changesets/013-rpc-payroll.sql b/db/core/changesets/013-rpc-payroll.sql index 5cacc4e..665919e 100644 --- a/db/core/changesets/013-rpc-payroll.sql +++ b/db/core/changesets/013-rpc-payroll.sql @@ -1636,6 +1636,371 @@ EXCEPTION WHEN OTHERS THEN END; $$; +--changeset panel:core-013i2-fn-payroll-week-list-open endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.fn_payroll_week_list_open(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_tid integer := (payload->>'tenant_id')::integer; + v_weeks jsonb; +BEGIN + SELECT COALESCE(jsonb_agg(to_jsonb(w) ORDER BY w.week_start DESC), '[]'::jsonb) + INTO v_weeks + FROM ( + SELECT pw.*, + (SELECT COALESCE(SUM(l.loan_discount), 0) + FROM payroll_week_lines l + JOIN payroll_sheets s ON s.id = l.sheet_id + WHERE s.week_id = pw.id) AS loan_recovery + FROM payroll_weeks pw + WHERE pw.tenant_id = v_tid AND pw.status != 'paid' + ) w; + RETURN core.rpc_ok( + jsonb_build_object('weeks', v_weeks), + format('Semanas abiertas listadas para tenant %s: %s semana(s)', v_tid, jsonb_array_length(v_weeks)), + jsonb_build_object('fn', 'fn_payroll_week_list_open', 'tenant_id', v_tid) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_payroll_week_list_open', SQLSTATE, SQLERRM); +END; +$$; + +--changeset panel:core-013i3-fn-payroll-week-csv endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.fn_payroll_week_csv(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_week_id bigint := (payload->>'week_id')::bigint; + v_csv text; +BEGIN + IF v_week_id IS NULL THEN + RETURN core.rpc_err('VALIDATION', 'fn_payroll_week_csv: week_id es obligatorio', + jsonb_build_object('fn', 'fn_payroll_week_csv', 'field', 'week_id')); + END IF; + SELECT string_agg(row_text, E'\n' ORDER BY ord) + INTO v_csv + FROM ( + SELECT 0 AS ord, 'Nombre,Apellido,Hoja,Proyecto,Dias,Jornal,Concepto,Cantidad,Unidad,Bruto,Descuentos,Prestamo,EtiquetaPrestamo,Requerido,APagar' AS row_text + UNION ALL + SELECT 1 AS ord, + concat_ws(',', + w.first_name, w.last_name_p, s.kind, COALESCE(p.name, ''), + l.days, l.daily_wage, COALESCE(l.concepto, ''), + COALESCE(l.qty_actual, 0) + COALESCE(l.qty_extra, 0), + COALESCE(l.unit_code, ''), l.gross, l.discounts, l.loan_discount, + replace(COALESCE(l.loan_label, ''), ',', ' '), + l.required_net, l.payable_net + ) + FROM payroll_week_lines l + JOIN payroll_sheets s ON s.id = l.sheet_id + JOIN workers w ON w.id = l.worker_id + LEFT JOIN projects p ON p.id = s.project_id + WHERE s.week_id = v_week_id + ORDER BY s.kind, p.name, w.last_name_p + ) q; + RETURN core.rpc_ok( + jsonb_build_object('csv', v_csv), + format('CSV de nómina generado para semana id=%s', v_week_id), + jsonb_build_object('fn', 'fn_payroll_week_csv', 'week_id', v_week_id) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_payroll_week_csv', SQLSTATE, SQLERRM); +END; +$$; + +--changeset panel:core-013i4-fn-loan-read endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.fn_loan_list(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_worker_id bigint := NULLIF(payload->>'worker_id', '')::bigint; + v_loans jsonb; +BEGIN + SELECT COALESCE(jsonb_agg(row_to_json(x)::jsonb ORDER BY x.id DESC), '[]'::jsonb) + INTO v_loans + FROM ( + SELECT l.*, w.first_name, w.last_name_p + FROM loans l + JOIN workers w ON w.id = l.worker_id + WHERE v_worker_id IS NULL OR l.worker_id = v_worker_id + ) x; + RETURN core.rpc_ok( + jsonb_build_object('loans', v_loans), + format('Préstamos listados: %s registro(s)', jsonb_array_length(v_loans)), + jsonb_build_object('fn', 'fn_loan_list', 'worker_id', v_worker_id, 'count', jsonb_array_length(v_loans)) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_loan_list', SQLSTATE, SQLERRM); +END; +$$; + +CREATE OR REPLACE FUNCTION core.fn_loan_get(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_id bigint := NULLIF(payload->>'id', '')::bigint; + v_loan jsonb; +BEGIN + IF v_id IS NULL THEN + RETURN core.rpc_err( + 'VALIDATION', + 'fn_loan_get: id es obligatorio', + jsonb_build_object('fn', 'fn_loan_get', 'field', 'id') + ); + END IF; + SELECT row_to_json(x)::jsonb INTO v_loan + FROM ( + SELECT l.*, w.first_name, w.middle_name, w.last_name_p, w.last_name_m + FROM loans l + JOIN workers w ON w.id = l.worker_id + WHERE l.id = v_id + ) x; + IF v_loan IS NULL THEN + RETURN core.rpc_err( + 'NOT_FOUND', + 'Préstamo no encontrado', + jsonb_build_object('fn', 'fn_loan_get', 'id', v_id) + ); + END IF; + RETURN core.rpc_ok( + jsonb_build_object('loan', v_loan), + format('Préstamo id=%s cargado', v_id), + jsonb_build_object('fn', 'fn_loan_get', 'id', v_id) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_loan_get', SQLSTATE, SQLERRM); +END; +$$; + +CREATE OR REPLACE FUNCTION core.fn_loan_payment_get(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_id bigint := NULLIF(payload->>'id', '')::bigint; + v_payment jsonb; +BEGIN + IF v_id IS NULL THEN + RETURN core.rpc_err( + 'VALIDATION', + 'fn_loan_payment_get: id es obligatorio', + jsonb_build_object('fn', 'fn_loan_payment_get', 'field', 'id') + ); + END IF; + SELECT row_to_json(x)::jsonb INTO v_payment + FROM ( + SELECT p.*, l.delivered, l.commission_pct, + w.first_name, w.middle_name, w.last_name_p, w.last_name_m, + wk.week_end + FROM loan_payments p + JOIN loans l ON l.id = p.loan_id + JOIN workers w ON w.id = l.worker_id + JOIN payroll_weeks wk ON wk.id = p.week_id + WHERE p.id = v_id + ) x; + IF v_payment IS NULL THEN + RETURN core.rpc_err( + 'NOT_FOUND', + 'Recibo no encontrado', + jsonb_build_object('fn', 'fn_loan_payment_get', 'id', v_id) + ); + END IF; + RETURN core.rpc_ok( + jsonb_build_object('payment', v_payment), + format('Pago de préstamo id=%s cargado', v_id), + jsonb_build_object('fn', 'fn_loan_payment_get', 'id', v_id) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_loan_payment_get', SQLSTATE, SQLERRM); +END; +$$; + +--changeset panel:core-013i5-fn-payroll-period endDelimiter:; splitStatements:true +CREATE OR REPLACE FUNCTION core.fn_payroll_period_create(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_project_id bigint := NULLIF(payload->>'project_id', '')::bigint; + v_week_start date := NULLIF(btrim(payload->>'week_start'), '')::date; + v_week_end date := NULLIF(btrim(payload->>'week_end'), '')::date; + v_extra jsonb := COALESCE(payload->'extra_discounts', '{}'::jsonb); + v_period_id bigint; + v_worker record; + v_att_n integer; + v_days integer; + v_gross numeric; + v_extra_disc numeric; + v_net numeric; +BEGIN + IF v_project_id IS NULL OR v_week_start IS NULL OR v_week_end IS NULL THEN + RETURN core.rpc_err( + 'VALIDATION', + 'fn_payroll_period_create: project_id, week_start y week_end son obligatorios', + jsonb_build_object('fn', 'fn_payroll_period_create') + ); + END IF; + INSERT INTO payroll_periods (project_id, week_start, week_end, status) + VALUES (v_project_id, v_week_start, v_week_end, 'draft') + RETURNING id INTO v_period_id; + FOR v_worker IN + SELECT w.id, w.daily_wage + FROM workers w + JOIN assignments a ON a.worker_id = w.id AND a.project_id = v_project_id AND a.active = true + WHERE w.status = 'activo' + LOOP + SELECT COUNT(*)::integer INTO v_att_n + FROM attendance + WHERE worker_id = v_worker.id + AND project_id = v_project_id + AND work_date BETWEEN v_week_start AND v_week_end + AND present = true; + v_days := v_att_n; + v_gross := v_days * v_worker.daily_wage; + v_extra_disc := COALESCE((v_extra->>v_worker.id::text)::numeric, 0); + v_net := v_gross - v_extra_disc; + INSERT INTO payroll_lines (period_id, worker_id, days, daily_wage, gross, discounts, loan_payment, net) + VALUES (v_period_id, v_worker.id, v_days, v_worker.daily_wage, v_gross, v_extra_disc, 0, v_net); + END LOOP; + RETURN core.rpc_created( + jsonb_build_object('id', v_period_id), + format('Periodo de nómina %s creado para proyecto %s', v_period_id, v_project_id), + jsonb_build_object('fn', 'fn_payroll_period_create', 'id', v_period_id, 'project_id', v_project_id) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_payroll_period_create', SQLSTATE, SQLERRM); +END; +$$; + +CREATE OR REPLACE FUNCTION core.fn_payroll_period_list(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_project_id bigint := NULLIF(payload->>'project_id', '')::bigint; + v_periods jsonb; +BEGIN + SELECT COALESCE(jsonb_agg(row_to_json(x)::jsonb ORDER BY x.id DESC), '[]'::jsonb) + INTO v_periods + FROM ( + SELECT pe.*, p.name AS project_name + FROM payroll_periods pe + JOIN projects p ON p.id = pe.project_id + WHERE v_project_id IS NULL OR pe.project_id = v_project_id + ) x; + RETURN core.rpc_ok( + jsonb_build_object('periods', v_periods), + format('Periodos de nómina listados: %s registro(s)', jsonb_array_length(v_periods)), + jsonb_build_object('fn', 'fn_payroll_period_list', 'project_id', v_project_id, 'count', jsonb_array_length(v_periods)) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_payroll_period_list', SQLSTATE, SQLERRM); +END; +$$; + +CREATE OR REPLACE FUNCTION core.fn_payroll_period_get(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_id bigint := NULLIF(payload->>'id', '')::bigint; + v_period jsonb; + v_lines jsonb; +BEGIN + IF v_id IS NULL THEN + RETURN core.rpc_err( + 'VALIDATION', + 'fn_payroll_period_get: id es obligatorio', + jsonb_build_object('fn', 'fn_payroll_period_get', 'field', 'id') + ); + END IF; + SELECT row_to_json(x)::jsonb INTO v_period + FROM ( + SELECT pe.*, p.name AS project_name + FROM payroll_periods pe + JOIN projects p ON p.id = pe.project_id + WHERE pe.id = v_id + ) x; + IF v_period IS NULL THEN + RETURN core.rpc_err( + 'NOT_FOUND', + 'Periodo no encontrado', + jsonb_build_object('fn', 'fn_payroll_period_get', 'id', v_id) + ); + END IF; + SELECT COALESCE(jsonb_agg(row_to_json(lx)::jsonb ORDER BY lx.last_name_p), '[]'::jsonb) + INTO v_lines + FROM ( + SELECT l.*, w.first_name, w.last_name_p, w.position + FROM payroll_lines l + JOIN workers w ON w.id = l.worker_id + WHERE l.period_id = v_id + ) lx; + RETURN core.rpc_ok( + jsonb_build_object('period', v_period, 'lines', v_lines), + format('Periodo de nómina %s cargado con %s línea(s)', v_id, jsonb_array_length(v_lines)), + jsonb_build_object('fn', 'fn_payroll_period_get', 'id', v_id, 'line_count', jsonb_array_length(v_lines)) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_payroll_period_get', SQLSTATE, SQLERRM); +END; +$$; + +CREATE OR REPLACE FUNCTION core.fn_payroll_period_patch(payload jsonb) +RETURNS jsonb +LANGUAGE plpgsql +SECURITY INVOKER +SET search_path = core +AS $$ +DECLARE + v_id bigint := NULLIF(payload->>'id', '')::bigint; + v_status text := btrim(COALESCE(payload->>'status', '')); +BEGIN + IF v_id IS NULL OR v_status = '' THEN + RETURN core.rpc_err( + 'VALIDATION', + 'fn_payroll_period_patch: id y status son obligatorios', + jsonb_build_object('fn', 'fn_payroll_period_patch') + ); + END IF; + IF NOT EXISTS (SELECT 1 FROM payroll_periods WHERE id = v_id) THEN + RETURN core.rpc_err( + 'NOT_FOUND', + 'Periodo no encontrado', + jsonb_build_object('fn', 'fn_payroll_period_patch', 'id', v_id) + ); + END IF; + UPDATE payroll_periods SET status = v_status WHERE id = v_id; + RETURN core.rpc_ok( + jsonb_build_object('ok', true, 'id', v_id, 'status', v_status), + format('Periodo de nómina %s actualizado a estado %s', v_id, v_status), + jsonb_build_object('fn', 'fn_payroll_period_patch', 'id', v_id, 'status', v_status) + ); +EXCEPTION WHEN OTHERS THEN + RETURN core.rpc_from_exception('fn_payroll_period_patch', SQLSTATE, SQLERRM); +END; +$$; + --changeset panel:core-013j-fn-payroll-grants endDelimiter:; splitStatements:true GRANT EXECUTE ON FUNCTION core.round_money(numeric) TO panels_core_app; GRANT EXECUTE ON FUNCTION core.fn_payroll_settings_get(jsonb) TO panels_core_app; @@ -1656,3 +2021,12 @@ GRANT EXECUTE ON FUNCTION core.fn_payroll_admin_line_add(jsonb) TO panels_core_a GRANT EXECUTE ON FUNCTION core.fn_payroll_admin_line_update(jsonb) TO panels_core_app; GRANT EXECUTE ON FUNCTION core.fn_payroll_admin_line_remove(jsonb) TO panels_core_app; GRANT EXECUTE ON FUNCTION core.fn_payroll_jornal_worker_add(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_payroll_week_list_open(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_payroll_week_csv(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_loan_list(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_loan_get(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_loan_payment_get(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_payroll_period_create(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_payroll_period_list(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_payroll_period_get(jsonb) TO panels_core_app; +GRANT EXECUTE ON FUNCTION core.fn_payroll_period_patch(jsonb) TO panels_core_app; From e92614acb9a15f80f00d710906a36c210a69e7de Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 3 Sep 2026 23:58:55 +0000 Subject: [PATCH 2/2] refactor(api): migrate main.ts CRUD routes to PostgreSQL RPC Replace all db.prepare() calls in main.ts with callCoreFn and the respondRpc/respondApiError HTTP envelope helpers. Company and project routes use the RPC response pattern; worker, budget, badge, catalog, and tenant settings endpoints call their matching core.fn_* functions. Adds app.onError(onAppError) for consistent error handling. Keeps companies.ts wrappers, excel.ts document storage, and auth/PDF/S3 unchanged. Co-authored-by: alberto.martinez --- api/main.ts | 1041 ++++++++++++++++++++++++++++----------------------- 1 file changed, 565 insertions(+), 476 deletions(-) diff --git a/api/main.ts b/api/main.ts index 0ac52b5..ba5a988 100644 --- a/api/main.ts +++ b/api/main.ts @@ -3,7 +3,6 @@ import { cors } from "hono/cors"; import { checklistFor, refreshPipeline, - nextProjectCode, isProjectStatus, projectById, projectMustBe, @@ -11,13 +10,19 @@ import { projectChecklistFor, companyChecklistFor, imssFlagsFor, - checklistItemsFor, - lastInsertId, pingCoreDb, tenantTimezone, getCoreDb, type Db, } from "./db.ts"; +import { callCoreFn } from "./rpc.ts"; +import { + respondRpc, + respondApiError, + onAppError, + routeLabel, + type RpcEnvelope, +} from "./http_errors.ts"; import { config } from "./config.ts"; import { clearSession, @@ -36,14 +41,13 @@ import { createTenant, getTenantDetail, issueTenantAdminAccess, listTenants, upd import { smtpConfigured, testSmtp } from "./mail.ts"; import { saveSmtpSettings, smtpPublicView } from "./smtp.ts"; import { decryptBytes } from "./docs_crypto.ts"; -import { importExcel, storeDocument, storeProjectDocument, storeCompanyDocument, findExisting, assign, buildImportTemplate } from "./excel.ts"; +import { importExcel, storeDocument, storeProjectDocument, storeCompanyDocument, buildImportTemplate } from "./excel.ts"; import { normalizeWorker, validateCurp, validateNss, validateRfc, validateWorkerFields, - fullName, formatNss, normUpper, titleCase, @@ -58,15 +62,48 @@ import { exportBudgetWorkbook, previewBudgetExcel, importBudgetExcel, - lineAmount, listBudget, } from "./budget.ts"; -import { companyDocKey, getObject, pingStorage, projectDocKey, workerDocKey } from "./storage.ts"; +import { badgeJobPdfKey, companyDocKey, getObject, pingStorage, projectDocKey, workerDocKey } from "./storage.ts"; import { cacheCore, cacheKeyCore } from "./cache.ts"; import { pingRedis } from "./redis.ts"; const app = new Hono<{ Variables: { user: AuthUser; db: Db } }>(); +app.onError(onAppError); + +type StorageDoc = { + storage_name: string; + iv: string; + original_name: string; + mime?: string; +}; + +async function rpcDocumentForDownload( + db: Db, + fn: string, + payload: Record, + route: string, +): Promise<{ doc?: StorageDoc; envelope: RpcEnvelope }> { + const envelope = await callCoreFn<{ document: StorageDoc }>(db, fn, payload, { route }); + if (!envelope.ok) return { envelope }; + const doc = envelope.data?.document; + if (!doc?.storage_name) { + return { + envelope: { + ok: false, + code: "NOT_FOUND", + layer: "db", + message: "Documento no encontrado", + context: { fn, ...payload }, + data: null, + errors: null, + }, + }; + } + return { doc, envelope }; +} + async function scopedCompany( db: Db, id: number, @@ -237,20 +274,16 @@ app.get("/v1/catalogs", ...requireCoreAuth, async (c) => { const db = c.get("db"); const user = c.get("user"); const tid = tenantScope(user); - // Catálogos de baja escritura (Fase 4e): TTL corto, sin invalidación - // explícita. La llave incluye tenant_id porque `companies` sí es dato - // de negocio por-tenant -- risk_levels/badge_themes/document_types son - // globales, pero se cachean juntos por simplicidad de esta única - // respuesta agregada. - const payload = await cacheCore(cacheKeyCore(tid, "catalogs"), 30, async () => ({ - risks: await db.prepare("SELECT * FROM risk_levels").all(), - themes: await db.prepare("SELECT * FROM badge_themes").all(), - document_types: await db.prepare("SELECT * FROM document_types").all(), - project_document_types: await db.prepare("SELECT * FROM project_document_types").all(), - company_document_types: await db.prepare("SELECT * FROM company_document_types").all(), - project_statuses: PROJECT_STATUS_CATALOG, - companies: await listCompanies(db, tid), - })); + const route = routeLabel(c); + const payload = await cacheCore(cacheKeyCore(tid, "catalogs"), 30, async () => { + const env = await callCoreFn>(db, "core.fn_catalogs", {}, { route }); + if (!env.ok) throw new Error(env.message); + return { + ...env.data, + project_statuses: PROJECT_STATUS_CATALOG, + companies: await listCompanies(db, tid), + }; + }); return c.json(payload); }); @@ -265,25 +298,44 @@ app.get("/v1/configuracion", ...requireCoreAuth, async (c) => { app.put("/v1/configuracion", ...requireCoreAuth, async (c) => { const user = c.get("user"); - if (user.role !== "tenant_admin") return c.json({ error: "Solo administradores" }, 403); + if (user.role !== "tenant_admin") { + return respondApiError(c, "FORBIDDEN", "Solo administradores pueden cambiar la configuración del tenant", { + route: routeLabel(c), + }); + } const db = c.get("db"); const tid = tenantScope(user); - if (tid == null) return c.json({ error: "Cuenta sin tenant" }, 400); + if (tid == null) { + return respondApiError(c, "VALIDATION", "Cuenta sin tenant asociada para guardar configuración", { + route: routeLabel(c), + }); + } const body = await c.req.json<{ timezone?: string }>(); const tz = (body.timezone ?? "").trim(); const valid = new Set(Intl.supportedValuesOf("timeZone")); - if (!tz || !valid.has(tz)) return c.json({ error: "Zona horaria inválida" }, 400); - await db.prepare( - `INSERT INTO tenant_settings (tenant_id, timezone, updated_at) VALUES (?, ?, now()) - ON CONFLICT (tenant_id) DO UPDATE SET timezone = excluded.timezone, updated_at = now()`, - ).run(tid, tz); - return c.json({ timezone: tz }); + if (!tz || !valid.has(tz)) { + return respondApiError(c, "VALIDATION", `Zona horaria inválida en ${routeLabel(c)}: ${tz || "(vacío)"}`, { + route: routeLabel(c), + timezone: tz || null, + }); + } + return respondRpc(c, await callCoreFn( + db, + "core.fn_tenant_settings_upsert", + { tenant_id: tid, timezone: tz }, + { route: routeLabel(c) }, + )); }); app.get("/v1/companies", ...requireCoreAuth, async (c) => { const db = c.get("db"); const tid = tenantScope(c.get("user")); - return c.json({ companies: await listCompanies(db, tid) }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_company_list", + { tenant_id: tid }, + { route: routeLabel(c) }, + )); }); app.post("/v1/companies", ...requireCoreAuth, async (c) => { @@ -291,8 +343,17 @@ app.post("/v1/companies", ...requireCoreAuth, async (c) => { const db = c.get("db"); const tid = tenantScope(c.get("user")); const result = await createSubcompany(db, body, tid); - if (result.error || !result.company) return c.json({ error: result.error }, 400); - return c.json({ company: result.company }, 201); + if (result.error || !result.company) { + const code = result.status === 404 ? "NOT_FOUND" : "VALIDATION"; + return respondApiError(c, code, result.error ?? "No se pudo crear la empresa", { route: routeLabel(c) }); + } + return respondRpc(c, { + ok: true, + code: "CREATED", + layer: "db", + message: `Empresa ${result.company.code} creada`, + data: { company: result.company }, + }); }); app.patch("/v1/companies/:id", ...requireCoreAuth, async (c) => { @@ -301,22 +362,50 @@ app.patch("/v1/companies/:id", ...requireCoreAuth, async (c) => { const db = c.get("db"); const tid = tenantScope(c.get("user")); const result = await updateCompany(db, id, body, tid); - if (result.error) return c.json({ error: result.error }, (result.status ?? 400) as 400 | 404); - return c.json({ company: result.company }); + if (result.error) { + const code = result.status === 404 ? "NOT_FOUND" : "VALIDATION"; + return respondApiError(c, code, result.error, { route: routeLabel(c), company_id: id }); + } + return respondRpc(c, { + ok: true, + code: "OK", + layer: "db", + message: `Empresa id=${id} actualizada`, + data: { company: result.company }, + }); }); app.get("/v1/companies/:id", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); - const company = await scopedCompany(db, id, tenantScope(c.get("user"))); - if (!company) return c.json({ error: "Empresa no encontrada" }, 404); - const documents = await db.prepare( - "SELECT id, type_code, original_name, mime, size_bytes, is_current, uploaded_at FROM company_documents WHERE company_id = ? ORDER BY uploaded_at DESC", - ).all(id); - return c.json({ - company, - documents, - checklist: await companyChecklistFor(db, id), + const tid = tenantScope(c.get("user")); + const route = routeLabel(c); + const companyEnv = await callCoreFn<{ company: Record }>(db, "core.fn_company_get", { id }, { route }); + if (!companyEnv.ok) return respondRpc(c, companyEnv); + const company = companyEnv.data?.company; + if (!company) { + return respondApiError(c, "NOT_FOUND", `Empresa id=${id} no encontrada en ${route}`, { route, company_id: id }); + } + if (tid != null && company.tenant_id != null && company.tenant_id !== tid) { + return respondApiError(c, "NOT_FOUND", `Empresa id=${id} no encontrada en ${route}`, { route, company_id: id }); + } + const docsEnv = await callCoreFn<{ documents: unknown[] }>( + db, + "core.fn_company_document_list", + { company_id: id }, + { route }, + ); + if (!docsEnv.ok) return respondRpc(c, docsEnv); + return respondRpc(c, { + ok: true, + code: "OK", + layer: "db", + message: `Empresa ${company.code} cargada`, + data: { + company, + documents: docsEnv.data?.documents ?? [], + checklist: await companyChecklistFor(db, id), + }, }); }); @@ -341,22 +430,19 @@ app.get("/v1/companies/:id/documents/:docId", ...requireCoreAuth, async (c) => { const companyId = Number(c.req.param("id")); const docId = Number(c.req.param("docId")); const db = c.get("db"); - if (!await scopedCompany(db, companyId, tenantScope(c.get("user")))) return c.json({ error: "Empresa no encontrada" }, 404); - const doc = await db.prepare( - "SELECT * FROM company_documents WHERE id=? AND company_id=?", - ).get(docId, companyId) as { - storage_name: string; - iv: string; - mime: string; - original_name: string; - } | undefined; - if (!doc) return c.json({ error: "Documento no encontrado" }, 404); + const route = routeLabel(c); + if (!await scopedCompany(db, companyId, tenantScope(c.get("user")))) { + return respondApiError(c, "NOT_FOUND", `Empresa id=${companyId} no encontrada en ${route}`, { route, company_id: companyId }); + } + const { doc, envelope } = await rpcDocumentForDownload( + db, + "core.fn_company_document_get", + { company_id: companyId, id: docId }, + route, + ); + if (!doc) return respondRpc(c, envelope); const enc = await getObject(companyDocKey(companyId, doc.storage_name)); const plain = await decryptBytes(doc.iv, enc); - // Documentos subidos por el usuario se sirven como adjunto, nunca inline: - // servir "inline" dejaría que el navegador renderice el Content-Type que - // el propio uploader eligió (ej. un .html disfrazado de "comprobante"), - // habilitando XSS almacenado dentro del origen autenticado de la app. c.header("Content-Type", "application/octet-stream"); c.header("X-Content-Type-Options", "nosniff"); c.header("Content-Disposition", `attachment; filename="${encodeURIComponent(doc.original_name)}"`); @@ -365,36 +451,20 @@ app.get("/v1/companies/:id/documents/:docId", ...requireCoreAuth, async (c) => { app.get("/v1/projects", ...requireCoreAuth, async (c) => { const db = c.get("db"); - const tid = tenantScope(c.get("user")); const status = (c.req.query("status") ?? "").trim(); const allowed = status ? status.split(",").map((s) => s.trim()).filter(isProjectStatus) : []; - const clauses: string[] = []; - const params: (string | number)[] = []; - if (tid != null) { - clauses.push("p.tenant_id = ?"); - params.push(tid); + if (status && !allowed.length) { + return respondApiError(c, "VALIDATION", `Ningún estado de proyecto válido en ${routeLabel(c)}: ${status}`, { + route: routeLabel(c), + status, + }); } - if (allowed.length) { - clauses.push(`p.status IN (${allowed.map(() => "?").join(",")})`); - params.push(...allowed); - } - const where = clauses.length ? `WHERE ${clauses.join(" AND ")}` : ""; - return c.json({ - projects: await db.prepare( - `SELECT p.*, c.name AS company_name, c.code AS company_code, - (SELECT COUNT(*) FROM assignments a WHERE a.project_id = p.id AND a.active = true) AS active_count, - (SELECT COUNT(*) FROM project_document_types t - WHERE t.required = true AND NOT EXISTS ( - SELECT 1 FROM project_documents d - WHERE d.project_id = p.id AND d.type_code = t.code AND d.is_current = true - )) AS missing_docs, - (SELECT COUNT(*) FROM budget_items b WHERE b.project_id = p.id) AS budget_count - FROM projects p - LEFT JOIN companies c ON c.id = p.company_id - ${where} - ORDER BY CASE p.status WHEN 'activo' THEN 0 WHEN 'pausado' THEN 1 WHEN 'concluido' THEN 2 ELSE 3 END, p.id`, - ).all(...params), - }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_project_list", + { status: allowed.length ? allowed.join(",") : undefined }, + { route: routeLabel(c) }, + )); }); type ProjectInput = { @@ -446,98 +516,129 @@ function projectFields(body: ProjectInput, cur?: Record) { app.post("/v1/projects", ...requireCoreAuth, async (c) => { const body = await c.req.json(); - if (!body.name?.trim()) return c.json({ error: "Nombre de proyecto obligatorio" }, 400); + if (!body.name?.trim()) { + return respondApiError(c, "VALIDATION", `Nombre de proyecto obligatorio en ${routeLabel(c)}`, { + route: routeLabel(c), + }); + } const status = body.status && isProjectStatus(body.status) ? body.status : "activo"; const db = c.get("db"); const fields = projectFields(body); - if (!fields.company_id) return c.json({ error: "Seleccione la empresa del proyecto" }, 400); + if (!fields.company_id) { + return respondApiError(c, "VALIDATION", `Seleccione la empresa del proyecto en ${routeLabel(c)}`, { + route: routeLabel(c), + }); + } if (!await resolveCompany(db, { company_id: fields.company_id })) { - return c.json({ error: "Empresa no encontrada" }, 400); + return respondApiError(c, "VALIDATION", `Empresa id=${fields.company_id} no encontrada en ${routeLabel(c)}`, { + route: routeLabel(c), + company_id: fields.company_id, + }); } if (fields.start_date && fields.end_date && fields.end_date < fields.start_date) { - return c.json({ error: "La fecha de término no puede ser anterior al inicio" }, 400); + return respondApiError(c, "VALIDATION", `La fecha de término no puede ser anterior al inicio en ${routeLabel(c)}`, { + route: routeLabel(c), + }); } - const code = await nextProjectCode(db); const tid = tenantScope(c.get("user")); - await db.prepare( - `INSERT INTO projects ( - code, name, address, theme_id, status, company_id, contract_amount, - start_date, end_date, resident_name, siroc, payroll_tax_pct, tenant_id - ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, - ).run( - code, - fields.name, - fields.address, - fields.theme_id, - status, - fields.company_id, - fields.contract_amount, - fields.start_date, - fields.end_date, - fields.resident_name, - fields.siroc, - fields.payroll_tax_pct, - tid, - ); - const id = await lastInsertId(db); - return c.json({ id, code, status }, 201); + return respondRpc(c, await callCoreFn( + db, + "core.fn_project_create", + { + tenant_id: tid, + name: fields.name, + address: fields.address, + theme_id: fields.theme_id, + status, + company_id: fields.company_id, + contract_amount: fields.contract_amount, + start_date: fields.start_date, + end_date: fields.end_date, + resident_name: fields.resident_name, + siroc: fields.siroc, + payroll_tax_pct: fields.payroll_tax_pct, + }, + { route: routeLabel(c) }, + )); }); app.patch("/v1/projects/:id", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const body = await c.req.json(); const db = c.get("db"); - // Con RLS activo (db/core/changesets/005-rls.sql), esta consulta ya no - // puede traer una fila de otro tenant: la conexión solo ve app.tenant_id. - const cur = await db.prepare("SELECT * FROM projects WHERE id = ?").get(id) as Record | undefined; - if (!cur) return c.json({ error: "Proyecto no encontrado" }, 404); + const route = routeLabel(c); + const curEnv = await callCoreFn<{ project: Record }>(db, "core.fn_project_get", { id }, { route }); + if (!curEnv.ok) return respondRpc(c, curEnv); + const cur = curEnv.data?.project; + if (!cur) { + return respondApiError(c, "NOT_FOUND", `Proyecto id=${id} no encontrado en ${route}`, { route, project_id: id }); + } const status = body.status && isProjectStatus(body.status) ? body.status : String(cur.status || "activo"); const fields = projectFields(body, cur); - if (!fields.company_id) return c.json({ error: "Seleccione la empresa del proyecto" }, 400); + if (!fields.company_id) { + return respondApiError(c, "VALIDATION", `Seleccione la empresa del proyecto en ${route}`, { route, project_id: id }); + } if (!await resolveCompany(db, { company_id: fields.company_id })) { - return c.json({ error: "Empresa no encontrada" }, 400); + return respondApiError(c, "VALIDATION", `Empresa id=${fields.company_id} no encontrada en ${route}`, { + route, + project_id: id, + company_id: fields.company_id, + }); } if (fields.start_date && fields.end_date && fields.end_date < fields.start_date) { - return c.json({ error: "La fecha de término no puede ser anterior al inicio" }, 400); + return respondApiError(c, "VALIDATION", `La fecha de término no puede ser anterior al inicio en ${route}`, { + route, + project_id: id, + }); } - await db.prepare( - `UPDATE projects SET name=?, address=?, theme_id=?, status=?, company_id=?, contract_amount=?, - start_date=?, end_date=?, resident_name=?, siroc=?, payroll_tax_pct=? WHERE id=?`, - ).run( - fields.name, - fields.address, - fields.theme_id, - status, - fields.company_id, - fields.contract_amount, - fields.start_date, - fields.end_date, - fields.resident_name, - fields.siroc, - fields.payroll_tax_pct, - id, - ); - return c.json({ ok: true, id, code: cur.code, status }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_project_update", + { + id, + name: fields.name, + address: fields.address, + theme_id: fields.theme_id, + status, + company_id: fields.company_id, + contract_amount: fields.contract_amount, + start_date: fields.start_date, + end_date: fields.end_date, + resident_name: fields.resident_name, + siroc: fields.siroc, + payroll_tax_pct: fields.payroll_tax_pct, + }, + { route }, + )); }); app.get("/v1/projects/:id", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); - const project = await db.prepare( - `SELECT p.*, c.name AS company_name, c.code AS company_code, - (SELECT COUNT(*) FROM assignments a WHERE a.project_id = p.id AND a.active = true) AS active_count - FROM projects p - LEFT JOIN companies c ON c.id = p.company_id - WHERE p.id = ?`, - ).get(id); - if (!project) return c.json({ error: "Proyecto no encontrado" }, 404); - const documents = await db.prepare( - "SELECT id, type_code, original_name, mime, size_bytes, is_current, uploaded_at FROM project_documents WHERE project_id = ? ORDER BY uploaded_at DESC", - ).all(id); - return c.json({ - project, - documents, - checklist: await projectChecklistFor(db, id), + const route = routeLabel(c); + const projectEnv = await callCoreFn<{ project: Record }>(db, "core.fn_project_get", { id }, { route }); + if (!projectEnv.ok) return respondRpc(c, projectEnv); + const project = projectEnv.data?.project; + if (!project) { + return respondApiError(c, "NOT_FOUND", `Proyecto id=${id} no encontrado en ${route}`, { route, project_id: id }); + } + const docsEnv = await callCoreFn<{ documents: unknown[] }>( + db, + "core.fn_project_document_list", + { project_id: id }, + { route }, + ); + if (!docsEnv.ok) return respondRpc(c, docsEnv); + return respondRpc(c, { + ok: true, + code: "OK", + layer: "db", + message: `Proyecto ${project.code} cargado`, + data: { + project, + documents: docsEnv.data?.documents ?? [], + checklist: await projectChecklistFor(db, id), + }, }); }); @@ -562,15 +663,14 @@ app.get("/v1/projects/:id/documents/:docId", ...requireCoreAuth, async (c) => { const projectId = Number(c.req.param("id")); const docId = Number(c.req.param("docId")); const db = c.get("db"); - const doc = await db.prepare( - "SELECT * FROM project_documents WHERE id=? AND project_id=?", - ).get(docId, projectId) as { - storage_name: string; - iv: string; - mime: string; - original_name: string; - } | undefined; - if (!doc) return c.json({ error: "Documento no encontrado" }, 404); + const route = routeLabel(c); + const { doc, envelope } = await rpcDocumentForDownload( + db, + "core.fn_project_document_get", + { project_id: projectId, id: docId }, + route, + ); + if (!doc) return respondRpc(c, envelope); const enc = await getObject(projectDocKey(projectId, doc.storage_name)); const plain = await decryptBytes(doc.iv, enc); c.header("Content-Type", "application/octet-stream"); @@ -644,18 +744,24 @@ app.post("/v1/projects/:id/budget/import", ...requireCoreAuth, async (c) => { app.post("/v1/projects/:id/budget/chapters", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const body = await c.req.json<{ name?: string; code?: string; parent_id?: number | null }>(); - if (!body.name?.trim()) return c.json({ error: "Nombre de capítulo obligatorio" }, 400); + if (!body.name?.trim()) { + return respondApiError(c, "VALIDATION", `Nombre de capítulo obligatorio en ${routeLabel(c)}`, { + route: routeLabel(c), + project_id: id, + }); + } const db = c.get("db"); - if (!await projectById(db, id)) return c.json({ error: "Proyecto no encontrado" }, 404); - const sort = (await db.prepare("SELECT COALESCE(MAX(sort_order), 0) + 1 AS n FROM budget_chapters WHERE project_id = ?").get(id) as { n: number }).n; - await db.prepare("INSERT INTO budget_chapters (project_id, parent_id, code, name, sort_order) VALUES (?, ?, ?, ?, ?)").run( - id, - body.parent_id || null, - (body.code || "").trim(), - body.name.trim(), - sort, - ); - return c.json({ id: await lastInsertId(db) }, 201); + return respondRpc(c, await callCoreFn( + db, + "core.fn_budget_chapter_create", + { + project_id: id, + parent_id: body.parent_id || null, + code: (body.code || "").trim(), + name: body.name.trim(), + }, + { route: routeLabel(c) }, + )); }); app.patch("/v1/projects/:id/budget/chapters/:cid", ...requireCoreAuth, async (c) => { @@ -663,25 +769,29 @@ app.patch("/v1/projects/:id/budget/chapters/:cid", ...requireCoreAuth, async (c) const cid = Number(c.req.param("cid")); const body = await c.req.json<{ name?: string; code?: string }>(); const db = c.get("db"); - const cur = await db.prepare("SELECT * FROM budget_chapters WHERE id = ? AND project_id = ?").get(cid, id) as { name: string; code: string } | undefined; - if (!cur) return c.json({ error: "Capítulo no encontrado" }, 404); - await db.prepare("UPDATE budget_chapters SET name = ?, code = ? WHERE id = ?").run( - (body.name ?? cur.name).trim(), - (body.code ?? cur.code).trim(), - cid, - ); - return c.json({ ok: true }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_budget_chapter_update", + { + project_id: id, + id: cid, + name: body.name, + code: body.code, + }, + { route: routeLabel(c) }, + )); }); app.delete("/v1/projects/:id/budget/chapters/:cid", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const cid = Number(c.req.param("cid")); const db = c.get("db"); - const found = await db.prepare("SELECT id FROM budget_chapters WHERE id = ? AND project_id = ?").get(cid, id); - if (!found) return c.json({ error: "Capítulo no encontrado" }, 404); - await db.prepare("UPDATE budget_items SET chapter_id = NULL WHERE project_id = ? AND chapter_id = ?").run(id, cid); - await db.prepare("DELETE FROM budget_chapters WHERE id = ? AND project_id = ?").run(cid, id); - return c.json({ ok: true }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_budget_chapter_delete", + { project_id: id, id: cid }, + { route: routeLabel(c) }, + )); }); type BudgetItemInput = { @@ -696,27 +806,27 @@ type BudgetItemInput = { app.post("/v1/projects/:id/budget/items", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const body = await c.req.json(); - if (!body.description?.trim()) return c.json({ error: "Descripción de partida obligatoria" }, 400); + if (!body.description?.trim()) { + return respondApiError(c, "VALIDATION", `Descripción de partida obligatoria en ${routeLabel(c)}`, { + route: routeLabel(c), + project_id: id, + }); + } const db = c.get("db"); - if (!await projectById(db, id)) return c.json({ error: "Proyecto no encontrado" }, 404); - const qty = Number(body.quantity ?? 0); - const price = Number(body.unit_price ?? 0); - const sort = (await db.prepare("SELECT COALESCE(MAX(sort_order), 0) + 1 AS n FROM budget_items WHERE project_id = ?").get(id) as { n: number }).n; - await db.prepare( - `INSERT INTO budget_items (project_id, chapter_id, code, description, unit, quantity, unit_price, amount, sort_order) - VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, - ).run( - id, - body.chapter_id || null, - (body.code || "").trim().toUpperCase(), - body.description.trim(), - (body.unit || "").trim().toUpperCase(), - qty, - price, - lineAmount(qty, price), - sort, - ); - return c.json({ id: await lastInsertId(db) }, 201); + return respondRpc(c, await callCoreFn( + db, + "core.fn_budget_item_create", + { + project_id: id, + chapter_id: body.chapter_id || null, + code: (body.code || "").trim(), + description: body.description.trim(), + unit: (body.unit || "").trim(), + quantity: body.quantity ?? 0, + unit_price: body.unit_price ?? 0, + }, + { route: routeLabel(c) }, + )); }); app.patch("/v1/projects/:id/budget/items/:iid", ...requireCoreAuth, async (c) => { @@ -724,39 +834,40 @@ app.patch("/v1/projects/:id/budget/items/:iid", ...requireCoreAuth, async (c) => const iid = Number(c.req.param("iid")); const body = await c.req.json(); const db = c.get("db"); - const cur = await db.prepare("SELECT * FROM budget_items WHERE id = ? AND project_id = ?").get(iid, id) as Record | undefined; - if (!cur) return c.json({ error: "Partida no encontrada" }, 404); - const qty = body.quantity !== undefined ? Number(body.quantity) : Number(cur.quantity); - const price = body.unit_price !== undefined ? Number(body.unit_price) : Number(cur.unit_price); - await db.prepare( - `UPDATE budget_items SET chapter_id=?, code=?, description=?, unit=?, quantity=?, unit_price=?, amount=? WHERE id=?`, - ).run( - body.chapter_id !== undefined ? (body.chapter_id || null) : cur.chapter_id, - (body.code ?? String(cur.code ?? "")).trim().toUpperCase(), - (body.description ?? String(cur.description ?? "")).trim(), - (body.unit ?? String(cur.unit ?? "")).trim().toUpperCase(), - qty, - price, - lineAmount(qty, price), - iid, - ); - return c.json({ ok: true }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_budget_item_update", + { + project_id: id, + id: iid, + chapter_id: body.chapter_id, + code: body.code, + description: body.description, + unit: body.unit, + quantity: body.quantity, + unit_price: body.unit_price, + }, + { route: routeLabel(c) }, + )); }); app.delete("/v1/projects/:id/budget/items/:iid", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const iid = Number(c.req.param("iid")); const db = c.get("db"); - const found = await db.prepare("SELECT id FROM budget_items WHERE id = ? AND project_id = ?").get(iid, id); - if (!found) return c.json({ error: "Partida no encontrada" }, 404); - await db.prepare("DELETE FROM budget_items WHERE id = ? AND project_id = ?").run(iid, id); - return c.json({ ok: true }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_budget_item_delete", + { project_id: id, id: iid }, + { route: routeLabel(c) }, + )); }); app.post("/v1/workers/validate", ...requireCoreAuth, async (c) => { const body = await c.req.json(); const errors = validateWorkerFields(body); const db = c.get("db"); + const route = routeLabel(c); const exclude = body.id ?? 0; const checks: { field: string; value: string; err: string | null }[] = [ { field: "curp", value: normUpper(body.curp), err: validateCurp(body.curp ?? "") }, @@ -765,9 +876,16 @@ app.post("/v1/workers/validate", ...requireCoreAuth, async (c) => { ]; for (const ch of checks) { if (ch.err) continue; - const row = await db.prepare( - `SELECT id, first_name, last_name_p FROM workers WHERE ${ch.field} = ? AND id != ?`, - ).get(ch.value, exclude) as { id: number; first_name: string; last_name_p: string } | undefined; + const env = await callCoreFn<{ worker?: { id: number; first_name: string; last_name_p: string } | null; found?: boolean }>( + db, + "core.fn_worker_find_existing", + { + [ch.field]: ch.value, + exclude_id: exclude, + }, + { route }, + ); + const row = env.ok && env.data?.found ? env.data.worker : undefined; if (row) { errors[ch.field] = `Este ${ch.field.toUpperCase()} ya pertenece a ${row.first_name} ${row.last_name_p}`; @@ -782,128 +900,36 @@ app.get("/v1/workers", ...requireCoreAuth, async (c) => { const q = (c.req.query("q") ?? "").trim(); const status = c.req.query("status"); const projectId = c.req.query("project_id"); - let sql = `SELECT w.*, r.label AS risk_label, r.color AS risk_color, r.text_color AS risk_text, - c.name AS company_name, c.kind AS company_kind, - ic.name AS imss_company_name, ic.code AS imss_company_code, ic.registro_patronal AS imss_registro_patronal, - (SELECT STRING_AGG(p.name, ', ') FROM assignments a JOIN projects p ON p.id = a.project_id - WHERE a.worker_id = w.id AND a.active = true) AS proyectos, - (SELECT STRING_AGG(p.code || '|' || replace(p.name, '|', '/'), ';;') FROM assignments a JOIN projects p ON p.id = a.project_id - WHERE a.worker_id = w.id AND a.active = true) AS proyecto_pairs, - (SELECT STRING_AGG(p.id::text, ',') FROM assignments a JOIN projects p ON p.id = a.project_id - WHERE a.worker_id = w.id AND a.active = true) AS project_ids, - (SELECT COUNT(*) FROM document_types t - WHERE t.required = true AND NOT EXISTS ( - SELECT 1 FROM documents d - WHERE d.worker_id = w.id AND d.type_code = t.code AND d.is_current = true - )) AS missing_docs, - (SELECT COALESCE(SUM(l.balance), 0) FROM loans l WHERE l.worker_id = w.id AND l.balance > 0) AS loan_balance, - CASE WHEN EXISTS (SELECT 1 FROM assignments a WHERE a.worker_id = w.id AND a.active = true) THEN true ELSE false END AS in_project - FROM workers w - JOIN risk_levels r ON r.code = w.risk_code - LEFT JOIN companies c ON c.id = w.company_id - LEFT JOIN companies ic ON ic.id = w.imss_company_id - WHERE 1=1`; - const params: (string | number)[] = []; - if (tid != null) { - sql += " AND w.tenant_id = ?"; - params.push(tid); - } - if (status) { - sql += " AND w.status = ?"; - params.push(status); - if (status === "activo") sql += " AND w.pipeline_status != 'baja'"; - } - if (projectId) { - sql += - " AND EXISTS (SELECT 1 FROM assignments a WHERE a.worker_id = w.id AND a.project_id = ? AND a.active = true)"; - params.push(Number(projectId)); - } - if (q) { - sql += - " AND (w.first_name ILIKE ? OR w.last_name_p ILIKE ? OR w.curp ILIKE ? OR w.rfc ILIKE ? OR w.nss ILIKE ?)"; - const like = `%${q}%`; - params.push(like, like, like, like, like); - } - sql += " ORDER BY w.last_name_p, w.first_name"; - const rows = await db.prepare(sql).all(...params) as Array & { - id: number; - status: string; - pipeline_status: string; - imss_status?: string; - in_project?: boolean; - }>; - const workers = []; - for (const row of rows) { - const flags = await imssFlagsFor(db, row.id, tid); - workers.push({ - ...row, - imss_status: flags.imss_status, - imss_ready: flags.imss_ready, - in_project_without_imss: flags.in_project_without_imss, - expediente_ok: flags.expediente_ok, - freshness_required: flags.freshness_required, - }); - } - const active = workers.filter((w) => w.status === "activo" && w.pipeline_status !== "baja"); - const withImss = active.filter((w) => w.imss_status === "alta").length; - const withoutImss = active.length - withImss; - const inProjectWithoutImss = active.filter((w) => w.in_project_without_imss).length; - const imssReady = active.filter((w) => w.imss_ready).length; - return c.json({ - workers, - imss_stats: { - with_imss: withImss, - without_imss: withoutImss, - in_project_without_imss: inProjectWithoutImss, - imss_ready_count: imssReady, + return respondRpc(c, await callCoreFn( + db, + "core.fn_worker_list", + { + tenant_id: tid, + q: q || undefined, + status: status || undefined, + project_id: projectId ? Number(projectId) : undefined, }, - }); + { route: routeLabel(c) }, + )); }); app.get("/v1/workers/:id", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); const tid = tenantScope(c.get("user")); - const worker = await db.prepare( - `SELECT w.*, r.label AS risk_label, r.color AS risk_color, r.text_color AS risk_text, - c.name AS company_name, c.kind AS company_kind, - ic.name AS imss_company_name, ic.code AS imss_company_code, ic.registro_patronal AS imss_registro_patronal - FROM workers w - JOIN risk_levels r ON r.code = w.risk_code - LEFT JOIN companies c ON c.id = w.company_id - LEFT JOIN companies ic ON ic.id = w.imss_company_id - WHERE w.id = ?`, - ).get(id); - if (!worker) return c.json({ error: "No encontrado" }, 404); - const documents = await db.prepare( - `SELECT id, type_code, original_name, mime, size_bytes, is_current, uploaded_at, - issued_at, expires_at, imss_company_id, imss_alta_at - FROM documents WHERE worker_id = ? ORDER BY uploaded_at DESC`, - ).all(id); - const assignments = await db.prepare( - `SELECT a.*, p.name AS project_name, p.code AS project_code FROM assignments a - JOIN projects p ON p.id = a.project_id WHERE a.worker_id = ? - ORDER BY a.active DESC, a.start_date DESC, a.id DESC`, - ).all(id); - const loans = await db.prepare("SELECT * FROM loans WHERE worker_id = ? ORDER BY id DESC").all(id); - const { items, freshness_required } = await checklistItemsFor(db, id, tid); - const flags = await imssFlagsFor(db, id, tid); - const docTypes = await db.prepare( - `SELECT code, label, required, validity_mode, freshness_days, requires_issued_at, requires_expires_at, category - FROM document_types ORDER BY required DESC, label`, - ).all(); - return c.json({ - worker: { ...worker, full_name: fullName(worker as never), ...flags }, - documents, - assignments, - loans, - checklist: items, - document_types: docTypes, - freshness_required, - }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_worker_get", + { id, tenant_id: tid }, + { route: routeLabel(c) }, + )); }); -function conflict(row: { id: number; first_name: string; last_name_p: string; curp: string; rfc: string; nss: string } | undefined, n: ReturnType, excludeId = 0) { +function conflict( + row: { id: number; first_name: string; last_name_p: string; curp: string; rfc: string; nss: string } | undefined, + n: ReturnType, + excludeId = 0, +) { if (row && row.id !== excludeId) { const field = row.curp === n.curp ? "CURP" : row.rfc === n.rfc ? "RFC" : "NSS"; return { @@ -917,97 +943,111 @@ function conflict(row: { id: number; first_name: string; last_name_p: string; cu return null; } +async function findExistingWorker( + db: Db, + curp: string, + rfc: string, + nss: string, + excludeId = 0, + route = "worker-duplicate-check", +) { + const env = await callCoreFn<{ + worker?: { id: number; first_name: string; last_name_p: string; curp: string; rfc: string; nss: string } | null; + found?: boolean; + }>(db, "core.fn_worker_find_existing", { curp, rfc, nss, exclude_id: excludeId }, { route }); + if (!env.ok || !env.data?.found) return undefined; + return env.data.worker; +} + app.post("/v1/workers", ...requireCoreAuth, async (c) => { const body = await c.req.json(); const errors = validateWorkerFields(body); const db = c.get("db"); - if (Object.keys(errors).length) return c.json({ errors }, 400); + const route = routeLabel(c); + if (Object.keys(errors).length) { + return respondApiError(c, "VALIDATION", `Datos del trabajador inválidos en ${route}`, { route }, errors); + } const n = normalizeWorker({ ...body, hire_type: "" }); - const existing = await findExisting(db, n.curp, n.rfc, n.nss); + const existing = await findExistingWorker(db, n.curp, n.rfc, n.nss, 0, route); const cf = conflict(existing, n); - if (cf) return c.json(cf.body, cf.status); + if (cf) { + return respondApiError(c, "CONFLICT", cf.body.error, { route, worker_id: cf.body.worker_id }); + } if (body.project_id) { const blocked = projectMustBe( await projectById(db, body.project_id), ["activo"], "Solo se asigna personal a proyectos activos", ); - if (blocked) return c.json({ error: blocked.error }, blocked.status); + if (blocked) { + const code = blocked.status === 404 ? "NOT_FOUND" : "VALIDATION"; + return respondApiError(c, code, blocked.error, { route, project_id: body.project_id }); + } } const tid = tenantScope(c.get("user")); - await db.prepare( - `INSERT INTO workers - (first_name, middle_name, last_name_p, last_name_m, curp, rfc, nss, phone, email, address, - blood_type, hire_type, company_id, position, risk_code, work_type, daily_wage, needs_badge, status, tenant_id) - VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)`, - ).run( - n.first_name, n.middle_name, n.last_name_p, n.last_name_m, n.curp, n.rfc, n.nss, - n.phone, n.email, n.address, n.blood_type, "", null, n.position, n.risk_code, - n.work_type, n.daily_wage, n.needs_badge, n.status, tid, - ); - const id = await lastInsertId(db); - if (body.project_id) { - await assign(db, id, body.project_id); - } - await refreshPipeline(db, id, tid); - return c.json({ id }, 201); + return respondRpc(c, await callCoreFn( + db, + "core.fn_worker_create", + { ...body, tenant_id: tid, project_id: body.project_id ?? null }, + { route }, + )); }); app.patch("/v1/workers/:id", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); const tid = tenantScope(c.get("user")); - const cur = await db.prepare("SELECT * FROM workers WHERE id = ?").get(id) as Record | undefined; - if (!cur) return c.json({ error: "No encontrado" }, 404); + const route = routeLabel(c); + const curEnv = await callCoreFn<{ worker: Record }>(db, "core.fn_worker_get", { id, tenant_id: tid }, { route }); + if (!curEnv.ok) return respondRpc(c, curEnv); + const cur = curEnv.data?.worker; + if (!cur) { + return respondApiError(c, "NOT_FOUND", `Trabajador id=${id} no encontrado en ${route}`, { route, worker_id: id }); + } const body = await c.req.json(); const merged = { ...cur, ...body } as WorkerInput; const errors = validateWorkerFields(merged); - if (Object.keys(errors).length) return c.json({ errors }, 400); - // Empresa/patrón solo cambia con alta/baja IMSS; no se sobrescribe desde el formulario. - const hireType = String(cur.hire_type || ""); - const companyId = (cur.company_id as number | null) ?? null; - const n = normalizeWorker({ ...merged, hire_type: hireType }); - const existing = await findExisting(db, n.curp, n.rfc, n.nss); + if (Object.keys(errors).length) { + return respondApiError(c, "VALIDATION", `Datos del trabajador inválidos en ${route}`, { route, worker_id: id }, errors); + } + const n = normalizeWorker({ ...merged, hire_type: String(cur.hire_type || "") }); + const existing = await findExistingWorker(db, n.curp, n.rfc, n.nss, id, route); const cf = conflict(existing, n, id); - if (cf) return c.json(cf.body, cf.status); - await db.prepare( - `UPDATE workers SET - first_name=?, middle_name=?, last_name_p=?, last_name_m=?, curp=?, rfc=?, nss=?, - phone=?, email=?, address=?, blood_type=?, hire_type=?, company_id=?, position=?, risk_code=?, - work_type=?, daily_wage=?, needs_badge=?, status=?, updated_at=now() - WHERE id=?`, - ).run( - n.first_name, n.middle_name, n.last_name_p, n.last_name_m, n.curp, n.rfc, n.nss, - n.phone, n.email, n.address, n.blood_type, hireType, companyId, n.position, n.risk_code, - n.work_type, n.daily_wage, n.needs_badge, n.status, id, - ); - await refreshPipeline(db, id, tid); - return c.json({ ok: true }); + if (cf) { + return respondApiError(c, "CONFLICT", cf.body.error, { route, worker_id: cf.body.worker_id }); + } + return respondRpc(c, await callCoreFn( + db, + "core.fn_worker_update", + { id, tenant_id: tid, ...body }, + { route }, + )); }); app.patch("/v1/workers/:id/pipeline", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const { pipeline_status } = await c.req.json<{ pipeline_status: string }>(); const allowed = ["incompleto", "listo_gafete", "impreso", "activo", "baja"]; - if (!allowed.includes(pipeline_status)) return c.json({ error: "Estado inválido" }, 400); + if (!allowed.includes(pipeline_status)) { + return respondApiError(c, "VALIDATION", `Estado de pipeline inválido en ${routeLabel(c)}: ${pipeline_status}`, { + route: routeLabel(c), + worker_id: id, + pipeline_status, + }); + } const db = c.get("db"); const tid = tenantScope(c.get("user")); - const prev = await db.prepare("SELECT status FROM workers WHERE id=?").get(id) as { status: string } | undefined; - if (pipeline_status === "baja") { - await db.prepare("UPDATE workers SET status='baja', pipeline_status='baja' WHERE id=?").run(id); - } else { - // Reactivar: si venía de baja, marca rehire para exigir docs frescos hasta nueva alta IMSS. - if (prev?.status === "baja") { - await db.prepare( - "UPDATE workers SET status='activo', last_rehire_at=current_date WHERE id=?", - ).run(id); - } else { - await db.prepare("UPDATE workers SET status='activo' WHERE id=?").run(id); - } - await refreshPipeline(db, id, tid); - } - const worker = await db.prepare("SELECT status, pipeline_status, last_rehire_at, imss_status FROM workers WHERE id=?").get(id); - return c.json({ ok: true, worker }); + const env = await callCoreFn( + db, + "core.fn_worker_set_pipeline", + { id, tenant_id: tid, pipeline_status }, + { route: routeLabel(c) }, + ); + if (!env.ok) return respondRpc(c, env); + return respondRpc(c, { + ...env, + data: { ok: true, ...(env.data as Record) }, + }); }); app.post("/v1/workers/:id/assign", ...requireCoreAuth, async (c) => { @@ -1015,20 +1055,35 @@ app.post("/v1/workers/:id/assign", ...requireCoreAuth, async (c) => { const { project_id, active } = await c.req.json<{ project_id: number; active?: boolean }>(); const db = c.get("db"); const tid = tenantScope(c.get("user")); + const route = routeLabel(c); if (active === false) { - await db.prepare( - "UPDATE assignments SET active=false, end_date=current_date WHERE worker_id=? AND project_id=?", - ).run(id, project_id); - } else { - const project = await projectById(db, project_id); - if (!project) return c.json({ error: "Proyecto no encontrado" }, 404); - if (project.status !== "activo") { - return c.json({ error: "Solo se asigna personal a proyectos activos" }, 400); - } - await assign(db, id, project_id); + return respondRpc(c, await callCoreFn( + db, + "core.fn_worker_unassign", + { worker_id: id, project_id, tenant_id: tid }, + { route }, + )); } - await refreshPipeline(db, id, tid); - return c.json({ ok: true }); + const project = await projectById(db, project_id); + if (!project) { + return respondApiError(c, "NOT_FOUND", `Proyecto id=${project_id} no encontrado en ${route}`, { + route, + project_id, + }); + } + if (project.status !== "activo") { + return respondApiError(c, "VALIDATION", `Solo se asigna personal a proyectos activos en ${route}`, { + route, + project_id, + status: project.status, + }); + } + return respondRpc(c, await callCoreFn( + db, + "core.fn_worker_assign", + { worker_id: id, project_id, tenant_id: tid }, + { route }, + )); }); app.get("/v1/workers/import/template", ...requireCoreAuth, async (c) => { @@ -1063,19 +1118,36 @@ app.post("/v1/workers/:id/documents", ...requireCoreAuth, async (c) => { const form = await c.req.formData(); const type = String(form.get("type") || ""); const file = form.get("file"); - if (!(file instanceof File) || !type) return c.json({ error: "type y file requeridos" }, 400); + if (!(file instanceof File) || !type) { + return respondApiError(c, "VALIDATION", `type y file son requeridos en ${routeLabel(c)}`, { + route: routeLabel(c), + worker_id: id, + }); + } const db = c.get("db"); - const exists = await db.prepare("SELECT id FROM workers WHERE id=?").get(id); - if (!exists) return c.json({ error: "No encontrado" }, 404); + const tid = tenantScope(c.get("user")); + const route = routeLabel(c); + const workerEnv = await callCoreFn(db, "core.fn_worker_get", { id, tenant_id: tid }, { route }); + if (!workerEnv.ok) return respondRpc(c, workerEnv); - const policy = await db.prepare( - `SELECT validity_mode, requires_issued_at, requires_expires_at FROM document_types WHERE code = ?`, - ).get(type) as { + const catEnv = await callCoreFn<{ document_types?: Array> }>( + db, + "core.fn_catalogs", + {}, + { route }, + ); + const policy = (catEnv.data?.document_types ?? []).find((t) => String(t.code) === type) as { validity_mode: string; requires_issued_at: boolean; requires_expires_at: boolean; } | undefined; - if (!policy) return c.json({ error: "Tipo de documento no válido" }, 400); + if (!policy) { + return respondApiError(c, "VALIDATION", `Tipo de documento no válido en ${route}: ${type}`, { + route, + worker_id: id, + type, + }); + } const issuedAt = String(form.get("issued_at") || "").trim() || null; const expiresAt = String(form.get("expires_at") || "").trim() || null; @@ -1084,19 +1156,19 @@ app.post("/v1/workers/:id/documents", ...requireCoreAuth, async (c) => { const imssBajaAt = String(form.get("imss_baja_at") || "").trim() || null; if (policy.requires_issued_at && !issuedAt) { - return c.json({ error: "Indique la fecha de emisión del documento" }, 400); + return respondApiError(c, "VALIDATION", `Indique la fecha de emisión del documento en ${route}`, { route, worker_id: id }); } if (policy.requires_expires_at && !expiresAt) { - return c.json({ error: "Indique la fecha de vigencia / vencimiento" }, 400); + return respondApiError(c, "VALIDATION", `Indique la fecha de vigencia / vencimiento en ${route}`, { route, worker_id: id }); } if (type === "alta_imss" && !imssCompanyId) { - return c.json({ error: "Seleccione la empresa patrón del alta IMSS" }, 400); + return respondApiError(c, "VALIDATION", `Seleccione la empresa patrón del alta IMSS en ${route}`, { route, worker_id: id }); } if (type === "alta_imss" && !imssAltaAt) { - return c.json({ error: "Indique la fecha de alta IMSS" }, 400); + return respondApiError(c, "VALIDATION", `Indique la fecha de alta IMSS en ${route}`, { route, worker_id: id }); } if (type === "baja_imss" && !imssBajaAt) { - return c.json({ error: "Indique la fecha de baja IMSS" }, 400); + return respondApiError(c, "VALIDATION", `Indique la fecha de baja IMSS en ${route}`, { route, worker_id: id }); } const bytes = new Uint8Array(await file.arrayBuffer()); @@ -1109,9 +1181,9 @@ app.post("/v1/workers/:id/documents", ...requireCoreAuth, async (c) => { imss_baja_at: imssBajaAt, }); } catch (error) { - return c.json({ error: error instanceof Error ? error.message : "No se pudo guardar" }, 400); + const message = error instanceof Error ? error.message : "No se pudo guardar el documento"; + return respondApiError(c, "VALIDATION", message, { route, worker_id: id }); } - const tid = tenantScope(c.get("user")); return c.json({ ok: true, checklist: await checklistFor(db, id, tid), @@ -1123,15 +1195,14 @@ app.get("/v1/workers/:id/documents/:docId", ...requireCoreAuth, async (c) => { const workerId = Number(c.req.param("id")); const docId = Number(c.req.param("docId")); const db = c.get("db"); - const doc = await db.prepare( - "SELECT * FROM documents WHERE id=? AND worker_id=?", - ).get(docId, workerId) as { - storage_name: string; - iv: string; - mime: string; - original_name: string; - } | undefined; - if (!doc) return c.json({ error: "Documento no encontrado" }, 404); + const route = routeLabel(c); + const { doc, envelope } = await rpcDocumentForDownload( + db, + "core.fn_worker_document_get", + { worker_id: workerId, id: docId }, + route, + ); + if (!doc) return respondRpc(c, envelope); const enc = await getObject(workerDocKey(workerId, doc.storage_name)); const plain = await decryptBytes(doc.iv, enc); c.header("Content-Type", "application/octet-stream"); @@ -1169,68 +1240,85 @@ app.post("/v1/projects/:id/badge-jobs", ...requireCoreAuth, async (c) => { const projectId = Number(c.req.param("id")); const body = await c.req.json<{ worker_ids?: number[] }>().catch(() => ({ worker_ids: [] as number[] })); const db = c.get("db"); + const user = c.get("user"); + const tid = tenantScope(user); + const route = routeLabel(c); const blocked = projectMustBe( await projectById(db, projectId), ["activo"], "Solo se generan gafetes de proyectos activos", ); - if (blocked) return c.json({ error: blocked.error }, blocked.status); + if (blocked) { + const code = blocked.status === 404 ? "NOT_FOUND" : "VALIDATION"; + return respondApiError(c, code, blocked.error, { route, project_id: projectId }); + } let ids = body.worker_ids ?? []; if (!ids.length) { - ids = (await db.prepare( - `SELECT w.id FROM workers w - JOIN assignments a ON a.worker_id = w.id AND a.project_id = ? AND a.active = true - WHERE w.status = 'activo' AND w.needs_badge = true - AND EXISTS (SELECT 1 FROM documents d WHERE d.worker_id=w.id AND d.type_code='foto' AND d.is_current=true)`, - ).all(projectId) as { id: number }[]).map((r) => r.id); + const listEnv = await callCoreFn<{ workers?: Array> }>( + db, + "core.fn_worker_list", + { tenant_id: tid, project_id: projectId, status: "activo" }, + { route }, + ); + if (!listEnv.ok) return respondRpc(c, listEnv); + ids = (listEnv.data?.workers ?? []) + .filter((w) => w.needs_badge && ["listo_gafete", "impreso", "activo"].includes(String(w.pipeline_status))) + .map((w) => Number(w.id)); + } + if (!ids.length) { + return respondApiError(c, "VALIDATION", `No hay personal activo con foto para imprimir en ${route}`, { + route, + project_id: projectId, + }); } - if (!ids.length) return c.json({ error: "No hay personal activo con foto para imprimir" }, 400); const bytes = await generateBadgePdf(db, projectId, ids); - const user = c.get("user"); - await db.prepare( - "INSERT INTO badge_jobs (project_id, status, created_by_id, created_by_name) VALUES (?, 'done', ?, ?)", - ).run(projectId, user.id, user.display_name); - const jobId = await lastInsertId(db); - const path = await saveJobPdf(bytes, jobId); - await db.prepare("UPDATE badge_jobs SET pdf_path=? WHERE id=?").run(path, jobId); - const ins = db.prepare("INSERT INTO badge_job_people (job_id, worker_id) VALUES (?, ?)"); - const tid = tenantScope(user); + const createEnv = await callCoreFn( + db, + "core.fn_badge_job_create", + { + project_id: projectId, + worker_ids: ids, + created_by_id: user.id, + created_by_name: user.display_name, + }, + { route }, + ); + if (!createEnv.ok) return respondRpc(c, createEnv); + const jobId = Number((createEnv.data as { id?: number })?.id); + await saveJobPdf(bytes, jobId); for (const wid of ids) { - await ins.run(jobId, wid); await refreshPipeline(db, wid, tid); } - return c.json({ id: jobId, count: ids.length }); + return respondRpc(c, { + ...createEnv, + data: { ...(createEnv.data as Record), count: ids.length }, + }); }); app.get("/v1/badge-jobs", ...requireCoreAuth, async (c) => { const db = c.get("db"); - const jobs = await db.prepare( - `SELECT j.*, p.name AS project_name, - (SELECT COUNT(*) FROM badge_job_people x WHERE x.job_id=j.id) AS people, - (SELECT COUNT(*) FROM badge_job_people x WHERE x.job_id=j.id AND x.delivered=true) AS delivered - FROM badge_jobs j JOIN projects p ON p.id=j.project_id ORDER BY j.id DESC`, - ).all(); - return c.json({ jobs }); + return respondRpc(c, await callCoreFn(db, "core.fn_badge_job_list", {}, { route: routeLabel(c) })); }); app.get("/v1/badge-jobs/:id", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); - const job = await db.prepare("SELECT * FROM badge_jobs WHERE id=?").get(id); - if (!job) return c.json({ error: "Job no encontrado" }, 404); - const people = await db.prepare( - `SELECT p.*, w.first_name, w.last_name_p, w.position FROM badge_job_people p - JOIN workers w ON w.id = p.worker_id WHERE p.job_id=?`, - ).all(id); - return c.json({ job, people }); + return respondRpc(c, await callCoreFn(db, "core.fn_badge_job_get", { job_id: id }, { route: routeLabel(c) })); }); app.get("/v1/badge-jobs/:id/pdf", ...requireCoreAuth, async (c) => { const id = Number(c.req.param("id")); const db = c.get("db"); - const job = await db.prepare("SELECT pdf_path FROM badge_jobs WHERE id=?").get(id) as { pdf_path: string } | undefined; - if (!job?.pdf_path) return c.json({ error: "PDF no encontrado" }, 404); - const bytes = await getObject(job.pdf_path); + const route = routeLabel(c); + const jobEnv = await callCoreFn<{ job?: { pdf_path?: string | null } }>( + db, + "core.fn_badge_job_get", + { job_id: id }, + { route }, + ); + if (!jobEnv.ok) return respondRpc(c, jobEnv); + const pdfPath = jobEnv.data?.job?.pdf_path || badgeJobPdfKey(id); + const bytes = await getObject(pdfPath); c.header("Content-Type", "application/pdf"); c.header("Content-Disposition", `attachment; filename="gafetes-${id}.pdf"`); return c.body(bytes.buffer as ArrayBuffer); @@ -1241,11 +1329,12 @@ app.patch("/v1/badge-jobs/:id/people/:workerId", ...requireCoreAuth, async (c) = const workerId = Number(c.req.param("workerId")); const { delivered } = await c.req.json<{ delivered: boolean }>(); const db = c.get("db"); - await db.prepare( - `UPDATE badge_job_people SET delivered=?, delivered_at=CASE WHEN ? THEN now() ELSE NULL END - WHERE job_id=? AND worker_id=?`, - ).run(delivered, delivered, jobId, workerId); - return c.json({ ok: true }); + return respondRpc(c, await callCoreFn( + db, + "core.fn_badge_job_set_delivered", + { job_id: jobId, worker_id: workerId, delivered }, + { route: routeLabel(c) }, + )); }); registerPayrollRoutes(app);