293 lines
17 KiB
Python
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()
|