Files
PANEL_BASES_ANEXO24/src/lib/server/controldesk-pg.ts
AlexeerCT 26eed5726c
Some checks failed
Aduanasoft/PANEL_BASES_ANEXO24/pipeline/head There was a failure building this commit
feat(alerts): enhance alert handling and client data enrichment in dashboard (#16)
- Updated the alert handling logic to include a 'notFound' flag for databases that are not found on the server, allowing for clearer UI representation.
- Enriched alert data with client contact information by integrating the `lookupAlertClientData` function, ensuring alerts display relevant client details.
- Improved the display of alert dates in the UI to reflect the new 'notFound' status, enhancing user experience and clarity.

This update improves the accuracy and usability of the dashboard alerts.

Reviewed-on: #16
Co-authored-by: AlexeerCT <acazares@aduanasoft.com.mx>
Co-committed-by: AlexeerCT <acazares@aduanasoft.com.mx>
2026-07-13 15:27:14 +00:00

1377 lines
49 KiB
TypeScript

/**
* Catálogo ControlDesk en PostgreSQL (esquema a24c; DDL lo aprovisiona otra aplicación).
* Devuelve columnas con alias en español para compatibilidad con la UI existente.
*/
import path from 'node:path';
import { env } from '$env/dynamic/private';
import { pgPool } from './db';
import { encryptSecret, decryptSecret } from './crypto';
function schemaName(): string {
const s = 'a24c';
if (!/^[a-zA-Z_][a-zA-Z0-9_]*$/.test(s)) return 'a24c';
return s;
}
function qNodes(): string {
const s = schemaName();
return `"${s.replace(/"/g, '""')}"."database_nodes"`;
}
function qUsers(): string {
const s = schemaName();
return `"${s.replace(/"/g, '""')}"."portal_users"`;
}
function qRestoreTargets(): string {
const s = schemaName();
return `"${s.replace(/"/g, '""')}"."restore_targets"`;
}
function qRestoreJobLogs(): string {
const s = schemaName();
return `"${s.replace(/"/g, '""')}"."restore_job_logs"`;
}
function qCloudRestoreStatus(): string {
const s = schemaName();
return `"${s.replace(/"/g, '""')}"."cloudrestore_status"`;
}
function qNodeLastRestore(): string {
const s = schemaName();
return `"${s.replace(/"/g, '""')}"."node_last_restore"`;
}
function isPgUndefinedTable(err: unknown): boolean {
return typeof err === 'object' && err !== null && (err as { code?: string }).code === '42P01';
}
function isPgUndefinedColumn(err: unknown): boolean {
return typeof err === 'object' && err !== null && (err as { code?: string }).code === '42703';
}
/** Crea cloudrestore_status si la BD existía antes de la migración 002 (idempotente). */
async function ensureCloudRestoreStatusTable(): Promise<void> {
await pgPool.query('CREATE SCHEMA IF NOT EXISTS a24c');
await pgPool.query(`
CREATE TABLE IF NOT EXISTS ${qCloudRestoreStatus()} (
id SERIAL PRIMARY KEY,
instance_key VARCHAR(120) NOT NULL UNIQUE DEFAULT 'default',
input_folder VARCHAR(500) NOT NULL,
host_name VARCHAR(255),
app_version VARCHAR(50),
reported_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
`);
// Carpeta de procesados reportada por el restaurador (opcional). Si no la reporta, el
// panel la deriva como hermana de input_folder (deriveSiblingFolder).
await pgPool.query(
`ALTER TABLE ${qCloudRestoreStatus()} ADD COLUMN IF NOT EXISTS processed_folder VARCHAR(500)`
);
}
/** Añade columnas de tamaño/ruta relativa a restore_job_logs si faltan (idempotente). */
async function ensureRestoreJobLogColumns(): Promise<void> {
for (const stmt of [
`ALTER TABLE ${qRestoreJobLogs()} ADD COLUMN IF NOT EXISTS size_bytes BIGINT`,
`ALTER TABLE ${qRestoreJobLogs()} ADD COLUMN IF NOT EXISTS rel_path VARCHAR(600)`
]) {
try {
await pgPool.query(stmt);
} catch (e) {
if (!isPgUndefinedTable(e)) throw e; // la tabla la crea a24c; si no existe aún, se ignora
}
}
}
/**
* Estado "último respaldo por nodo": una fila por nodo, con el restaurador donde vive el
* archivo. Llaveada por nodo (NO por la asignación actual) para que el tracking del último
* respaldo persista aunque el nodo se reasigne a otro restaurador. Idempotente.
*/
async function ensureNodeLastRestoreTable(): Promise<void> {
await pgPool.query('CREATE SCHEMA IF NOT EXISTS a24c');
await pgPool.query(`
CREATE TABLE IF NOT EXISTS ${qNodeLastRestore()} (
database_node_id INTEGER PRIMARY KEY,
restore_target_id INTEGER,
db_name VARCHAR(255),
node_key VARCHAR(255),
filename VARCHAR(500) NOT NULL,
rel_path VARCHAR(600),
size_bytes BIGINT,
restored_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
`);
await pgPool.query(
`CREATE INDEX IF NOT EXISTS idx_a24c_node_last_restore_target
ON ${qNodeLastRestore()} (restore_target_id)`
);
}
/** Crea restore_targets y filas Alfa/Omega/Gamma si faltan (migración 001, idempotente). */
async function ensureRestoreTargetsSchema(): Promise<void> {
await pgPool.query('CREATE SCHEMA IF NOT EXISTS a24c');
await pgPool.query(`
CREATE TABLE IF NOT EXISTS ${qRestoreTargets()} (
id SERIAL PRIMARY KEY,
name VARCHAR(120) NOT NULL UNIQUE,
server_ip VARCHAR(255),
sql_username VARCHAR(128),
sql_password_encrypted TEXT,
data_folder VARCHAR(500),
ssh_host VARCHAR(255),
ssh_port INTEGER DEFAULT 22,
ssh_username VARCHAR(128),
ssh_password_encrypted TEXT,
remote_inbox_path VARCHAR(500),
notes TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
)
`);
await pgPool.query(`
INSERT INTO ${qRestoreTargets()} (name) VALUES ('Alfa'), ('Omega'), ('Gamma')
ON CONFLICT (name) DO NOTHING
`);
for (const stmt of [
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS ssh_host VARCHAR(255)`,
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS ssh_port INTEGER DEFAULT 22`,
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS ssh_username VARCHAR(128)`,
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS ssh_password_encrypted TEXT`,
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS remote_inbox_path VARCHAR(500)`,
// Características de hardware opcionales (para distribuir bases por capacidad).
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS os VARCHAR(50)`,
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS ram_gb INTEGER`,
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS disk_gb INTEGER`,
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS location VARCHAR(255)`,
`ALTER TABLE ${qRestoreTargets()} ADD COLUMN IF NOT EXISTS size_category VARCHAR(20)`
]) {
await pgPool.query(stmt);
}
try {
await pgPool.query(`
ALTER TABLE ${qNodes()}
ADD COLUMN IF NOT EXISTS restore_target_id INTEGER
REFERENCES ${qRestoreTargets()} (id) ON DELETE SET NULL
`);
} catch {
/* database_nodes puede no existir aún en algunos entornos */
}
}
const ROW_DATABASE_NODE = `
id AS "ID",
node_subnode_key AS "NodoSubNodo",
is_active AS "Activo",
rfc AS "RFC",
legal_name AS "Nombre",
branch_name AS "Sucursal",
notification_email AS "CorreoNotificacion",
server_name AS "ServerName",
database_name AS "BDName",
restore_target_id AS "RestoreTargetId",
anexo24c_aviso_fecha AS "Anexo24CAvisoFecha"
`;
const ROW_PORTAL_USER = `
id AS "ID",
database_node_id AS "IDNodoSubNodo",
is_authority_client AS "ClienteAutoridad",
full_name AS "Nombre",
username AS "Usuario",
bd_shelter AS "BD_Shelter"
`;
export async function listDatabaseNodes(): Promise<any[]> {
const sql = `SELECT ${ROW_DATABASE_NODE} FROM ${qNodes()}`;
const r = await pgPool.query(sql);
return r.rows;
}
export async function listDatabaseNodesActive(): Promise<any[]> {
const sql = `
SELECT ${ROW_DATABASE_NODE}
FROM ${qNodes()}
WHERE is_active = 1
AND database_name IS NOT NULL
AND TRIM(database_name) <> ''
ORDER BY node_subnode_key
`;
const r = await pgPool.query(sql);
return r.rows;
}
/** Nodos activos con `sql_password` para conectar a SQL Server (no exponer al cliente). */
export async function listDatabaseNodesForMssql(): Promise<any[]> {
const where = `
WHERE is_active = 1
AND database_name IS NOT NULL
AND TRIM(database_name) <> ''
ORDER BY node_subnode_key
`;
try {
const r = await pgPool.query(
`SELECT ${ROW_DATABASE_NODE}, sql_password AS "sql_password" FROM ${qNodes()} ${where}`
);
return r.rows;
} catch (e: any) {
if (e?.code === '42703') {
const r = await pgPool.query(`SELECT ${ROW_DATABASE_NODE} FROM ${qNodes()} ${where}`);
return r.rows.map((row: any) => ({ ...row, sql_password: null }));
}
throw e;
}
}
export async function listClientsCatalog(): Promise<any[]> {
const sql = `
SELECT
id AS "ID",
legal_name AS "Nombre",
node_subnode_key AS "NodoSubNodo",
notification_email AS "CorreoNotificacion",
is_active AS "Activo",
database_name AS "BDName"
FROM ${qNodes()}
`;
const r = await pgPool.query(sql);
return r.rows;
}
export async function listPortalUsers(): Promise<any[]> {
const sql = `SELECT ${ROW_PORTAL_USER} FROM ${qUsers()} ORDER BY id`;
const r = await pgPool.query(sql);
return r.rows;
}
/**
* Usuarios del portal (cliente/autoridad) con los datos de su nodo, para el
* reporte de asesores. NO incluye contraseñas: `portal_users.password_hash`
* es un hash bcrypt irreversible y nunca se expone.
*/
export async function listPortalUsersWithNode(): Promise<any[]> {
const sql = `
SELECT
pu.id AS "ID",
pu.is_authority_client AS "ClienteAutoridad",
pu.full_name AS "Nombre",
pu.username AS "Usuario",
pu.bd_shelter AS "BD_Shelter",
dn.node_subnode_key AS "NodoSubNodo",
dn.legal_name AS "Cliente",
dn.rfc AS "RFC",
dn.database_name AS "BDName",
dn.is_active AS "NodoActivo"
FROM ${qUsers()} pu
LEFT JOIN ${qNodes()} dn ON pu.database_node_id = dn.id
ORDER BY dn.node_subnode_key, pu.is_authority_client, pu.username
`;
const r = await pgPool.query(sql);
return r.rows;
}
export async function lookupNodeByNodoOrBdName(nodoName: string): Promise<any | null> {
const sql = `
SELECT
legal_name AS "Nombre",
node_subnode_key AS "NodoSubNodo",
rfc AS "RFC"
FROM ${qNodes()}
WHERE LOWER(TRIM(node_subnode_key)) = LOWER(TRIM($1::text))
OR LOWER(TRIM(database_name)) = LOWER(TRIM($1::text))
LIMIT 1
`;
const r = await pgPool.query(sql, [nodoName]);
return r.rows[0] ?? null;
}
/**
* Asocia el nombre de archivo de respaldo (sin extensión) a una fila de `database_nodes`.
* Cubre:
* - igualdad sin mayúsculas;
* - el primer segmento separado por punto como token de nodo
* (ej. `08037natm001.KNOWNWORLD.000` → nodo `08037NATM001`; el sufijo
* `.KNOWNWORLD` = nombre lógico de la BD origen y `.000` = secuencia se ignoran);
* - prefijos con separadores `.`, `_` o `-` (ej. `NODO001_full_20240414` → `NODO001`).
*/
export function matchNodeRowFromBackupStem(stem: string, nodes: any[]): any | null {
const raw = String(stem ?? '').trim();
if (!raw || !nodes?.length) return null;
const key = raw.toLowerCase();
// Candidatos para igualdad exacta: el stem completo y su primer segmento (token de nodo).
const head = key.split('.')[0];
const exactKeys = head && head !== key ? [key, head] : [key];
const nodoOf = (row: any) => String(row.NodoSubNodo ?? '').trim().toLowerCase();
const bdOf = (row: any) => String(row.BDName ?? '').trim().toLowerCase();
for (const row of nodes) {
const n = nodoOf(row);
const b = bdOf(row);
for (const k of exactKeys) {
if (n && n === k) return row;
if (b && b === k) return row;
}
}
// Coincidencia por prefijo: el nodo/BD seguido de un separador (`.`, `_` o `-`).
const startsWithToken = (token: string) =>
key.startsWith(`${token}.`) || key.startsWith(`${token}_`) || key.startsWith(`${token}-`);
type Cand = { row: any; len: number };
const cands: Cand[] = [];
for (const row of nodes) {
const n = nodoOf(row);
const b = bdOf(row);
if (n && (key === n || startsWithToken(n))) {
cands.push({ row, len: n.length });
}
if (b && b !== n && (key === b || startsWithToken(b))) {
cands.push({ row, len: b.length });
}
}
if (!cands.length) return null;
// Ganar el match más específico (token más largo) para no confundir nodos con prefijo común.
cands.sort((a, b) => b.len - a.len);
return cands[0].row;
}
/** Datos de contacto para alertas (equivalente a la consulta previa sobre Usuarios/BasesDeDatos). */
export async function lookupAlertClientData(nodoName: string): Promise<any | null> {
// El "Cliente" de la alerta es el nombre de la tabla de bases de datos
// (database_nodes.legal_name), NO el full_name del usuario del portal. Se resuelve el nodo por
// database_name o node_subnode_key; así también funciona para bases sin usuario asociado.
const sql = `
SELECT
dn.legal_name AS "Nombre",
dn.notification_email AS "CorreoNotificacion",
dn.node_subnode_key AS "NodoSubNodo"
FROM ${qNodes()} dn
WHERE LOWER(TRIM(dn.database_name)) = LOWER(TRIM($1::text))
OR LOWER(TRIM(dn.node_subnode_key)) = LOWER(TRIM($1::text))
LIMIT 1
`;
const r = await pgPool.query(sql, [nodoName]);
return r.rows[0] ?? null;
}
export async function updateNodeActive(id: number, activo: boolean): Promise<void> {
await pgPool.query(`UPDATE ${qNodes()} SET is_active = $1 WHERE id = $2`, [activo ? 1 : 0, id]);
}
export async function updateNodeLegalName(id: number, nombre: string): Promise<void> {
await pgPool.query(`UPDATE ${qNodes()} SET legal_name = $1 WHERE id = $2`, [nombre, id]);
}
export async function insertDatabaseNode(row: {
nodoSubNodo: string;
rfc: string;
nombre: string;
sucursal: string;
correo: string;
serverName: string;
bdName: string;
activo: number;
restoreTargetId?: number | null;
anexo24CAvisoFecha?: string | null;
}): Promise<void> {
// server_name y sql_password se derivan del servidor de restauración asignado (su IP y su
// credencial SQL cifrada); si el target aún no tiene IP, se conserva el serverName recibido.
await pgPool.query(
`
INSERT INTO ${qNodes()} (
node_subnode_key, rfc, legal_name, branch_name,
notification_email, server_name, database_name, is_active, restore_target_id,
sql_password, anexo24c_aviso_fecha
) VALUES (
$1, $2, $3, $4, $5,
COALESCE((SELECT server_ip FROM ${qRestoreTargets()} WHERE id = $9), $6),
$7, $8, $9,
(SELECT sql_password_encrypted FROM ${qRestoreTargets()} WHERE id = $9),
$10
)
`,
[
row.nodoSubNodo,
row.rfc,
row.nombre,
row.sucursal,
row.correo,
row.serverName,
row.bdName,
row.activo,
row.restoreTargetId ?? null,
row.anexo24CAvisoFecha ?? null
]
);
}
export async function updateDatabaseNode(
id: number,
row: {
nodoSubNodo: string;
rfc: string;
nombre: string;
sucursal: string;
correo: string;
serverName: string;
bdName: string;
activo: number;
restoreTargetId?: number | null;
anexo24CAvisoFecha?: string | null;
}
): Promise<void> {
await pgPool.query(
`
UPDATE ${qNodes()} SET
node_subnode_key = $1,
rfc = $2,
legal_name = $3,
branch_name = $4,
notification_email = $5,
server_name = COALESCE((SELECT server_ip FROM ${qRestoreTargets()} WHERE id = $9), $6),
database_name = $7,
is_active = $8,
restore_target_id = $9,
sql_password = (SELECT sql_password_encrypted FROM ${qRestoreTargets()} WHERE id = $9),
anexo24c_aviso_fecha = $10
WHERE id = $11
`,
[
row.nodoSubNodo,
row.rfc,
row.nombre,
row.sucursal,
row.correo,
row.serverName,
row.bdName,
row.activo,
row.restoreTargetId ?? null,
row.anexo24CAvisoFecha ?? null,
id
]
);
}
export async function deleteDatabaseNode(id: number): Promise<void> {
await pgPool.query(`DELETE FROM ${qNodes()} WHERE id = $1`, [id]);
}
export async function insertPortalUser(row: {
databaseNodeId: number;
isAuthorityClient: number;
fullName: string;
username: string;
passwordHash: string;
bdShelter: string | null;
}): Promise<void> {
await pgPool.query(
`
INSERT INTO ${qUsers()} (
database_node_id, is_authority_client, full_name, username, password_hash, bd_shelter
) VALUES ($1, $2, $3, $4, $5, $6)
`,
[
row.databaseNodeId,
row.isAuthorityClient,
row.fullName,
row.username,
row.passwordHash,
row.bdShelter
]
);
}
export async function updatePortalUser(
id: number,
row: {
databaseNodeId: number;
isAuthorityClient: number;
fullName: string;
username: string;
bdShelter: string | null;
passwordHash?: string;
}
): Promise<void> {
if (row.passwordHash !== undefined) {
await pgPool.query(
`
UPDATE ${qUsers()} SET
database_node_id = $1,
is_authority_client = $2,
full_name = $3,
username = $4,
password_hash = $5,
bd_shelter = $6
WHERE id = $7
`,
[
row.databaseNodeId,
row.isAuthorityClient,
row.fullName,
row.username,
row.passwordHash,
row.bdShelter,
id
]
);
} else {
await pgPool.query(
`
UPDATE ${qUsers()} SET
database_node_id = $1,
is_authority_client = $2,
full_name = $3,
username = $4,
bd_shelter = $5
WHERE id = $6
`,
[row.databaseNodeId, row.isAuthorityClient, row.fullName, row.username, row.bdShelter, id]
);
}
}
export async function deletePortalUser(id: number): Promise<void> {
await pgPool.query(`DELETE FROM ${qUsers()} WHERE id = $1`, [id]);
}
// ============================================================================
// Servidores de restauración (restore_targets) — integración CloudRestoreAS.
// La contraseña SQL se cifra con AES-256-GCM antes de persistir y solo se
// descifra al entregarla al servicio CloudRestoreAS por el endpoint con token.
// ============================================================================
/** Servidor de restauración sin las contraseñas (seguro para la UI). */
export interface RestoreTarget {
id: number;
name: string;
server_ip: string;
sql_username: string;
data_folder: string;
ssh_host: string;
ssh_port: number;
ssh_username: string;
remote_inbox_path: string;
notes: string | null;
// Características de hardware opcionales (NULL = sin capturar).
os: string | null;
ram_gb: number | null;
disk_gb: number | null;
location: string | null;
size_category: string | null;
}
/** Datos de alta/edición. Las contraseñas en texto plano; se cifran aquí. */
export interface RestoreTargetInput {
name: string;
server_ip: string;
sql_username: string;
sql_password?: string; // opcional en edición: si se omite, no se cambia
data_folder: string;
ssh_host: string;
ssh_port: number;
ssh_username: string;
ssh_password?: string; // opcional en edición: si se omite, no se cambia
remote_inbox_path: string;
notes?: string | null;
// Características de hardware opcionales.
os?: string | null;
ram_gb?: number | null;
disk_gb?: number | null;
location?: string | null;
size_category?: string | null;
}
const ROW_RESTORE_TARGET = `
id, name, server_ip, sql_username, data_folder,
ssh_host, ssh_port, ssh_username, remote_inbox_path, notes,
os, ram_gb, disk_gb, location, size_category
`;
async function queryRestoreTargets(): Promise<RestoreTarget[]> {
const r = await pgPool.query(
`SELECT ${ROW_RESTORE_TARGET} FROM ${qRestoreTargets()} ORDER BY name`
);
return r.rows as RestoreTarget[];
}
/** Lista los servidores de restauración sin exponer la contraseña (alimenta los radios). */
export async function listRestoreTargets(): Promise<RestoreTarget[]> {
try {
return await queryRestoreTargets();
} catch (e) {
if (!isPgUndefinedTable(e) && !isPgUndefinedColumn(e)) throw e;
await ensureRestoreTargetsSchema();
return await queryRestoreTargets();
}
}
export async function getRestoreTargetById(id: number): Promise<RestoreTarget | null> {
const r = await pgPool.query(
`SELECT ${ROW_RESTORE_TARGET} FROM ${qRestoreTargets()} WHERE id = $1`,
[id]
);
return (r.rows[0] as RestoreTarget) ?? null;
}
/**
* Servidor de restauración por id con la contraseña SQL descifrada. Uso exclusivo del servidor
* (conectar a SQL Server para depurar bases duplicadas); NUNCA se expone al cliente.
*/
export async function getRestoreTargetWithPasswordById(
id: number
): Promise<(RestoreTarget & { sql_password: string }) | null> {
const r = await pgPool.query(
`SELECT ${ROW_RESTORE_TARGET}, sql_password_encrypted FROM ${qRestoreTargets()} WHERE id = $1`,
[id]
);
const row = r.rows[0];
if (!row) return null;
const { sql_password_encrypted, ...rest } = row;
const sql_password = sql_password_encrypted ? decryptSecret(sql_password_encrypted) : '';
return { ...(rest as RestoreTarget), sql_password };
}
/**
* Servidor de restauración por id con AMBAS contraseñas descifradas (SQL para BACKUP/RESTORE, SSH
* para SFTP). Uso exclusivo del servidor (mover bases duplicadas); NUNCA se expone al cliente.
*/
export async function getRestoreTargetWithSecretsById(
id: number
): Promise<(RestoreTarget & { sql_password: string; ssh_password: string }) | null> {
const r = await pgPool.query(
`SELECT ${ROW_RESTORE_TARGET}, sql_password_encrypted, ssh_password_encrypted
FROM ${qRestoreTargets()} WHERE id = $1`,
[id]
);
const row = r.rows[0];
if (!row) return null;
const { sql_password_encrypted, ssh_password_encrypted, ...rest } = row;
return {
...(rest as RestoreTarget),
sql_password: sql_password_encrypted ? decryptSecret(sql_password_encrypted) : '',
ssh_password: ssh_password_encrypted ? decryptSecret(ssh_password_encrypted) : ''
};
}
/**
* Devuelve el servidor de restauración ASIGNADO a una base de datos, con la contraseña
* descifrada. La base se identifica por database_name o node_subnode_key (lo que CloudRestoreAS
* resuelve del nombre de archivo). Uso exclusivo del endpoint servicio-a-servicio.
* Retorna null si la base no existe o no tiene servidor asignado.
*/
export async function getRestoreTargetForDatabase(
dbNameOrNodo: string,
instanceName?: string | null
): Promise<(RestoreTarget & { sql_password: string; ssh_password: string }) | null> {
const instanceFilter = (instanceName ?? '').trim();
const r = await pgPool.query(
`SELECT rt.id, rt.name, rt.server_ip, rt.sql_username, rt.data_folder,
rt.ssh_host, rt.ssh_port, rt.ssh_username, rt.remote_inbox_path, rt.notes,
rt.sql_password_encrypted, rt.ssh_password_encrypted
FROM ${qNodes()} dn
JOIN ${qRestoreTargets()} rt ON rt.id = dn.restore_target_id
WHERE (LOWER(TRIM(dn.database_name)) = LOWER(TRIM($1::text))
OR LOWER(TRIM(dn.node_subnode_key)) = LOWER(TRIM($1::text)))
AND (
$2::text = ''
OR LOWER(TRIM(rt.name)) = LOWER(TRIM($2::text))
)
LIMIT 1`,
[dbNameOrNodo, instanceFilter]
);
const row = r.rows[0];
if (!row) return null;
const { sql_password_encrypted, ssh_password_encrypted, ...rest } = row;
// Si el servidor aún no tiene contraseñas configuradas se devuelven vacías:
// el cliente (panel_client) las rechaza por campo requerido y difiere el job.
const sql_password = sql_password_encrypted ? decryptSecret(sql_password_encrypted) : '';
const ssh_password = ssh_password_encrypted ? decryptSecret(ssh_password_encrypted) : '';
return { ...(rest as RestoreTarget), sql_password, ssh_password };
}
export type RouteAction = 'restore_local' | 'forward';
export type RestoreTargetWithRoute = RestoreTarget & {
sql_password: string;
ssh_password: string;
input_folder: string | null;
};
/**
* Servidor asignado a una base, con contraseñas descifradas y carpeta de entrada
* reportada por CloudRestoreAS (cloudrestore_status), si existe.
*/
export async function getRestoreTargetWithInputFolderForDatabase(
dbNameOrNodo: string
): Promise<RestoreTargetWithRoute | null> {
const r = await pgPool.query(
`SELECT rt.id, rt.name, rt.server_ip, rt.sql_username, rt.data_folder,
rt.ssh_host, rt.ssh_port, rt.ssh_username, rt.remote_inbox_path, rt.notes,
rt.sql_password_encrypted, rt.ssh_password_encrypted,
cs.input_folder
FROM ${qNodes()} dn
JOIN ${qRestoreTargets()} rt ON rt.id = dn.restore_target_id
LEFT JOIN ${qCloudRestoreStatus()} cs
ON LOWER(TRIM(cs.instance_key)) = LOWER(TRIM(rt.name))
WHERE (LOWER(TRIM(dn.database_name)) = LOWER(TRIM($1::text))
OR LOWER(TRIM(dn.node_subnode_key)) = LOWER(TRIM($1::text)))
LIMIT 1`,
[dbNameOrNodo]
);
const row = r.rows[0];
if (!row) return null;
const { sql_password_encrypted, ssh_password_encrypted, input_folder, ...rest } = row;
const sql_password = sql_password_encrypted ? decryptSecret(sql_password_encrypted) : '';
const ssh_password = ssh_password_encrypted ? decryptSecret(ssh_password_encrypted) : '';
const folder =
input_folder && String(input_folder).trim() ? String(input_folder).trim() : null;
return { ...(rest as RestoreTarget), sql_password, ssh_password, input_folder: folder };
}
export interface ResolveRouteResult {
action: RouteAction;
db_name: string;
node_key: string;
target: RestoreTargetWithRoute;
}
/**
* Resuelve nodo/destino desde el nombre del archivo ZIP y determina si el CRA debe
* restaurar localmente o reenviar el ZIP por SFTP.
*/
export async function resolveRouteForFilename(
filename: string,
instanceName?: string | null
): Promise<ResolveRouteResult | null> {
const stem = path.parse(filename).name;
const nodes = await listDatabaseNodesActive();
const matched = matchNodeRowFromBackupStem(stem, nodes);
if (!matched) return null;
const dbName = String(matched.BDName ?? matched.NodoSubNodo ?? '').trim();
if (!dbName) return null;
const target = await getRestoreTargetWithInputFolderForDatabase(dbName);
if (!target) return null;
const instance = (instanceName ?? '').trim();
let action: RouteAction;
if (!instance) {
action = 'forward';
} else if (target.name.trim().toLowerCase() === instance.toLowerCase()) {
action = 'restore_local';
} else {
action = 'forward';
}
return {
action,
db_name: dbName,
node_key: String(matched.NodoSubNodo ?? stem).trim(),
target
};
}
export async function createRestoreTarget(input: RestoreTargetInput): Promise<number> {
if (!input.sql_password) {
throw new Error('La contraseña SQL es obligatoria al crear un servidor.');
}
const cols = [
'name', 'server_ip', 'sql_username', 'data_folder',
'ssh_host', 'ssh_port', 'ssh_username', 'remote_inbox_path', 'notes',
'os', 'ram_gb', 'disk_gb', 'location', 'size_category',
'sql_password_encrypted'
];
const params: unknown[] = [
input.name, input.server_ip, input.sql_username, input.data_folder,
input.ssh_host, input.ssh_port, input.ssh_username, input.remote_inbox_path,
input.notes ?? null,
input.os ?? null, input.ram_gb ?? null, input.disk_gb ?? null,
input.location ?? null, input.size_category ?? null,
encryptSecret(input.sql_password)
];
if (input.ssh_password) {
cols.push('ssh_password_encrypted');
params.push(encryptSecret(input.ssh_password));
}
const placeholders = params.map((_, i) => `$${i + 1}`).join(', ');
const r = await pgPool.query(
`INSERT INTO ${qRestoreTargets()} (${cols.join(', ')}) VALUES (${placeholders}) RETURNING id`,
params
);
return r.rows[0].id as number;
}
/** Actualiza un servidor. Cada contraseña (SQL/SSH) solo se reescribe si viene definida. */
export async function updateRestoreTarget(id: number, input: RestoreTargetInput): Promise<void> {
const sets = [
'name = $1', 'server_ip = $2', 'sql_username = $3', 'data_folder = $4',
'ssh_host = $5', 'ssh_port = $6', 'ssh_username = $7', 'remote_inbox_path = $8',
'notes = $9', 'os = $10', 'ram_gb = $11', 'disk_gb = $12', 'location = $13',
'size_category = $14', 'updated_at = now()'
];
const params: unknown[] = [
input.name, input.server_ip, input.sql_username, input.data_folder,
input.ssh_host, input.ssh_port, input.ssh_username, input.remote_inbox_path,
input.notes ?? null,
input.os ?? null, input.ram_gb ?? null, input.disk_gb ?? null,
input.location ?? null, input.size_category ?? null
];
let p = params.length;
if (input.sql_password) {
params.push(encryptSecret(input.sql_password));
sets.push(`sql_password_encrypted = $${++p}`);
}
if (input.ssh_password) {
params.push(encryptSecret(input.ssh_password));
sets.push(`ssh_password_encrypted = $${++p}`);
}
params.push(id);
await pgPool.query(
`UPDATE ${qRestoreTargets()} SET ${sets.join(', ')} WHERE id = $${params.length}`,
params
);
}
export async function deleteRestoreTarget(id: number): Promise<void> {
await pgPool.query(`DELETE FROM ${qRestoreTargets()} WHERE id = $1`, [id]);
}
// ============================================================================
// Asignación masiva de nodos a servidores de restauración.
// ============================================================================
/** Nodo con su asignación de restaurador, para el checklist/distribución (sin secretos). */
export interface AssignmentNode {
ID: number;
NodoSubNodo: string;
Nombre: string;
BDName: string | null;
ServerName: string | null;
Activo: number;
RestoreTargetId: number | null;
}
/** Todos los nodos con su asignación actual (base del checklist y de la distribución). */
export async function listNodesForAssignment(): Promise<AssignmentNode[]> {
const r = await pgPool.query(
`
SELECT
id AS "ID",
node_subnode_key AS "NodoSubNodo",
legal_name AS "Nombre",
database_name AS "BDName",
server_name AS "ServerName",
is_active AS "Activo",
restore_target_id AS "RestoreTargetId"
FROM ${qNodes()}
ORDER BY node_subnode_key
`
);
return r.rows as AssignmentNode[];
}
/**
* Guarda el checklist manual de un restaurador: los `checkedNodeIds` quedan asignados a
* `targetId` (reasignando desde donde estuvieran y derivando server_name y sql_password del
* target); los nodos que estaban en este restaurador y ya NO vienen marcados quedan sin asignar
* (restore_target_id y sql_password en NULL). Los nodos de OTROS restauradores no marcados no se
* tocan. Transacción con prepared statements.
*/
export async function assignNodesToRestoreTarget(
targetId: number,
checkedNodeIds: number[]
): Promise<void> {
const ids = Array.from(new Set(checkedNodeIds.filter((n) => Number.isInteger(n) && n > 0)));
const client = await pgPool.connect();
try {
await client.query('BEGIN');
// 1) Asignar/reasignar los marcados (server_name = IP del target; sql_password = su credencial cifrada).
await client.query(
`
UPDATE ${qNodes()}
SET restore_target_id = $1,
server_name = COALESCE((SELECT server_ip FROM ${qRestoreTargets()} WHERE id = $1), server_name),
sql_password = (SELECT sql_password_encrypted FROM ${qRestoreTargets()} WHERE id = $1)
WHERE id = ANY($2::int[])
`,
[targetId, ids]
);
// 2) Quitar de este restaurador los desmarcados (restore_target_id y sql_password a NULL); no toca otros targets.
await client.query(
`
UPDATE ${qNodes()}
SET restore_target_id = NULL,
sql_password = NULL
WHERE restore_target_id = $1
AND NOT (id = ANY($2::int[]))
`,
[targetId, ids]
);
await client.query('COMMIT');
} catch (e) {
await client.query('ROLLBACK');
throw e;
} finally {
client.release();
}
}
/**
* Aplica un conjunto de asignaciones (distribución global o automática): cada par fija el
* restore_target_id del nodo (y deriva server_name y sql_password del restaurador si no es NULL).
* Pares con targetId NULL dejan el nodo sin asignar (conservando server_name, sql_password a NULL).
* Transacción.
*/
export async function applyNodeAssignments(
pairs: { nodeId: number; targetId: number | null }[]
): Promise<void> {
const clean = pairs.filter((p) => Number.isInteger(p.nodeId) && p.nodeId > 0);
if (clean.length === 0) return;
const client = await pgPool.connect();
try {
await client.query('BEGIN');
for (const { nodeId, targetId } of clean) {
await client.query(
`
UPDATE ${qNodes()}
SET restore_target_id = $2,
server_name = CASE
WHEN $2::int IS NULL THEN server_name
ELSE COALESCE((SELECT server_ip FROM ${qRestoreTargets()} WHERE id = $2), server_name)
END,
sql_password = CASE
WHEN $2::int IS NULL THEN NULL
ELSE (SELECT sql_password_encrypted FROM ${qRestoreTargets()} WHERE id = $2)
END
WHERE id = $1
`,
[nodeId, targetId]
);
}
await client.query('COMMIT');
} catch (e) {
await client.query('ROLLBACK');
throw e;
} finally {
client.release();
}
}
// ============================================================================
// Estado de CloudRestoreAS (carpeta de entrada reportada por el servicio).
// ============================================================================
export interface CloudRestoreStatus {
id: number;
instance_key: string;
input_folder: string;
processed_folder: string | null;
host_name: string | null;
app_version: string | null;
reported_at: Date;
}
/** Estados reportados por cada instancia CloudRestoreAS (Alfa/Omega/Gamma). */
export async function listCloudRestoreStatuses(): Promise<CloudRestoreStatus[]> {
try {
await ensureCloudRestoreStatusTable();
const r = await pgPool.query(
`
SELECT id, instance_key, input_folder, processed_folder, host_name, app_version, reported_at
FROM ${qCloudRestoreStatus()}
ORDER BY instance_key
`
);
return r.rows as CloudRestoreStatus[];
} catch (e) {
if (isPgUndefinedTable(e)) return [];
throw e;
}
}
/** @deprecated Usar listCloudRestoreStatuses. Mantiene compatibilidad con instancia default. */
export async function getCloudRestoreStatus(): Promise<CloudRestoreStatus | null> {
const rows = await listCloudRestoreStatuses();
return rows.find((r) => r.instance_key === 'default') ?? rows[0] ?? null;
}
/** UPSERT del estado reportado por CloudRestoreAS (solo vía API servicio). */
export async function upsertCloudRestoreStatus(row: {
inputFolder: string;
processedFolder?: string | null;
hostName: string | null;
appVersion: string | null;
instanceKey?: string;
}): Promise<void> {
await ensureCloudRestoreStatusTable();
const key = row.instanceKey ?? 'default';
await pgPool.query(
`
INSERT INTO ${qCloudRestoreStatus()} (
instance_key, input_folder, processed_folder, host_name, app_version, reported_at
) VALUES ($1, $2, $3, $4, $5, now())
ON CONFLICT (instance_key) DO UPDATE SET
input_folder = EXCLUDED.input_folder,
processed_folder = EXCLUDED.processed_folder,
host_name = EXCLUDED.host_name,
app_version = EXCLUDED.app_version,
reported_at = now()
`,
[key, row.inputFolder, row.processedFolder ?? null, row.hostName, row.appVersion]
);
}
/** Inserta un registro de bitácora reportado por CloudRestoreAS. */
export async function insertRestoreJobLog(row: {
filename: string;
restoreTargetId: number | null;
dbName: string | null;
status: 'completed' | 'failed' | 'forwarded';
durationMs: number | null;
errorMessage: string | null;
sizeBytes?: number | null;
relPath?: string | null;
}): Promise<void> {
await ensureRestoreJobLogColumns();
await pgPool.query(
`
INSERT INTO ${qRestoreJobLogs()} (
filename, restore_target_id, db_name, status, duration_ms, error_message,
size_bytes, rel_path
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
`,
[
row.filename,
row.restoreTargetId,
row.dbName,
row.status,
row.durationMs,
row.errorMessage,
row.sizeBytes ?? null,
row.relPath ?? null
]
);
}
/** Resumen de bitácora por servidor: conteos por status (últimos 30 días) y última exitosa. */
export interface RestoreJobLogSummary {
restore_target_id: number;
completed: number;
failed: number;
forwarded: number;
last_completed_at: Date | null;
}
export async function listRestoreJobLogSummaries(): Promise<RestoreJobLogSummary[]> {
try {
const r = await pgPool.query(
`
SELECT
restore_target_id,
COUNT(*) FILTER (
WHERE status = 'completed' AND restored_at >= now() - INTERVAL '30 days'
)::int AS completed,
COUNT(*) FILTER (
WHERE status = 'failed' AND restored_at >= now() - INTERVAL '30 days'
)::int AS failed,
COUNT(*) FILTER (
WHERE status = 'forwarded' AND restored_at >= now() - INTERVAL '30 days'
)::int AS forwarded,
MAX(restored_at) FILTER (WHERE status = 'completed') AS last_completed_at
FROM ${qRestoreJobLogs()}
WHERE restore_target_id IS NOT NULL
GROUP BY restore_target_id
`
);
return r.rows as RestoreJobLogSummary[];
} catch (e) {
if (isPgUndefinedTable(e)) return [];
throw e;
}
}
/** Últimas restauraciones de un servidor (para el modal de bitácora del panel). */
export interface RestoreJobLogRow {
id: number;
filename: string;
db_name: string | null;
status: string;
duration_ms: number | null;
error_message: string | null;
restored_at: Date;
}
export async function listRecentRestoreJobLogs(
targetId: number,
limit = 20
): Promise<RestoreJobLogRow[]> {
const capped = Math.min(Math.max(1, Math.trunc(limit)), 100);
try {
const r = await pgPool.query(
`
SELECT id, filename, db_name, status, duration_ms, error_message, restored_at
FROM ${qRestoreJobLogs()}
WHERE restore_target_id = $1
ORDER BY restored_at DESC
LIMIT $2
`,
[targetId, capped]
);
return r.rows as RestoreJobLogRow[];
} catch (e) {
if (isPgUndefinedTable(e)) return [];
throw e;
}
}
// ============================================================================
// Inventario "último respaldo por nodo" (node_last_restore) y restores fallidos.
// ============================================================================
export interface NodeLastRestoreRow {
database_node_id: number;
restore_target_id: number | null;
server_name: string | null; // restaurador (Alfa/Omega/Gamma) donde vive el archivo
node_key: string | null; // NodoSubNodo
client_name: string | null; // legal_name
db_name: string | null;
filename: string;
rel_path: string | null;
size_bytes: number | null;
restored_at: Date;
}
/** Último respaldo restaurado por nodo. Sobrevive a la reasignación del restaurador. */
export async function listNodeLastRestore(): Promise<NodeLastRestoreRow[]> {
try {
await ensureNodeLastRestoreTable();
const r = await pgPool.query(
`
SELECT
nlr.database_node_id,
nlr.restore_target_id,
rt.name AS server_name,
COALESCE(nlr.node_key, dn.node_subnode_key) AS node_key,
dn.legal_name AS client_name,
nlr.db_name,
nlr.filename,
nlr.rel_path,
nlr.size_bytes,
nlr.restored_at
FROM ${qNodeLastRestore()} nlr
LEFT JOIN ${qRestoreTargets()} rt ON rt.id = nlr.restore_target_id
LEFT JOIN ${qNodes()} dn ON dn.id = nlr.database_node_id
ORDER BY nlr.restored_at DESC
`
);
return r.rows as NodeLastRestoreRow[];
} catch (e) {
if (isPgUndefinedTable(e)) return [];
throw e;
}
}
/**
* UPSERT del último respaldo de un nodo. Solo actualiza si el nuevo `restored_at` es igual o
* más reciente, para no retroceder el tracking. Llaveada por nodo → persiste ante reasignación.
*/
export async function upsertNodeLastRestore(row: {
databaseNodeId: number;
restoreTargetId: number | null;
dbName: string | null;
nodeKey: string | null;
filename: string;
relPath: string | null;
sizeBytes: number | null;
restoredAt?: Date | null;
}): Promise<void> {
await ensureNodeLastRestoreTable();
await pgPool.query(
`
INSERT INTO ${qNodeLastRestore()} AS nlr (
database_node_id, restore_target_id, db_name, node_key,
filename, rel_path, size_bytes, restored_at, updated_at
) VALUES ($1, $2, $3, $4, $5, $6, $7, COALESCE($8, now()), now())
ON CONFLICT (database_node_id) DO UPDATE SET
restore_target_id = EXCLUDED.restore_target_id,
db_name = EXCLUDED.db_name,
node_key = EXCLUDED.node_key,
filename = EXCLUDED.filename,
rel_path = EXCLUDED.rel_path,
size_bytes = EXCLUDED.size_bytes,
restored_at = EXCLUDED.restored_at,
updated_at = now()
WHERE EXCLUDED.restored_at >= nlr.restored_at
`,
[
row.databaseNodeId,
row.restoreTargetId,
row.dbName,
row.nodeKey,
row.filename,
row.relPath,
row.sizeBytes,
row.restoredAt ?? null
]
);
}
/**
* Resuelve el database_node_id de un respaldo por db_name o por el nombre de archivo.
* Reutiliza matchNodeRowFromBackupStem para el fallback por stem del filename.
*/
export async function findNodeIdForBackup(
filename: string,
dbName: string | null
): Promise<{ id: number; nodeKey: string | null } | null> {
if (dbName) {
const r = await pgPool.query(
`SELECT id, node_subnode_key FROM ${qNodes()}
WHERE LOWER(TRIM(database_name)) = LOWER(TRIM($1))
OR LOWER(TRIM(node_subnode_key)) = LOWER(TRIM($1))
LIMIT 1`,
[dbName]
);
if (r.rows.length) {
return { id: Number(r.rows[0].id), nodeKey: r.rows[0].node_subnode_key ?? null };
}
}
const stem = path.parse(filename).name;
const nodes = await pgPool.query(
`SELECT id AS "ID", node_subnode_key AS "NodoSubNodo", database_name AS "BDName"
FROM ${qNodes()}`
);
const match = matchNodeRowFromBackupStem(stem, nodes.rows);
if (match && match.ID != null) {
return { id: Number(match.ID), nodeKey: (match.NodoSubNodo as string) ?? null };
}
return null;
}
/** Restores fallidos (todos los servidores) para el panel de fallidos. */
export interface FailedRestoreRow {
id: number;
restore_target_id: number | null;
server_name: string | null;
filename: string;
db_name: string | null;
error_message: string | null;
rel_path: string | null;
size_bytes: number | null;
restored_at: Date;
}
export async function listFailedRestoreJobLogs(limit = 100): Promise<FailedRestoreRow[]> {
const capped = Math.min(Math.max(1, Math.trunc(limit)), 500);
try {
await ensureRestoreJobLogColumns();
const r = await pgPool.query(
`
SELECT
jl.id, jl.restore_target_id, rt.name AS server_name,
jl.filename, jl.db_name, jl.error_message, jl.rel_path, jl.size_bytes, jl.restored_at
FROM ${qRestoreJobLogs()} jl
LEFT JOIN ${qRestoreTargets()} rt ON rt.id = jl.restore_target_id
WHERE jl.status = 'failed'
ORDER BY jl.restored_at DESC
LIMIT $1
`,
[capped]
);
return r.rows as FailedRestoreRow[];
} catch (e) {
if (isPgUndefinedTable(e)) return [];
throw e;
}
}
/** Restores completados/reenviados (todos los servidores) para el panel de restaurados. */
export interface RestoredRestoreRow {
id: number;
restore_target_id: number | null;
server_name: string | null;
node_key: string | null;
client_name: string | null;
db_name: string | null;
filename: string;
rel_path: string | null;
size_bytes: number | null;
restored_at: Date;
}
/**
* Restores exitosos (`completed`/`forwarded`) desde restore_job_logs — lo que reporta
* CloudRestoreAS vía POST /api/restore/job-result. Resuelve node_key/client_name por nombre
* de base contra database_nodes (LATERAL … LIMIT 1 para no duplicar la fila del log si dos
* nodos comparten database_name). Cuando no hay match, ambos quedan null y la UI cae a db_name.
*/
export async function listRestoredRestoreJobLogs(limit = 200): Promise<RestoredRestoreRow[]> {
const capped = Math.min(Math.max(1, Math.trunc(limit)), 500);
try {
await ensureRestoreJobLogColumns();
const r = await pgPool.query(
`
SELECT
jl.id, jl.restore_target_id, rt.name AS server_name,
dn.node_subnode_key AS node_key, dn.legal_name AS client_name,
jl.filename, jl.db_name, jl.rel_path, jl.size_bytes, jl.restored_at
FROM ${qRestoreJobLogs()} jl
LEFT JOIN ${qRestoreTargets()} rt ON rt.id = jl.restore_target_id
LEFT JOIN LATERAL (
SELECT n.node_subnode_key, n.legal_name
FROM ${qNodes()} n
WHERE LOWER(TRIM(n.database_name)) = LOWER(TRIM(jl.db_name))
LIMIT 1
) dn ON true
WHERE jl.status IN ('completed', 'forwarded')
ORDER BY jl.restored_at DESC
LIMIT $1
`,
[capped]
);
return r.rows as RestoredRestoreRow[];
} catch (e) {
if (isPgUndefinedTable(e)) return [];
throw e;
}
}
/** Carpetas de un restaurador para descarga por filesystem (input/processed reportados). */
export interface RestoreTargetDownload {
id: number;
name: string;
input_folder: string | null;
processed_folder: string | null;
}
export async function getRestoreTargetForDownload(
targetId: number
): Promise<RestoreTargetDownload | null> {
await ensureRestoreTargetsSchema();
await ensureCloudRestoreStatusTable();
const r = await pgPool.query(
`
SELECT rt.id, rt.name, cs.input_folder, cs.processed_folder
FROM ${qRestoreTargets()} rt
LEFT JOIN ${qCloudRestoreStatus()} cs
ON LOWER(TRIM(cs.instance_key)) = LOWER(TRIM(rt.name))
WHERE rt.id = $1
`,
[targetId]
);
if (!r.rows.length) return null;
const row = r.rows[0];
return {
id: Number(row.id),
name: row.name,
input_folder: row.input_folder ?? null,
processed_folder: row.processed_folder ?? null
};
}