Files
clientflow_backend/scripts/audit_operational_consistency.py

293 lines
17 KiB
Python

#!/usr/bin/env python3
"""Read-only audit of operator-action consistency.
The script deliberately refuses any database that is not transaction read-only.
It is intended to be run against the isolated production snapshot before and
after read-model/UI changes, with identical detection rules.
"""
from __future__ import annotations
import argparse
import json
from collections import Counter, defaultdict
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from sqlalchemy import text
from app.db import engine
from app.operations_service import get_operations_summary
from app.opportunity_next_action_service import get_opportunity_next_actions
from app.opportunity_service import list_opportunities
from app.work_center_action_policy import canonical_action_code, reconstructed_review_required
CLOSED_STATUSES = {"closed", "lost", "won", "completed", "cancelled", "canceled", "spam", "archived"}
FOLLOWUP_CODES = {
"CALL_CUSTOMER", "CONFIRM_DELIVERY", "FOLLOW_UP_QUOTE", "FOLLOW_UP_PROFORMA",
"FOLLOW_UP_PAYMENT", "FOLLOW_UP_CUSTOMER_REVIEW", "FOLLOW_UP_GENERIC",
"RECOVER_OPPORTUNITY", "REVIEW_NURTURE",
}
RESPONSE_CODES = {"SEND_INFO", "SUPPORT", "SEND_QUOTE", "SEND_PROFORMA", "SEND_INVOICE"}
FINANCIAL_CODES = {"SEND_INVOICE", "CREATE_JASMIN_INVOICE", "CONFIRM_PAYMENT", "SEND_PROFORMA"}
SEVERITY = {
"ACTION_SOURCE_DIVERGENCE": "CRITICAL",
"DETAIL_ACTION_DIVERGENCE": "CRITICAL",
"STALE_PENDING_TASK": "HIGH",
"MULTIPLE_COMPETING_PENDING_TASKS": "HIGH",
"FOLLOWUP_SATISFIED_BY_CUSTOMER_INBOUND": "HIGH",
"RESPONSE_TASK_ALREADY_SATISFIED": "HIGH",
"CLOSED_PROCESS_WITH_PENDING_TASK": "HIGH",
"ACTION_BLOCKER_CONTRADICTION": "CRITICAL",
"DOCUMENT_ACTION_WITHOUT_REQUIRED_DOCUMENT": "CRITICAL",
"SCHEDULED_FOLLOWUP_SOURCE_DIVERGENCE": "HIGH",
"OPERATIONS_FALSE_POSITIVE": "HIGH",
"OPERATIONS_FALSE_NEGATIVE": "HIGH",
"STAGE_ACTION_CONTRADICTION": "MEDIUM",
"DUPLICATE_CURRENT_OBLIGATION": "HIGH",
"TASK_STATUS_TIMELINE_CONTRADICTION": "HIGH",
}
def _s(value: Any) -> str:
return str(value or "").strip()
def _dt(value: Any):
if isinstance(value, datetime):
return value if value.tzinfo else value.replace(tzinfo=timezone.utc)
if not value:
return None
try:
parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00"))
return parsed if parsed.tzinfo else parsed.replace(tzinfo=timezone.utc)
except ValueError:
return None
def _jsonable(value: Any):
if isinstance(value, datetime):
return value.isoformat()
if isinstance(value, dict):
return {k: _jsonable(v) for k, v in value.items()}
if isinstance(value, (list, tuple)):
return [_jsonable(v) for v in value]
return value
def _detail_primary_action(opp: dict, decision: dict, pending: list[dict]) -> tuple[str, str]:
"""Mirror the current detail-page precedence without rendering HTML."""
metadata = opp.get("metadata") if isinstance(opp.get("metadata"), dict) else {}
if reconstructed_review_required(metadata):
return "REVIEW_RECONSTRUCTED_PROCESS", "explicit_reconstructed_review"
lifecycle = _s(opp.get("lifecycle_state")).lower()
call = next((t for t in pending if canonical_action_code(t.get("action_code")) == "CALL_CUSTOMER"), None)
if call:
due = _dt(call.get("due_at"))
lifecycle = "follow_up_due" if due and due <= datetime.now(timezone.utc) else "scheduled_follow_up"
elif lifecycle == "scheduled_follow_up":
lifecycle = "active"
lifecycle_task = next((t for t in pending if canonical_action_code(t.get("action_code")) in FOLLOWUP_CODES), None)
if lifecycle_task and lifecycle in {"follow_up_due", "scheduled_follow_up", "recovery", "nurture"}:
return canonical_action_code(lifecycle_task.get("action_code")), "lifecycle_task"
decision_code = canonical_action_code(decision.get("action_code"))
if decision_code:
return decision_code, "opportunity_decision"
if pending:
return canonical_action_code(pending[0].get("action_code")), "pending_task"
return "", "none"
def collect() -> dict:
with engine.connect() as conn:
identity = conn.execute(text(
"SELECT current_database(), current_user, current_setting('transaction_read_only')"
)).one()
if identity[2] != "on":
raise RuntimeError(f"audit requires transaction_read_only=on, got {identity!r}")
task_rows = conn.execute(text("""
SELECT t.id::text, t.opportunity_id::text, t.action_code, t.action,
t.status, t.due_at, t.created_at, t.updated_at, t.done_at,
t.source_system, COALESCE(t.metadata, '{}'::jsonb) AS metadata
FROM tasks t WHERE t.opportunity_id IS NOT NULL
ORDER BY t.opportunity_id, t.created_at, t.id
""")).mappings().all()
timeline_rows = conn.execute(text("""
SELECT o.id::text AS opportunity_id,
max(m.created_at) FILTER (WHERE m.direction='inbound') AS latest_inbound,
max(m.created_at) FILTER (WHERE m.direction='outbound') AS latest_outbound
FROM opportunities o
LEFT JOIN messages m ON m.conversation_id=o.conversation_id
AND m.source_system IN ('chatwoot','chatwoot_backfill')
GROUP BY o.id
""")).mappings().all()
recon_rows = conn.execute(text("""
SELECT opportunity_id::text, status, suggested_action, source_system,
created_at, COALESCE(payload, '{}'::jsonb) AS payload
FROM reconciliation_items WHERE opportunity_id IS NOT NULL
ORDER BY opportunity_id, created_at DESC
""")).mappings().all()
document_rows = conn.execute(text("""
SELECT opportunity_id::text, to_jsonb(commercial_documents) AS document
FROM commercial_documents WHERE opportunity_id IS NOT NULL
ORDER BY opportunity_id, created_at DESC
""")).mappings().all()
opportunities = list_opportunities(status="all", limit=1000)
ids = [_s(o.get("id")) for o in opportunities]
decisions = get_opportunity_next_actions(ids)
operations = get_operations_summary(limit=2000)
canonical = {}
all_operation_items = []
for key in ("work_items", "waiting_items", "backlog_items", "not_current_items"):
for item in operations.get(key, []):
all_operation_items.append(item)
if item.get("opportunity_id"):
canonical[_s(item.get("opportunity_id"))] = item
tasks_by_opp: dict[str, list[dict]] = defaultdict(list)
for row in task_rows:
tasks_by_opp[_s(row["opportunity_id"])].append(dict(row))
timelines = {_s(r["opportunity_id"]): dict(r) for r in timeline_rows}
recon_by_opp: dict[str, list[dict]] = defaultdict(list)
for row in recon_rows:
recon_by_opp[_s(row["opportunity_id"])].append(dict(row))
docs_by_opp: dict[str, list[dict]] = defaultdict(list)
for row in document_rows:
docs_by_opp[_s(row["opportunity_id"])].append(dict(row["document"] or {}))
findings = []
records = []
now = datetime.now(timezone.utc)
def add(family: str, oid: str, detail: str, **evidence):
findings.append({
"family": family, "severity": SEVERITY[family], "opportunity_id": oid,
"detail": detail, "evidence": _jsonable(evidence),
})
for opp in opportunities:
oid = _s(opp.get("id"))
decision = decisions.get(oid) or {}
decision_code = canonical_action_code(decision.get("action_code"))
item = canonical.get(oid)
operation_code = canonical_action_code(item.get("current_action_code")) if item else ""
all_tasks = tasks_by_opp.get(oid, [])
pending = [t for t in all_tasks if _s(t.get("status")).lower() == "pending"]
detail_code, detail_source = _detail_primary_action(opp, decision, pending)
timeline = timelines.get(oid, {})
latest_in = _dt(timeline.get("latest_inbound"))
latest_out = _dt(timeline.get("latest_outbound"))
stale_ids = {_s(r.get("task_id") or r.get("id")) for r in (item or {}).get("stale_task_refs", [])}
docs = docs_by_opp.get(oid, [])
recons = recon_by_opp.get(oid, [])
active_call = next((t for t in pending if canonical_action_code(t.get("action_code")) == "CALL_CUSTOMER"), None)
if item and operation_code != decision_code and operation_code != "CALL_CUSTOMER":
add("ACTION_SOURCE_DIVERGENCE", oid, "Operations action differs from OpportunityDecision", decision=decision_code, operations=operation_code)
if detail_code and decision_code and detail_code != decision_code and detail_source == "pending_task":
add("DETAIL_ACTION_DIVERGENCE", oid, "Detail primary action is overridden by a pending task", decision=decision_code, detail_action=detail_code)
for task in pending:
if _s(task.get("id")) in stale_ids:
add("STALE_PENDING_TASK", oid, "Canonical projection marks pending task as stale", task_id=task["id"], task_action=task["action_code"], current_action=operation_code or decision_code)
codes = [canonical_action_code(t.get("action_code")) for t in pending]
if len(set(codes)) > 1:
add("MULTIPLE_COMPETING_PENDING_TASKS", oid, "Opportunity has multiple distinct pending action codes", task_codes=codes, current_action=operation_code or decision_code)
duplicates = {code: count for code, count in Counter(codes).items() if code and count > 1}
if duplicates:
add("DUPLICATE_CURRENT_OBLIGATION", oid, "Duplicate pending action codes", duplicates=duplicates)
for task in pending:
code = canonical_action_code(task.get("action_code"))
created = _dt(task.get("created_at"))
if code.startswith("FOLLOW_UP_") and latest_in and created and latest_in > created:
add("FOLLOWUP_SATISFIED_BY_CUSTOMER_INBOUND", oid, "Customer inbound is later than pending follow-up", task_id=task["id"], task_action=code, task_created=created, latest_inbound=latest_in)
add("TASK_STATUS_TIMELINE_CONTRADICTION", oid, "Pending follow-up predates later customer inbound", task_id=task["id"], event="inbound")
if code in RESPONSE_CODES and latest_out and created and latest_out > created:
add("RESPONSE_TASK_ALREADY_SATISFIED", oid, "Public outbound is later than pending response task", task_id=task["id"], task_action=code, task_created=created, latest_outbound=latest_out)
add("TASK_STATUS_TIMELINE_CONTRADICTION", oid, "Pending response task predates later public outbound", task_id=task["id"], event="outbound")
status = _s(opp.get("status")).lower()
stage = _s(opp.get("stage")).lower()
lifecycle = _s(opp.get("lifecycle_state")).lower()
if pending and ({status, stage, lifecycle} & CLOSED_STATUSES):
add("CLOSED_PROCESS_WITH_PENDING_TASK", oid, "Closed/inactive process has pending tasks", status=status, stage=stage, lifecycle=lifecycle, task_codes=codes)
if decision_code in {"RECONCILE_DOCUMENTS", "VALIDATE_FISCAL_CUSTOMER"}:
contradicted = [c for c in codes if c in FINANCIAL_CODES or c in {"CREATE_JASMIN_QUOTE", "SEND_QUOTE"}]
if contradicted:
add("ACTION_BLOCKER_CONTRADICTION", oid, "Pending task proposes downstream action while decision requires prerequisite", blocker=decision_code, downstream=contradicted, detail_action=detail_code)
if detail_code in FINANCIAL_CODES and decision_code in {"RECONCILE_DOCUMENTS", "VALIDATE_FISCAL_CUSTOMER"}:
add("DOCUMENT_ACTION_WITHOUT_REQUIRED_DOCUMENT", oid, "Detail exposes document action despite structured prerequisite", detail_action=detail_code, prerequisite=decision_code)
scheduled_mirror = lifecycle == "scheduled_follow_up" or bool(opp.get("next_follow_up_at"))
if scheduled_mirror != bool(active_call):
add("SCHEDULED_FOLLOWUP_SOURCE_DIVERGENCE", oid, "Compatibility schedule fields disagree with active CALL_CUSTOMER task", lifecycle=lifecycle, next_follow_up_at=opp.get("next_follow_up_at"), active_call=bool(active_call))
if item and _s(item.get("operational_queue")) in {"do_now", "review", "exception"}:
refs = item.get("eligibility", {}).get("obligation_source_refs", [])
if not refs and item.get("source") != "outbox" and decision_code in {"NO_ACTION", "NOT_FOUND", "FOLLOW_UP"}:
add("OPERATIONS_FALSE_POSITIVE", oid, "Actionable Operations item lacks a strong explicit obligation", action=operation_code, queue=item.get("operational_queue"))
if item and _s(item.get("operational_queue")) in {"waiting", "backlog", "not_current"} and pending:
non_stale = [t for t in pending if _s(t.get("id")) not in stale_ids]
if non_stale and not (operation_code == "CALL_CUSTOMER" and _dt(active_call.get("due_at") if active_call else None) and _dt(active_call.get("due_at")) > now):
add("OPERATIONS_FALSE_NEGATIVE", oid, "Demoted Operations item retains explicit non-stale pending obligation", action=operation_code, queue=item.get("operational_queue"), task_codes=[t["action_code"] for t in non_stale])
if detail_code != decision_code and _s(opp.get("stage")).upper() in {"INVOICE_REQUESTED", "PROFORMA_SENT", "QUOTE_SENT", "PAYMENT_CONFIRMED"}:
add("STAGE_ACTION_CONTRADICTION", oid, "Presentation follows stage/task rather than structured decision", stage=opp.get("stage"), detail_action=detail_code, decision=decision_code)
records.append(_jsonable({
"opportunity_id": oid, "title": opp.get("title"), "customer": opp.get("linked_customer_name") or opp.get("customer_name"),
"status": opp.get("status"), "commercial_stage": opp.get("stage"), "lifecycle_state": opp.get("lifecycle_state"),
"next_follow_up_at": opp.get("next_follow_up_at"), "decision_action": decision_code,
"decision_reason": decision.get("reason"), "decision_can_execute": decision.get("can_execute"),
"operations_action": operation_code or None, "operational_queue": (item or {}).get("operational_queue"),
"eligibility_reason": (item or {}).get("eligibility_reason_code"), "pending_tasks": pending,
"stale_task_refs": (item or {}).get("stale_task_refs", []), "latest_public_inbound": latest_in,
"latest_public_outbound": latest_out, "active_call_customer_task": active_call,
"financial_state": decision.get("financial_state"), "document_state": [dict(d) for d in docs],
"payment_state": decision.get("financial_state"), "reconciliation_state": recons,
"primary_detail_action": detail_code, "primary_detail_source": detail_source,
}))
counts = Counter(f["family"] for f in findings)
severity_counts = Counter(f["severity"] for f in findings)
ranked = sorted(findings, key=lambda f: ({"CRITICAL": 0, "HIGH": 1, "MEDIUM": 2, "LOW": 3}[f["severity"]], f["family"], f["opportunity_id"]))
return _jsonable({
"generated_at": datetime.now(timezone.utc),
"database": {"name": identity[0], "user": identity[1], "transaction_read_only": identity[2]},
"summary": {
"opportunities": len(opportunities), "tasks": len(task_rows),
"canonical_total": operations.get("canonical_count"), "work_queue_total": operations.get("counts", {}).get("work_queue_total"),
"families": dict(sorted(counts.items())), "severities": dict(severity_counts),
},
"top_20": ranked[:20], "findings": ranked, "opportunities": records,
"operations_diagnostics": operations.get("projection_metrics", {}),
})
def render_text(report: dict) -> str:
lines = [
"ClientFlow operational consistency audit",
f"Generated: {report['generated_at']}",
f"Database: {report['database']}",
"",
"Summary",
json.dumps(report["summary"], ensure_ascii=False, indent=2),
"",
"Top 20 findings",
]
for finding in report["top_20"]:
lines.append(f"[{finding['severity']}] {finding['family']} {finding['opportunity_id']}: {finding['detail']} | {json.dumps(finding['evidence'], ensure_ascii=False)}")
return "\n".join(lines) + "\n"
def main() -> None:
parser = argparse.ArgumentParser()
parser.add_argument("--json", required=True)
parser.add_argument("--text", required=True)
args = parser.parse_args()
report = collect()
Path(args.json).write_text(json.dumps(report, ensure_ascii=False, indent=2), encoding="utf-8")
Path(args.text).write_text(render_text(report), encoding="utf-8")
print(json.dumps(report["summary"], ensure_ascii=False, indent=2))
if __name__ == "__main__":
main()