diff --git a/app/admin/__init__.py b/app/admin/__init__.py new file mode 100644 index 0000000..87c5f3e --- /dev/null +++ b/app/admin/__init__.py @@ -0,0 +1 @@ +"""Restricted administrative reporting modules.""" diff --git a/app/admin/archive_bundle.py b/app/admin/archive_bundle.py new file mode 100644 index 0000000..2f39d84 --- /dev/null +++ b/app/admin/archive_bundle.py @@ -0,0 +1,254 @@ +"""Portable, checksummed Project CT archive bundles. Never contacts vTiger.""" + +from __future__ import annotations + +import hashlib +import json +import os +import socket +import tempfile +import zipfile +from datetime import datetime, timezone +from typing import Any, Iterable + +from psycopg2.extras import RealDictCursor + +from app.core.database import get_db_connection, release_db_connection + +BUNDLE_SCHEMA_VERSION = 1 +ENTRY_NAMES = ("versions.jsonl", "records.jsonl", "relations.jsonl", "checkpoints.jsonl", "files.jsonl") + + +def file_sha256(path: str) -> str: + digest = hashlib.sha256() + with open(path, "rb") as source: + for chunk in iter(lambda: source.read(1024 * 1024), b""): + digest.update(chunk) + return digest.hexdigest() + + +def _json_line(value: Any) -> bytes: + return (json.dumps(value, ensure_ascii=False, sort_keys=True, default=str) + "\n").encode("utf-8") + + +def _write_cursor_entry(archive: zipfile.ZipFile, name: str, cursor, batch_size: int = 500) -> tuple[int, str]: + digest = hashlib.sha256() + count = 0 + with archive.open(name, "w") as target: + while True: + rows = cursor.fetchmany(batch_size) + if not rows: + break + for row in rows: + raw = _json_line(dict(row)) + target.write(raw) + digest.update(raw) + count += 1 + return count, digest.hexdigest() + + +def create_archive_bundle(through_version_id: int, initiated_by: int | None = None) -> tuple[str, dict]: + fd, path = tempfile.mkstemp(prefix="project-ct-", suffix=".zip") + os.close(fd) + conn = get_db_connection() + manifest = { + "schema_version": BUNDLE_SCHEMA_VERSION, + "created_at": datetime.now(timezone.utc).isoformat(), + "source_instance": socket.gethostname(), + "through_version_id": through_version_id, + "entries": {}, + "files": {}, + "contains_files": True, + "document_policy": "metadata_and_original_bytes", + } + try: + with zipfile.ZipFile(path, "w", compression=zipfile.ZIP_DEFLATED, compresslevel=6) as archive: + queries = { + "versions.jsonl": ("SELECT * FROM vtiger_archive_versions WHERE id<=%s ORDER BY id", (through_version_id,)), + "records.jsonl": ("SELECT * FROM vtiger_archive_records WHERE version_id<=%s ORDER BY module,vtiger_id,revision_no", (through_version_id,)), + "relations.jsonl": ("SELECT * FROM vtiger_archive_relations WHERE version_id<=%s ORDER BY id", (through_version_id,)), + "checkpoints.jsonl": ( + """SELECT module,MAX(source_modified_at) AS last_modified_at,MAX(vtiger_id) AS last_vtiger_id, + %s::bigint AS last_successful_version_id,NOW() AS updated_at + FROM vtiger_archive_records WHERE version_id<=%s GROUP BY module ORDER BY module""", + (through_version_id, through_version_id), + ), + "files.jsonl": ( + """SELECT id,version_id,source_module,document_vtiger_id,resource_vtiger_id,filename,content_type, + size_bytes,content_sha256,archived_at + FROM vtiger_archive_files WHERE version_id<=%s ORDER BY id""", + (through_version_id,), + ), + } + for index, (name, (sql, params)) in enumerate(queries.items()): + cursor = conn.cursor(name=f"project_ct_export_{index}", cursor_factory=RealDictCursor) + cursor.itersize = 500 + cursor.execute(sql, params) + count, digest = _write_cursor_entry(archive, name, cursor) + cursor.close() + manifest["entries"][name] = {"count": count, "sha256": digest} + file_rows = conn.cursor(cursor_factory=RealDictCursor) + file_rows.execute( + """SELECT DISTINCT ON (content_sha256) content_sha256,content,size_bytes + FROM vtiger_archive_files WHERE version_id<=%s ORDER BY content_sha256,id""", + (through_version_id,), + ) + for row in file_rows: + content = bytes(row["content"]) + digest = hashlib.sha256(content).hexdigest() + if digest != row["content_sha256"]: + raise ValueError(f"Lokal filchecksum er ugyldig: {row['content_sha256']}") + entry_name = f"files/{digest}" + archive.writestr(entry_name, content) + manifest["files"][digest] = {"entry": entry_name, "size": len(content), "sha256": digest} + file_rows.close() + archive.writestr("manifest.json", json.dumps(manifest, ensure_ascii=False, sort_keys=True, indent=2, default=str)) + manifest["bundle_sha256"] = file_sha256(path) + return path, manifest + except Exception: + if os.path.exists(path): + os.unlink(path) + raise + finally: + release_db_connection(conn) + + +def _verify_entry(archive: zipfile.ZipFile, manifest: dict, name: str) -> None: + expected = ((manifest.get("entries") or {}).get(name) or {}).get("sha256") + if not expected: + raise ValueError(f"Manifest mangler checksum for {name}") + digest = hashlib.sha256() + with archive.open(name) as source: + for chunk in iter(lambda: source.read(1024 * 1024), b""): + digest.update(chunk) + if digest.hexdigest() != expected: + raise ValueError(f"Checksumfejl i {name}") + + +def _rows(archive: zipfile.ZipFile, name: str) -> Iterable[dict]: + with archive.open(name) as source: + for line_number, raw in enumerate(source, 1): + if not raw.strip(): + continue + try: + yield json.loads(raw) + except json.JSONDecodeError as exc: + raise ValueError(f"Ugyldig JSON i {name}, linje {line_number}") from exc + + +def import_archive_bundle(path: str, initiated_by: int | None = None) -> dict: + bundle_digest = file_sha256(path) + conn = get_db_connection() + counts = {"versions": 0, "records": 0, "records_skipped": 0, "relations": 0, "checkpoints": 0} + try: + with zipfile.ZipFile(path) as archive: + try: + manifest = json.loads(archive.read("manifest.json")) + except (KeyError, json.JSONDecodeError) as exc: + raise ValueError("Filen er ikke en gyldig Project CT-arkivpakke") from exc + if manifest.get("schema_version") != BUNDLE_SCHEMA_VERSION: + raise ValueError(f"Ikke-understøttet arkivformat: {manifest.get('schema_version')}") + if manifest.get("contains_files") is not True: + raise ValueError("Pakken mangler originale vTiger-dokumentfiler") + for name in ENTRY_NAMES: + _verify_entry(archive, manifest, name) + for digest, info in (manifest.get("files") or {}).items(): + entry_name = info.get("entry") + if not entry_name or entry_name not in archive.namelist(): + raise ValueError(f"Pakken mangler dokumentfil {digest}") + actual = hashlib.sha256(archive.read(entry_name)).hexdigest() + if actual != digest or actual != info.get("sha256"): + raise ValueError(f"Checksumfejl i dokumentfil {digest}") + + with conn.cursor(cursor_factory=RealDictCursor) as cursor: + version_map: dict[int, int] = {} + for row in _rows(archive, "versions.jsonl"): + cursor.execute( + """INSERT INTO vtiger_archive_versions + (sync_kind,status,started_at,completed_at,source_cutoff,module_counts,warnings, + critical_errors,control_report,control_approved_at,raw_export_sha256) + VALUES (%s,%s,%s,%s,%s,%s::jsonb,%s::jsonb,%s::jsonb,%s::jsonb,%s,%s) RETURNING id""", + (row["sync_kind"], row["status"], row["started_at"], row.get("completed_at"), row["source_cutoff"], + json.dumps(row.get("module_counts") or {}), json.dumps(row.get("warnings") or []), + json.dumps(row.get("critical_errors") or []), json.dumps(row.get("control_report") or {}), + row.get("control_approved_at"), row.get("raw_export_sha256")), + ) + version_map[int(row["id"])] = int(cursor.fetchone()["id"]) + counts["versions"] += 1 + + for row in _rows(archive, "records.jsonl"): + cursor.execute( + "SELECT id FROM vtiger_archive_records WHERE module=%s AND vtiger_id=%s AND payload_sha256=%s", + (row["module"], row["vtiger_id"], row["payload_sha256"]), + ) + if cursor.fetchone(): + counts["records_skipped"] += 1 + continue + cursor.execute( + "SELECT COALESCE(MAX(revision_no),0)+1 AS revision FROM vtiger_archive_records WHERE module=%s AND vtiger_id=%s", + (row["module"], row["vtiger_id"]), + ) + revision = int(cursor.fetchone()["revision"]) + cursor.execute( + """INSERT INTO vtiger_archive_records + (version_id,module,vtiger_id,revision_no,source_created_at,source_modified_at,is_deleted,payload,payload_sha256,archived_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s::jsonb,%s,%s)""", + (version_map[int(row["version_id"])], row["module"], row["vtiger_id"], revision, + row.get("source_created_at"), row.get("source_modified_at"), row.get("is_deleted", False), + json.dumps(row["payload"], ensure_ascii=False), row["payload_sha256"], row.get("archived_at")), + ) + counts["records"] += 1 + + counts["files"] = 0 + counts["files_skipped"] = 0 + for row in _rows(archive, "files.jsonl"): + cursor.execute( + """SELECT id FROM vtiger_archive_files + WHERE source_module=%s AND document_vtiger_id=%s AND resource_vtiger_id=%s AND content_sha256=%s""", + (row.get("source_module") or "Documents", row["document_vtiger_id"], row["resource_vtiger_id"], row["content_sha256"]), + ) + if cursor.fetchone(): + counts["files_skipped"] += 1 + continue + info = manifest["files"][row["content_sha256"]] + content = archive.read(info["entry"]) + cursor.execute( + """INSERT INTO vtiger_archive_files + (version_id,source_module,document_vtiger_id,resource_vtiger_id,filename,content_type,size_bytes, + content_sha256,content,archived_at) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s)""", + (version_map[int(row["version_id"])], row.get("source_module") or "Documents", row["document_vtiger_id"], row["resource_vtiger_id"], + row["filename"], row.get("content_type"), row["size_bytes"], row["content_sha256"], content, + row.get("archived_at")), + ) + counts["files"] += 1 + + for row in _rows(archive, "relations.jsonl"): + cursor.execute( + """INSERT INTO vtiger_archive_relations + (version_id,source_module,source_vtiger_id,field_name,target_vtiger_id,target_module) + VALUES (%s,%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING""", + (version_map[int(row["version_id"])], row["source_module"], row["source_vtiger_id"], + row["field_name"], row["target_vtiger_id"], row.get("target_module")), + ) + counts["relations"] += cursor.rowcount + + for row in _rows(archive, "checkpoints.jsonl"): + mapped_version = version_map.get(int(row["last_successful_version_id"])) if row.get("last_successful_version_id") else None + cursor.execute( + """INSERT INTO vtiger_archive_checkpoints + (module,last_modified_at,last_vtiger_id,last_successful_version_id,updated_at) + VALUES (%s,%s,%s,%s,%s) ON CONFLICT(module) DO UPDATE SET + last_modified_at=GREATEST(vtiger_archive_checkpoints.last_modified_at,EXCLUDED.last_modified_at), + last_vtiger_id=EXCLUDED.last_vtiger_id,last_successful_version_id=EXCLUDED.last_successful_version_id, + updated_at=NOW()""", + (row["module"], row.get("last_modified_at"), row.get("last_vtiger_id"), mapped_version, row.get("updated_at")), + ) + counts["checkpoints"] += 1 + conn.commit() + return {"bundle_sha256": bundle_digest, "manifest": manifest, "counts": counts, "contacted_vtiger": False} + except Exception: + conn.rollback() + raise + finally: + release_db_connection(conn) diff --git a/app/admin/hub_impact.html b/app/admin/hub_impact.html new file mode 100644 index 0000000..eebb4c5 --- /dev/null +++ b/app/admin/hub_impact.html @@ -0,0 +1,60 @@ +{% extends "shared/frontend/base.html" %} +{% block title %}Project CT · Hub Impact{% endblock %} +{% block extra_css %} + +{% endblock %} +{% block content %} +
+
+
Fortrolig · kun superadmin

Project CT · Hub Impact

Dokumentér effekten af Hub mod et permanent vTiger-arkiv.

+
Indlæser arkivstatus…
+
+ +
+

Ny reproducerbar sammenligning

Perioderne skal have præcis samme antal kalenderdage.
Download pakke til prod
+
+
+
+
+
+
+ +
+
+ +
+ +
+

Hvad Hub effektiviserer og optimerer

Positive og negative ændringer vises ensartet. Tallene sammenligner lige lange perioder og er dokumentation—ikke en skjult vægtet score.
+

Pr. medarbejder

KildeMedarbejderE-mailTimerRegistreringerSagerOrdrerTimer/dag
+

Afvigelser udeladt fra KPI

Datakvalitet

+
+ +

Gemte rapporter

+

Flytning mellem test og produktion

