185 lines
7.0 KiB
Python
185 lines
7.0 KiB
Python
"""Esquema y gestión de la base de datos SQLite."""
|
|
|
|
import sqlite3
|
|
import threading
|
|
from pathlib import Path
|
|
from typing import Optional
|
|
from datetime import datetime
|
|
from ..constants import DB_PATH
|
|
|
|
# Espera máxima (s) por un lock de SQLite antes de fallar. Con WAL + busy_timeout los
|
|
# lectores (UI cada 5s) y escritores (workers + file_watcher) dejan de chocar con
|
|
# "database is locked": esperan en vez de fallar de inmediato.
|
|
_SQLITE_TIMEOUT_SECONDS = 30.0
|
|
_SQLITE_BUSY_TIMEOUT_MS = 5000
|
|
|
|
|
|
class DatabaseManager:
|
|
"""Gestor de la base de datos SQLite."""
|
|
|
|
def __init__(self, db_path: Optional[Path] = None):
|
|
"""
|
|
Inicializa el gestor de base de datos.
|
|
|
|
Args:
|
|
db_path: Ruta a la base de datos (usa DB_PATH por defecto)
|
|
"""
|
|
self.db_path = db_path or DB_PATH
|
|
self.db_path.parent.mkdir(parents=True, exist_ok=True)
|
|
# Serializa las escrituras del propio proceso (varios workers + file_watcher) para
|
|
# eliminar las colisiones write-write; WAL cubre la concurrencia lectura/escritura.
|
|
self._write_lock = threading.Lock()
|
|
self._initialize_schema()
|
|
|
|
def get_connection(self) -> sqlite3.Connection:
|
|
"""Obtiene una conexión a SQLite con protecciones de concurrencia (WAL + timeouts)."""
|
|
conn = sqlite3.connect(
|
|
str(self.db_path), check_same_thread=False, timeout=_SQLITE_TIMEOUT_SECONDS
|
|
)
|
|
conn.row_factory = sqlite3.Row
|
|
# WAL: lectores concurrentes con un escritor (persistente en el archivo, idempotente).
|
|
conn.execute("PRAGMA journal_mode=WAL")
|
|
conn.execute(f"PRAGMA busy_timeout={_SQLITE_BUSY_TIMEOUT_MS}")
|
|
conn.execute("PRAGMA synchronous=NORMAL")
|
|
return conn
|
|
|
|
def _initialize_schema(self):
|
|
"""Crea las tablas si no existen."""
|
|
conn = self.get_connection()
|
|
try:
|
|
cursor = conn.cursor()
|
|
|
|
# Tabla de jobs
|
|
cursor.execute("""
|
|
CREATE TABLE IF NOT EXISTS jobs (
|
|
job_id TEXT PRIMARY KEY,
|
|
created_at TEXT NOT NULL,
|
|
updated_at TEXT NOT NULL,
|
|
source_path TEXT NOT NULL,
|
|
source_name TEXT NOT NULL,
|
|
source_hash TEXT,
|
|
node_name TEXT,
|
|
db_name TEXT,
|
|
status TEXT NOT NULL,
|
|
attempts INTEGER DEFAULT 0,
|
|
last_error TEXT,
|
|
started_at TEXT,
|
|
finished_at TEXT,
|
|
total_ms INTEGER,
|
|
extract_ms INTEGER,
|
|
restore_ms INTEGER,
|
|
filelist_ms INTEGER,
|
|
purged_at TEXT
|
|
)
|
|
""")
|
|
|
|
# Tabla de pasos de jobs
|
|
cursor.execute("""
|
|
CREATE TABLE IF NOT EXISTS job_steps (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
job_id TEXT NOT NULL,
|
|
step TEXT NOT NULL,
|
|
started_at TEXT NOT NULL,
|
|
finished_at TEXT,
|
|
duration_ms INTEGER,
|
|
exit_code INTEGER,
|
|
stdout TEXT,
|
|
stderr TEXT,
|
|
error TEXT,
|
|
FOREIGN KEY (job_id) REFERENCES jobs(job_id)
|
|
)
|
|
""")
|
|
|
|
# Tabla de nodos (mapeo node_name -> db_name)
|
|
cursor.execute("""
|
|
CREATE TABLE IF NOT EXISTS nodes (
|
|
node_name TEXT PRIMARY KEY,
|
|
db_name TEXT NOT NULL,
|
|
active INTEGER DEFAULT 1,
|
|
notes TEXT,
|
|
updated_at TEXT NOT NULL
|
|
)
|
|
""")
|
|
|
|
# Tabla de eventos
|
|
cursor.execute("""
|
|
CREATE TABLE IF NOT EXISTS events (
|
|
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
|
created_at TEXT NOT NULL,
|
|
level TEXT NOT NULL,
|
|
job_id TEXT,
|
|
message TEXT NOT NULL,
|
|
FOREIGN KEY (job_id) REFERENCES jobs(job_id)
|
|
)
|
|
""")
|
|
|
|
# Tabla de configuración
|
|
cursor.execute("""
|
|
CREATE TABLE IF NOT EXISTS config (
|
|
key TEXT PRIMARY KEY,
|
|
value TEXT NOT NULL,
|
|
updated_at TEXT NOT NULL
|
|
)
|
|
""")
|
|
|
|
# Migración idempotente: columna purged_at para marcar los respaldos que la
|
|
# retención ya borró del disco. Se conserva la fila (historial/stats/UI) y solo
|
|
# se anota que su archivo físico dejó de existir. Backward-compatible: columna
|
|
# nullable que las versiones anteriores de la app simplemente ignoran.
|
|
cursor.execute("PRAGMA table_info(jobs)")
|
|
job_columns = {row[1] for row in cursor.fetchall()}
|
|
if "purged_at" not in job_columns:
|
|
cursor.execute("ALTER TABLE jobs ADD COLUMN purged_at TEXT")
|
|
|
|
# Índices
|
|
cursor.execute("CREATE INDEX IF NOT EXISTS idx_jobs_status ON jobs(status)")
|
|
cursor.execute("CREATE INDEX IF NOT EXISTS idx_jobs_created ON jobs(created_at)")
|
|
# Soporta la retención por nodo (GROUP BY node_name + filtro por finished_at).
|
|
cursor.execute(
|
|
"CREATE INDEX IF NOT EXISTS idx_jobs_node_finished "
|
|
"ON jobs(node_name, finished_at)"
|
|
)
|
|
cursor.execute("CREATE INDEX IF NOT EXISTS idx_events_created ON events(created_at)")
|
|
cursor.execute("CREATE INDEX IF NOT EXISTS idx_events_job ON events(job_id)")
|
|
cursor.execute("CREATE INDEX IF NOT EXISTS idx_job_steps_job ON job_steps(job_id)")
|
|
|
|
conn.commit()
|
|
finally:
|
|
conn.close()
|
|
|
|
def execute(self, query: str, params: tuple = ()) -> sqlite3.Cursor:
|
|
"""Ejecuta una escritura, serializada por lock de proceso y con conexión cerrada."""
|
|
with self._write_lock:
|
|
conn = self.get_connection()
|
|
try:
|
|
cursor = conn.cursor()
|
|
cursor.execute(query, params)
|
|
conn.commit()
|
|
return cursor
|
|
finally:
|
|
conn.close()
|
|
|
|
def fetchone(self, query: str, params: tuple = ()) -> Optional[sqlite3.Row]:
|
|
"""Ejecuta una consulta y devuelve una fila."""
|
|
conn = self.get_connection()
|
|
try:
|
|
cursor = conn.cursor()
|
|
cursor.execute(query, params)
|
|
return cursor.fetchone()
|
|
finally:
|
|
conn.close()
|
|
|
|
def fetchall(self, query: str, params: tuple = ()) -> list[sqlite3.Row]:
|
|
"""Ejecuta una consulta y devuelve todas las filas."""
|
|
conn = self.get_connection()
|
|
try:
|
|
cursor = conn.cursor()
|
|
cursor.execute(query, params)
|
|
return cursor.fetchall()
|
|
finally:
|
|
conn.close()
|
|
|
|
|
|
# Instancia global
|
|
db = DatabaseManager()
|