251 lines
8.6 KiB
Python
251 lines
8.6 KiB
Python
from __future__ import annotations
|
|
|
|
from datetime import datetime, timedelta, timezone
|
|
from pathlib import Path
|
|
from typing import Any
|
|
import json
|
|
import os
|
|
import uuid
|
|
|
|
from sqlalchemy import cast, func, or_, Text
|
|
|
|
from .database import SessionLocal
|
|
from .models import SoftwareJob
|
|
|
|
PRIVACY_DELETION_LOG = Path(os.getenv("PRIVACY_DELETION_LOG", "/app/data/logs/privacy-deletion-audit.log"))
|
|
APP_LOG_DIR = Path(os.getenv("APP_LOG_DIR", "/app/data/logs"))
|
|
|
|
IMPLEMENTED_RETENTION_KEYS = {"am_job_payloads", "am_diagnostic_logs"}
|
|
PROTECTED_LOG_FILES = {
|
|
"privacy-policy-audit.log",
|
|
"privacy-deletion-audit.log",
|
|
}
|
|
ACTIVE_LOG_FILES = {"errors.log", "ldap.log", "software-callback-debug.log"}
|
|
TERMINAL_JOB_STATUSES = {"success", "failed", "partial", "timeout", "cancelled"}
|
|
|
|
|
|
def _utcnow() -> datetime:
|
|
return datetime.now(timezone.utc)
|
|
|
|
|
|
def append_deletion_audit(
|
|
action: str,
|
|
changed_by: str,
|
|
*,
|
|
category_key: str,
|
|
retention_days: int | None,
|
|
result: dict[str, Any],
|
|
run_id: str | None = None,
|
|
) -> str:
|
|
run_id = run_id or uuid.uuid4().hex
|
|
PRIVACY_DELETION_LOG.parent.mkdir(parents=True, exist_ok=True)
|
|
record = {
|
|
"timestamp_utc": _utcnow().isoformat(),
|
|
"run_id": run_id,
|
|
"action": action,
|
|
"changed_by": changed_by or "unknown",
|
|
"category_key": category_key,
|
|
"retention_days": retention_days,
|
|
"result": result,
|
|
}
|
|
with PRIVACY_DELETION_LOG.open("a", encoding="utf-8") as handle:
|
|
handle.write(json.dumps(record, ensure_ascii=False, sort_keys=True) + "\n")
|
|
return run_id
|
|
|
|
|
|
def deletion_audit_tail(limit: int = 100) -> list[dict[str, Any]]:
|
|
if not PRIVACY_DELETION_LOG.exists():
|
|
return []
|
|
try:
|
|
lines = PRIVACY_DELETION_LOG.read_text(encoding="utf-8", errors="replace").splitlines()[-max(1, limit):]
|
|
except OSError:
|
|
return []
|
|
rows: list[dict[str, Any]] = []
|
|
for line in reversed(lines):
|
|
try:
|
|
rows.append(json.loads(line))
|
|
except json.JSONDecodeError:
|
|
continue
|
|
return rows
|
|
|
|
|
|
def _job_cutoff_expression():
|
|
return func.coalesce(SoftwareJob.finished_at, SoftwareJob.callback_completed_at, SoftwareJob.created_at)
|
|
|
|
|
|
def _job_payload_query(db, cutoff: datetime):
|
|
content_present = or_(
|
|
SoftwareJob.command_preview.isnot(None),
|
|
SoftwareJob.resolved_script.isnot(None),
|
|
SoftwareJob.meshctrl_stdout.isnot(None),
|
|
SoftwareJob.meshctrl_stderr.isnot(None),
|
|
SoftwareJob.message.isnot(None),
|
|
SoftwareJob.log_text != "",
|
|
cast(SoftwareJob.result_data, Text) != "{}",
|
|
)
|
|
return db.query(SoftwareJob).filter(
|
|
SoftwareJob.status.in_(TERMINAL_JOB_STATUSES),
|
|
_job_cutoff_expression() < cutoff,
|
|
content_present,
|
|
)
|
|
|
|
|
|
def _check_job_payloads(retention_days: int) -> dict[str, Any]:
|
|
cutoff = datetime.utcnow() - timedelta(days=retention_days)
|
|
db = SessionLocal()
|
|
try:
|
|
query = _job_payload_query(db, cutoff)
|
|
count = query.count()
|
|
oldest = query.with_entities(func.min(_job_cutoff_expression())).scalar()
|
|
# Approximate payload size without transferring the payloads to Python.
|
|
byte_expression = (
|
|
func.length(func.coalesce(SoftwareJob.command_preview, ""))
|
|
+ func.length(func.coalesce(SoftwareJob.resolved_script, ""))
|
|
+ func.length(func.coalesce(SoftwareJob.meshctrl_stdout, ""))
|
|
+ func.length(func.coalesce(SoftwareJob.meshctrl_stderr, ""))
|
|
+ func.length(func.coalesce(SoftwareJob.message, ""))
|
|
+ func.length(func.coalesce(SoftwareJob.log_text, ""))
|
|
+ func.length(func.coalesce(cast(SoftwareJob.result_data, Text), ""))
|
|
)
|
|
approx_bytes = query.with_entities(func.coalesce(func.sum(byte_expression), 0)).scalar() or 0
|
|
return {
|
|
"supported": True,
|
|
"category_key": "am_job_payloads",
|
|
"cutoff_utc": cutoff.isoformat(),
|
|
"matched_records": int(count),
|
|
"approx_bytes": int(approx_bytes),
|
|
"oldest_record_utc": oldest.isoformat() if oldest else None,
|
|
"operation": "clear_job_payloads",
|
|
}
|
|
finally:
|
|
db.close()
|
|
|
|
|
|
def _delete_job_payloads(retention_days: int) -> dict[str, Any]:
|
|
preview = _check_job_payloads(retention_days)
|
|
cutoff = datetime.utcnow() - timedelta(days=retention_days)
|
|
db = SessionLocal()
|
|
try:
|
|
updated = _job_payload_query(db, cutoff).update(
|
|
{
|
|
SoftwareJob.command_preview: None,
|
|
SoftwareJob.resolved_script: None,
|
|
SoftwareJob.meshctrl_stdout: None,
|
|
SoftwareJob.meshctrl_stderr: None,
|
|
SoftwareJob.message: None,
|
|
SoftwareJob.log_text: "",
|
|
SoftwareJob.result_data: {},
|
|
},
|
|
synchronize_session=False,
|
|
)
|
|
db.commit()
|
|
preview["deleted_or_anonymised_records"] = int(updated)
|
|
preview["status"] = "success"
|
|
return preview
|
|
except Exception:
|
|
db.rollback()
|
|
raise
|
|
finally:
|
|
db.close()
|
|
|
|
|
|
def _eligible_log_files(cutoff: datetime) -> list[Path]:
|
|
if not APP_LOG_DIR.exists():
|
|
return []
|
|
result: list[Path] = []
|
|
cutoff_timestamp = cutoff.timestamp()
|
|
for path in APP_LOG_DIR.rglob("*"):
|
|
try:
|
|
if not path.is_file() or path.name in PROTECTED_LOG_FILES:
|
|
continue
|
|
if path.stat().st_mtime < cutoff_timestamp:
|
|
result.append(path)
|
|
except OSError:
|
|
continue
|
|
return result
|
|
|
|
|
|
def _check_diagnostic_logs(retention_days: int) -> dict[str, Any]:
|
|
cutoff = _utcnow() - timedelta(days=retention_days)
|
|
files = _eligible_log_files(cutoff)
|
|
total_bytes = 0
|
|
oldest: datetime | None = None
|
|
for path in files:
|
|
try:
|
|
stat = path.stat()
|
|
except OSError:
|
|
continue
|
|
total_bytes += stat.st_size
|
|
modified = datetime.fromtimestamp(stat.st_mtime, tz=timezone.utc)
|
|
oldest = modified if oldest is None or modified < oldest else oldest
|
|
return {
|
|
"supported": True,
|
|
"category_key": "am_diagnostic_logs",
|
|
"cutoff_utc": cutoff.isoformat(),
|
|
"matched_records": len(files),
|
|
"approx_bytes": int(total_bytes),
|
|
"oldest_record_utc": oldest.isoformat() if oldest else None,
|
|
"operation": "delete_log_files",
|
|
}
|
|
|
|
|
|
def _delete_diagnostic_logs(retention_days: int) -> dict[str, Any]:
|
|
cutoff = _utcnow() - timedelta(days=retention_days)
|
|
files = _eligible_log_files(cutoff)
|
|
deleted = 0
|
|
failed = 0
|
|
freed_bytes = 0
|
|
for path in files:
|
|
try:
|
|
size = path.stat().st_size
|
|
if path.name in ACTIVE_LOG_FILES:
|
|
# Keep the inode used by the active FileHandler and clear only its content.
|
|
path.write_text("", encoding="utf-8")
|
|
else:
|
|
path.unlink()
|
|
deleted += 1
|
|
freed_bytes += size
|
|
except OSError:
|
|
failed += 1
|
|
return {
|
|
"supported": True,
|
|
"category_key": "am_diagnostic_logs",
|
|
"cutoff_utc": cutoff.isoformat(),
|
|
"matched_records": len(files),
|
|
"deleted_or_anonymised_records": deleted,
|
|
"failed_records": failed,
|
|
"approx_bytes": int(freed_bytes),
|
|
"operation": "delete_log_files",
|
|
"status": "success" if failed == 0 else "partial",
|
|
}
|
|
|
|
|
|
def check_retention_category(category_key: str, retention_days: int | None) -> dict[str, Any]:
|
|
if category_key not in IMPLEMENTED_RETENTION_KEYS:
|
|
return {
|
|
"supported": False,
|
|
"category_key": category_key,
|
|
"matched_records": 0,
|
|
"reason": "No executable retention handler is implemented for this category.",
|
|
}
|
|
if retention_days is None or retention_days < 1:
|
|
return {
|
|
"supported": False,
|
|
"category_key": category_key,
|
|
"matched_records": 0,
|
|
"reason": "A retention period of at least one day is required.",
|
|
}
|
|
if category_key == "am_job_payloads":
|
|
return _check_job_payloads(retention_days)
|
|
return _check_diagnostic_logs(retention_days)
|
|
|
|
|
|
def delete_retention_category(category_key: str, retention_days: int | None) -> dict[str, Any]:
|
|
if category_key not in IMPLEMENTED_RETENTION_KEYS:
|
|
raise ValueError("No executable retention handler is implemented for this category.")
|
|
if retention_days is None or retention_days < 1:
|
|
raise ValueError("A retention period of at least one day is required.")
|
|
if category_key == "am_job_payloads":
|
|
return _delete_job_payloads(retention_days)
|
|
return _delete_diagnostic_logs(retention_days)
|