feature/bitacora-correccion-filtro-tenant
This commit is contained in:
@@ -2,12 +2,20 @@
|
||||
Audit Log Events
|
||||
"""
|
||||
|
||||
import logging
|
||||
|
||||
from sqlalchemy import event, inspect
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from api.v1.modules.a76.general_catalogs.company.models import Company
|
||||
from core.database import rls_company_var, rls_tenant_var
|
||||
|
||||
from .services.service import AuditService
|
||||
from .utils.serialization import serialize_for_json
|
||||
from core.context import get_user_context
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def register_audit_listeners(models_to_audit):
|
||||
"""
|
||||
@@ -30,11 +38,37 @@ def _get_current_username():
|
||||
or context.get("sub")
|
||||
or "System"
|
||||
)
|
||||
except:
|
||||
except Exception:
|
||||
pass
|
||||
return "System"
|
||||
|
||||
|
||||
def _resolve_audit_company_tenant(session: Session, target) -> tuple:
|
||||
"""
|
||||
company_id / tenant_id desde la fila ORM, ContextVars RLS (petición HTTP),
|
||||
o ``Company.tenant_id`` por ``company_id``.
|
||||
"""
|
||||
company_id = getattr(target, "company_id", None)
|
||||
if company_id is None and getattr(target, "__tablename__", None) == "company":
|
||||
company_id = getattr(target, "id", None)
|
||||
if company_id is None:
|
||||
company_id = rls_company_var.get()
|
||||
|
||||
tenant_id = getattr(target, "tenant_id", None)
|
||||
if tenant_id is None:
|
||||
tenant_id = rls_tenant_var.get()
|
||||
if tenant_id is None and company_id is not None:
|
||||
row = (
|
||||
session.query(Company.tenant_id)
|
||||
.filter(Company.id == company_id, Company.deleted_at.is_(None))
|
||||
.first()
|
||||
)
|
||||
if row:
|
||||
tenant_id = int(row[0])
|
||||
|
||||
return company_id, tenant_id
|
||||
|
||||
|
||||
def after_insert_listener(mapper, connection, target):
|
||||
"""
|
||||
Listener for INSERT operations
|
||||
@@ -42,12 +76,19 @@ def after_insert_listener(mapper, connection, target):
|
||||
table_name = target.__tablename__
|
||||
record_data = {c.name: getattr(target, c.name) for c in mapper.columns}
|
||||
username = _get_current_username()
|
||||
company_id = getattr(target, "company_id", None) or getattr(target, "id", None)
|
||||
tenant_id = getattr(target, "tenant_id", None)
|
||||
|
||||
# Create a session bound to the connection
|
||||
session = Session(bind=connection)
|
||||
try:
|
||||
company_id, tenant_id = _resolve_audit_company_tenant(session, target)
|
||||
if company_id is None or tenant_id is None:
|
||||
logger.debug(
|
||||
"Audit skip INSERT %s: missing company_id=%s tenant_id=%s",
|
||||
table_name,
|
||||
company_id,
|
||||
tenant_id,
|
||||
)
|
||||
return
|
||||
|
||||
AuditService.log_crud_operation(
|
||||
db=session,
|
||||
table_name=table_name,
|
||||
@@ -59,7 +100,7 @@ def after_insert_listener(mapper, connection, target):
|
||||
tenant_id=tenant_id,
|
||||
)
|
||||
except Exception as e:
|
||||
print(f"Error logging insert: {e}")
|
||||
logger.warning("Error logging insert for %s: %s", table_name, e)
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
@@ -87,11 +128,19 @@ def after_update_listener(mapper, connection, target):
|
||||
|
||||
record_data = {c.name: getattr(target, c.name) for c in mapper.columns}
|
||||
username = _get_current_username()
|
||||
company_id = getattr(target, "company_id", None) or getattr(target, "id", None)
|
||||
tenant_id = getattr(target, "tenant_id", None)
|
||||
|
||||
session = Session(bind=connection)
|
||||
try:
|
||||
company_id, tenant_id = _resolve_audit_company_tenant(session, target)
|
||||
if company_id is None or tenant_id is None:
|
||||
logger.debug(
|
||||
"Audit skip UPDATE %s: missing company_id=%s tenant_id=%s",
|
||||
table_name,
|
||||
company_id,
|
||||
tenant_id,
|
||||
)
|
||||
return
|
||||
|
||||
AuditService.log_crud_operation(
|
||||
db=session,
|
||||
table_name=table_name,
|
||||
@@ -105,7 +154,7 @@ def after_update_listener(mapper, connection, target):
|
||||
tenant_id=tenant_id,
|
||||
)
|
||||
except Exception as e:
|
||||
print(f"Error logging update: {e}")
|
||||
logger.warning("Error logging update for %s: %s", table_name, e)
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
@@ -117,11 +166,19 @@ def after_delete_listener(mapper, connection, target):
|
||||
table_name = target.__tablename__
|
||||
record_data = {c.name: getattr(target, c.name) for c in mapper.columns}
|
||||
username = _get_current_username()
|
||||
company_id = getattr(target, "company_id", None) or getattr(target, "id", None)
|
||||
tenant_id = getattr(target, "tenant_id", None)
|
||||
|
||||
session = Session(bind=connection)
|
||||
try:
|
||||
company_id, tenant_id = _resolve_audit_company_tenant(session, target)
|
||||
if company_id is None or tenant_id is None:
|
||||
logger.debug(
|
||||
"Audit skip DELETE %s: missing company_id=%s tenant_id=%s",
|
||||
table_name,
|
||||
company_id,
|
||||
tenant_id,
|
||||
)
|
||||
return
|
||||
|
||||
AuditService.log_crud_operation(
|
||||
db=session,
|
||||
table_name=table_name,
|
||||
@@ -133,6 +190,6 @@ def after_delete_listener(mapper, connection, target):
|
||||
tenant_id=tenant_id,
|
||||
)
|
||||
except Exception as e:
|
||||
print(f"Error logging delete: {e}")
|
||||
logger.warning("Error logging delete for %s: %s", table_name, e)
|
||||
finally:
|
||||
session.close()
|
||||
|
||||
@@ -25,6 +25,19 @@ from api.v1.modules.a76.general_catalogs.company.models import Company
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
||||
def _audit_scope_tenant_id(db: Session, company_id: int) -> int:
|
||||
"""Tenant_id de la fila ``Company`` para filtrar ``audit_logs`` (alineado con lo persistido)."""
|
||||
row = (
|
||||
db.query(Company.tenant_id)
|
||||
.filter(Company.id == company_id, Company.deleted_at.is_(None))
|
||||
.first()
|
||||
)
|
||||
if not row:
|
||||
raise HTTPException(status_code=404, detail="Company not found")
|
||||
return int(row[0])
|
||||
|
||||
|
||||
_SEGMENT_LABELS = {
|
||||
"tenants": "Espacio",
|
||||
"companies": "Companias",
|
||||
@@ -189,16 +202,17 @@ async def get_bitacora(
|
||||
"""
|
||||
Bitácora por compañía. Requiere permiso ``audit_logs.view``.
|
||||
"""
|
||||
tenant_id = validate_access_to_resource(
|
||||
validate_access_to_resource(
|
||||
db,
|
||||
company_id,
|
||||
current_user,
|
||||
required_permissions=["audit_logs.view"],
|
||||
)
|
||||
scope_tenant_id = _audit_scope_tenant_id(db, company_id)
|
||||
|
||||
query = db.query(AuditLog).filter(
|
||||
AuditLog.company_id == company_id,
|
||||
AuditLog.tenant_id == tenant_id,
|
||||
AuditLog.tenant_id == scope_tenant_id,
|
||||
)
|
||||
|
||||
# Filters
|
||||
@@ -250,18 +264,19 @@ async def get_procedures(
|
||||
"""
|
||||
Lista de procedimientos para filtros (alcance compañía). Requiere ``audit_logs.view``.
|
||||
"""
|
||||
tenant_id = validate_access_to_resource(
|
||||
validate_access_to_resource(
|
||||
db,
|
||||
company_id,
|
||||
current_user,
|
||||
required_permissions=["audit_logs.view"],
|
||||
)
|
||||
scope_tenant_id = _audit_scope_tenant_id(db, company_id)
|
||||
|
||||
results = (
|
||||
db.query(AuditLog.procedure)
|
||||
.filter(
|
||||
AuditLog.company_id == company_id,
|
||||
AuditLog.tenant_id == tenant_id,
|
||||
AuditLog.tenant_id == scope_tenant_id,
|
||||
)
|
||||
.distinct()
|
||||
.order_by(AuditLog.procedure)
|
||||
@@ -279,19 +294,20 @@ async def get_audit_detail(
|
||||
"""
|
||||
Detalle de un registro de bitácora. Requiere ``audit_logs.view``.
|
||||
"""
|
||||
tenant_id = validate_access_to_resource(
|
||||
validate_access_to_resource(
|
||||
db,
|
||||
company_id,
|
||||
current_user,
|
||||
required_permissions=["audit_logs.view"],
|
||||
)
|
||||
scope_tenant_id = _audit_scope_tenant_id(db, company_id)
|
||||
|
||||
log = (
|
||||
db.query(AuditLog)
|
||||
.filter(
|
||||
AuditLog.spec_id == spec_id,
|
||||
AuditLog.company_id == company_id,
|
||||
AuditLog.tenant_id == tenant_id,
|
||||
AuditLog.tenant_id == scope_tenant_id,
|
||||
)
|
||||
.first()
|
||||
)
|
||||
@@ -316,13 +332,14 @@ async def list_tenant_files(
|
||||
if not should_ensure_s3_bucket():
|
||||
raise HTTPException(status_code=400, detail="S3 storage is disabled")
|
||||
|
||||
tenant_id = validate_access_to_resource(
|
||||
validate_access_to_resource(
|
||||
db,
|
||||
company_id,
|
||||
current_user,
|
||||
required_permissions=["audit_logs.view"],
|
||||
)
|
||||
tenant_prefix = _tenant_prefix(tenant_id)
|
||||
scope_tenant_id = _audit_scope_tenant_id(db, company_id)
|
||||
tenant_prefix = _tenant_prefix(scope_tenant_id)
|
||||
rel_path = _normalize_relative_path(path)
|
||||
list_prefix = f"{tenant_prefix}{rel_path}/" if rel_path else tenant_prefix
|
||||
|
||||
@@ -342,7 +359,7 @@ async def list_tenant_files(
|
||||
continuation_token=continuation_token,
|
||||
)
|
||||
|
||||
company_names = _companies_map(db, tenant_id)
|
||||
company_names = _companies_map(db, scope_tenant_id)
|
||||
folders: List[AuditFileFolderItem] = []
|
||||
for prefix in data.get("prefixes", []):
|
||||
rel = _relative_from_tenant_prefix(prefix, tenant_prefix)
|
||||
@@ -410,13 +427,14 @@ async def download_tenant_file(
|
||||
if not should_ensure_s3_bucket():
|
||||
raise HTTPException(status_code=400, detail="S3 storage is disabled")
|
||||
|
||||
tenant_id = validate_access_to_resource(
|
||||
validate_access_to_resource(
|
||||
db,
|
||||
company_id,
|
||||
current_user,
|
||||
required_permissions=["audit_logs.view"],
|
||||
)
|
||||
tenant_prefix = _tenant_prefix(tenant_id)
|
||||
scope_tenant_id = _audit_scope_tenant_id(db, company_id)
|
||||
tenant_prefix = _tenant_prefix(scope_tenant_id)
|
||||
rel_path = _normalize_relative_path(path)
|
||||
if not rel_path or rel_path.endswith("/"):
|
||||
raise HTTPException(status_code=400, detail="A file path is required")
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"""
|
||||
Audit Log Service
|
||||
"""
|
||||
import logging
|
||||
from datetime import datetime, timedelta, date, time
|
||||
from decimal import Decimal
|
||||
import uuid
|
||||
@@ -9,8 +10,8 @@ from typing import Optional, List, Dict, Any
|
||||
from sqlalchemy.orm import Session
|
||||
from ..models import AuditLog
|
||||
from .core import AuditMapper, ReferenceGenerator
|
||||
from core.security import verify_token # keep if needed or simpler just remove if unused
|
||||
# We don't need security import here anymore as context is passed explicitly or handled by events
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _make_json_safe(obj: Any) -> Any:
|
||||
@@ -188,30 +189,35 @@ class AuditService:
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def log_login(db: Session, username: str, ip_address: str = None, user_agent: str = None):
|
||||
def log_login(
|
||||
db: Session,
|
||||
username: str,
|
||||
ip_address: str = None,
|
||||
user_agent: str = None,
|
||||
company_id: Optional[int] = None,
|
||||
tenant_id: Optional[int] = None,
|
||||
):
|
||||
"""
|
||||
``audit_logs`` exige ``tenant_id`` y ``company_id``. El login vía Hub no
|
||||
define compañía activa; sin ambos argumentos no se inserta fila (antes fallaba NOT NULL).
|
||||
"""
|
||||
if company_id is None or tenant_id is None:
|
||||
return None
|
||||
|
||||
# Prevent duplicate login logs (debounce 5 seconds)
|
||||
# This handles cases where frontend might submit twice or redirects trigger re-auth
|
||||
try:
|
||||
# Timezone handling: Use UTC for consistency
|
||||
now = datetime.now(pytz.UTC)
|
||||
|
||||
five_seconds_ago = now - timedelta(seconds=5)
|
||||
|
||||
# Check for recent login from same user
|
||||
existing = db.query(AuditLog).filter(
|
||||
AuditLog.username == username,
|
||||
AuditLog.operation_type == "LOGIN",
|
||||
# Compare against timestamp (timezone aware)
|
||||
AuditLog.timestamp >= five_seconds_ago
|
||||
AuditLog.company_id == company_id,
|
||||
AuditLog.tenant_id == tenant_id,
|
||||
AuditLog.timestamp >= five_seconds_ago,
|
||||
).first()
|
||||
|
||||
if existing:
|
||||
return existing
|
||||
|
||||
except Exception as e:
|
||||
import traceback
|
||||
traceback.print_exc()
|
||||
logger.warning("Login audit debounce query failed: %s", e)
|
||||
|
||||
return AuditService.create_audit_log(
|
||||
db=db,
|
||||
@@ -221,16 +227,24 @@ class AuditService:
|
||||
username=username,
|
||||
operation_type="LOGIN",
|
||||
ip_address=ip_address,
|
||||
user_agent=user_agent
|
||||
user_agent=user_agent,
|
||||
company_id=company_id,
|
||||
tenant_id=tenant_id,
|
||||
)
|
||||
|
||||
@staticmethod
|
||||
def log_logout(db: Session, username: str, ip_address: str = None, user_agent: str = None):
|
||||
"""
|
||||
Registra un evento de cierre de sesión
|
||||
"""
|
||||
def log_logout(
|
||||
db: Session,
|
||||
username: str,
|
||||
ip_address: str = None,
|
||||
user_agent: str = None,
|
||||
company_id: Optional[int] = None,
|
||||
tenant_id: Optional[int] = None,
|
||||
):
|
||||
"""Misma condición que ``log_login``: sin alcance compañía/tenant no se escribe."""
|
||||
if company_id is None or tenant_id is None:
|
||||
return None
|
||||
try:
|
||||
# Reutilizamos create_audit_log para mantener consistencia
|
||||
AuditService.create_audit_log(
|
||||
db=db,
|
||||
reference="LOGOUT",
|
||||
@@ -239,8 +253,9 @@ class AuditService:
|
||||
username=username,
|
||||
operation_type="LOGOUT",
|
||||
ip_address=ip_address,
|
||||
user_agent=user_agent
|
||||
user_agent=user_agent,
|
||||
company_id=company_id,
|
||||
tenant_id=tenant_id,
|
||||
)
|
||||
except Exception as e:
|
||||
# No re-lanzamos la excepción para no interrumpir el flujo de logout
|
||||
pass
|
||||
logger.warning("Logout audit insert failed: %s", e)
|
||||
|
||||
@@ -136,8 +136,13 @@ def get_core_db(request: Request = None) -> Generator[Session, None, None]:
|
||||
FastAPI inyecta ``Request`` automáticamente; los llamadores existentes que
|
||||
escriben ``db: Session = Depends(get_core_db)`` siguen funcionando sin
|
||||
cambios porque ``Request`` se resuelve en la capa de dependencia.
|
||||
|
||||
Replica el mismo ``(tenant_id, company_id)`` en ContextVars para código que
|
||||
comparte la transacción sin la misma instancia de sesión (p. ej. listeners).
|
||||
"""
|
||||
tenant_id, company_id = _extract_rls_context(request)
|
||||
token_t = rls_tenant_var.set(tenant_id)
|
||||
token_c = rls_company_var.set(company_id)
|
||||
db = CoreSessionLocal()
|
||||
db.info[RLS_TENANT_KEY] = tenant_id
|
||||
db.info[RLS_COMPANY_KEY] = company_id
|
||||
@@ -145,11 +150,16 @@ def get_core_db(request: Request = None) -> Generator[Session, None, None]:
|
||||
yield db
|
||||
finally:
|
||||
db.close()
|
||||
rls_tenant_var.reset(token_t)
|
||||
rls_company_var.reset(token_c)
|
||||
|
||||
|
||||
async def get_async_core_db(request: Request = None) -> AsyncGenerator[AsyncSession, None]:
|
||||
"""Dependency async para obtener sesión con contexto RLS."""
|
||||
tenant_id, company_id = _extract_rls_context(request)
|
||||
token_t = rls_tenant_var.set(tenant_id)
|
||||
token_c = rls_company_var.set(company_id)
|
||||
try:
|
||||
async with AsyncCoreSessionLocal() as session:
|
||||
session.info[RLS_TENANT_KEY] = tenant_id
|
||||
session.info[RLS_COMPANY_KEY] = company_id
|
||||
@@ -157,6 +167,9 @@ async def get_async_core_db(request: Request = None) -> AsyncGenerator[AsyncSess
|
||||
yield session
|
||||
finally:
|
||||
await session.close()
|
||||
finally:
|
||||
rls_tenant_var.reset(token_t)
|
||||
rls_company_var.reset(token_c)
|
||||
|
||||
|
||||
@contextmanager
|
||||
|
||||
Reference in New Issue
Block a user