Checksummede arkivpakker. Import læser kun ZIP-filen og kontakter aldrig vTiger.
Ingen overførsler endnu.
+
+{% endblock %} +{% block extra_js %} + +{% endblock %} diff --git a/app/admin/hub_impact.py b/app/admin/hub_impact.py new file mode 100644 index 0000000..1a0916f --- /dev/null +++ b/app/admin/hub_impact.py @@ -0,0 +1,235 @@ +"""Reproducible vTiger-versus-Hub impact metrics for Project CT.""" + +from __future__ import annotations + +import hashlib +import json +from collections import defaultdict +from datetime import date, datetime, timedelta +from decimal import Decimal, InvalidOperation +from typing import Any + +from app.core.database import execute_insert, execute_query, execute_query_single +from app.services.vtiger_service import VTigerService + +DATE_FIELDS = ("worked_date", "date", "createdtime", "created_at", "modifiedtime") +EMAIL_FIELDS = ("email1", "email", "email_address", "user_email") +USER_REF_FIELDS = ("assigned_user_id", "creator", "created_by", "userid", "user_id") + + +def _date(value: Any) -> date | None: + if isinstance(value, datetime): + return value.date() + if isinstance(value, date): + return value + try: + return datetime.fromisoformat(str(value or "").strip().replace("Z", "+00:00")).date() + except (ValueError, TypeError): + return None + + +def _first(payload: dict, fields: tuple[str, ...]): + return next((payload.get(field) for field in fields if payload.get(field) not in (None, "")), None) + + +def _hours(payload: dict) -> Decimal: + return VTigerService._extract_timelog_hours(payload) + + +def _archive_snapshot(version_id: int) -> list[dict[str, Any]]: + return execute_query( + """SELECT DISTINCT ON (module,vtiger_id) module,vtiger_id,payload,is_deleted + FROM vtiger_archive_records WHERE version_id <= %s + ORDER BY module,vtiger_id,revision_no DESC""", (version_id,), + ) or [] + + +def _hub_users() -> tuple[dict[int, dict], dict[str, dict]]: + rows = execute_query("SELECT user_id,email,full_name,username FROM users") or [] + by_id = {int(row["user_id"]): dict(row) for row in rows} + by_email = {str(row.get("email") or "").strip().lower(): dict(row) for row in rows if row.get("email")} + return by_id, by_email + + +def _employee_bucket(store: dict, key: str, name: str, email: str | None): + if key not in store: + store[key] = {"key": key, "name": name, "email": email, "hours": Decimal("0"), + "time_entries": 0, "cases": 0, "orders": 0} + return store[key] + + +def _finalize(metrics: dict[str, Any], start: date, end: date) -> dict[str, Any]: + days = (end - start).days + 1 + workdays = max(sum(1 for offset in range(days) if (start + timedelta(days=offset)).weekday() < 5), 1) + employees = [] + for item in metrics["employees"].values(): + normalized = dict(item) + normalized["hours"] = round(float(item["hours"]), 2) + normalized["productivity"] = { + "hours_per_workday": round(normalized["hours"] / workdays, 2), + "registrations_per_workday": round(normalized["time_entries"] / workdays, 2), + "cases_per_workday": round(normalized["cases"] / workdays, 2), + "orders_per_workday": round(normalized["orders"] / workdays, 2), + "deliveries_per_workday": round((normalized["time_entries"] + normalized["cases"] + normalized["orders"]) / workdays, 2), + } + employees.append(normalized) + employees.sort(key=lambda row: (row.get("name") or "").lower()) + totals = {key: sum(float(row[key]) for row in employees) for key in ("hours", "time_entries", "cases", "orders")} + totals["hours"] = round(totals["hours"], 2) + totals["time_entries"] = int(totals["time_entries"]) + totals["cases"] = int(totals["cases"]) + totals["orders"] = int(totals["orders"]) + totals["productivity"] = { + "hours_per_workday": round(totals["hours"] / workdays, 2), + "registrations_per_workday": round(totals["time_entries"] / workdays, 2), + "cases_per_workday": round(totals["cases"] / workdays, 2), + "orders_per_workday": round(totals["orders"] / workdays, 2), + "deliveries_per_workday": round((totals["time_entries"] + totals["cases"] + totals["orders"]) / workdays, 2), + "hours_per_case_or_order": round(totals["hours"] / max(totals["cases"] + totals["orders"], 1), 2), + } + return {"totals": totals, "employees": employees, "anomalies": metrics["anomalies"], + "data_quality": metrics["data_quality"], "workdays": workdays} + + +def _vtiger_metrics(records: list[dict], start: date, end: date, hub_email_users: dict[str, dict]) -> dict: + active = [row for row in records if not row.get("is_deleted")] + vt_users = {} + for row in active: + if row["module"] != "Users": + continue + payload = row["payload"] or {} + email = str(_first(payload, EMAIL_FIELDS) or "").strip().lower() + vt_users[row["vtiger_id"]] = { + "email": email or None, + "name": str(payload.get("first_name") or "") + " " + str(payload.get("last_name") or payload.get("user_name") or row["vtiger_id"]), + } + result = {"employees": {}, "anomalies": [], "data_quality": {"unmatched_employees": [], "missing_dates": 0}} + module_metric = {"Timelog": "time_entries", "Cases": "cases", "SalesOrder": "orders"} + for row in active: + metric = module_metric.get(row["module"]) + if not metric: + continue + payload = row["payload"] or {} + occurred = _date(_first(payload, DATE_FIELDS)) + if not occurred: + result["data_quality"]["missing_dates"] += 1 + continue + if occurred < start or occurred > end: + continue + user_ref = str(_first(payload, USER_REF_FIELDS) or "") + vt_user = vt_users.get(user_ref, {"email": None, "name": user_ref or "Ukendt vTiger-bruger"}) + email = vt_user.get("email") + matched = hub_email_users.get(email) if email else None + key = f"hub:{matched['user_id']}" if matched else f"vtiger:{user_ref or 'unknown'}" + name = str(matched.get("full_name") or matched.get("username")) if matched else str(vt_user.get("name") or key).strip() + bucket = _employee_bucket(result["employees"], key, name, email) + if not matched and key not in result["data_quality"]["unmatched_employees"]: + result["data_quality"]["unmatched_employees"].append(key) + if metric == "time_entries": + hours = _hours(payload) + if hours > 16: + result["anomalies"].append({"source": "vtiger", "type": "over_16_hours", "id": row["vtiger_id"], "hours": float(hours)}) + continue + bucket["hours"] += hours + bucket[metric] += 1 + return result + + +def _hub_metrics(start: date, end: date, hub_users: dict[int, dict]) -> dict: + result = {"employees": {}, "anomalies": [], "data_quality": {"unmatched_employees": [], "missing_dates": 0}} + times = execute_query( + """SELECT id,medarbejder_id,worked_date,start_tid,created_at,aktiv_timer, + COALESCE(faktisk_tid_min::numeric/60,approved_hours,original_hours,0) AS hours + FROM tmodule_times WHERE vtiger_id IS NULL + AND COALESCE(worked_date,start_tid::date,created_at::date) BETWEEN %s AND %s""", (start, end), + ) or [] + for row in times: + user = hub_users.get(int(row["medarbejder_id"])) if row.get("medarbejder_id") else None + key = f"hub:{user['user_id']}" if user else "hub:unknown" + bucket = _employee_bucket(result["employees"], key, + str((user or {}).get("full_name") or (user or {}).get("username") or "Ukendt Hub-bruger"), + (user or {}).get("email")) + if not user and key not in result["data_quality"]["unmatched_employees"]: + result["data_quality"]["unmatched_employees"].append(key) + hours = Decimal(str(row.get("hours") or 0)) + if row.get("aktiv_timer") or hours > 16: + result["anomalies"].append({"source": "hub", "type": "active_timer" if row.get("aktiv_timer") else "over_16_hours", + "id": row["id"], "hours": float(hours)}) + continue + bucket["hours"] += hours + bucket["time_entries"] += 1 + cases = execute_query( + """SELECT id,created_by_user_id,created_at FROM sag_sager + WHERE deleted_at IS NULL AND created_at::date BETWEEN %s AND %s""", (start, end), + ) or [] + orders = execute_query( + "SELECT id,created_by,order_date FROM tmodule_orders WHERE order_date BETWEEN %s AND %s", (start, end), + ) or [] + for rows, metric, user_field in ((cases, "cases", "created_by_user_id"), (orders, "orders", "created_by")): + for row in rows: + user = hub_users.get(int(row[user_field])) if row.get(user_field) else None + key = f"hub:{user['user_id']}" if user else "hub:unknown" + bucket = _employee_bucket(result["employees"], key, + str((user or {}).get("full_name") or (user or {}).get("username") or "Ukendt Hub-bruger"), + (user or {}).get("email")) + if not user and key not in result["data_quality"]["unmatched_employees"]: + result["data_quality"]["unmatched_employees"].append(key) + bucket[metric] += 1 + return result + + +def build_impact_report(version_id: int, vtiger_from: date, vtiger_to: date, hub_from: date, hub_to: date) -> dict: + vt_days = (vtiger_to - vtiger_from).days + 1 + hub_days = (hub_to - hub_from).days + 1 + if vt_days <= 0 or hub_days <= 0 or vt_days != hub_days: + raise ValueError("vTiger- og Hub-perioderne skal være gyldige og lige lange") + version = execute_query_single( + "SELECT id,status,source_cutoff FROM vtiger_archive_versions WHERE id=%s", (version_id,), + ) + if not version or version.get("status") != "completed": + raise ValueError("Vælg en fuldført arkivversion") + by_id, by_email = _hub_users() + vtiger = _finalize(_vtiger_metrics(_archive_snapshot(version_id), vtiger_from, vtiger_to, by_email), vtiger_from, vtiger_to) + hub = _finalize(_hub_metrics(hub_from, hub_to, by_id), hub_from, hub_to) + keys = ("hours", "time_entries", "cases", "orders") + delta = {key: round(float(hub["totals"][key]) - float(vtiger["totals"][key]), 2) for key in keys} + percent = { + key: (round(delta[key] / float(vtiger["totals"][key]) * 100, 1) if float(vtiger["totals"][key]) else None) + for key in keys + } + productivity_keys = ("hours_per_workday", "registrations_per_workday", "cases_per_workday", + "orders_per_workday", "deliveries_per_workday", "hours_per_case_or_order") + productivity_delta = {} + for key in productivity_keys: + old, new = float(vtiger["totals"]["productivity"][key]), float(hub["totals"]["productivity"][key]) + productivity_delta[key] = {"delta": round(new - old, 2), "percent": round((new - old) / old * 100, 1) if old else None} + evidence = [] + for key, label in (("deliveries_per_workday", "samlet registreret output pr. arbejdsdag"), + ("cases_per_workday", "sager pr. arbejdsdag"), + ("orders_per_workday", "ordrer pr. arbejdsdag"), + ("registrations_per_workday", "tidsregistreringer pr. arbejdsdag")): + change = productivity_delta[key] + if change["percent"] is not None: + direction = "mere" if change["percent"] >= 0 else "mindre" + evidence.append(f"Hub registrerer {abs(change['percent']):.1f}% {direction} {label}.") + return { + "archive_version": {"id": version_id, "source_cutoff": version.get("source_cutoff")}, + "periods": {"vtiger": {"from": vtiger_from, "to": vtiger_to, "days": vt_days}, + "hub": {"from": hub_from, "to": hub_to, "days": hub_days}}, + "vtiger": vtiger, "hub": hub, + "comparison": delta, + "improvement": {"percent": percent, "productivity": productivity_delta, "evidence": evidence}, + } + + +def save_impact_report(result: dict, generated_by: int) -> int: + canonical = json.dumps(result, ensure_ascii=False, sort_keys=True, separators=(",", ":"), default=str) + periods = result["periods"] + return int(execute_insert( + """INSERT INTO hub_impact_reports + (archive_version_id,vtiger_from,vtiger_to,hub_from,hub_to,result,result_sha256,generated_by) + VALUES (%s,%s,%s,%s,%s,%s::jsonb,%s,%s) RETURNING id""", + (result["archive_version"]["id"], periods["vtiger"]["from"], periods["vtiger"]["to"], + periods["hub"]["from"], periods["hub"]["to"], canonical, + hashlib.sha256(canonical.encode()).hexdigest(), generated_by), + )) diff --git a/app/admin/router.py b/app/admin/router.py new file mode 100644 index 0000000..4d5645b --- /dev/null +++ b/app/admin/router.py @@ -0,0 +1,344 @@ +"""Superadmin-only Project CT archive endpoints.""" + +from __future__ import annotations + +import hashlib +import io +import json +import os +import tempfile +from datetime import date +from fastapi import APIRouter, Depends, File, HTTPException, Query, UploadFile, status +from fastapi.responses import FileResponse, Response +from starlette.background import BackgroundTask +from pydantic import BaseModel + +from app.admin.vtiger_archive import ARCHIVE_MODULES, run_archive_sync, termination_readiness +from app.admin.hub_impact import build_impact_report, save_impact_report +from app.admin.archive_bundle import create_archive_bundle, file_sha256, import_archive_bundle +from app.core.auth_dependencies import get_current_user +from app.core.database import execute_query, execute_query_single + +router = APIRouter() + + +class ImpactReportRequest(BaseModel): + archive_version_id: int + vtiger_from: date + vtiger_to: date + hub_from: date + hub_to: date + + +def require_hidden_superadmin(current_user: dict = Depends(get_current_user)) -> dict: + # This feature must not disclose its existence to ordinary users. + if not current_user.get("is_superadmin"): + raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Not found") + return current_user + + +@router.get("/admin/vtiger-archive/status") +async def archive_status(current_user: dict = Depends(require_hidden_superadmin)): + versions = execute_query( + """SELECT id,sync_kind,status,started_at,completed_at,source_cutoff,module_counts, + critical_errors,control_report,control_approved_at,raw_export_sha256 + FROM vtiger_archive_versions ORDER BY id DESC LIMIT 25""" + ) or [] + checkpoints = execute_query( + "SELECT * FROM vtiger_archive_checkpoints ORDER BY module" + ) or [] + transfers = execute_query( + """SELECT id,direction,bundle_sha256,through_version_id,source_instance,status,counts, + error_message,started_at,completed_at + FROM vtiger_archive_transfers ORDER BY id DESC LIMIT 25""" + ) or [] + file_archive = execute_query( + """SELECT source_module,COUNT(*)::integer AS files,SUM(size_bytes)::bigint AS size_bytes + FROM vtiger_archive_files GROUP BY source_module ORDER BY source_module""" + ) or [] + return {"versions": versions, "checkpoints": checkpoints, "transfers": transfers, + "file_archive": file_archive, "readiness": termination_readiness()} + + +@router.post("/admin/vtiger-archive/sync/{sync_kind}") +async def archive_sync(sync_kind: str, current_user: dict = Depends(require_hidden_superadmin)): + return await run_archive_sync(sync_kind, int(current_user["id"])) + + +@router.post("/admin/vtiger-archive/versions/{version_id}/approve-control") +async def approve_control(version_id: int, current_user: dict = Depends(require_hidden_superadmin)): + version = execute_query_single( + "SELECT * FROM vtiger_archive_versions WHERE id=%s", (version_id,), + ) + if not version: + raise HTTPException(status_code=404, detail="Arkivversionen findes ikke") + if version.get("status") != "completed" or version.get("critical_errors"): + raise HTTPException(status_code=409, detail="En fejlet eller ukomplet kontrolrapport kan ikke godkendes") + counts = version.get("module_counts") or {} + missing = [module for module in ARCHIVE_MODULES if module not in counts] + if missing: + raise HTTPException(status_code=409, detail=f"Kontrolrapport mangler moduler: {', '.join(missing)}") + execute_query( + """UPDATE vtiger_archive_versions SET control_approved_at=NOW(),control_approved_by=%s + WHERE id=%s""", (current_user["id"], version_id), fetch=False, + ) + return {"approved": True, "version_id": version_id, "readiness": termination_readiness()} + + +@router.get("/admin/vtiger-archive/versions/{version_id}/raw-export") +async def raw_export( + version_id: int, + module: str | None = Query(None), + current_user: dict = Depends(require_hidden_superadmin), +): + version = execute_query_single("SELECT id FROM vtiger_archive_versions WHERE id=%s", (version_id,)) + if not version: + raise HTTPException(status_code=404, detail="Arkivversionen findes ikke") + params: list[object] = [version_id] + module_filter = "" + if module: + module_filter = "AND module=%s" + params.append(module) + rows = execute_query( + f"""SELECT DISTINCT ON (module,vtiger_id) module,vtiger_id,revision_no, + source_created_at,source_modified_at,is_deleted,payload,payload_sha256 + FROM vtiger_archive_records WHERE version_id <= %s {module_filter} + ORDER BY module,vtiger_id,revision_no DESC""", tuple(params), + ) or [] + relations = execute_query( + """SELECT source_module,source_vtiger_id,field_name,target_vtiger_id,target_module + FROM vtiger_archive_relations WHERE version_id <= %s + ORDER BY source_module,source_vtiger_id,field_name,target_vtiger_id""", (version_id,), + ) or [] + version_meta = execute_query_single( + """SELECT id,sync_kind,status,started_at,completed_at,source_cutoff,module_counts, + warnings,critical_errors,control_report,control_approved_at + FROM vtiger_archive_versions WHERE id=%s""", (version_id,), + ) or {} + export_lines = [{"record_type": "archive_version", "data": dict(version_meta)}] + export_lines.extend({"record_type": "entity", "data": dict(row)} for row in rows) + export_lines.extend({"record_type": "relation", "data": dict(row)} for row in relations) + body = "\n".join(json.dumps(item, ensure_ascii=False, default=str, sort_keys=True) for item in export_lines) + "\n" + digest = hashlib.sha256(body.encode("utf-8")).hexdigest() + if not module: + execute_query( + "UPDATE vtiger_archive_versions SET raw_export_sha256=%s WHERE id=%s", + (digest, version_id), fetch=False, + ) + filename = f"project-ct-vtiger-v{version_id}{'-' + module if module else ''}.jsonl" + return Response( + content=body, + media_type="application/x-ndjson", + headers={"Content-Disposition": f'attachment; filename="{filename}"', "X-Content-SHA256": digest}, + ) + + +@router.get("/admin/vtiger-archive/readiness") +async def readiness(current_user: dict = Depends(require_hidden_superadmin)): + return termination_readiness() + + +@router.get("/admin/vtiger-archive/bundle.zip") +async def export_archive_bundle( + through_version_id: int | None = Query(None), + current_user: dict = Depends(require_hidden_superadmin), +): + if through_version_id is None: + latest = execute_query_single( + "SELECT id FROM vtiger_archive_versions WHERE status='completed' ORDER BY id DESC LIMIT 1" + ) + if not latest: + raise HTTPException(status_code=409, detail="Der findes ingen fuldført arkivversion") + through_version_id = int(latest["id"]) + selected = execute_query_single( + "SELECT id,status FROM vtiger_archive_versions WHERE id=%s", (through_version_id,), + ) + if not selected: + raise HTTPException(status_code=404, detail="Arkivversionen findes ikke") + if selected.get("status") != "completed": + raise HTTPException(status_code=409, detail="Kun en fuldført arkivversion kan eksporteres til produktion") + transfer_id = execute_query_single( + """INSERT INTO vtiger_archive_transfers(direction,through_version_id,status,initiated_by) + VALUES ('export',%s,'running',%s) RETURNING id""", (through_version_id, current_user["id"]), + )["id"] + try: + path, manifest = create_archive_bundle(through_version_id, int(current_user["id"])) + counts = {name: info["count"] for name, info in manifest["entries"].items()} + execute_query( + """UPDATE vtiger_archive_transfers SET status='completed',completed_at=NOW(), + bundle_sha256=%s,source_instance=%s,counts=%s::jsonb WHERE id=%s""", + (manifest["bundle_sha256"], manifest["source_instance"], json.dumps(counts), transfer_id), fetch=False, + ) + return FileResponse( + path, media_type="application/zip", + filename=f"project-ct-vtiger-archive-v{through_version_id}.zip", + background=BackgroundTask(lambda: os.path.exists(path) and os.unlink(path)), + headers={"X-Archive-SHA256": manifest["bundle_sha256"], "X-vTiger-Contacted": "false"}, + ) + except Exception as exc: + execute_query( + "UPDATE vtiger_archive_transfers SET status='failed',completed_at=NOW(),error_message=%s WHERE id=%s", + (str(exc), transfer_id), fetch=False, + ) + raise HTTPException(status_code=500, detail=f"Arkivpakken kunne ikke oprettes: {exc}") from exc + + +@router.post("/admin/vtiger-archive/import-bundle") +async def import_bundle( + file: UploadFile = File(...), + current_user: dict = Depends(require_hidden_superadmin), +): + if not str(file.filename or "").lower().endswith(".zip"): + raise HTTPException(status_code=400, detail="Vælg en Project CT .zip-arkivpakke") + fd, path = tempfile.mkstemp(prefix="project-ct-upload-", suffix=".zip") + os.close(fd) + transfer_id = execute_query_single( + """INSERT INTO vtiger_archive_transfers(direction,status,initiated_by) + VALUES ('import','running',%s) RETURNING id""", (current_user["id"],), + )["id"] + try: + size = 0 + with open(path, "wb") as target: + while chunk := await file.read(1024 * 1024): + size += len(chunk) + if size > 20 * 1024 * 1024 * 1024: + raise ValueError("Arkivpakken må højst fylde 20 GB") + target.write(chunk) + uploaded_sha256 = file_sha256(path) + previous = execute_query_single( + """SELECT id,through_version_id,counts FROM vtiger_archive_transfers + WHERE direction='import' AND status='completed' AND bundle_sha256=%s + ORDER BY id DESC LIMIT 1""", (uploaded_sha256,), + ) + if previous: + execute_query( + """UPDATE vtiger_archive_transfers SET status='completed',completed_at=NOW(),bundle_sha256=%s, + through_version_id=%s,counts=%s::jsonb WHERE id=%s""", + (uploaded_sha256, previous.get("through_version_id"), json.dumps(previous.get("counts") or {}), transfer_id), + fetch=False, + ) + return { + "bundle_sha256": uploaded_sha256, + "counts": previous.get("counts") or {}, + "contacted_vtiger": False, + "already_imported": True, + "original_transfer_id": previous["id"], + "transfer_id": transfer_id, + } + result = import_archive_bundle(path, int(current_user["id"])) + execute_query( + """UPDATE vtiger_archive_transfers SET status='completed',completed_at=NOW(),bundle_sha256=%s, + source_instance=%s,through_version_id=%s,counts=%s::jsonb WHERE id=%s""", + (result["bundle_sha256"], result["manifest"].get("source_instance"), + result["manifest"].get("through_version_id"), json.dumps(result["counts"]), transfer_id), fetch=False, + ) + return {**result, "transfer_id": transfer_id} + except ValueError as exc: + execute_query( + "UPDATE vtiger_archive_transfers SET status='failed',completed_at=NOW(),error_message=%s WHERE id=%s", + (str(exc), transfer_id), fetch=False, + ) + raise HTTPException(status_code=400, detail=str(exc)) from exc + except Exception as exc: + execute_query( + "UPDATE vtiger_archive_transfers SET status='failed',completed_at=NOW(),error_message=%s WHERE id=%s", + (str(exc), transfer_id), fetch=False, + ) + raise HTTPException(status_code=500, detail=f"Arkivpakken kunne ikke importeres: {exc}") from exc + finally: + if os.path.exists(path): + os.unlink(path) + + +@router.get("/admin/hub-impact/options") +async def impact_options(current_user: dict = Depends(require_hidden_superadmin)): + versions = execute_query( + """SELECT id,sync_kind,completed_at,source_cutoff,module_counts,control_report,control_approved_at + FROM vtiger_archive_versions WHERE status='completed' ORDER BY id DESC""" + ) or [] + reports = execute_query( + """SELECT id,archive_version_id,vtiger_from,vtiger_to,hub_from,hub_to,result_sha256,generated_at + FROM hub_impact_reports ORDER BY id DESC LIMIT 30""" + ) or [] + return {"versions": versions, "reports": reports, "readiness": termination_readiness()} + + +@router.post("/admin/hub-impact/reports") +async def create_impact_report(payload: ImpactReportRequest, current_user: dict = Depends(require_hidden_superadmin)): + try: + result = build_impact_report(payload.archive_version_id, payload.vtiger_from, payload.vtiger_to, + payload.hub_from, payload.hub_to) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + report_id = save_impact_report(result, int(current_user["id"])) + return {"id": report_id, "result": result} + + +@router.get("/admin/hub-impact/reports/{report_id}") +async def get_impact_report(report_id: int, current_user: dict = Depends(require_hidden_superadmin)): + row = execute_query_single("SELECT * FROM hub_impact_reports WHERE id=%s", (report_id,)) + if not row: + raise HTTPException(status_code=404, detail="Rapporten findes ikke") + return row + + +def _impact_workbook(report: dict) -> bytes: + from openpyxl import Workbook + from openpyxl.styles import Font, PatternFill + wb = Workbook() + overview = wb.active + overview.title = "Overblik" + overview.append(["Project CT · Hub Impact", "vTiger", "Hub", "Forskel", "% ændring"]) + for cell in overview[1]: + cell.font = Font(bold=True, color="FFFFFF") + cell.fill = PatternFill("solid", fgColor="0F4C75") + labels = {"hours": "Timer", "time_entries": "Tidsregistreringer", "cases": "Sager", "orders": "Ordrer"} + for key, label in labels.items(): + overview.append([label, report["vtiger"]["totals"][key], report["hub"]["totals"][key], report["comparison"][key], (report.get("improvement") or {}).get("percent", {}).get(key)]) + overview.append([]) + overview.append(["Arkivversion", report["archive_version"]["id"]]) + overview.append(["vTiger-periode", f"{report['periods']['vtiger']['from']} – {report['periods']['vtiger']['to']}"]) + overview.append(["Hub-periode", f"{report['periods']['hub']['from']} – {report['periods']['hub']['to']}"]) + overview.freeze_panes = "A2" + overview.column_dimensions["A"].width = 28 + for col in "BCDE": overview.column_dimensions[col].width = 18 + evidence = wb.create_sheet("Effekt") + evidence.append(["Dokumenteret ændring"]) + for line in (report.get("improvement") or {}).get("evidence", []): + evidence.append([line]) + evidence.column_dimensions["A"].width = 90 + + employees = wb.create_sheet("Pr medarbejder") + employees.append(["Kilde", "Medarbejder", "E-mail", "Timer", "Registreringer", "Sager", "Ordrer", "Timer/arbejdsdag"]) + for source in ("vtiger", "hub"): + for row in report[source]["employees"]: + employees.append([source, row["name"], row.get("email"), row["hours"], row["time_entries"], row["cases"], row["orders"], row["productivity"]["hours_per_workday"]]) + for cell in employees[1]: cell.font = Font(bold=True) + employees.freeze_panes = "A2" + employees.auto_filter.ref = employees.dimensions + + anomalies = wb.create_sheet("Afvigelser") + anomalies.append(["Kilde", "Type", "ID", "Timer"]) + for source in ("vtiger", "hub"): + for row in report[source]["anomalies"]: + anomalies.append([row.get("source"), row.get("type"), row.get("id"), row.get("hours")]) + quality = wb.create_sheet("Datakvalitet") + quality.append(["Kilde", "Manglende datoer", "Ikke matchede medarbejdere"]) + for source in ("vtiger", "hub"): + data = report[source]["data_quality"] + quality.append([source, data["missing_dates"], ", ".join(data["unmatched_employees"])]) + output = io.BytesIO() + wb.save(output) + return output.getvalue() + + +@router.get("/admin/hub-impact/reports/{report_id}/export.xlsx") +async def export_impact_report(report_id: int, current_user: dict = Depends(require_hidden_superadmin)): + row = execute_query_single("SELECT result FROM hub_impact_reports WHERE id=%s", (report_id,)) + if not row: + raise HTTPException(status_code=404, detail="Rapporten findes ikke") + return Response( + content=_impact_workbook(row["result"]), + media_type="application/vnd.openxmlformats-officedocument.spreadsheetml.sheet", + headers={"Content-Disposition": f'attachment; filename="project-ct-impact-{report_id}.xlsx"'}, + ) diff --git a/app/admin/views.py b/app/admin/views.py new file mode 100644 index 0000000..c33fd5f --- /dev/null +++ b/app/admin/views.py @@ -0,0 +1,13 @@ +from fastapi import APIRouter, Depends, Request +from fastapi.responses import HTMLResponse +from fastapi.templating import Jinja2Templates + +from app.admin.router import require_hidden_superadmin + +router = APIRouter() +templates = Jinja2Templates(directory="app") + + +@router.get("/admin/hub-impact", response_class=HTMLResponse) +async def hub_impact_page(request: Request, current_user: dict = Depends(require_hidden_superadmin)): + return templates.TemplateResponse("admin/hub_impact.html", {"request": request, "current_user": current_user}) diff --git a/app/admin/vtiger_archive.py b/app/admin/vtiger_archive.py new file mode 100644 index 0000000..7a15880 --- /dev/null +++ b/app/admin/vtiger_archive.py @@ -0,0 +1,378 @@ +"""Permanent append-only vTiger archive used by Project CT.""" + +from __future__ import annotations + +import hashlib +import asyncio +import json +import logging +import re +from datetime import datetime, timezone +from typing import Any, Iterable + +from app.core.database import execute_insert, execute_query, execute_query_single +from app.services.vtiger_service import VTigerService, get_vtiger_service + +logger = logging.getLogger(__name__) + +# Document metadata and original file bytes are both archived before vTiger is retired. +ARCHIVE_MODULES = ( + "Users", "Groups", "Roles", "Currency", "Accounts", "Contacts", "Leads", "Potentials", "Cases", "Project", + "ProjectMilestone", "ModComments", "Timelog", "Calendar", "Events", "Emails", + "SalesOrder", "Quotes", "Invoice", "PurchaseOrder", "Vendors", "ServiceContracts", + "Subscription", "Products", "Services", "Assets", "Documents", +) +RELATION_FIELD_RE = re.compile(r"(?:^|_)(?:id|ids|account|contact|parent|related|assigned_user)$", re.I) +MODULE_BATCH_SIZES = {"Emails": 20, "Invoice": 50, "SalesOrder": 50, "Subscription": 50} + + +def _transient_vtiger_failure(service) -> bool: + status = service.last_query_status + error_text = json.dumps(service.last_query_error or {}).upper() + return status is None or status == 429 or (isinstance(status, int) and status >= 500) or any( + token in error_text for token in ("TOO_MANY_REQUESTS", "TIMEOUT", "TEMPORAR") + ) + + +def _canonical_payload(record: dict[str, Any]) -> tuple[str, str]: + raw = json.dumps(record, ensure_ascii=False, sort_keys=True, separators=(",", ":"), default=str) + return raw, hashlib.sha256(raw.encode("utf-8")).hexdigest() + + +def _parse_source_time(value: Any): + text = str(value or "").strip() + if not text: + return None + try: + parsed = datetime.fromisoformat(text.replace("Z", "+00:00")) + return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc) + except ValueError: + return None + + +def _relation_values(payload: dict[str, Any]) -> Iterable[tuple[str, str]]: + for field, value in payload.items(): + if not value or not RELATION_FIELD_RE.search(str(field)): + continue + values = value if isinstance(value, list) else [value] + for candidate in values: + candidate = str(candidate or "").strip() + if re.fullmatch(r"\d+x\d+", candidate): + yield str(field), candidate + + +def _target_module(vtiger_id: str) -> str | None: + prefix = str(vtiger_id).split("x", 1)[0] + row = execute_query_single( + """SELECT module FROM vtiger_archive_records + WHERE split_part(vtiger_id,'x',1)=%s ORDER BY id DESC LIMIT 1""", (prefix,), + ) + return str(row["module"]) if row else None + + +def archive_record(version_id: int, module: str, record: dict[str, Any]) -> bool: + vtiger_id = str(record.get("id") or "").strip() + if not vtiger_id: + raise ValueError(f"{module}-post mangler id") + raw, digest = _canonical_payload(record) + previous = execute_query_single( + """SELECT revision_no, payload_sha256 FROM vtiger_archive_records + WHERE module=%s AND vtiger_id=%s ORDER BY revision_no DESC LIMIT 1""", + (module, vtiger_id), + ) + if previous and previous.get("payload_sha256") == digest: + return False + revision = int((previous or {}).get("revision_no") or 0) + 1 + execute_query( + """INSERT INTO vtiger_archive_records + (version_id,module,vtiger_id,revision_no,source_created_at,source_modified_at,is_deleted,payload,payload_sha256) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s::jsonb,%s)""", + (version_id, module, vtiger_id, revision, _parse_source_time(record.get("createdtime")), + _parse_source_time(record.get("modifiedtime")), bool(record.get("deleted")), raw, digest), + fetch=False, + ) + for field, target in _relation_values(record): + execute_query( + """INSERT INTO vtiger_archive_relations + (version_id,source_module,source_vtiger_id,field_name,target_vtiger_id,target_module) + VALUES (%s,%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING""", + (version_id, module, vtiger_id, field, target, _target_module(target)), fetch=False, + ) + return True + + +async def archive_document_files(version_id: int) -> dict[str, int]: + """Archive immutable Document and Email bytes via vTiger's files_retrieve API.""" + documents = execute_query( + """SELECT DISTINCT ON (module,vtiger_id) module,vtiger_id,payload + FROM vtiger_archive_records + WHERE module IN ('Documents','Emails') AND version_id<=%s AND is_deleted=false + ORDER BY module,vtiger_id,revision_no DESC""", (version_id,), + ) or [] + stats = {"entities": len(documents), "resources": 0, "archived": 0, "unchanged": 0, "errors": 0} + # A small hard ceiling keeps this one-time preservation fast without flooding vTiger. + semaphore = asyncio.Semaphore(6) + jobs = [] + + async def archive_resource(document: dict, payload: dict, resource_id: str) -> None: + async with semaphore: + previous_resource = execute_query_single( + """SELECT id FROM vtiger_archive_files + WHERE source_module=%s AND document_vtiger_id=%s AND resource_vtiger_id=%s + ORDER BY id DESC LIMIT 1""", + (document["module"], document["vtiger_id"], resource_id), + ) + if previous_resource: + stats["unchanged"] += 1 + return + service = VTigerService() + result = None + for attempt in range(8): + result = await service.retrieve_file(resource_id) + if result or not _transient_vtiger_failure(service): + break + await asyncio.sleep(min(1.5 * (2 ** attempt), 30)) + if not result: + stats["errors"] += 1 + return + content = result.get("content") or b"" + digest = hashlib.sha256(content).hexdigest() + filename = str(result.get("filename") or payload.get("filename") or resource_id) + execute_query( + """INSERT INTO vtiger_archive_files + (version_id,source_module,document_vtiger_id,resource_vtiger_id,filename,content_type,size_bytes,content_sha256,content) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s) ON CONFLICT DO NOTHING""", + (version_id, document["module"], document["vtiger_id"], resource_id, filename, + result.get("filetype") or payload.get("filetype"), len(content), digest, content), + fetch=False, + ) + stats["archived"] += 1 + await asyncio.sleep(0.35) + + for document in documents: + payload = document.get("payload") or {} + resource_ids = re.findall(r"\d+x\d+", str(payload.get("imageattachmentids") or "")) + for resource_id in resource_ids: + stats["resources"] += 1 + jobs.append(archive_resource(document, payload, resource_id)) + if jobs: + await asyncio.gather(*jobs) + return stats + + +async def _fetch_module(module: str, modified_after=None) -> list[dict[str, Any]]: + service = get_vtiger_service() + records: list[dict[str, Any]] = [] + seen_ids: set[str] = set() + offset = 0 + batch_size = MODULE_BATCH_SIZES.get(module, 100) + while True: + clauses = [] + if modified_after: + clauses.append(f"modifiedtime >= '{modified_after.isoformat()}'") + where = f" WHERE {' AND '.join(clauses)}" if clauses else "" + # This vTiger deployment paginates reliably with LIMIT offset,count. + # ID-watermarks skip records in modules whose entity IDs are compared as text. + query = f"SELECT * FROM {module}{where} ORDER BY id ASC LIMIT {offset}, {batch_size};" + batch = [] + for attempt in range(8): + batch = await service.query(query) + if not _transient_vtiger_failure(service): + break + await asyncio.sleep(min(1.5 * (2 ** attempt), 30)) + if not batch: + if service.last_query_error or (service.last_query_status not in (None, 200)): + raise RuntimeError(f"vTiger kunne ikke hente {module}: {service.last_query_error or service.last_query_status}") + break + fresh = [row for row in batch if str(row.get("id") or "") not in seen_ids] + if not fresh: + break + records.extend(fresh) + seen_ids.update(str(row.get("id") or "") for row in fresh) + # Payload-heavy modules may return fewer than the requested 100 records + # due to a response-size cap even though later offsets still exist. + offset += len(batch) + await asyncio.sleep(0.35) + return records + + +async def _source_module_count(module: str) -> int: + service = get_vtiger_service() + rows = [] + for attempt in range(8): + rows = await service.query(f"SELECT count(*) FROM {module};") + if not _transient_vtiger_failure(service): + break + await asyncio.sleep(min(1.5 * (2 ** attempt), 30)) + if not rows or rows[0].get("count") is None: + raise RuntimeError( + f"vTiger kunne ikke kontrollere antal poster i {module}: " + f"{service.last_query_error or service.last_query_status}" + ) + return int(rows[0]["count"]) + + +async def _close_relation_targets(version_id: int) -> dict[str, int]: + """Fetch referenced records omitted from normal module lists (notably inactive users).""" + unresolved = execute_query( + """SELECT DISTINCT r.target_vtiger_id + FROM vtiger_archive_relations r + WHERE r.version_id <= %s AND NOT EXISTS ( + SELECT 1 FROM vtiger_archive_records t + WHERE t.vtiger_id=r.target_vtiger_id AND t.version_id <= %s + ) ORDER BY r.target_vtiger_id""", (version_id, version_id), + ) or [] + prefix_rows = execute_query( + """SELECT split_part(vtiger_id,'x',1) AS prefix,MIN(module) AS module + FROM vtiger_archive_records GROUP BY 1""" + ) or [] + module_by_prefix = {str(row["prefix"]): str(row["module"]) for row in prefix_rows} + module_by_prefix.update({"19": "Users", "21": "Currency", "53": "Roles"}) + service = get_vtiger_service() + stats = {"requested": 0, "archived": 0, "not_found": 0, "errors": 0} + for item in unresolved: + target_id = str(item.get("target_vtiger_id") or "") + if not re.fullmatch(r"\d+x\d+", target_id): + continue + module = module_by_prefix.get(target_id.split("x", 1)[0]) + if not module: + continue + stats["requested"] += 1 + rows = [] + for attempt in range(8): + rows = await service.query(f"SELECT * FROM {module} WHERE id='{target_id}' LIMIT 1;") + if not _transient_vtiger_failure(service): + break + await asyncio.sleep(min(1.5 * (2 ** attempt), 30)) + if rows: + stats["archived"] += int(archive_record(version_id, module, rows[0])) + elif _transient_vtiger_failure(service): + stats["errors"] += 1 + else: + stats["not_found"] += 1 + await asyncio.sleep(0.35) + return stats + + +async def run_archive_sync(sync_kind: str, initiated_by: int | None = None) -> dict[str, Any]: + if sync_kind not in {"full", "incremental", "final"}: + raise ValueError("sync_kind skal være full, incremental eller final") + running = execute_query_single( + """SELECT id,started_at FROM vtiger_archive_versions + WHERE status='running' ORDER BY id DESC LIMIT 1""" + ) + if running: + age = datetime.now(timezone.utc) - running["started_at"] + if age.total_seconds() < 12 * 3600: + raise RuntimeError(f"Arkivsync version {running['id']} kører allerede") + execute_query( + """UPDATE vtiger_archive_versions SET status='failed',completed_at=NOW(), + critical_errors='[{"error":"Stale sync blev lukket automatisk"}]'::jsonb WHERE id=%s""", + (running["id"],), fetch=False, + ) + version_id = int(execute_insert( + """INSERT INTO vtiger_archive_versions(sync_kind,initiated_by,source_cutoff) + VALUES (%s,%s,NOW()) RETURNING id""", (sync_kind, initiated_by), + )) + counts: dict[str, Any] = {} + errors: list[dict[str, str]] = [] + try: + for module in ARCHIVE_MODULES: + await asyncio.sleep(0.45) + checkpoint = execute_query_single( + "SELECT last_modified_at FROM vtiger_archive_checkpoints WHERE module=%s", (module,), + ) + watermark = checkpoint.get("last_modified_at") if checkpoint and sync_kind == "incremental" else None + try: + source_total = await _source_module_count(module) if sync_kind in {"full", "final"} else None + records = await _fetch_module(module, watermark) + inserted = sum(1 for record in records if archive_record(version_id, module, record)) + counts[module] = {"fetched": len(records), "new_revisions": inserted} + if source_total is not None: + counts[module]["source_total"] = source_total + if len(records) != source_total: + raise RuntimeError(f"Paritetsfejl: vTiger={source_total}, hentet={len(records)}") + source_times = [parsed for parsed in (_parse_source_time(r.get("modifiedtime")) for r in records) if parsed] + newest = max(source_times, default=None) + execute_query( + """INSERT INTO vtiger_archive_checkpoints(module,last_modified_at,last_vtiger_id,last_successful_version_id) + VALUES (%s,%s,%s,%s) ON CONFLICT(module) DO UPDATE SET + last_modified_at=COALESCE(EXCLUDED.last_modified_at,vtiger_archive_checkpoints.last_modified_at), + last_vtiger_id=EXCLUDED.last_vtiger_id,last_successful_version_id=EXCLUDED.last_successful_version_id, + updated_at=NOW()""", + (module, newest, str(records[-1].get("id")) if records else None, version_id), fetch=False, + ) + except Exception as exc: + logger.exception("Project CT sync failed for %s", module) + errors.append({"module": module, "error": str(exc)}) + document_files = await archive_document_files(version_id) if sync_kind in {"full", "final"} else { + "entities": 0, "resources": 0, "archived": 0, "unchanged": 0, "errors": 0, + } + if document_files.get("errors"): + errors.append({ + "module": "DocumentFiles", + "error": f"{document_files['errors']} dokumentfiler kunne ikke arkiveres", + }) + relation_backfill = await _close_relation_targets(version_id) + if relation_backfill.get("errors"): + errors.append({ + "module": "relations", + "error": f"{relation_backfill['errors']} relationer kunne ikke kontrolleres mod vTiger", + }) + archived_counts = execute_query( + """SELECT module,COUNT(*) AS records FROM ( + SELECT DISTINCT ON (module,vtiger_id) module,vtiger_id,is_deleted + FROM vtiger_archive_records WHERE version_id <= %s + ORDER BY module,vtiger_id,revision_no DESC + ) snapshot WHERE is_deleted=false GROUP BY module ORDER BY module""", (version_id,), + ) or [] + unresolved = execute_query_single( + """SELECT COUNT(*)::integer AS count FROM vtiger_archive_relations r + WHERE r.version_id <= %s AND NOT EXISTS ( + SELECT 1 FROM vtiger_archive_records t + WHERE t.vtiger_id=r.target_vtiger_id AND t.version_id <= %s + )""", (version_id, version_id), + ) or {"count": 0} + document_quality = execute_query_single( + """SELECT COUNT(*)::integer AS total, + COUNT(*) FILTER (WHERE COALESCE(payload->>'filename',payload->>'notes_title','')='')::integer AS missing_name + FROM (SELECT DISTINCT ON (vtiger_id) vtiger_id,payload,is_deleted + FROM vtiger_archive_records WHERE module='Documents' AND version_id <= %s + ORDER BY vtiger_id,revision_no DESC) d WHERE is_deleted=false""", (version_id,), + ) or {"total": 0, "missing_name": 0} + report = { + "modules_expected": len(ARCHIVE_MODULES), "modules_completed": len(counts), + "missing_modules": [module for module in ARCHIVE_MODULES if module not in counts], + "module_sync": counts, "archived_snapshot_counts": archived_counts, + "relations": {"unresolved": int(unresolved.get("count") or 0)}, + "relation_backfill": relation_backfill, + "documents": {**document_quality, "files": document_files}, "critical_errors": errors, + } + status = "failed" if errors else "completed" + execute_query( + """UPDATE vtiger_archive_versions SET status=%s,completed_at=NOW(),module_counts=%s::jsonb, + critical_errors=%s::jsonb,control_report=%s::jsonb WHERE id=%s""", + (status, json.dumps(counts), json.dumps(errors), json.dumps(report), version_id), fetch=False, + ) + return {"version_id": version_id, "status": status, "counts": counts, "critical_errors": errors} + except Exception as exc: + execute_query( + "UPDATE vtiger_archive_versions SET status='failed',completed_at=NOW(),critical_errors=%s::jsonb WHERE id=%s", + (json.dumps([{"error": str(exc)}]), version_id), fetch=False, + ) + raise + + +def termination_readiness() -> dict[str, Any]: + final = execute_query_single( + """SELECT id,status,completed_at,critical_errors,control_approved_at + FROM vtiger_archive_versions WHERE sync_kind='final' ORDER BY id DESC LIMIT 1""" + ) + reasons = [] + if not final or final.get("status") != "completed": + reasons.append("Final sync mangler eller er ikke fuldført") + if final and final.get("critical_errors"): + reasons.append("Final sync har kritiske fejl") + if not final or not final.get("control_approved_at"): + reasons.append("Kontrolrapporten er ikke godkendt") + return {"ready": not reasons, "reasons": reasons, "final_version": final} diff --git a/app/billing/backend/supplier_invoices.py b/app/billing/backend/supplier_invoices.py index ebd7232..1c5758f 100644 --- a/app/billing/backend/supplier_invoices.py +++ b/app/billing/backend/supplier_invoices.py @@ -16,6 +16,7 @@ from app.services.economic_service import get_economic_service from app.services.ollama_service import ollama_service from app.services.template_service import template_service from app.services.invoice2data_service import get_invoice2data_service +from app.modules.internet_connections.backend.change_case_service import ensure_external_change_case import logging import os import re @@ -26,15 +27,11 @@ router = APIRouter() _PURCHASE_CASE_TYPE = "indkøb" _INTERNET_CASE_RELEVANT_CHANGE_FIELDS = { - "address", - "service_address", - "monthly_cost", - "technology", - "connection_type", - "circuit_number", - "speed_mbps", - "download_mbps", - "upload_mbps", + "address", "service_address", "monthly_cost", "sales_price", "technology", + "connection_type", "circuit_number", "provider_reference", "provider", "vendor_id", + "speed_mbps", "download_mbps", "upload_mbps", "status", "sla_subscription_id", + "sla_price", "sla_status", "cidr", "contract_number", "range_added", "range_removed", + "range_monthly_cost", "range_sales_price", } SUPPLIER_STATUS_V2 = ("modtaget", "godkendt", "betalt", "afvist") @@ -267,90 +264,19 @@ def _ensure_internet_change_case( owner_customer_id: Optional[int], changes: Dict[str, Dict[str, object]], ) -> Optional[int]: - # Customer ownership, initial activation and internal classification are - # bookkeeping outcomes of a successful import, not operational incidents. - # Only create cases for changes that can affect delivery or billing. - relevant_changes = { - field: change - for field, change in (changes or {}).items() - if field in _INTERNET_CASE_RELEVANT_CHANGE_FIELDS - } - if not relevant_changes: - return None - - title = f"Internet ændring {reference or connection_name} - faktura {invoice_number}" - existing = execute_query_single( - """ - SELECT id - FROM sag_sager - WHERE deleted_at IS NULL - AND titel = %s - ORDER BY id DESC - LIMIT 1 - """, - (title,), + outcome = ensure_external_change_case( + connection_id=connection_id, + source_type="globalconnect_invoice", + source_key=str(invoice_number), + source_label=f"GlobalConnect faktura {invoice_number}", + source_url="/billing/supplier-invoices", + reference=reference, + connection_name=connection_name, + provider="GlobalConnect", + owner_customer_id=owner_customer_id, + changes=changes, ) - if existing: - return int(existing["id"]) - - assigned_group_id = _resolve_group_id_by_name_tokens(["økonomi", "okonomi", "economic"]) - try: - case_customer_id = owner_customer_id or _resolve_procurement_customer_id() - except Exception as exc: - logger.warning( - "Skipping automatic internet change case for connection %s on invoice %s: could not resolve case customer (%s)", - connection_id, - invoice_number, - exc, - ) - return None - change_lines = "\n".join( - f"- {field}: {change.get('from')} -> {change.get('to')}" - for field, change in relevant_changes.items() - ) - description = ( - "Automatisk oprettet ved import af internetfaktura.\n" - f"Forbindelse: {connection_name}\n" - f"Reference: {reference or '-'}\n" - f"Faktura: {invoice_number}\n" - "Registrerede ændringer:\n" - f"{change_lines}" - ) - - try: - try: - row = execute_query_single( - """ - INSERT INTO sag_sager ( - titel, beskrivelse, type, status, customer_id, assigned_group_id, created_by_user_id - ) - VALUES (%s, %s, %s, %s, %s, %s, %s) - RETURNING id - """, - (title, description, _PURCHASE_CASE_TYPE, "åben", case_customer_id, assigned_group_id, 1), - ) - except Exception as insert_error: - if 'column "type"' not in str(insert_error): - raise - row = execute_query_single( - """ - INSERT INTO sag_sager ( - titel, beskrivelse, status, customer_id, assigned_group_id, created_by_user_id - ) - VALUES (%s, %s, %s, %s, %s, %s) - RETURNING id - """, - (title, description, "åben", case_customer_id, assigned_group_id, 1), - ) - except Exception as exc: - logger.warning( - "Skipping automatic internet change case for connection %s on invoice %s: %s", - connection_id, - invoice_number, - exc, - ) - return None - return int(row["id"]) if row and row.get("id") else None + return outcome.get("case_id") def _ensure_case_for_supplier_invoice( @@ -1573,12 +1499,16 @@ def _upsert_globalconnect_ip_range(connection_id: int, line: Dict, invoice_numbe {"cidr": cidr, "changes": changes}, ) reference = display_reference or cidr + connection_owner = execute_query_single( + "SELECT customer_id FROM internet_connections_connections WHERE id=%s", + (connection_id,), + ) or {} _ensure_internet_change_case( connection_id=connection_id, invoice_number=invoice_number, reference=reference, connection_name=f"IP-range {cidr}", - owner_customer_id=matched_customer["id"] if matched_customer else None, + owner_customer_id=connection_owner.get("customer_id"), changes=changes, ) return range_id @@ -1623,6 +1553,28 @@ def _upsert_globalconnect_ip_range(connection_id: int, line: Dict, invoice_numbe f"IP-range {cidr} oprettet fra faktura {invoice_number}", {"cidr": cidr, "provider_reference": params[1], "contract_number": params[2]}, ) + # A range added to an existing connection is a commercial change. When the + # connection itself was created by this invoice, it is merely initial data. + created_with_invoice = execute_query_single( + """SELECT 1 FROM internet_connections_history + WHERE connection_id=%s AND event_type='connection_created_from_supplier_invoice' + AND details->>'invoice_number'=%s LIMIT 1""", + (connection_id, invoice_number), + ) + if not created_with_invoice: + connection = execute_query_single( + """SELECT name, circuit_number, customer_id, provider + FROM internet_connections_connections WHERE id=%s""", + (connection_id,), + ) or {} + _ensure_internet_change_case( + connection_id=connection_id, + invoice_number=invoice_number, + reference=str(connection.get("circuit_number") or display_reference or cidr), + connection_name=str(connection.get("name") or f"IP-range {cidr}"), + owner_customer_id=connection.get("customer_id"), + changes={"range_added": {"from": None, "to": cidr}}, + ) return range_id @@ -1998,6 +1950,22 @@ def _sync_globalconnect_extraction_to_internet_impl(extraction_row: Dict, simula connection_groups=len(grouped_connections), ip_range_candidates=len(ip_range_lines), ) + case_creation_errors = [] + if not simulate: + try: + case_creation_errors = execute_query( + """SELECT connection_id, last_error AS error + FROM internet_connection_change_cases + WHERE source_type='globalconnect_invoice' AND source_key=%s + AND last_error IS NOT NULL + ORDER BY connection_id""", + (str(invoice_number),), + ) or [] + except Exception as exc: + # The invoice still completes during staggered migration rollout. + logger.warning("Could not load internet change-case control report: %s", exc) + verification["change_case_errors"] = case_creation_errors + verification["requires_manual_review"] = bool(case_creation_errors) or bool(verification.get("requires_manual_review")) return { "skipped": False, @@ -2012,6 +1980,7 @@ def _sync_globalconnect_extraction_to_internet_impl(extraction_row: Dict, simula "line_audit": line_audit, "skipped_items": [entry for entry in line_audit if entry["status"] == "skipped"], "verification": verification, + "change_case_errors": case_creation_errors, } diff --git a/app/core/config.py b/app/core/config.py index f69daa9..f54a7b2 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -225,6 +225,8 @@ class Settings(BaseSettings): ARCHIVED_VTIGER_SYNC_INTERVAL_MINUTES: int = 30 ARCHIVED_VTIGER_SYNC_LIMIT: int = 5000 ARCHIVED_VTIGER_SYNC_INCLUDE_MESSAGES: bool = False + PROJECT_CT_ARCHIVE_SYNC_ENABLED: bool = True + PROJECT_CT_ARCHIVE_SYNC_INTERVAL_MINUTES: int = 360 # Backup System Configuration BACKUP_ENABLED: bool = True diff --git a/app/jobs/project_ct_archive_sync.py b/app/jobs/project_ct_archive_sync.py new file mode 100644 index 0000000..8faf9f3 --- /dev/null +++ b/app/jobs/project_ct_archive_sync.py @@ -0,0 +1,20 @@ +"""Scheduled permanent Project CT vTiger archive sync.""" + +import logging + +from app.admin.vtiger_archive import run_archive_sync + +logger = logging.getLogger(__name__) + + +async def run_project_ct_archive_sync() -> None: + try: + result = await run_archive_sync("incremental") + logger.info("Project CT incremental archive sync completed: %s", result) + except RuntimeError as exc: + if "kører allerede" in str(exc): + logger.info("Project CT scheduled sync skipped: %s", exc) + return + logger.exception("Project CT scheduled archive sync failed") + except Exception: + logger.exception("Project CT scheduled archive sync failed") diff --git a/app/models/schemas.py b/app/models/schemas.py index 662b72d..cb2005d 100644 --- a/app/models/schemas.py +++ b/app/models/schemas.py @@ -3,7 +3,7 @@ Pydantic Models and Schemas """ from enum import Enum -from pydantic import BaseModel, ConfigDict +from pydantic import BaseModel, ConfigDict, Field from typing import Optional, List from datetime import datetime, date diff --git a/app/modules/internet_connections/backend/change_case_service.py b/app/modules/internet_connections/backend/change_case_service.py new file mode 100644 index 0000000..9cc691d --- /dev/null +++ b/app/modules/internet_connections/backend/change_case_service.py @@ -0,0 +1,168 @@ +"""Create deduplicated procurement cases for externally detected internet changes.""" + +from __future__ import annotations + +import json +import logging +from typing import Any, Mapping, Optional + +from app.core.config import settings +from app.core.database import execute_query, execute_query_single + +logger = logging.getLogger(__name__) + +RELEVANT_FIELDS = { + "address", "service_address", "monthly_cost", "sales_price", "technology", + "connection_type", "circuit_number", "provider_reference", "provider", "vendor_id", + "speed_mbps", "download_mbps", "upload_mbps", "status", "sla_subscription_id", + "sla_price", "sla_status", "ip_range", "cidr", "contract_number", "range_added", + "range_removed", "range_monthly_cost", "range_sales_price", +} + +FIELD_LABELS = { + "address": "Adresse", "service_address": "Serviceadresse", "monthly_cost": "Indkøbspris", + "sales_price": "Salgspris", "technology": "Teknologi", "connection_type": "Forbindelsestype", + "circuit_number": "Kredsløbsnummer", "provider_reference": "Leverandørreference", + "provider": "Leverandør", "vendor_id": "Leverandør", "speed_mbps": "Hastighed", + "download_mbps": "Download", "upload_mbps": "Upload", "status": "Status", + "sla_subscription_id": "SLA-aftale", "sla_price": "SLA-pris", "sla_status": "SLA-status", + "ip_range": "IP-range", "cidr": "IP-range", "contract_number": "Kontraktnummer", + "range_added": "Nyt IP-range", "range_removed": "Fjernet IP-range", + "range_monthly_cost": "IP-range indkøbspris", "range_sales_price": "IP-range salgspris", +} + + +def filter_relevant_changes(changes: Mapping[str, Any] | None) -> dict[str, dict[str, Any]]: + filtered: dict[str, dict[str, Any]] = {} + for field, raw in (changes or {}).items(): + if field not in RELEVANT_FIELDS: + continue + change = raw if isinstance(raw, Mapping) else {"from": None, "to": raw} + before, after = change.get("from"), change.get("to") + if str(before or "").strip() == str(after or "").strip(): + continue + filtered[field] = {"from": before, "to": after} + return filtered + + +def _procurement_customer_id() -> int: + configured = getattr(settings, "PROCUREMENT_CASE_CUSTOMER_ID", None) + if configured: + row = execute_query_single("SELECT id FROM customers WHERE id=%s AND is_active=true", (configured,)) + if row: + return int(row["id"]) + row = execute_query_single( + """SELECT id FROM customers WHERE is_active=true AND LOWER(name) LIKE %s + ORDER BY CASE WHEN LOWER(name) LIKE %s THEN 0 ELSE 1 END, id LIMIT 1""", + ("%bmc%", "%bmc networks%"), + ) + if not row: + raise ValueError("BMC's interne indkøbskunde blev ikke fundet") + return int(row["id"]) + + +def _economy_group_id() -> Optional[int]: + row = execute_query_single( + """SELECT id FROM groups WHERE LOWER(name) LIKE ANY(%s) + ORDER BY id LIMIT 1""", (["%økonomi%", "%okonomi%", "%economic%"],), + ) + return int(row["id"]) if row else None + + +def _render_description( + *, connection_id: int, connection_name: str, reference: str, provider: str, + source_label: str, source_url: Optional[str], changes: Mapping[str, Mapping[str, Any]], +) -> str: + lines = [ + "Automatisk oprettet efter en ekstern ændring af en internetforbindelse.", "", + f"Forbindelse: {connection_name}", f"Reference: {reference or '-'}", + f"Leverandør: {provider or '-'}", f"Kilde: {source_label}", + f"Link til forbindelse: /economy/internet-connections/{connection_id}", + ] + if source_url: + lines.append(f"Link til kilde: {source_url}") + lines.extend(["", "Registrerede ændringer:"]) + for field, change in changes.items(): + lines.append(f"- {FIELD_LABELS.get(field, field)}: {change.get('from')} → {change.get('to')}") + return "\n".join(lines) + + +def ensure_external_change_case( + *, connection_id: int, source_type: str, source_key: str, source_label: str, + changes: Mapping[str, Any], connection_name: str = "Internetforbindelse", + reference: str = "", provider: str = "", owner_customer_id: Optional[int] = None, + source_url: Optional[str] = None, +) -> dict[str, Any]: + relevant = filter_relevant_changes(changes) + if not relevant: + return {"case_id": None, "created": False, "changes": {}, "error": None} + source_type = str(source_type or "external").strip().lower() + source_key = str(source_key or "").strip() + if not source_key: + return {"case_id": None, "created": False, "changes": relevant, "error": "Kilden mangler en stabil nøgle"} + + try: + existing = execute_query_single( + """SELECT id, sag_id, changes FROM internet_connection_change_cases + WHERE connection_id=%s AND source_type=%s AND source_key=%s""", + (connection_id, source_type, source_key), + ) + merged = dict((existing or {}).get("changes") or {}) + merged.update(relevant) + case_customer_id = int(owner_customer_id) if owner_customer_id else _procurement_customer_id() + title = f"Internetændring {reference or connection_name} · {source_label}"[:255] + description = _render_description( + connection_id=connection_id, connection_name=connection_name, reference=reference, + provider=provider, source_label=source_label, source_url=source_url, changes=merged, + ) + case_id = int(existing["sag_id"]) if existing and existing.get("sag_id") else None + created = False + if case_id: + execute_query( + "UPDATE sag_sager SET beskrivelse=%s, updated_at=NOW() WHERE id=%s AND deleted_at IS NULL", + (description, case_id), fetch=False, + ) + else: + row = execute_query_single( + """INSERT INTO sag_sager + (titel, beskrivelse, type, status, customer_id, assigned_group_id, created_by_user_id) + VALUES (%s,%s,'indkøb','åben',%s,%s,1) RETURNING id""", + (title, description, case_customer_id, _economy_group_id()), + ) + case_id = int(row["id"]) + created = True + + if existing: + execute_query( + """UPDATE internet_connection_change_cases SET sag_id=%s, source_label=%s, + source_url=%s, changes=%s::jsonb, last_error=NULL, updated_at=NOW() WHERE id=%s""", + (case_id, source_label, source_url, json.dumps(merged, ensure_ascii=False), existing["id"]), fetch=False, + ) + else: + execute_query( + """INSERT INTO internet_connection_change_cases + (connection_id,source_type,source_key,source_label,source_url,sag_id,changes) + VALUES (%s,%s,%s,%s,%s,%s,%s::jsonb)""", + (connection_id, source_type, source_key, source_label, source_url, case_id, + json.dumps(merged, ensure_ascii=False)), fetch=False, + ) + return {"case_id": case_id, "created": created, "changes": merged, "error": None} + except Exception as exc: + logger.warning("Could not create internet change case for connection %s: %s", connection_id, exc) + # Keep the import operational, but persist a visible control item whenever + # the audit table itself is available. + try: + execute_query( + """INSERT INTO internet_connection_change_cases + (connection_id,source_type,source_key,source_label,source_url,changes,last_error) + VALUES (%s,%s,%s,%s,%s,%s::jsonb,%s) + ON CONFLICT (connection_id,source_type,source_key) DO UPDATE + SET changes=internet_connection_change_cases.changes || EXCLUDED.changes, + last_error=EXCLUDED.last_error, updated_at=NOW()""", + (connection_id, source_type, source_key, source_label, source_url, + json.dumps(relevant, ensure_ascii=False), str(exc)), + fetch=False, + ) + except Exception: + logger.exception("Could not persist failed internet change-case audit") + return {"case_id": None, "created": False, "changes": relevant, "error": str(exc)} diff --git a/app/modules/internet_connections/backend/router.py b/app/modules/internet_connections/backend/router.py index cad15ae..a453fb3 100644 --- a/app/modules/internet_connections/backend/router.py +++ b/app/modules/internet_connections/backend/router.py @@ -1,4 +1,5 @@ import ipaddress +import hashlib import io import logging import re @@ -26,6 +27,7 @@ from app.modules.internet_connections.backend.provisioning_utils import ( build_network_product_profile, summarize_subscription_network_requirements, ) +from app.modules.internet_connections.backend.change_case_service import ensure_external_change_case from app.services.ollama_service import ollama_service logger = logging.getLogger(__name__) @@ -1862,15 +1864,21 @@ async def import_ip_nordic_connections(file: UploadFile = File(...), commit: boo filename = str(file.filename or "") if not filename.lower().endswith(".xlsx"): raise HTTPException(status_code=400, detail="Vælg en .xlsx-fil fra IP Nordic") - items = _parse_ip_nordic_xlsx(await file.read()) + content = await file.read() + import_key = hashlib.sha256(content).hexdigest() + items = _parse_ip_nordic_xlsx(content) created_count = 0 + updated_count = 0 skipped_count = 0 + change_case_ids: set[int] = set() + case_errors: List[Dict[str, Any]] = [] preview_items: List[Dict[str, Any]] = [] for item in items: existing = execute_query_single( """ - SELECT id, name, address + SELECT id, name, address, customer_id, monthly_cost, sales_price, + provider, status, technology, connection_type, circuit_number FROM internet_connections_connections WHERE deleted_at IS NULL AND LOWER(COALESCE(provider, '')) = LOWER(%s) @@ -1882,10 +1890,51 @@ async def import_ip_nordic_connections(file: UploadFile = File(...), commit: boo """, ("IP Nordic", _normalize_service_location(item["address"])), ) - action = "skip" if existing else "create" + changes: Dict[str, Dict[str, Any]] = {} + if existing: + desired = { + "monthly_cost": item["monthly_cost"], + "sales_price": item["sales_price"], + } + for field, after in desired.items(): + before = existing.get(field) + if Decimal(str(before or 0)) != Decimal(str(after or 0)): + changes[field] = {"from": before, "to": after} + action = "update" if changes else ("skip" if existing else "create") connection_id = int(existing["id"]) if existing else None - if commit and not existing: + if commit and existing and changes: + execute_query( + """UPDATE internet_connections_connections + SET monthly_cost=%s, sales_price=%s, updated_at=CURRENT_TIMESTAMP + WHERE id=%s""", + (item["monthly_cost"], item["sales_price"], connection_id), + fetch=False, + ) + _create_history_entry( + connection_id, + "ip_nordic_import_changed", + f"Opdateret fra IP Nordic-filen {filename}", + {"source_file": filename, "import_key": import_key, "changes": changes}, + ) + outcome = ensure_external_change_case( + connection_id=connection_id, + source_type="ip_nordic_spreadsheet", + source_key=import_key, + source_label=f"IP Nordic import {filename}", + source_url="/economy/internet-connections", + changes=changes, + connection_name=str(existing.get("name") or f"IP Nordic · {item['address']}"), + reference=str(existing.get("circuit_number") or item["address"]), + provider="IP Nordic", + owner_customer_id=existing.get("customer_id"), + ) + if outcome.get("case_id"): + change_case_ids.add(int(outcome["case_id"])) + if outcome.get("error"): + case_errors.append({"connection_id": connection_id, "address": item["address"], "error": outcome["error"]}) + updated_count += 1 + elif commit and not existing: notes = ( f"Importeret fra {filename}. Leverandørens firmanr.: {item['company_number']}. " f"Rapporteret firma: {item['reported_company']}. {item['line_count']} regnearkslinje(r) samlet. " @@ -1930,6 +1979,7 @@ async def import_ip_nordic_connections(file: UploadFile = File(...), commit: boo "sales_price": float(item["sales_price"]), "monthly_cost": float(item["monthly_cost"]), "line_count": item["line_count"], + "changes": changes, }) return { @@ -1938,9 +1988,14 @@ async def import_ip_nordic_connections(file: UploadFile = File(...), commit: boo "items": preview_items, "total": len(preview_items), "create_count": sum(1 for item in preview_items if item["action"] == "create"), - "existing_count": sum(1 for item in preview_items if item["action"] == "skip"), + "existing_count": sum(1 for item in preview_items if item["action"] in {"skip", "update"}), "created_count": created_count, + "updated_count": updated_count, "skipped_count": skipped_count, + "change_case_ids": sorted(change_case_ids), + "case_errors": case_errors, + "requires_manual_review": bool(case_errors), + "import_key": import_key, "customer_auto_assignment": False, } @@ -2182,6 +2237,27 @@ async def get_connection(connection_id: int): return connection +@router.get("/internet-connections/{connection_id:int}/cases", response_model=List[dict]) +async def list_connection_cases(connection_id: int): + """Direct case history only; this never infers or changes the connection's customer.""" + if not execute_query_single( + "SELECT id FROM internet_connections_connections WHERE id=%s AND deleted_at IS NULL", (connection_id,), + ): + raise HTTPException(status_code=404, detail="Connection not found") + return execute_query( + """SELECT s.id,s.titel,s.status,s.template_key,s.customer_id,s.updated_at,s.created_at, + c.name AS customer_name, + COALESCE(NULLIF(TRIM(u.full_name),''),u.username) AS responsible_name, + link.created_at AS linked_at + FROM sag_internet_connections link + JOIN sag_sager s ON s.id=link.sag_id AND s.deleted_at IS NULL + LEFT JOIN customers c ON c.id=s.customer_id + LEFT JOIN users u ON u.user_id=s.ansvarlig_bruger_id + WHERE link.connection_id=%s + ORDER BY s.updated_at DESC NULLS LAST,s.id DESC""", (connection_id,), + ) or [] + + @router.get("/internet-connections/{connection_id:int}/cross-field-ports") async def get_connection_cross_field_ports(connection_id: int): """Find cross-field ports related by customer or service-location address.""" diff --git a/app/modules/internet_connections/templates/detail.html b/app/modules/internet_connections/templates/detail.html index 286da7a..045315a 100644 --- a/app/modules/internet_connections/templates/detail.html +++ b/app/modules/internet_connections/templates/detail.html @@ -325,12 +325,16 @@ + + Opret sag + @@ -674,6 +678,14 @@ +
+
+
Sager på forbindelsen
Kun direkte tilknyttede sager.
+ Opret sag +
+
+
+
Historik
@@ -1531,6 +1543,7 @@ pricingHistoryResponse, contractsResponse, historyResponse, + casesResponse, crossFieldPortsResponse, ] = await Promise.all([ fetch(`/api/v1/internet-connections/${connectionId}`), @@ -1541,6 +1554,7 @@ fetch(`/api/v1/internet-connections/${connectionId}/pricing/history`), fetch('/api/v1/internet-connections/contracts'), fetch(`/api/v1/internet-connections/${connectionId}/history`), + fetch(`/api/v1/internet-connections/${connectionId}/cases`), fetch(`/api/v1/internet-connections/${connectionId}/cross-field-ports`), ]); @@ -1558,12 +1572,14 @@ const pricingHistory = await safeJson(pricingHistoryResponse, []); const contracts = await safeJson(contractsResponse, []); const history = await safeJson(historyResponse, []); + const cases = await safeJson(casesResponse, []); const crossFieldPorts = await safeJson(crossFieldPortsResponse, { items: [], summary: {} }); const failedSections = [ [rangesResponse, 'IP-ranges'], [addressesResponse, 'IP-adresser'], [summaryResponse, 'IP-oversigt'], [pricingResponse, 'priser'], [pricingHistoryResponse, 'prishistorik'], [contractsResponse, 'kontrakter'], [historyResponse, 'historik'], [crossFieldPortsResponse, 'krydsfelt'], + [casesResponse, 'sager'], ].filter(([response]) => !response.ok).map(([, label]) => label); if (failedSections.length) { const feedback = document.getElementById('detailSaveFeedback'); @@ -1585,10 +1601,27 @@ renderAddresses(currentAddresses); renderPricing(pricing, pricingHistory); renderHistory(history); + renderConnectionCases(cases); renderRelationGrid(connection); renderBmcnetChildren(connection, currentBmcnetChildren); renderContractsOverview(contracts); renderCrossFieldPorts(crossFieldPorts); + const createCaseUrl = `/sag/new?internet_connection_id=${encodeURIComponent(connectionId)}`; + document.getElementById('createCaseForConnectionBtn').href = createCaseUrl; + document.getElementById('createCaseForConnectionPanelBtn').href = createCaseUrl; + } + + function renderConnectionCases(cases) { + const list = document.getElementById('connectionCasesList'); + if (!Array.isArray(cases) || !cases.length) { + list.innerHTML = '
Ingen sager er knyttet direkte til forbindelsen endnu.
'; + return; + } + list.innerHTML = cases.map(item => ` + +
SAG-${item.id} · ${escapeHtml(item.titel || 'Uden titel')}
${escapeHtml(item.customer_name || 'Ingen kunde')} · ${escapeHtml(item.responsible_name || 'Ikke tildelt')} · ændret ${formatDateTime(item.updated_at)}
+ ${escapeHtml(item.status || '-')} +
`).join(''); } function renderCrossFieldPorts(payload) { diff --git a/app/modules/sag/backend/router.py b/app/modules/sag/backend/router.py index 8dada3b..c43b963 100644 --- a/app/modules/sag/backend/router.py +++ b/app/modules/sag/backend/router.py @@ -1030,6 +1030,14 @@ async def create_sag(request: Request, data: dict): if contact_id and contact_id not in contact_ids: contact_ids.append(contact_id) telefoni_opkald_id = _coerce_optional_int(data.get("telefoni_opkald_id"), "telefoni_opkald_id") + raw_connection_ids = data.get("internet_connection_ids") or [] + if not isinstance(raw_connection_ids, list): + raise HTTPException(status_code=400, detail="internet_connection_ids skal være en liste") + internet_connection_ids = [] + for raw_connection_id in raw_connection_ids: + connection_id = _coerce_optional_int(raw_connection_id, "internet_connection_id") + if connection_id and connection_id not in internet_connection_ids: + internet_connection_ids.append(connection_id) if pipeline is not None and not isinstance(pipeline, dict): raise HTTPException(status_code=400, detail="pipeline skal være et objekt") if not isinstance(order_items, list): @@ -1134,6 +1142,23 @@ async def create_sag(request: Request, data: dict): (contact_id, data.get("customer_id")), ) + if internet_connection_ids: + cursor.execute( + """SELECT id FROM internet_connections_connections + WHERE id = ANY(%s) AND deleted_at IS NULL""", + (internet_connection_ids,), + ) + existing_connection_ids = {int(row["id"]) for row in cursor.fetchall()} + missing_connection_ids = [item for item in internet_connection_ids if item not in existing_connection_ids] + if missing_connection_ids: + raise HTTPException(status_code=400, detail=f"Internetforbindelsen findes ikke: {missing_connection_ids[0]}") + for connection_id in internet_connection_ids: + cursor.execute( + """INSERT INTO sag_internet_connections (sag_id, connection_id, linked_by_user_id) + VALUES (%s,%s,%s) ON CONFLICT DO NOTHING""", + (result["id"], connection_id, current_user_id), + ) + if telefoni_opkald_id: cursor.execute( """UPDATE telefoni_opkald @@ -1176,6 +1201,45 @@ async def create_sag(request: Request, data: dict): logger.error("❌ Error creating case: %s", e) raise HTTPException(status_code=500, detail="Failed to create case") + +@router.get("/sag/{sag_id}/internet-connections") +async def list_sag_internet_connections(sag_id: int): + return execute_query( + """SELECT link.connection_id,link.created_at,ic.name,ic.circuit_number,ic.provider_reference, + ic.address,ic.customer_id + FROM sag_internet_connections link + JOIN internet_connections_connections ic ON ic.id=link.connection_id AND ic.deleted_at IS NULL + WHERE link.sag_id=%s ORDER BY link.created_at DESC""", (sag_id,), + ) or [] + + +@router.post("/sag/{sag_id}/internet-connections", dependencies=[Depends(case_edit_access)]) +async def link_sag_internet_connection(sag_id: int, request: Request, data: dict): + connection_id = _coerce_optional_int(data.get("connection_id"), "connection_id") + if not connection_id: + raise HTTPException(status_code=400, detail="connection_id er påkrævet") + if not execute_query_single("SELECT id FROM sag_sager WHERE id=%s AND deleted_at IS NULL", (sag_id,)): + raise HTTPException(status_code=404, detail="Sagen findes ikke") + if not execute_query_single( + "SELECT id FROM internet_connections_connections WHERE id=%s AND deleted_at IS NULL", (connection_id,), + ): + raise HTTPException(status_code=404, detail="Internetforbindelsen findes ikke") + execute_query( + """INSERT INTO sag_internet_connections (sag_id,connection_id,linked_by_user_id) + VALUES (%s,%s,%s) ON CONFLICT DO NOTHING""", + (sag_id, connection_id, _get_user_id_from_request(request)), fetch=False, + ) + return {"linked": True, "sag_id": sag_id, "connection_id": connection_id} + + +@router.delete("/sag/{sag_id}/internet-connections/{connection_id}", dependencies=[Depends(case_edit_access)]) +async def unlink_sag_internet_connection(sag_id: int, connection_id: int): + execute_query( + "DELETE FROM sag_internet_connections WHERE sag_id=%s AND connection_id=%s", + (sag_id, connection_id), fetch=False, + ) + return {"unlinked": True} + @router.get("/sag/{sag_id:int}") async def get_sag(sag_id: int): """Get a specific case.""" diff --git a/app/modules/sag/backend/solutions.py b/app/modules/sag/backend/solutions.py index 6339e7e..4c88ba2 100644 --- a/app/modules/sag/backend/solutions.py +++ b/app/modules/sag/backend/solutions.py @@ -1,5 +1,6 @@ import json import logging +import re from typing import Optional from fastapi import APIRouter, Depends, HTTPException, Query, Request @@ -7,6 +8,7 @@ from fastapi import APIRouter, Depends, HTTPException, Query, Request from app.core.auth_dependencies import get_current_user, require_any_permission from app.core.database import execute_query, execute_query_single from app.models.schemas import Solution, SolutionCreate, SolutionUpdate +from app.services.case_analysis_service import CaseAnalysisService logger = logging.getLogger(__name__) router = APIRouter() @@ -16,6 +18,13 @@ APPROVAL_STATUSES = {"draft", "pending", "approved", "outdated", "rejected"} RESULT_ALIASES = {"resolved": "Løst", "partial": "Delvist", "unresolved": "Ej løst", "løst": "Løst", "delvist": "Delvist", "workaround": "Workaround", "ej løst": "Ej løst"} TYPE_ALIASES = {"standard": "Support", "permanent": "Support", "external": "Ekstern", "support": "Support", "drift": "Drift", "konsulent": "Konsulent", "infrastruktur": "Infrastruktur", "workaround": "Support"} +SECRET_PATTERNS = ( + (re.compile(r"(?i)\b(password|passwd|kodeord|api[_ -]?key|secret|token)\s*[:=]\s*([^\s,;]+)"), r"\1: [FJERNET]"), + (re.compile(r"(?i)\bBearer\s+[A-Za-z0-9._~+\-/]+=*"), "Bearer [FJERNET]"), + (re.compile(r"-----BEGIN [A-Z ]*PRIVATE KEY-----.*?-----END [A-Z ]*PRIVATE KEY-----", re.S), "[PRIVAT NØGLE FJERNET]"), + (re.compile(r"(?i)(?:postgres|mysql|mongodb(?:\+srv)?)://[^\s]+"), "[DATABASEFORBINDELSE FJERNET]"), +) + def _user_id(current_user: dict) -> Optional[int]: value = current_user.get("id") or current_user.get("user_id") @@ -58,6 +67,43 @@ def _normalize_payload(data: dict) -> dict: return data +def _redact_sensitive(text: str) -> tuple[str, list[str]]: + cleaned = str(text or "") + warnings = [] + for pattern, replacement in SECRET_PATTERNS: + cleaned, count = pattern.subn(replacement, cleaned) + if count: + warnings.append(f"{count} mulig(e) hemmelighed(er) blev fjernet før AI-behandling") + return cleaned, warnings + + +def _knowledge_tokens(text: str) -> list[str]: + ignored = {"eller", "ikke", "med", "den", "det", "der", "som", "for", "fra", "til", "har", "kan", "skal", "sag", "sagen", "test", "viden"} + tokens, seen = [], set() + for token in re.findall(r"[A-Za-zÀ-ÿ0-9_.-]{3,}", str(text or "").lower()): + if token in ignored or token in seen: + continue + seen.add(token) + tokens.append(token) + return tokens[:24] + + +def _normalize_source_refs(raw_refs, allowed_refs: set[str], sag_id: int) -> list[str]: + if isinstance(raw_refs, str): + raw_refs = re.split(r"[,;\n]+", raw_refs) + normalized = [f"Sag {sag_id}"] + for raw in raw_refs or []: + value = str(raw or "").strip().strip("[]") + match = re.search(r"(?i)\b(sag|kommentar|mail|tid|artikel)\s*#?\s*(\d+)\b", value) + if not match: + continue + kind = match.group(1).capitalize() + canonical = f"{kind} {int(match.group(2))}" + if canonical in allowed_refs and canonical not in normalized: + normalized.append(canonical) + return normalized + + def _version_solution(solution: dict, user_id: Optional[int], change_note: Optional[str] = None) -> None: version_row = execute_query_single("SELECT COALESCE(MAX(version_number), 0) + 1 AS next_version FROM sag_solution_versions WHERE solution_id = %s", (solution["id"],)) or {"next_version": 1} snapshot = dict(solution) @@ -70,15 +116,47 @@ def _version_solution(solution: dict, user_id: Optional[int], change_note: Optio ) +def _case_knowledge_context(sag_id: int, limit: int = 6) -> tuple[dict, list[dict]]: + case = execute_query_single( + """SELECT s.id,s.titel,s.beskrivelse,s.status, + (SELECT sk.customer_id FROM sag_kunder sk WHERE sk.sag_id=s.id AND sk.deleted_at IS NULL ORDER BY sk.id LIMIT 1) AS customer_id, + COALESCE((SELECT string_agg(t.name,' ') FROM entity_tags et JOIN tags t ON t.id=et.tag_id WHERE et.entity_type='case' AND et.entity_id=s.id),'') AS tag_text + FROM sag_sager s WHERE s.id=%s AND s.deleted_at IS NULL""", + (sag_id,), + ) + if not case: + raise HTTPException(status_code=404, detail="Sagen findes ikke") + tokens = _knowledge_tokens(" ".join([str(case.get("titel") or ""), str(case.get("beskrivelse") or ""), str(case.get("tag_text") or "")])) + if not tokens: + return case, [] + search_query = " OR ".join(tokens) + articles = execute_query( + """SELECT ka.id,ka.title,ka.summary,ka.problem,ka.root_cause,ka.solution,ka.workaround, + ka.visibility,ka.customer_id,ka.tags,ka.products,ka.sag_id,ka.updated_at, + ts_rank_cd(ka.search_document,websearch_to_tsquery('simple',%s)) AS relevance + FROM knowledge_articles ka + WHERE ka.status='published' AND ka.archived_at IS NULL + AND (ka.visibility IN ('general','internal') OR (ka.visibility='customer' AND ka.customer_id=%s)) + AND ka.search_document @@ websearch_to_tsquery('simple',%s) + ORDER BY relevance DESC,ka.updated_at DESC LIMIT %s""", + (search_query, case.get("customer_id"), search_query, limit), + ) or [] + if articles: + top_relevance = float(articles[0].get("relevance") or 0) + relative_floor = max(0.05, top_relevance * 0.35) + articles = [row for row in articles if float(row.get("relevance") or 0) >= relative_floor] + return case, articles + + @router.get("/sag/{sag_id}/solution", response_model=Optional[Solution]) async def get_solution(sag_id: int, _current_user: dict = Depends(get_current_user)): - result = execute_query("SELECT * FROM sag_solutions WHERE sag_id = %s", (sag_id,)) + result = execute_query("SELECT * FROM sag_solutions WHERE sag_id = %s AND deleted_at IS NULL", (sag_id,)) return result[0] if result else None @router.get("/sag/{sag_id}/solution/versions") async def get_solution_versions(sag_id: int, _current_user: dict = Depends(get_current_user)): - solution = execute_query_single("SELECT id FROM sag_solutions WHERE sag_id = %s", (sag_id,)) + solution = execute_query_single("SELECT id FROM sag_solutions WHERE sag_id = %s AND deleted_at IS NULL", (sag_id,)) if not solution: return {"items": [], "total": 0} items = execute_query( @@ -95,8 +173,10 @@ async def get_solution_versions(sag_id: int, _current_user: dict = Depends(get_c async def create_solution(sag_id: int, solution: SolutionCreate, current_user: dict = Depends(get_current_user)): if not execute_query_single("SELECT id FROM sag_sager WHERE id=%s AND deleted_at IS NULL", (sag_id,)): raise HTTPException(status_code=404, detail="Sagen findes ikke") - if execute_query_single("SELECT id FROM sag_solutions WHERE sag_id=%s", (sag_id,)): - raise HTTPException(status_code=409, detail="Der findes allerede en løsning på sagen") + existing = execute_query_single("SELECT id, deleted_at FROM sag_solutions WHERE sag_id=%s", (sag_id,)) + if existing: + detail = "Sagen har en arkiveret løsning, som skal gendannes" if existing.get("deleted_at") else "Der findes allerede en løsning på sagen" + raise HTTPException(status_code=409, detail=detail) data = _normalize_payload(solution.model_dump(exclude={"sag_id", "created_by_user_id"})) user_id = _user_id(current_user) result = execute_query( @@ -113,7 +193,7 @@ async def create_solution(sag_id: int, solution: SolutionCreate, current_user: d @router.patch("/sag/{sag_id}/solution", response_model=Solution, dependencies=[Depends(case_edit_access)]) async def update_solution(sag_id: int, updates: SolutionUpdate, current_user: dict = Depends(get_current_user)): - if not execute_query_single("SELECT id FROM sag_solutions WHERE sag_id=%s", (sag_id,)): + if not execute_query_single("SELECT id FROM sag_solutions WHERE sag_id=%s AND deleted_at IS NULL", (sag_id,)): raise HTTPException(status_code=404, detail="Løsningen findes ikke") data = _normalize_payload(updates.model_dump(exclude_unset=True)) change_note = data.pop("change_note", None) @@ -137,7 +217,7 @@ async def update_solution(sag_id: int, updates: SolutionUpdate, current_user: di @router.post("/sag/{sag_id}/solution/workflow", dependencies=[Depends(case_edit_access)]) async def solution_workflow(sag_id: int, request: Request, current_user: dict = Depends(get_current_user)): action = str((await request.json()).get("action") or "").lower() - solution = execute_query_single("SELECT * FROM sag_solutions WHERE sag_id=%s", (sag_id,)) + solution = execute_query_single("SELECT * FROM sag_solutions WHERE sag_id=%s AND deleted_at IS NULL", (sag_id,)) if not solution: raise HTTPException(status_code=404, detail="Løsningen findes ikke") user_id = _user_id(current_user) @@ -161,11 +241,15 @@ async def solution_workflow(sag_id: int, request: Request, current_user: dict = @router.post("/sag/{sag_id}/solution/publish", dependencies=[Depends(case_edit_access)]) async def publish_solution(sag_id: int, current_user: dict = Depends(get_current_user)): - solution = execute_query_single("SELECT * FROM sag_solutions WHERE sag_id=%s", (sag_id,)) + solution = execute_query_single("SELECT * FROM sag_solutions WHERE sag_id=%s AND deleted_at IS NULL", (sag_id,)) if not solution: raise HTTPException(status_code=404, detail="Løsningen findes ikke") if solution.get("approval_status") != "approved": raise HTTPException(status_code=409, detail="Løsningen skal godkendes før udgivelse") + publish_text = "\n".join(str(solution.get(field) or "") for field in ("title", "problem", "root_cause", "investigation", "description", "workaround")) + _clean_publish_text, secret_warnings = _redact_sensitive(publish_text) + if secret_warnings: + raise HTTPException(status_code=422, detail="Løsningen indeholder muligvis password, token eller anden hemmelighed. Fjern det før udgivelse.") customer_id = None if solution.get("visibility") == "customer": customer = execute_query_single("SELECT customer_id FROM sag_kunder WHERE sag_id=%s AND deleted_at IS NULL ORDER BY id LIMIT 1", (sag_id,)) @@ -184,12 +268,197 @@ async def publish_solution(sag_id: int, current_user: dict = Depends(get_current investigation=EXCLUDED.investigation,solution=EXCLUDED.solution,workaround=EXCLUDED.workaround, visibility=EXCLUDED.visibility,tags=EXCLUDED.tags,products=EXCLUDED.products,status='published', version_number=knowledge_articles.version_number+1,published_by_user_id=EXCLUDED.published_by_user_id, - reviewed_at=NOW(),updated_at=NOW() RETURNING *""", + reviewed_at=NOW(),updated_at=NOW(),archived_at=NULL,archived_by_user_id=NULL RETURNING *""", (solution["id"], sag_id, customer_id, solution["title"], summary, solution.get("problem"), solution.get("root_cause"), solution.get("investigation"), description, solution.get("workaround"), solution.get("visibility", "internal"), json.dumps(solution.get("tags") or []), json.dumps(solution.get("products") or []), _user_id(current_user)), ) return article +@router.get("/sag/{sag_id}/knowledge-suggestions") +async def case_knowledge_suggestions( + sag_id: int, + limit: int = Query(6, ge=1, le=20), + _current_user: dict = Depends(get_current_user), +): + _case, articles = _case_knowledge_context(sag_id, limit) + items = [] + for article in articles: + item = dict(article) + item["reason"] = "Matcher sagens titel, beskrivelse eller tags" + item["relevance_percent"] = min(99, max(1, round(float(item.get("relevance") or 0) * 100))) + items.append(item) + return {"items": items, "total": len(items), "source": "approved_knowledge_only"} + + +@router.post("/sag/{sag_id}/solution/ai-draft", dependencies=[Depends(case_edit_access)]) +async def generate_solution_ai_draft(sag_id: int, current_user: dict = Depends(get_current_user)): + case, articles = _case_knowledge_context(sag_id, 5) + comments = execute_query( + """SELECT id,forfatter,indhold,created_at FROM sag_kommentarer + WHERE sag_id=%s AND deleted_at IS NULL ORDER BY created_at DESC LIMIT 60""", + (sag_id,), + ) or [] + emails = execute_query( + """SELECT DISTINCT e.id,e.subject,e.sender_email,e.recipient_email,e.body_text,e.received_date + FROM email_messages e LEFT JOIN sag_emails se ON se.email_id=e.id + WHERE e.deleted_at IS NULL AND (se.sag_id=%s OR e.linked_case_id=%s) + ORDER BY e.received_date DESC NULLS LAST LIMIT 25""", + (sag_id, sag_id), + ) or [] + time_entries = execute_query( + """SELECT id,description,original_hours,worked_date FROM tmodule_times + WHERE sag_id=%s ORDER BY worked_date DESC NULLS LAST LIMIT 30""", + (sag_id,), + ) or [] + + sections = [f"SAG #{sag_id}\nTitel: {case.get('titel') or ''}\nBeskrivelse: {case.get('beskrivelse') or ''}"] + if comments: + sections.append("KOMMENTARER:\n" + "\n".join(f"[Kommentar {row['id']}] {row.get('forfatter') or 'Ukendt'}: {row.get('indhold') or ''}" for row in reversed(comments))) + if emails: + sections.append("MAILS:\n" + "\n".join(f"[Mail {row['id']}] {row.get('subject') or '(uden emne)'}: {row.get('body_text') or ''}" for row in reversed(emails))) + if time_entries: + sections.append("TIDSNOTER:\n" + "\n".join(f"[Tid {row['id']}] {row.get('description') or ''}" for row in time_entries)) + if articles: + sections.append("GODKENDT VIDEN:\n" + "\n".join(f"[Artikel {row['id']}] {row['title']}\nProblem: {row.get('problem') or ''}\nLøsning: {row.get('solution') or ''}" for row in articles)) + context, warnings = _redact_sensitive("\n\n".join(sections)) + context = context[:24000] + prompt = """Du laver et UDKAST til dokumentation af en IT-supportløsning. +Returner KUN gyldig JSON med felterne title, problem, root_cause, investigation, description, workaround, solution_type, result, tags, products, confidence, uncertainties og source_refs. +Regler: +- Brug kun fakta fra den vedlagte sag og de nummererede kilder. +- Opfind aldrig en årsag eller handling. Sæt feltet tomt og beskriv manglen i uncertainties. +- description er den endelige løsning. Hvis sagen endnu ikke dokumenterer en løsning, skal description være tom. +- source_refs skal indeholde de præcise kilder, fx "Kommentar 12" eller "Artikel 4". +- Medtag aldrig passwords, tokens, API-nøgler eller andre hemmeligheder. +- Skriv kort, teknisk og professionelt på dansk. +- confidence er et tal mellem 0 og 1. +""" + service = CaseAnalysisService() + # A full case history is materially larger than QuickCreate input. Keep the + # global QuickCreate timeout unchanged, but allow the local model time to + # produce a source-grounded solution draft. + service.ai_timeout = max(service.ai_timeout, 60) + result = await service._call_ollama(prompt, context) + if not result: + warnings.append("Den lokale AI-tjeneste er utilgængelig; der vises en sikker grundkladde uden AI-genererede konklusioner") + result = { + "title": case.get("titel") or "", + "problem": case.get("beskrivelse") or "", + "root_cause": "", + "investigation": "", + "description": "", + "workaround": "", + "solution_type": "Support", + "result": "Ej løst", + "tags": _knowledge_tokens(case.get("tag_text") or "")[:10], + "products": [], + "confidence": 0.2, + "uncertainties": ["Årsag og endelig løsning skal udfyldes af en medarbejder"], + "source_refs": [f"Sag {sag_id}"], + } + allowed_refs = {f"Kommentar {row['id']}" for row in comments} | {f"Mail {row['id']}" for row in emails} | {f"Tid {row['id']}" for row in time_entries} | {f"Artikel {row['id']}" for row in articles} | {f"Sag {sag_id}"} + source_refs = _normalize_source_refs(result.get("source_refs", []), allowed_refs, sag_id) + draft = { + "title": str(result.get("title") or case.get("titel") or "")[:255], + "problem": str(result.get("problem") or ""), + "root_cause": str(result.get("root_cause") or ""), + "investigation": str(result.get("investigation") or ""), + "description": str(result.get("description") or ""), + "workaround": str(result.get("workaround") or ""), + "solution_type": TYPE_ALIASES.get(str(result.get("solution_type") or "Support").casefold(), "Support"), + "result": RESULT_ALIASES.get(str(result.get("result") or "Ej løst").casefold(), "Ej løst"), + "tags": _clean_list(result.get("tags")), + "products": _clean_list(result.get("products")), + "confidence": max(0.0, min(1.0, float(result.get("confidence") or 0))), + "uncertainties": _clean_list(result.get("uncertainties")), + "source_refs": source_refs, + } + if not draft["description"]: + warnings.append("Sagen indeholder ikke en sikkert dokumenteret endelig løsning") + if not source_refs: + warnings.append("AI-udkastet indeholdt ingen gyldige kildehenvisninger") + return {"draft": draft, "warnings": warnings, "model": service.ollama_model, "review_required": True} + + +@router.get("/solutions") +async def list_solutions( + q: str = Query("", max_length=200), + approval_status: Optional[str] = None, + visibility: Optional[str] = None, + include_archived: bool = False, + limit: int = Query(50, ge=1, le=200), + offset: int = Query(0, ge=0), + _current_user: dict = Depends(get_current_user), +): + where = ["s.deleted_at IS NOT NULL" if include_archived else "s.deleted_at IS NULL"] + params: list = [] + if approval_status: + if approval_status not in APPROVAL_STATUSES: + raise HTTPException(status_code=422, detail="Ugyldig godkendelsesstatus") + where.append("s.approval_status=%s") + params.append(approval_status) + if visibility: + if visibility not in VISIBILITIES: + raise HTTPException(status_code=422, detail="Ugyldig synlighed") + where.append("s.visibility=%s") + params.append(visibility) + term = q.strip() + if term: + where.append("(s.title ILIKE %s OR s.description ILIKE %s OR s.problem ILIKE %s OR sg.titel ILIKE %s OR c.name ILIKE %s OR s.tags::text ILIKE %s OR s.products::text ILIKE %s)") + like = f"%{term}%" + params.extend([like] * 7) + where_sql = " AND ".join(where) + count = execute_query_single( + f"""SELECT COUNT(*) AS total FROM sag_solutions s + JOIN sag_sager sg ON sg.id=s.sag_id + LEFT JOIN LATERAL ( + SELECT cu.name FROM sag_kunder sk JOIN customers cu ON cu.id=sk.customer_id + WHERE sk.sag_id=s.sag_id AND sk.deleted_at IS NULL ORDER BY sk.id LIMIT 1 + ) c ON TRUE WHERE {where_sql}""", + tuple(params), + ) or {"total": 0} + items = execute_query( + f"""SELECT s.*, sg.titel AS case_title, sg.status AS case_status, c.name AS customer_name, + COALESCE(creator.full_name,creator.username) AS created_by, + COALESCE(updater.full_name,updater.username) AS updated_by, + ka.id AS article_id, ka.status AS article_status + FROM sag_solutions s JOIN sag_sager sg ON sg.id=s.sag_id + LEFT JOIN LATERAL ( + SELECT cu.name FROM sag_kunder sk JOIN customers cu ON cu.id=sk.customer_id + WHERE sk.sag_id=s.sag_id AND sk.deleted_at IS NULL ORDER BY sk.id LIMIT 1 + ) c ON TRUE + LEFT JOIN users creator ON creator.user_id=s.created_by_user_id + LEFT JOIN users updater ON updater.user_id=s.updated_by_user_id + LEFT JOIN knowledge_articles ka ON ka.solution_id=s.id + WHERE {where_sql} ORDER BY s.updated_at DESC,s.id DESC LIMIT %s OFFSET %s""", + tuple(params + [limit, offset]), + ) or [] + return {"items": items, "total": int(count["total"]), "limit": limit, "offset": offset, "archived": include_archived} + + +@router.delete("/solutions/{solution_id}", dependencies=[Depends(case_edit_access)]) +async def archive_solution(solution_id: int, current_user: dict = Depends(get_current_user)): + solution = execute_query_single("SELECT * FROM sag_solutions WHERE id=%s AND deleted_at IS NULL", (solution_id,)) + if not solution: + raise HTTPException(status_code=404, detail="Løsningen findes ikke eller er allerede arkiveret") + user_id = _user_id(current_user) + execute_query("UPDATE knowledge_articles SET status='archived',archived_at=NOW(),archived_by_user_id=%s,updated_at=NOW() WHERE solution_id=%s", (user_id, solution_id)) + archived = execute_query_single("UPDATE sag_solutions SET deleted_at=NOW(),deleted_by_user_id=%s,updated_at=NOW() WHERE id=%s RETURNING *", (user_id, solution_id)) + _version_solution(archived, user_id, "Løsning arkiveret") + return {"status": "archived", "id": solution_id, "sag_id": solution["sag_id"]} + + +@router.post("/solutions/{solution_id}/restore", dependencies=[Depends(case_edit_access)]) +async def restore_solution(solution_id: int, current_user: dict = Depends(get_current_user)): + solution = execute_query_single("SELECT * FROM sag_solutions WHERE id=%s AND deleted_at IS NOT NULL", (solution_id,)) + if not solution: + raise HTTPException(status_code=404, detail="Den arkiverede løsning findes ikke") + user_id = _user_id(current_user) + restored = execute_query_single("UPDATE sag_solutions SET deleted_at=NULL,deleted_by_user_id=NULL,updated_by_user_id=%s,updated_at=NOW() WHERE id=%s RETURNING *", (user_id, solution_id)) + _version_solution(restored, user_id, "Løsning gendannet") + return {"status": "restored", "id": solution_id, "sag_id": solution["sag_id"]} + + @router.get("/knowledge/articles") async def search_knowledge_articles(q: str = Query("", max_length=200), customer_id: Optional[int] = None, limit: int = Query(25, ge=1, le=100), offset: int = Query(0, ge=0), _current_user: dict = Depends(get_current_user)): term = q.strip() diff --git a/app/modules/sag/frontend/views.py b/app/modules/sag/frontend/views.py index ae25279..7ab36f4 100644 --- a/app/modules/sag/frontend/views.py +++ b/app/modules/sag/frontend/views.py @@ -19,6 +19,12 @@ async def knowledge_index(request: Request): return templates.TemplateResponse("modules/sag/templates/knowledge_index.html", {"request": request}) +@router.get("/solutions", response_class=HTMLResponse) +async def solutions_management(request: Request): + """Central management surface for case solutions and their publication flow.""" + return templates.TemplateResponse("modules/sag/templates/solutions_management.html", {"request": request}) + + @router.get("/knowledge/{article_id:int}", response_class=HTMLResponse) async def knowledge_detail(request: Request, article_id: int): article = execute_query( @@ -857,7 +863,7 @@ async def sag_detaljer(request: Request, sag_id: int): comments = execute_query(comments_query, (sag_id,)) # Fetch Solution - solution_query = "SELECT * FROM sag_solutions WHERE sag_id = %s" + solution_query = "SELECT * FROM sag_solutions WHERE sag_id = %s AND deleted_at IS NULL" solution_res = execute_query(solution_query, (sag_id,)) solution = solution_res[0] if solution_res else None @@ -1218,7 +1224,7 @@ async def sag_detaljer_v3(request: Request, sag_id: int): comments = execute_query(comments_query, (sag_id,)) # Fetch Solution - solution_query = "SELECT * FROM sag_solutions WHERE sag_id = %s" + solution_query = "SELECT * FROM sag_solutions WHERE sag_id = %s AND deleted_at IS NULL" solution_res = execute_query(solution_query, (sag_id,)) solution = solution_res[0] if solution_res else None diff --git a/app/modules/sag/templates/create.html b/app/modules/sag/templates/create.html index 067a9c8..215654c 100644 --- a/app/modules/sag/templates/create.html +++ b/app/modules/sag/templates/create.html @@ -199,6 +199,7 @@
+

@@ -386,7 +387,7 @@ let customerContactsLoadToken = 0; let successAlertTimeout; let orderLineCounter = 0; - let telefoniPrefill = { contactId: null, title: null, callId: null, customerId: null, description: null }; + let telefoniPrefill = { contactId: null, title: null, callId: null, customerId: null, description: null, internetConnectionId: null }; let topAlertLoadToken = 0; function escapeTopAlertHtml(value) { @@ -843,6 +844,7 @@ const callIdRaw = params.get('telefoni_opkald_id'); const customerIdRaw = params.get('customer_id'); const descriptionRaw = params.get('description'); + const internetConnectionIdRaw = params.get('internet_connection_id'); const contactId = contactIdRaw ? parseInt(contactIdRaw) : null; const customerId = customerIdRaw ? parseInt(customerIdRaw) : null; @@ -851,6 +853,8 @@ telefoniPrefill.title = titleRaw ? String(titleRaw) : null; telefoniPrefill.callId = callIdRaw ? String(callIdRaw) : null; telefoniPrefill.description = descriptionRaw ? String(descriptionRaw) : null; + const internetConnectionId = internetConnectionIdRaw ? parseInt(internetConnectionIdRaw) : null; + telefoniPrefill.internetConnectionId = Number.isFinite(internetConnectionId) ? internetConnectionId : null; } async function applyTelefoniPrefill() { @@ -898,6 +902,23 @@ console.error('Telefoni prefill failed', e); } } + + if (telefoniPrefill.internetConnectionId) { + const panel = document.getElementById('internetConnectionPrefill'); + try { + const res = await fetch(`/api/v1/internet-connections/${telefoniPrefill.internetConnectionId}`, { credentials: 'include' }); + if (!res.ok) throw new Error('Forbindelsen findes ikke'); + const connection = await res.json(); + const reference = connection.circuit_number || connection.provider_reference || `#${telefoniPrefill.internetConnectionId}`; + panel.innerHTML = `Internetforbindelse knyttes til sagen: ${escapeTopAlertHtml(connection.name || 'Internetforbindelse')} · ${escapeTopAlertHtml(reference)}${connection.address ? ` · ${escapeTopAlertHtml(connection.address)}` : ''}. Åbn forbindelse
Kunden vælges stadig manuelt og ændres ikke på forbindelsen.
`; + panel.classList.remove('d-none'); + const titleInput = document.getElementById('titel'); + if (titleInput && !titleInput.value.trim()) titleInput.value = `Vedr. forbindelse ${reference}`; + } catch (error) { + panel.textContent = `Internetforbindelse #${telefoniPrefill.internetConnectionId} kunne ikke hentes.`; + panel.className = 'alert alert-warning mb-4'; + } + } } function removeContact(id) { @@ -1346,6 +1367,9 @@ contact_ids: Object.keys(selectedContacts).map(id => parseInt(id)).filter(Number.isFinite), telefoni_opkald_id: telefoniPrefill.callId ? parseInt(telefoniPrefill.callId) : null }; + if (telefoniPrefill.internetConnectionId) { + data.internet_connection_ids = [telefoniPrefill.internetConnectionId]; + } if (data.type === 'pipeline') { data.pipeline = { diff --git a/app/modules/sag/templates/detail_v3.html b/app/modules/sag/templates/detail_v3.html index 4034e3d..076b75b 100644 --- a/app/modules/sag/templates/detail_v3.html +++ b/app/modules/sag/templates/detail_v3.html @@ -8362,6 +8362,13 @@
+
+
+
AI og relateret viden
Forslag bygger kun på godkendte vidensartikler og gemmes aldrig automatisk.
+ +
+
Finder relateret viden…
+
{% if is_nextcloud %}
@@ -8420,6 +8427,7 @@
Løsning
+ Alle løsninger Vidensdatabase {% if solution %}{% endif %}
@@ -10568,17 +10576,20 @@ async function registerSolutionFromComment(content) { const firstLine = content.split('\n').find((line) => line.trim()) || ''; const title = firstLine.slice(0, 120) || 'Løsning fra kommentar'; + const hasExistingSolution = Boolean(typeof existingCaseSolution !== 'undefined' && existingCaseSolution); + const existingDescription = hasExistingSolution ? String(existingCaseSolution.description || '').trim() : ''; const response = await fetch('/api/v1/sag/{{ case.id }}/solution', { - method: 'POST', + method: hasExistingSolution ? 'PATCH' : 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify({ sag_id: {{ case.id }}, - title, + title: hasExistingSolution ? existingCaseSolution.title : title, solution_type: 'Support', result: 'Løst', - description: content, - created_by_user_id: 1 + description: existingDescription ? `${existingDescription}\n\nSupplerende løsning:\n${content}` : content, + approval_status: 'draft', + change_note: hasExistingSolution ? 'Suppleret fra kommentar' : 'Oprettet fra kommentar' }) }); @@ -13752,6 +13763,7 @@