1714 lines
85 KiB
Python
Executable File
1714 lines
85 KiB
Python
Executable File
#!/usr/bin/env python3
|
|
"""Forensic READ-ONLY audit of ClientFlow Work Center and associations.
|
|
|
|
Target: ClientFlow v4928.1.5.130+
|
|
|
|
This auditor explains *why* a Work Center action is or is not coherent. It
|
|
separates proven defects from policy ambiguity and expected behaviour.
|
|
|
|
It compares:
|
|
1. the normalized item shown in /operations;
|
|
2. the stored pending task;
|
|
3. the central opportunity next-action decision;
|
|
4. explicit blockers (association, fiscal conflict, reconstructed review);
|
|
5. active operation links and open reconciliation candidates;
|
|
6. the task materialization and visibility policies.
|
|
|
|
No INSERT/UPDATE/DELETE is executed. There is deliberately no --apply option.
|
|
All SQL owned by this script runs in a PostgreSQL READ ONLY transaction.
|
|
|
|
Examples:
|
|
python scripts/audit_work_center_forensics.py --limit 500 --opportunity-limit 1000
|
|
python scripts/audit_work_center_forensics.py --focus NOLTIA --focus S00323
|
|
python scripts/audit_work_center_forensics.py --fail-on confirmed_high
|
|
python scripts/audit_work_center_forensics.py --self-test
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import argparse
|
|
import csv
|
|
import json
|
|
import re
|
|
import sys
|
|
from collections import Counter, defaultdict
|
|
from dataclasses import asdict, dataclass, field
|
|
from datetime import date, datetime, timezone
|
|
from pathlib import Path
|
|
from typing import Any, Dict, Iterable, List, Mapping, Optional, Sequence, Tuple
|
|
|
|
|
|
def _detect_project_root() -> Path:
|
|
script_path = Path(__file__).resolve()
|
|
candidates = [script_path.parent.parent, Path.cwd().resolve(), *script_path.parents]
|
|
seen: set[str] = set()
|
|
for candidate in candidates:
|
|
key = str(candidate)
|
|
if key in seen:
|
|
continue
|
|
seen.add(key)
|
|
if (candidate / "app" / "db.py").is_file():
|
|
return candidate
|
|
return script_path.parent.parent
|
|
|
|
|
|
PROJECT_ROOT = _detect_project_root()
|
|
if str(PROJECT_ROOT) not in sys.path:
|
|
sys.path.insert(0, str(PROJECT_ROOT))
|
|
|
|
from sqlalchemy import text # noqa: E402
|
|
|
|
|
|
OPEN_RECONCILIATION_STATUSES = {"open", "needs_review", "conflict"}
|
|
TERMINAL_STAGES = {"WON", "LOST", "NO_INTEREST", "CLOSED"}
|
|
DONE_TASK_STATUSES = {"done", "completed", "concluida", "concluída", "ignored", "ignorada"}
|
|
ASSOCIATION_ACTIONS = {
|
|
"ASSOCIATE_OPPORTUNITY", "REVIEW_ASSOCIATION", "LINK_DOCUMENT", "RECONCILE_DOCUMENTS",
|
|
}
|
|
RECONSTRUCTED_REVIEW_ACTIONS = {"REVIEW_RECONSTRUCTED_PROCESS", "REVIEW_MANUALLY"}
|
|
UI_RECONSTRUCTED_SENSITIVE_ACTIONS = {
|
|
# Exact allow-list used by app.operations_service v130.
|
|
"SEND_PROFORMA", "SEND_INVOICE", "CONFIRM_PAYMENT",
|
|
"CONFIRM_PAYMENT_AND_PREPARE_SHIPMENT", "PREPARE_ORDER", "CREATE_SHIPMENT",
|
|
}
|
|
OPERATIONAL_SENSITIVE_ACTIONS = UI_RECONSTRUCTED_SENSITIVE_ACTIONS | {
|
|
"VALIDATE_PHYSICAL_ORDER", "SHIP_ORDER", "CLOSE_OPPORTUNITY",
|
|
}
|
|
ACTION_ALIASES = {
|
|
"SHIP_ORDER": "CREATE_SHIPMENT",
|
|
"RECONCILE_DOCUMENTS": "ASSOCIATE_OPPORTUNITY",
|
|
"REVIEW_ASSOCIATION": "ASSOCIATE_OPPORTUNITY",
|
|
"LINK_DOCUMENT": "ASSOCIATE_OPPORTUNITY",
|
|
"REVIEW": "REVIEW_MANUALLY",
|
|
}
|
|
NON_TASK_ACTIONS = {"", "NO_ACTION", "WAIT_PRODUCTION", "NOT_FOUND", "AUDIT_ERROR"}
|
|
CATEGORY_ORDER = {"confirmed": 0, "needs_review": 1, "policy": 2, "expected": 3}
|
|
SEVERITY_ORDER = {"critical": 0, "high": 1, "medium": 2, "low": 3, "info": 4}
|
|
CONFIDENCE_ORDER = {"high": 0, "medium": 1, "low": 2}
|
|
|
|
|
|
@dataclass
|
|
class Finding:
|
|
category: str
|
|
severity: str
|
|
confidence: str
|
|
actionable: bool
|
|
code: str
|
|
title: str
|
|
explanation: str
|
|
recommendation: str
|
|
root_cause: str = ""
|
|
opportunity_id: str = ""
|
|
opportunity_title: str = ""
|
|
customer_name: str = ""
|
|
work_item_id: str = ""
|
|
task_id: str = ""
|
|
association_item_id: str = ""
|
|
external_ref: str = ""
|
|
displayed_action: str = ""
|
|
stored_action: str = ""
|
|
central_action: str = ""
|
|
effective_action: str = ""
|
|
evidence: Dict[str, Any] = field(default_factory=dict)
|
|
proposed_next_step: Dict[str, Any] = field(default_factory=dict)
|
|
|
|
|
|
@dataclass
|
|
class AssociationAssessment:
|
|
item_id: str
|
|
source_system: str
|
|
external_type: str
|
|
external_ref: str
|
|
status: str
|
|
assigned_opportunity_id: str
|
|
exact_link_opportunity_ids: List[str]
|
|
match_kind: str
|
|
classification: str
|
|
blocks_opportunity: bool
|
|
confidence: str
|
|
identity: Dict[str, Any] = field(default_factory=dict)
|
|
candidate: Dict[str, Any] = field(default_factory=dict)
|
|
matching_links: List[Dict[str, Any]] = field(default_factory=list)
|
|
|
|
|
|
@dataclass
|
|
class ReconstructedAssessment:
|
|
record_mode: str
|
|
is_legacy_record: bool
|
|
explicit_required: bool
|
|
explicit_cleared: bool
|
|
pending_review_task_ids: List[str]
|
|
completed_review_task_ids: List[str]
|
|
state: str
|
|
evidence: Dict[str, Any] = field(default_factory=dict)
|
|
|
|
|
|
@dataclass
|
|
class ActionChain:
|
|
opportunity_id: str
|
|
central_action: str
|
|
central_can_execute: Optional[bool]
|
|
explicit_blocker: str
|
|
effective_action: str
|
|
blocker_reason: str
|
|
association_blocker_ids: List[str]
|
|
reconstructed_state: str
|
|
materializable: bool
|
|
visible_actions: List[str]
|
|
pending_actions: List[str]
|
|
|
|
|
|
def _s(value: Any) -> str:
|
|
return str(value or "").strip()
|
|
|
|
|
|
def _upper(value: Any) -> str:
|
|
return _s(value).upper()
|
|
|
|
|
|
def _lower(value: Any) -> str:
|
|
return _s(value).lower()
|
|
|
|
|
|
def _as_dict(value: Any) -> Dict[str, Any]:
|
|
if isinstance(value, dict):
|
|
return value
|
|
if isinstance(value, str) and value.strip():
|
|
try:
|
|
parsed = json.loads(value)
|
|
return parsed if isinstance(parsed, dict) else {}
|
|
except Exception:
|
|
return {}
|
|
return {}
|
|
|
|
|
|
def _json_default(value: Any) -> str:
|
|
if isinstance(value, (datetime, date)):
|
|
return value.isoformat()
|
|
return str(value)
|
|
|
|
|
|
def _canonical_action(value: Any) -> str:
|
|
code = _upper(value)
|
|
return ACTION_ALIASES.get(code, code)
|
|
|
|
|
|
def _actions_match(left: Any, right: Any) -> bool:
|
|
a, b = _canonical_action(left), _canonical_action(right)
|
|
if not a or not b:
|
|
return False
|
|
if a == b:
|
|
return True
|
|
if {a, b} <= {"REVIEW_MANUALLY", "REVIEW_RECONSTRUCTED_PROCESS"}:
|
|
return True
|
|
if a == "FOLLOW_UP" and b.startswith("FOLLOW_UP_"):
|
|
return True
|
|
if b == "FOLLOW_UP" and a.startswith("FOLLOW_UP_"):
|
|
return True
|
|
return False
|
|
|
|
|
|
def _norm_text(value: Any) -> str:
|
|
value = _s(value).casefold()
|
|
value = re.sub(r"[^0-9a-zà-ÿ]+", " ", value)
|
|
return " ".join(value.split())
|
|
|
|
|
|
def _norm_email(value: Any) -> str:
|
|
return _s(value).casefold()
|
|
|
|
|
|
def _norm_tax_id(value: Any) -> str:
|
|
return "".join(ch for ch in _s(value) if ch.isdigit())
|
|
|
|
|
|
def _record_mode(metadata: Any) -> str:
|
|
return _lower(_as_dict(metadata).get("clientflow_record_mode"))
|
|
|
|
|
|
def _is_legacy_record(metadata: Any) -> bool:
|
|
return _record_mode(metadata) in {
|
|
"reconstructed_invoice_review", "historical_reconstructed", "legacy_review",
|
|
}
|
|
|
|
|
|
def _review_type(task: Mapping[str, Any]) -> str:
|
|
metadata = _as_dict(task.get("metadata"))
|
|
return _lower(metadata.get("review_type") or metadata.get("task_type"))
|
|
|
|
|
|
def _link_system(source_system: Any) -> str:
|
|
source = _lower(source_system)
|
|
if source.startswith("odoo"):
|
|
return "odoo"
|
|
if source.startswith("jasmin"):
|
|
return "jasmin"
|
|
if source.startswith("packlink"):
|
|
return "packlink"
|
|
return source
|
|
|
|
|
|
def _reference_variants(value: Any, *, system: str = "", external_type: str = "") -> set[str]:
|
|
raw = _s(value)
|
|
if not raw:
|
|
return set()
|
|
variants = {raw.casefold()}
|
|
compact = re.sub(r"\s+", "", raw).casefold()
|
|
variants.add(compact)
|
|
system = _lower(system)
|
|
external_type = _lower(external_type)
|
|
# Odoo sale orders are often represented both by numeric database id and
|
|
# human name S00325. Never infer one from the other. Only normalize zeros in
|
|
# an existing S-number, which is a presentation-equivalent reference.
|
|
if system == "odoo" or "sale" in external_type:
|
|
match = re.fullmatch(r"s0*(\d+)", compact)
|
|
if match:
|
|
number = int(match.group(1))
|
|
variants.add(f"s{number}")
|
|
variants.add(f"s{number:05d}")
|
|
return {variant for variant in variants if variant}
|
|
|
|
|
|
def _candidate_refs(item: Mapping[str, Any]) -> set[str]:
|
|
system = _link_system(item.get("source_system"))
|
|
external_type = _s(item.get("external_type"))
|
|
values = [item.get("external_id"), item.get("document_number")]
|
|
payload = _as_dict(item.get("payload"))
|
|
record = payload.get("record") if isinstance(payload.get("record"), dict) else payload
|
|
for key in ("id", "external_id", "name", "number", "document_number", "order_name"):
|
|
values.append(record.get(key))
|
|
refs: set[str] = set()
|
|
for value in values:
|
|
refs.update(_reference_variants(value, system=system, external_type=external_type))
|
|
return refs
|
|
|
|
|
|
def _link_refs(link: Mapping[str, Any]) -> set[str]:
|
|
system = _lower(link.get("system"))
|
|
external_type = _s(link.get("external_type"))
|
|
refs: set[str] = set()
|
|
for value in (link.get("external_id"), link.get("external_name")):
|
|
refs.update(_reference_variants(value, system=system, external_type=external_type))
|
|
return refs
|
|
|
|
|
|
def _display_ref(item: Mapping[str, Any]) -> str:
|
|
return _s(item.get("document_number") or item.get("external_id"))
|
|
|
|
|
|
def _truthy_metadata(metadata: Mapping[str, Any], keys: Iterable[str]) -> bool:
|
|
for key in keys:
|
|
value = metadata.get(key)
|
|
if value is True:
|
|
return True
|
|
if isinstance(value, str) and value.strip().casefold() in {"true", "1", "yes", "sim", "required"}:
|
|
return True
|
|
return False
|
|
|
|
|
|
def _falsey_metadata(metadata: Mapping[str, Any], keys: Iterable[str]) -> bool:
|
|
for key in keys:
|
|
value = metadata.get(key)
|
|
if value is False:
|
|
return True
|
|
if isinstance(value, str) and value.strip().casefold() in {"false", "0", "no", "não", "nao", "cleared", "validated"}:
|
|
return True
|
|
return False
|
|
|
|
|
|
def _focus_match(focus: Sequence[str], *values: Any) -> bool:
|
|
if not focus:
|
|
return True
|
|
haystack = " ".join(_s(value) for value in values).casefold()
|
|
return any(token.casefold() in haystack for token in focus if token)
|
|
|
|
|
|
def _severity_key(category: str, severity: str) -> str:
|
|
if category == "confirmed" and severity in {"critical", "high"}:
|
|
return "confirmed_high"
|
|
return severity
|
|
|
|
|
|
class ForensicAuditor:
|
|
def __init__(
|
|
self,
|
|
*,
|
|
limit: int,
|
|
opportunity_limit: int,
|
|
focus: Sequence[str],
|
|
include_clean: bool,
|
|
) -> None:
|
|
self.limit = min(max(int(limit), 1), 5000)
|
|
self.opportunity_limit = min(max(int(opportunity_limit), 1), 10000)
|
|
self.focus = [_s(value) for value in focus if _s(value)]
|
|
self.include_clean = bool(include_clean)
|
|
|
|
self.work_items: List[Dict[str, Any]] = []
|
|
self.operations_counts: Dict[str, Any] = {}
|
|
self.opportunities: Dict[str, Dict[str, Any]] = {}
|
|
self.tasks_by_opportunity: Dict[str, List[Dict[str, Any]]] = defaultdict(list)
|
|
self.pending_tasks_by_opportunity: Dict[str, List[Dict[str, Any]]] = defaultdict(list)
|
|
self.all_pending_tasks: List[Dict[str, Any]] = []
|
|
self.reconciliation_items: List[Dict[str, Any]] = []
|
|
self.operation_links: List[Dict[str, Any]] = []
|
|
self.documents_by_opportunity: Dict[str, List[Dict[str, Any]]] = defaultdict(list)
|
|
self.task_events_by_task: Dict[str, List[Dict[str, Any]]] = defaultdict(list)
|
|
self.next_actions: Dict[str, Dict[str, Any]] = {}
|
|
self.operation_snapshots: Dict[str, Dict[str, Any]] = {}
|
|
self.materialized_actions: set[str] = set()
|
|
|
|
self.association_assessments: List[AssociationAssessment] = []
|
|
self.associations_by_opportunity: Dict[str, List[AssociationAssessment]] = defaultdict(list)
|
|
self.reconstructed: Dict[str, ReconstructedAssessment] = {}
|
|
self.action_chains: Dict[str, ActionChain] = {}
|
|
self.visibility: Dict[str, Dict[str, Any]] = {}
|
|
self.findings: List[Finding] = []
|
|
self.dossiers: Dict[str, Dict[str, Any]] = {}
|
|
self.summary: Dict[str, Any] = {}
|
|
|
|
def add(self, finding: Finding) -> None:
|
|
if self.focus and not _focus_match(
|
|
self.focus,
|
|
finding.opportunity_id,
|
|
finding.opportunity_title,
|
|
finding.customer_name,
|
|
finding.external_ref,
|
|
finding.code,
|
|
):
|
|
return
|
|
self.findings.append(finding)
|
|
|
|
def load(self) -> None:
|
|
try:
|
|
from app.db import engine
|
|
from app.operations_service import get_operations_summary
|
|
from app.opportunity_next_action_service import get_opportunity_next_action
|
|
from app.operation_service import get_operation_snapshot
|
|
from app.opportunity_action_task_materializer import MATERIALIZED_ACTIONS, ACTION_ALIASES as MATERIALIZER_ALIASES
|
|
except ModuleNotFoundError as exc:
|
|
raise SystemExit(
|
|
"Não foi possível importar o backend ClientFlow. "
|
|
f"Raiz detetada: {PROJECT_ROOT}. Instale em <backend>/scripts/."
|
|
) from exc
|
|
|
|
operations = get_operations_summary(limit=self.limit)
|
|
self.work_items = [dict(item) for item in operations.get("work_items", [])]
|
|
self.operations_counts = dict(operations.get("counts", {}))
|
|
self.materialized_actions = {
|
|
_canonical_action(MATERIALIZER_ALIASES.get(_upper(code), _upper(code)))
|
|
for code in MATERIALIZED_ACTIONS
|
|
}
|
|
|
|
with engine.connect() as conn:
|
|
tx = conn.begin()
|
|
try:
|
|
conn.execute(text("SET TRANSACTION READ ONLY"))
|
|
required = {
|
|
"opportunities", "tasks", "reconciliation_items", "operation_links",
|
|
"commercial_documents", "customers", "task_events",
|
|
}
|
|
existing = set(conn.execute(text("""
|
|
SELECT table_name
|
|
FROM information_schema.tables
|
|
WHERE table_schema = current_schema()
|
|
AND table_name = ANY(:tables)
|
|
"""), {"tables": sorted(required)}).scalars().all())
|
|
missing = sorted(required - existing)
|
|
if missing:
|
|
raise RuntimeError("Tabelas necessárias em falta: " + ", ".join(missing))
|
|
|
|
opportunity_rows = conn.execute(text("""
|
|
SELECT
|
|
o.id::text, o.title, o.customer_name, o.customer_email,
|
|
o.stage, o.status, o.value_amount, o.currency,
|
|
o.local_customer_id::text, o.lifecycle_state,
|
|
o.next_follow_up_at, o.follow_up_attempts,
|
|
o.last_customer_activity_at, o.last_operator_activity_at,
|
|
o.last_commercial_activity_at, o.metadata,
|
|
o.created_at, o.updated_at,
|
|
c.name AS fiscal_customer_name,
|
|
c.email AS fiscal_customer_email,
|
|
c.tax_id AS fiscal_customer_tax_id,
|
|
c.street_name AS fiscal_customer_street,
|
|
c.postal_zone AS fiscal_customer_postal_zone,
|
|
c.city_name AS fiscal_customer_city
|
|
FROM opportunities o
|
|
LEFT JOIN customers c ON c.id = o.local_customer_id
|
|
WHERE COALESCE(o.status, 'open') = 'open'
|
|
AND COALESCE(o.stage, '') <> ALL(:terminal_stages)
|
|
ORDER BY o.updated_at DESC, o.created_at DESC
|
|
LIMIT :opportunity_limit
|
|
"""), {
|
|
"terminal_stages": sorted(TERMINAL_STAGES),
|
|
"opportunity_limit": self.opportunity_limit,
|
|
}).mappings().all()
|
|
self.opportunities = {str(row["id"]): dict(row) for row in opportunity_rows}
|
|
|
|
relevant_ids = set(self.opportunities)
|
|
relevant_ids.update(
|
|
_s(item.get("opportunity_id")) for item in self.work_items
|
|
if _s(item.get("opportunity_id"))
|
|
)
|
|
ids = sorted(relevant_ids)
|
|
|
|
task_rows = conn.execute(text("""
|
|
SELECT
|
|
t.id::text, t.opportunity_id::text, t.action_code,
|
|
t.action, t.note, t.route, t.priority, t.status,
|
|
t.due_at, t.source_system, t.metadata,
|
|
t.conversation_id, t.contact_id, t.customer_id,
|
|
t.message_id::text, t.raw_event_id::text,
|
|
t.idempotency_key, t.action_required, t.safe_to_post,
|
|
t.created_at, t.updated_at, t.done_at,
|
|
o.title AS opportunity_title,
|
|
o.stage AS opportunity_stage,
|
|
o.metadata AS opportunity_metadata,
|
|
COALESCE(c.name, o.customer_name, '') AS fiscal_customer_name,
|
|
COALESCE(c.email, o.customer_email, '') AS fiscal_customer_email,
|
|
COALESCE(c.tax_id, '') AS fiscal_customer_tax_id,
|
|
COALESCE(m.clean_body, m.raw_body, re.payload->>'content', '') AS request_text
|
|
FROM tasks t
|
|
LEFT JOIN opportunities o ON o.id = t.opportunity_id
|
|
LEFT JOIN customers c ON c.id = o.local_customer_id
|
|
LEFT JOIN messages m ON m.id = t.message_id
|
|
LEFT JOIN raw_events re ON re.id = t.raw_event_id
|
|
WHERE t.opportunity_id::text = ANY(:ids)
|
|
OR t.status = 'pending'
|
|
ORDER BY t.created_at DESC
|
|
"""), {"ids": ids or ["00000000-0000-0000-0000-000000000000"]}).mappings().all()
|
|
for row in task_rows:
|
|
task = dict(row)
|
|
oid = _s(task.get("opportunity_id"))
|
|
if oid:
|
|
self.tasks_by_opportunity[oid].append(task)
|
|
if _lower(task.get("status")) == "pending":
|
|
self.pending_tasks_by_opportunity[oid].append(task)
|
|
if _lower(task.get("status")) == "pending":
|
|
self.all_pending_tasks.append(task)
|
|
|
|
task_ids = [_s(row.get("id")) for row in task_rows if _s(row.get("id"))]
|
|
if task_ids:
|
|
event_rows = conn.execute(text("""
|
|
SELECT id::text, task_id::text, event_type, payload, created_by, created_at
|
|
FROM task_events
|
|
WHERE task_id::text = ANY(:task_ids)
|
|
ORDER BY created_at DESC
|
|
"""), {"task_ids": task_ids}).mappings().all()
|
|
for row in event_rows:
|
|
event = dict(row)
|
|
self.task_events_by_task[_s(event.get("task_id"))].append(event)
|
|
|
|
recon_rows = conn.execute(text("""
|
|
SELECT
|
|
ri.id::text, ri.opportunity_id::text, ri.customer_id::text,
|
|
ri.source_system, ri.external_type, ri.external_id,
|
|
ri.document_number, ri.title, ri.customer_name,
|
|
ri.customer_email, ri.customer_tax_id, ri.amount, ri.currency,
|
|
NULL::text AS suggested_stage, ri.suggested_action, ri.confidence,
|
|
ri.status, ri.payload, ri.resolution_note,
|
|
ri.created_at, ri.updated_at
|
|
FROM reconciliation_items ri
|
|
WHERE ri.status IN ('open','needs_review','conflict')
|
|
OR ri.opportunity_id::text = ANY(:ids)
|
|
ORDER BY ri.updated_at DESC, ri.created_at DESC
|
|
"""), {"ids": ids or ["00000000-0000-0000-0000-000000000000"]}).mappings().all()
|
|
self.reconciliation_items = [dict(row) for row in recon_rows]
|
|
|
|
link_rows = conn.execute(text("""
|
|
SELECT
|
|
ol.id::text, ol.opportunity_id::text, ol.system,
|
|
ol.external_type, ol.external_id, ol.external_name,
|
|
ol.status, ol.payload, ol.created_at, ol.updated_at
|
|
FROM operation_links ol
|
|
WHERE ol.opportunity_id::text = ANY(:ids)
|
|
OR COALESCE(ol.status, '') NOT IN ('ignored','deleted')
|
|
"""), {"ids": ids or ["00000000-0000-0000-0000-000000000000"]}).mappings().all()
|
|
self.operation_links = [dict(row) for row in link_rows]
|
|
|
|
document_rows = conn.execute(text("""
|
|
SELECT
|
|
cd.id::text, cd.opportunity_id::text, cd.system,
|
|
cd.document_kind, cd.external_id, cd.document_number,
|
|
cd.status, cd.total_amount, cd.currency, cd.role,
|
|
cd.is_active, cd.is_primary, cd.document_date,
|
|
cd.created_at, cd.updated_at
|
|
FROM commercial_documents cd
|
|
WHERE cd.opportunity_id::text = ANY(:ids)
|
|
ORDER BY cd.created_at DESC
|
|
"""), {"ids": ids or ["00000000-0000-0000-0000-000000000000"]}).mappings().all()
|
|
for row in document_rows:
|
|
doc = dict(row)
|
|
self.documents_by_opportunity[_s(doc.get("opportunity_id"))].append(doc)
|
|
finally:
|
|
tx.rollback()
|
|
|
|
for opportunity_id in sorted(self.opportunities):
|
|
try:
|
|
self.next_actions[opportunity_id] = dict(get_opportunity_next_action(opportunity_id))
|
|
except Exception as exc:
|
|
self.next_actions[opportunity_id] = {
|
|
"action_code": "AUDIT_ERROR",
|
|
"label": "Erro ao calcular próxima ação",
|
|
"can_execute": False,
|
|
"reason_if_blocked": str(exc),
|
|
}
|
|
try:
|
|
snapshot = get_operation_snapshot(opportunity_id)
|
|
self.operation_snapshots[opportunity_id] = dict(snapshot or {})
|
|
except Exception as exc:
|
|
self.operation_snapshots[opportunity_id] = {"error": str(exc), "cards": [], "links": []}
|
|
|
|
def _identity_assessment(self, item: Mapping[str, Any], opportunity_id: str) -> Dict[str, Any]:
|
|
opportunity = self.opportunities.get(opportunity_id, {})
|
|
candidate_nif = _norm_tax_id(item.get("customer_tax_id"))
|
|
opportunity_nif = _norm_tax_id(opportunity.get("fiscal_customer_tax_id"))
|
|
candidate_email = _norm_email(item.get("customer_email"))
|
|
opportunity_emails = {
|
|
_norm_email(opportunity.get("fiscal_customer_email")),
|
|
_norm_email(opportunity.get("customer_email")),
|
|
} - {""}
|
|
candidate_name = _norm_text(item.get("customer_name"))
|
|
opportunity_names = {
|
|
_norm_text(opportunity.get("fiscal_customer_name")),
|
|
_norm_text(opportunity.get("customer_name")),
|
|
} - {""}
|
|
nif_state = "unknown"
|
|
if candidate_nif and opportunity_nif:
|
|
nif_state = "match" if candidate_nif == opportunity_nif else "mismatch"
|
|
email_state = "unknown"
|
|
if candidate_email and opportunity_emails:
|
|
email_state = "match" if candidate_email in opportunity_emails else "mismatch"
|
|
name_state = "unknown"
|
|
if candidate_name and opportunity_names:
|
|
exact = candidate_name in opportunity_names
|
|
token_overlap = max(
|
|
(len(set(candidate_name.split()) & set(name.split())) for name in opportunity_names),
|
|
default=0,
|
|
)
|
|
name_state = "match" if exact or token_overlap >= 2 else "mismatch"
|
|
return {
|
|
"nif": nif_state,
|
|
"email": email_state,
|
|
"name": name_state,
|
|
"candidate": {
|
|
"nif": candidate_nif,
|
|
"email": candidate_email,
|
|
"name": _s(item.get("customer_name")),
|
|
},
|
|
"opportunity": {
|
|
"nif": opportunity_nif,
|
|
"emails": sorted(opportunity_emails),
|
|
"names": sorted(opportunity_names),
|
|
},
|
|
}
|
|
|
|
def assess_associations(self) -> None:
|
|
active_links = [
|
|
link for link in self.operation_links
|
|
if _lower(link.get("status")) not in {"ignored", "deleted"}
|
|
]
|
|
for item in self.reconciliation_items:
|
|
if _lower(item.get("status")) not in OPEN_RECONCILIATION_STATUSES:
|
|
continue
|
|
system = _link_system(item.get("source_system"))
|
|
refs = _candidate_refs(item)
|
|
matches: List[Dict[str, Any]] = []
|
|
for link in active_links:
|
|
if _lower(link.get("system")) != system:
|
|
continue
|
|
if refs & _link_refs(link):
|
|
matches.append(link)
|
|
linked_ids = sorted({_s(link.get("opportunity_id")) for link in matches if _s(link.get("opportunity_id"))})
|
|
assigned = _s(item.get("opportunity_id"))
|
|
|
|
classification = "unresolved_unassigned"
|
|
match_kind = "none"
|
|
blocks = False
|
|
confidence = "medium"
|
|
identity: Dict[str, Any] = {}
|
|
|
|
if len(linked_ids) > 1:
|
|
classification = "conflict_multiple_exact_links"
|
|
match_kind = "exact_reference"
|
|
blocks = True
|
|
confidence = "high"
|
|
elif len(linked_ids) == 1:
|
|
exact = linked_ids[0]
|
|
match_kind = "exact_reference"
|
|
identity = self._identity_assessment(item, exact)
|
|
if assigned and assigned != exact:
|
|
classification = "wrong_assigned_opportunity"
|
|
blocks = True
|
|
confidence = "high"
|
|
else:
|
|
classification = "stale_already_linked"
|
|
blocks = False
|
|
confidence = "high"
|
|
elif assigned:
|
|
identity = self._identity_assessment(item, assigned)
|
|
if identity.get("nif") == "mismatch":
|
|
classification = "assigned_identity_conflict"
|
|
blocks = True
|
|
confidence = "high"
|
|
else:
|
|
classification = "assigned_unresolved"
|
|
blocks = True
|
|
confidence = "high" if identity.get("nif") == "match" else "medium"
|
|
else:
|
|
classification = "unresolved_unassigned"
|
|
blocks = False
|
|
confidence = "medium"
|
|
|
|
assessment = AssociationAssessment(
|
|
item_id=_s(item.get("id")),
|
|
source_system=system,
|
|
external_type=_s(item.get("external_type")),
|
|
external_ref=_display_ref(item),
|
|
status=_lower(item.get("status")),
|
|
assigned_opportunity_id=assigned,
|
|
exact_link_opportunity_ids=linked_ids,
|
|
match_kind=match_kind,
|
|
classification=classification,
|
|
blocks_opportunity=blocks,
|
|
confidence=confidence,
|
|
identity=identity,
|
|
candidate=dict(item),
|
|
matching_links=matches,
|
|
)
|
|
self.association_assessments.append(assessment)
|
|
target_ids = set(linked_ids)
|
|
if assigned:
|
|
target_ids.add(assigned)
|
|
for oid in target_ids:
|
|
self.associations_by_opportunity[oid].append(assessment)
|
|
|
|
common = self._common(assigned or (linked_ids[0] if len(linked_ids) == 1 else ""))
|
|
external_ref = assessment.external_ref
|
|
if classification == "conflict_multiple_exact_links":
|
|
self.add(Finding(
|
|
category="confirmed", severity="critical", confidence="high", actionable=True,
|
|
code="EXTERNAL_REFERENCE_LINKED_TO_MULTIPLE_OPPORTUNITIES",
|
|
title="A mesma evidência externa está ligada a várias oportunidades",
|
|
explanation="A referência externa tem mais de uma ligação ativa; nenhuma ação posterior é segura até ficar uma ligação única.",
|
|
recommendation="Rever as ligações e preservar apenas a oportunidade correta.",
|
|
root_cause="association_integrity",
|
|
association_item_id=assessment.item_id, external_ref=external_ref,
|
|
effective_action="ASSOCIATE_OPPORTUNITY",
|
|
evidence={"opportunity_ids": linked_ids, "matching_links": matches, "candidate_refs": sorted(refs)},
|
|
proposed_next_step={"action": "manual_resolve_duplicate_external_link", "automatic": False},
|
|
**common,
|
|
))
|
|
elif classification == "wrong_assigned_opportunity":
|
|
self.add(Finding(
|
|
category="confirmed", severity="critical", confidence="high", actionable=True,
|
|
code="ASSOCIATION_ASSIGNED_TO_WRONG_OPPORTUNITY",
|
|
title="O candidato aponta para uma oportunidade diferente da ligação externa existente",
|
|
explanation="A referência externa já tem uma ligação exata, mas o candidato está atribuído a outra oportunidade.",
|
|
recommendation="Bloquear a ação e corrigir manualmente a associação.",
|
|
root_cause="association_integrity",
|
|
association_item_id=assessment.item_id, external_ref=external_ref,
|
|
effective_action="ASSOCIATE_OPPORTUNITY",
|
|
evidence={"assigned_opportunity_id": assigned, "exact_link_opportunity_id": linked_ids[0], "identity": identity},
|
|
proposed_next_step={"action": "review_wrong_association", "automatic": False},
|
|
**common,
|
|
))
|
|
elif classification == "stale_already_linked":
|
|
exact = linked_ids[0]
|
|
self.add(Finding(
|
|
category="confirmed", severity="high", confidence="high", actionable=True,
|
|
code="STALE_OPEN_CANDIDATE_ALREADY_EXACTLY_LINKED",
|
|
title="Candidato aberto apesar de a evidência já estar ligada exatamente uma vez",
|
|
explanation="A ligação externa já existe numa única oportunidade. O candidato aberto é redundante e não deve bloquear o processo.",
|
|
recommendation="Resolver como já ligado, sem alterar fase, documentos ou tasks.",
|
|
root_cause="reconciliation_stale_candidate",
|
|
association_item_id=assessment.item_id, external_ref=external_ref,
|
|
effective_action="RESOLVE_AS_EXISTING_LINK",
|
|
evidence={"exact_opportunity_id": exact, "matching_link": matches[0], "identity": identity},
|
|
proposed_next_step={"action": "resolve_as_existing_link", "automatic": True, "guard": "exactly_one_active_link"},
|
|
**self._common(exact),
|
|
))
|
|
elif classification == "assigned_identity_conflict":
|
|
self.add(Finding(
|
|
category="confirmed", severity="critical", confidence="high", actionable=True,
|
|
code="ASSOCIATION_IDENTITY_CONFLICT",
|
|
title="O candidato atribuído tem NIF diferente do cliente fiscal",
|
|
explanation="Existe evidência fiscal contraditória; a associação não pode ser aceite automaticamente.",
|
|
recommendation="Validar cliente e documento antes de qualquer ação financeira ou logística.",
|
|
root_cause="identity_conflict",
|
|
association_item_id=assessment.item_id, external_ref=external_ref,
|
|
effective_action="ASSOCIATE_OPPORTUNITY",
|
|
evidence={"identity": identity},
|
|
proposed_next_step={"action": "manual_identity_review", "automatic": False},
|
|
**common,
|
|
))
|
|
elif classification == "assigned_unresolved":
|
|
self.add(Finding(
|
|
category="needs_review", severity="high", confidence=confidence, actionable=True,
|
|
code="ASSIGNED_ASSOCIATION_REQUIRES_DECISION",
|
|
title="Existe candidato atribuído sem ligação externa confirmada",
|
|
explanation="O candidato está associado à oportunidade no staging, mas ainda não existe uma ligação externa exata.",
|
|
recommendation="Confirmar a compra/processo e decidir ligar, criar nova oportunidade ou ignorar.",
|
|
root_cause="association_pending_decision",
|
|
association_item_id=assessment.item_id, external_ref=external_ref,
|
|
effective_action="ASSOCIATE_OPPORTUNITY",
|
|
evidence={"identity": identity, "candidate_refs": sorted(refs)},
|
|
proposed_next_step={"action": "operator_association_decision", "automatic": False},
|
|
**common,
|
|
))
|
|
|
|
def assess_reconstructed(self) -> None:
|
|
required_keys = {
|
|
"historical_evidence_review_required",
|
|
"reconstructed_process_review_required",
|
|
"reconstructed_review_required",
|
|
}
|
|
cleared_keys = {
|
|
"historical_evidence_review_required",
|
|
"reconstructed_process_review_required",
|
|
"reconstructed_review_required",
|
|
}
|
|
validated_keys = {
|
|
"reconstructed_process_validated",
|
|
"historical_evidence_validated",
|
|
"reconstructed_review_validated",
|
|
}
|
|
for oid, opportunity in self.opportunities.items():
|
|
metadata = _as_dict(opportunity.get("metadata"))
|
|
tasks = self.tasks_by_opportunity.get(oid, [])
|
|
pending_reviews: List[str] = []
|
|
completed_reviews: List[str] = []
|
|
for task in tasks:
|
|
code = _upper(task.get("action_code"))
|
|
review_type = _review_type(task)
|
|
is_review = code == "REVIEW_RECONSTRUCTED_PROCESS" or (
|
|
code == "REVIEW_MANUALLY" and review_type in {
|
|
"reconstructed_process", "historical_evidence_review",
|
|
"historical_evidence_review_required",
|
|
}
|
|
)
|
|
if not is_review:
|
|
continue
|
|
if _lower(task.get("status")) == "pending":
|
|
pending_reviews.append(_s(task.get("id")))
|
|
elif _lower(task.get("status")) in DONE_TASK_STATUSES:
|
|
completed_reviews.append(_s(task.get("id")))
|
|
|
|
explicit_required = _truthy_metadata(metadata, required_keys) or bool(pending_reviews)
|
|
explicit_cleared = _truthy_metadata(metadata, validated_keys) or _falsey_metadata(metadata, cleared_keys) or bool(completed_reviews)
|
|
legacy = _is_legacy_record(metadata)
|
|
if explicit_required and not explicit_cleared:
|
|
state = "required"
|
|
elif explicit_cleared:
|
|
state = "cleared"
|
|
elif legacy:
|
|
state = "legacy_without_explicit_state"
|
|
else:
|
|
state = "not_applicable"
|
|
assessment = ReconstructedAssessment(
|
|
record_mode=_record_mode(metadata),
|
|
is_legacy_record=legacy,
|
|
explicit_required=explicit_required,
|
|
explicit_cleared=explicit_cleared,
|
|
pending_review_task_ids=pending_reviews,
|
|
completed_review_task_ids=completed_reviews,
|
|
state=state,
|
|
evidence={
|
|
"required_metadata": {key: metadata.get(key) for key in sorted(required_keys) if key in metadata},
|
|
"validated_metadata": {key: metadata.get(key) for key in sorted(validated_keys) if key in metadata},
|
|
},
|
|
)
|
|
self.reconstructed[oid] = assessment
|
|
|
|
central_code = _upper(self.next_actions.get(oid, {}).get("action_code"))
|
|
if state == "legacy_without_explicit_state" and central_code in OPERATIONAL_SENSITIVE_ACTIONS:
|
|
self.add(Finding(
|
|
category="policy", severity="medium", confidence="high", actionable=False,
|
|
code="LEGACY_RECORD_WITHOUT_EXPLICIT_REVIEW_STATE",
|
|
title="Registo reconstruído sem estado explícito de revisão",
|
|
explanation=(
|
|
"O metadata identifica um processo reconstruído, mas não diz se a revisão é obrigatória ou já foi concluída. "
|
|
"O sistema usa apenas uma lista parcial de ações para decidir quando mostrar a revisão."
|
|
),
|
|
recommendation="Adicionar um estado persistido required/validated e deixar de inferir apenas por clientflow_record_mode.",
|
|
root_cause="reconstructed_review_policy",
|
|
central_action=central_code,
|
|
effective_action=central_code,
|
|
evidence={"assessment": asdict(assessment), "ui_sensitive_allowlist": sorted(UI_RECONSTRUCTED_SENSITIVE_ACTIONS)},
|
|
proposed_next_step={"action": "define_explicit_reconstructed_review_state", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
if central_code in OPERATIONAL_SENSITIVE_ACTIONS and central_code not in UI_RECONSTRUCTED_SENSITIVE_ACTIONS:
|
|
self.add(Finding(
|
|
category="policy", severity="high", confidence="high", actionable=True,
|
|
code="RECONSTRUCTED_UI_GATE_ALLOWLIST_GAP",
|
|
title="A política de revisão reconstruída trata ações sensíveis de forma diferente",
|
|
explanation=(
|
|
f"A ação {central_code} é operacionalmente sensível, mas não pertence à lista usada pela UI para substituir a ação por revisão. "
|
|
"Outras ações sensíveis são bloqueadas apenas por estarem em clientflow_record_mode reconstruído."
|
|
),
|
|
recommendation="Decidir uma política única e aplicá-la no motor central, não numa lista local da página.",
|
|
root_cause="duplicated_action_policy",
|
|
central_action=central_code,
|
|
effective_action=central_code,
|
|
evidence={"assessment": asdict(assessment), "missing_from_ui_allowlist": central_code},
|
|
proposed_next_step={"action": "centralize_reconstructed_gate", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
def _explicit_association_blockers(self, oid: str) -> List[AssociationAssessment]:
|
|
return [
|
|
assessment for assessment in self.associations_by_opportunity.get(oid, [])
|
|
if assessment.blocks_opportunity and assessment.classification in {
|
|
"wrong_assigned_opportunity", "assigned_identity_conflict", "assigned_unresolved",
|
|
"conflict_multiple_exact_links",
|
|
}
|
|
]
|
|
|
|
def build_action_chains(self) -> None:
|
|
visible_by_opportunity: Dict[str, List[Dict[str, Any]]] = defaultdict(list)
|
|
for item in self.work_items:
|
|
oid = _s(item.get("opportunity_id"))
|
|
if oid:
|
|
visible_by_opportunity[oid].append(item)
|
|
|
|
for oid, opportunity in self.opportunities.items():
|
|
central = self.next_actions.get(oid, {})
|
|
central_code = _canonical_action(central.get("action_code"))
|
|
assoc_blockers = self._explicit_association_blockers(oid)
|
|
reconstructed = self.reconstructed.get(oid) or ReconstructedAssessment("", False, False, False, [], [], "not_applicable")
|
|
metadata = _as_dict(opportunity.get("metadata"))
|
|
fiscal_conflict = bool(
|
|
opportunity.get("has_nif_conflict") or metadata.get("has_nif_conflict")
|
|
or metadata.get("fiscal_conflict")
|
|
)
|
|
|
|
blocker = ""
|
|
reason = "central_decision"
|
|
effective = central_code
|
|
if assoc_blockers:
|
|
blocker = "association"
|
|
reason = "open_association_requires_decision"
|
|
effective = "ASSOCIATE_OPPORTUNITY"
|
|
elif reconstructed.state == "required":
|
|
blocker = "reconstructed_review"
|
|
reason = "explicit_reconstructed_review_required"
|
|
effective = "REVIEW_RECONSTRUCTED_PROCESS"
|
|
elif fiscal_conflict:
|
|
blocker = "fiscal_conflict"
|
|
reason = "fiscal_or_nif_conflict"
|
|
effective = "REVIEW_MANUALLY"
|
|
|
|
pending = self.pending_tasks_by_opportunity.get(oid, [])
|
|
visible = visible_by_opportunity.get(oid, [])
|
|
chain = ActionChain(
|
|
opportunity_id=oid,
|
|
central_action=central_code,
|
|
central_can_execute=central.get("can_execute"),
|
|
explicit_blocker=blocker,
|
|
effective_action=effective,
|
|
blocker_reason=reason,
|
|
association_blocker_ids=[assessment.item_id for assessment in assoc_blockers],
|
|
reconstructed_state=reconstructed.state,
|
|
materializable=_canonical_action(effective) in self.materialized_actions,
|
|
visible_actions=[_canonical_action(item.get("action_code")) for item in visible],
|
|
pending_actions=[_canonical_action(task.get("action_code")) for task in pending],
|
|
)
|
|
self.action_chains[oid] = chain
|
|
|
|
def _common(self, opportunity_id: str) -> Dict[str, Any]:
|
|
opportunity = self.opportunities.get(opportunity_id, {})
|
|
return {
|
|
"opportunity_id": opportunity_id,
|
|
"opportunity_title": _s(opportunity.get("title")),
|
|
"customer_name": _s(
|
|
opportunity.get("fiscal_customer_name") or opportunity.get("customer_name")
|
|
),
|
|
}
|
|
|
|
def _visible_item_for_task(self, task_id: str) -> Optional[Dict[str, Any]]:
|
|
return next(
|
|
(item for item in self.work_items if _lower(item.get("source")) == "task" and _s(item.get("id")) == task_id),
|
|
None,
|
|
)
|
|
|
|
def assess_visibility(self) -> None:
|
|
now = datetime.now(timezone.utc)
|
|
for task in self.all_pending_tasks:
|
|
task_id = _s(task.get("id"))
|
|
visible = self._visible_item_for_task(task_id)
|
|
code = _upper(task.get("action_code"))
|
|
due_at = task.get("due_at")
|
|
reason = "visible" if visible else "eligible_but_not_visible"
|
|
expected_hidden = False
|
|
if not visible and code.startswith("FOLLOW_UP_") and isinstance(due_at, datetime) and due_at > now:
|
|
reason = "scheduled_future_followup"
|
|
expected_hidden = True
|
|
elif not visible and _lower(task.get("status")) != "pending":
|
|
reason = "not_pending"
|
|
expected_hidden = True
|
|
|
|
self.visibility[task_id] = {
|
|
"visible": bool(visible),
|
|
"reason": reason,
|
|
"expected_hidden": expected_hidden,
|
|
"task": task,
|
|
"visible_item": visible,
|
|
}
|
|
|
|
if visible or expected_hidden:
|
|
continue
|
|
oid = _s(task.get("opportunity_id"))
|
|
chain = self.action_chains.get(oid)
|
|
severity = "high" if chain and _actions_match(code, chain.effective_action) else "medium"
|
|
self.add(Finding(
|
|
category="confirmed" if severity == "high" else "needs_review",
|
|
severity=severity,
|
|
confidence="high" if severity == "high" else "medium",
|
|
actionable=True,
|
|
code="PENDING_TASK_ELIGIBLE_BUT_NOT_VISIBLE",
|
|
title="Task pendente e vencida/atual não aparece no Centro de Trabalho",
|
|
explanation=(
|
|
"A task cumpre o filtro SQL base (pending e não é follow-up futuro), mas não foi devolvida na fila normalizada. "
|
|
"Pode estar a ser removida por regra de ruído ou ficar fora do limite/ordenação."
|
|
),
|
|
recommendation="Executar o diagnóstico de visibilidade e identificar a regra exata antes de recriar a task.",
|
|
root_cause="work_center_visibility",
|
|
task_id=task_id,
|
|
stored_action=code,
|
|
central_action=chain.central_action if chain else "",
|
|
effective_action=chain.effective_action if chain else "",
|
|
evidence={
|
|
"due_at": due_at,
|
|
"route": task.get("route"),
|
|
"priority": task.get("priority"),
|
|
"source_system": task.get("source_system"),
|
|
"operations_limit": self.limit,
|
|
"visible_queue_size": len(self.work_items),
|
|
},
|
|
proposed_next_step={"action": "trace_work_center_filter", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
def audit_visible_actions(self) -> None:
|
|
for item in self.work_items:
|
|
if _lower(item.get("source")) != "task":
|
|
continue
|
|
task_id = _s(item.get("id"))
|
|
oid = _s(item.get("opportunity_id"))
|
|
if not oid or oid not in self.opportunities:
|
|
continue
|
|
chain = self.action_chains.get(oid)
|
|
if not chain:
|
|
continue
|
|
displayed = _canonical_action(item.get("action_code"))
|
|
stored = _canonical_action(item.get("original_action_code") or item.get("action_code"))
|
|
central = chain.central_action
|
|
effective = chain.effective_action
|
|
ui_override = displayed != stored
|
|
|
|
if not _actions_match(displayed, effective):
|
|
confirmed_blocker = bool(chain.explicit_blocker)
|
|
downstream_sensitive = stored in OPERATIONAL_SENSITIVE_ACTIONS
|
|
category = "confirmed" if confirmed_blocker else "needs_review"
|
|
severity = "critical" if confirmed_blocker and downstream_sensitive else ("high" if confirmed_blocker else "medium")
|
|
confidence = "high" if confirmed_blocker else "medium"
|
|
self.add(Finding(
|
|
category=category, severity=severity, confidence=confidence, actionable=True,
|
|
code="VISIBLE_ACTION_DIFFERS_FROM_FIRST_EXPLICITLY_SAFE_ACTION",
|
|
title="A ação visível não corresponde à primeira ação suportada pelas evidências",
|
|
explanation=(
|
|
f"O Centro mostra {displayed}, enquanto a cadeia de decisão indica {effective}. "
|
|
+ ("Existe um bloqueio explícito comprovado." if confirmed_blocker else "Não existe bloqueio explícito; a diferença pode ser uma política ou task antiga.")
|
|
),
|
|
recommendation="Preservar apenas uma primeira ação e explicar o bloqueio/razão no cartão.",
|
|
root_cause="action_precedence",
|
|
work_item_id=task_id, task_id=task_id,
|
|
displayed_action=displayed, stored_action=stored,
|
|
central_action=central, effective_action=effective,
|
|
evidence={"action_chain": asdict(chain), "ui_override": ui_override, "item": item},
|
|
proposed_next_step={"action": "review_or_reclassify_task", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
if ui_override:
|
|
supported = False
|
|
if displayed == "ASSOCIATE_OPPORTUNITY" and chain.explicit_blocker == "association":
|
|
supported = True
|
|
if displayed == "REVIEW_RECONSTRUCTED_PROCESS" and chain.explicit_blocker == "reconstructed_review":
|
|
supported = True
|
|
if not supported:
|
|
self.add(Finding(
|
|
category="policy", severity="medium", confidence="high", actionable=True,
|
|
code="UI_ACTION_OVERRIDE_WITHOUT_MATCHING_EXPLICIT_BLOCKER",
|
|
title="A UI substituiu a ação persistida sem existir o bloqueio explícito correspondente",
|
|
explanation=f"A task guarda {stored}, mas a fila mostra {displayed}; a regra de apresentação não está alinhada com o estado persistido.",
|
|
recommendation="Mover a precedência para o motor central e persistir o estado que justifica o override.",
|
|
root_cause="ui_only_action_override",
|
|
work_item_id=task_id, task_id=task_id,
|
|
displayed_action=displayed, stored_action=stored,
|
|
central_action=central, effective_action=effective,
|
|
evidence={"action_chain": asdict(chain), "item": item},
|
|
proposed_next_step={"action": "centralize_ui_override_policy", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
central_payload = self.next_actions.get(oid, {})
|
|
if central_payload.get("can_execute") is False and _actions_match(displayed, central):
|
|
self.add(Finding(
|
|
category="confirmed", severity="high", confidence="high", actionable=True,
|
|
code="CENTRAL_BLOCKED_ACTION_VISIBLE",
|
|
title="Ação que o motor central declara bloqueada aparece como trabalho executável",
|
|
explanation="O mesmo action_code é mostrado, mas can_execute=false no motor central.",
|
|
recommendation="Mostrar o motivo e impedir conclusão até resolver o bloqueio.",
|
|
root_cause="central_block_not_enforced",
|
|
work_item_id=task_id, task_id=task_id,
|
|
displayed_action=displayed, stored_action=stored,
|
|
central_action=central, effective_action=effective,
|
|
evidence={"reason_if_blocked": central_payload.get("reason_if_blocked"), "central": central_payload},
|
|
proposed_next_step={"action": "enforce_central_can_execute", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
def audit_task_coherence(self) -> None:
|
|
for oid, opportunity in self.opportunities.items():
|
|
pending = self.pending_tasks_by_opportunity.get(oid, [])
|
|
chain = self.action_chains.get(oid)
|
|
if not chain:
|
|
continue
|
|
pending_codes = [_canonical_action(task.get("action_code")) for task in pending]
|
|
effective = chain.effective_action
|
|
|
|
if len(pending) > 1:
|
|
competing = [code for code in pending_codes if not _actions_match(code, effective)]
|
|
severity = "high" if competing else "medium"
|
|
self.add(Finding(
|
|
category="confirmed" if competing else "needs_review",
|
|
severity=severity, confidence="high", actionable=True,
|
|
code="MULTIPLE_PENDING_TASKS_REQUIRE_ORDERING",
|
|
title="A oportunidade tem várias tasks pendentes",
|
|
explanation=(
|
|
"Existem várias ações humanas em paralelo. "
|
|
+ ("Pelo menos uma não corresponde à primeira ação segura." if competing else "As tasks parecem equivalentes/duplicadas.")
|
|
),
|
|
recommendation="Definir a primeira ação e resolver/adiar as restantes após revisão.",
|
|
root_cause="task_competition",
|
|
central_action=chain.central_action, effective_action=effective,
|
|
evidence={"task_ids": [_s(task.get("id")) for task in pending], "pending_codes": pending_codes, "competing": competing},
|
|
proposed_next_step={"action": "review_competing_tasks", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
matching_pending = [task for task in pending if _actions_match(task.get("action_code"), effective)]
|
|
matching_visible = [
|
|
item for item in self.work_items
|
|
if _s(item.get("opportunity_id")) == oid and _actions_match(item.get("action_code"), effective)
|
|
]
|
|
|
|
if effective in NON_TASK_ACTIONS:
|
|
for task in pending:
|
|
code = _canonical_action(task.get("action_code"))
|
|
if code.startswith("FOLLOW_UP_") and self.visibility.get(_s(task.get("id")), {}).get("reason") == "scheduled_future_followup":
|
|
continue
|
|
self.add(Finding(
|
|
category="needs_review", severity="medium", confidence="medium", actionable=True,
|
|
code="PENDING_TASK_DURING_NON_ACTION_STATE",
|
|
title="Existe task pendente enquanto o motor indica espera/sem ação",
|
|
explanation=f"A decisão efetiva é {effective}, mas existe a task {code}.",
|
|
recommendation="Confirmar se a task é futura, obsoleta ou uma exceção de negócio.",
|
|
root_cause="stale_or_exception_task",
|
|
task_id=_s(task.get("id")), stored_action=code,
|
|
central_action=chain.central_action, effective_action=effective,
|
|
evidence={"task": task, "action_chain": asdict(chain)},
|
|
proposed_next_step={"action": "review_task_during_wait_state", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
continue
|
|
|
|
if chain.materializable and chain.central_can_execute is not False and not matching_pending and not matching_visible:
|
|
self.add(Finding(
|
|
category="confirmed", severity="high", confidence="high", actionable=True,
|
|
code="MATERIALIZABLE_EFFECTIVE_ACTION_HAS_NO_TASK",
|
|
title="A primeira ação segura é materializável, mas não existe task",
|
|
explanation=f"{effective} pertence ao conjunto real MATERIALIZED_ACTIONS e está executável, mas não existe task pendente nem item visível.",
|
|
recommendation="Materializar uma única task idempotente após confirmar que não existe bloqueio anterior.",
|
|
root_cause="task_materialization_gap",
|
|
central_action=chain.central_action, effective_action=effective,
|
|
evidence={"materialized_actions": sorted(self.materialized_actions), "action_chain": asdict(chain)},
|
|
proposed_next_step={"action": "materialize_effective_action", "action_code": effective, "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
elif not chain.materializable and effective not in NON_TASK_ACTIONS and not matching_pending and not matching_visible:
|
|
# This is a policy fact, not a bug: several central actions are
|
|
# intentionally not materialized automatically.
|
|
self.add(Finding(
|
|
category="policy", severity="info", confidence="high", actionable=False,
|
|
code="CENTRAL_ACTION_NOT_COVERED_BY_TASK_MATERIALIZER",
|
|
title="Ação central sem task porque não pertence à política de materialização",
|
|
explanation=f"A decisão é {effective}, mas o materializador atual só cobre {', '.join(sorted(self.materialized_actions))}.",
|
|
recommendation="Decidir se a ação deve permanecer apenas no pipeline ou passar a gerar task.",
|
|
root_cause="task_materialization_policy",
|
|
central_action=chain.central_action, effective_action=effective,
|
|
evidence={"materialized_actions": sorted(self.materialized_actions), "action_chain": asdict(chain)},
|
|
proposed_next_step={"action": "review_materialization_policy", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
def audit_association_surface(self) -> None:
|
|
visible_by_oid: Dict[str, List[Dict[str, Any]]] = defaultdict(list)
|
|
for item in self.work_items:
|
|
oid = _s(item.get("opportunity_id"))
|
|
if oid:
|
|
visible_by_oid[oid].append(item)
|
|
for oid, assessments in self.associations_by_opportunity.items():
|
|
blockers = [assessment for assessment in assessments if assessment.blocks_opportunity]
|
|
if not blockers:
|
|
continue
|
|
pending = self.pending_tasks_by_opportunity.get(oid, [])
|
|
visible = visible_by_oid.get(oid, [])
|
|
has_association_task = any(_canonical_action(task.get("action_code")) == "ASSOCIATE_OPPORTUNITY" for task in pending)
|
|
has_association_visible = any(_canonical_action(item.get("action_code")) == "ASSOCIATE_OPPORTUNITY" for item in visible)
|
|
if not has_association_task and not has_association_visible:
|
|
self.add(Finding(
|
|
category="confirmed", severity="high", confidence="high", actionable=True,
|
|
code="EXPLICIT_ASSOCIATION_BLOCKER_NOT_SURFACED",
|
|
title="Associação atribuída bloqueia a oportunidade mas não aparece no Centro de Trabalho",
|
|
explanation="Existe um candidato atribuído/conflituoso que exige decisão, sem task nem cartão de associação.",
|
|
recommendation="Materializar uma task de associação ou mostrar o candidato diretamente na fila de revisão.",
|
|
root_cause="association_not_materialized",
|
|
association_item_id=blockers[0].item_id,
|
|
external_ref=blockers[0].external_ref,
|
|
central_action=self.action_chains.get(oid).central_action if oid in self.action_chains else "",
|
|
effective_action="ASSOCIATE_OPPORTUNITY",
|
|
evidence={"blockers": [asdict(assessment) for assessment in blockers]},
|
|
proposed_next_step={"action": "surface_association_decision", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
for oid, pending in self.pending_tasks_by_opportunity.items():
|
|
assoc_tasks = [task for task in pending if _canonical_action(task.get("action_code")) == "ASSOCIATE_OPPORTUNITY"]
|
|
if not assoc_tasks:
|
|
continue
|
|
blockers = self._explicit_association_blockers(oid)
|
|
stale = [a for a in self.associations_by_opportunity.get(oid, []) if a.classification == "stale_already_linked"]
|
|
if not blockers:
|
|
category = "confirmed" if stale else "needs_review"
|
|
severity = "high" if stale else "medium"
|
|
self.add(Finding(
|
|
category=category, severity=severity, confidence="high" if stale else "medium", actionable=True,
|
|
code="ASSOCIATION_TASK_WITHOUT_UNRESOLVED_BLOCKER",
|
|
title="Existe task de associação sem candidato bloqueante por resolver",
|
|
explanation=(
|
|
"Os candidatos relacionados já estão exatamente ligados." if stale
|
|
else "Não foi encontrado candidato atribuído/conflituoso que justifique a task."
|
|
),
|
|
recommendation="Resolver candidatos obsoletos e depois rever/ignorar a task de associação.",
|
|
root_cause="stale_association_task",
|
|
task_id=_s(assoc_tasks[0].get("id")), stored_action="ASSOCIATE_OPPORTUNITY",
|
|
effective_action=self.action_chains.get(oid).effective_action if oid in self.action_chains else "",
|
|
evidence={"task_ids": [_s(task.get("id")) for task in assoc_tasks], "stale_candidates": [asdict(a) for a in stale]},
|
|
proposed_next_step={"action": "review_stale_association_task", "automatic": False},
|
|
**self._common(oid),
|
|
))
|
|
|
|
def build_dossiers(self) -> None:
|
|
finding_oids = {finding.opportunity_id for finding in self.findings if finding.opportunity_id}
|
|
visible_oids = {_s(item.get("opportunity_id")) for item in self.work_items if _s(item.get("opportunity_id"))}
|
|
candidate_oids = set(self.associations_by_opportunity)
|
|
selected = set(self.opportunities) if self.include_clean else finding_oids | visible_oids | candidate_oids
|
|
if self.focus:
|
|
selected = {
|
|
oid for oid in selected
|
|
if _focus_match(
|
|
self.focus,
|
|
oid,
|
|
self.opportunities.get(oid, {}).get("title"),
|
|
self.opportunities.get(oid, {}).get("customer_name"),
|
|
self.opportunities.get(oid, {}).get("fiscal_customer_name"),
|
|
*[a.external_ref for a in self.associations_by_opportunity.get(oid, [])],
|
|
)
|
|
}
|
|
|
|
for oid in sorted(selected):
|
|
opportunity = self.opportunities.get(oid, {})
|
|
snapshot = self.operation_snapshots.get(oid, {})
|
|
cards = []
|
|
for card in snapshot.get("cards", []) or []:
|
|
cards.append({
|
|
"system": card.get("system"),
|
|
"external_type": card.get("external_type"),
|
|
"external_id": card.get("external_id"),
|
|
"status": card.get("status"),
|
|
"title": card.get("title"),
|
|
"subtitle": card.get("subtitle"),
|
|
"payload": card.get("payload"),
|
|
})
|
|
self.dossiers[oid] = {
|
|
"opportunity": opportunity,
|
|
"action_chain": asdict(self.action_chains.get(oid)) if oid in self.action_chains else {},
|
|
"central_decision": self.next_actions.get(oid, {}),
|
|
"reconstructed_review": asdict(self.reconstructed.get(oid)) if oid in self.reconstructed else {},
|
|
"visible_work_items": [item for item in self.work_items if _s(item.get("opportunity_id")) == oid],
|
|
"pending_tasks": self.pending_tasks_by_opportunity.get(oid, []),
|
|
"recent_tasks": self.tasks_by_opportunity.get(oid, [])[:20],
|
|
"task_visibility": [
|
|
self.visibility.get(_s(task.get("id")), {})
|
|
for task in self.pending_tasks_by_opportunity.get(oid, [])
|
|
],
|
|
"association_assessments": [asdict(a) for a in self.associations_by_opportunity.get(oid, [])],
|
|
"documents": self.documents_by_opportunity.get(oid, []),
|
|
"operation_snapshot": {"cards": cards, "links": snapshot.get("links", []), "error": snapshot.get("error")},
|
|
"findings": [asdict(f) for f in self.findings if f.opportunity_id == oid],
|
|
}
|
|
|
|
def audit(self) -> None:
|
|
self.assess_associations()
|
|
self.assess_reconstructed()
|
|
self.build_action_chains()
|
|
self.assess_visibility()
|
|
self.audit_visible_actions()
|
|
self.audit_task_coherence()
|
|
self.audit_association_surface()
|
|
|
|
unique: Dict[Tuple[str, str, str, str, str], Finding] = {}
|
|
for finding in self.findings:
|
|
key = (
|
|
finding.code, finding.opportunity_id, finding.task_id,
|
|
finding.association_item_id, finding.effective_action,
|
|
)
|
|
unique.setdefault(key, finding)
|
|
self.findings = sorted(
|
|
unique.values(),
|
|
key=lambda finding: (
|
|
CATEGORY_ORDER.get(finding.category, 99),
|
|
SEVERITY_ORDER.get(finding.severity, 99),
|
|
CONFIDENCE_ORDER.get(finding.confidence, 99),
|
|
finding.customer_name.casefold(), finding.code,
|
|
),
|
|
)
|
|
self.build_dossiers()
|
|
|
|
categories = Counter(f.category for f in self.findings)
|
|
severities = Counter(f.severity for f in self.findings)
|
|
actionable = Counter("actionable" if f.actionable else "informational" for f in self.findings)
|
|
root_causes = Counter(f.root_cause or "other" for f in self.findings)
|
|
self.summary = {
|
|
"generated_at": datetime.now(timezone.utc).isoformat(),
|
|
"project_root": str(PROJECT_ROOT),
|
|
"read_only": True,
|
|
"work_items_visible": len(self.work_items),
|
|
"operations_counts": self.operations_counts,
|
|
"open_opportunities_audited": len(self.opportunities),
|
|
"pending_tasks_audited": len(self.all_pending_tasks),
|
|
"open_reconciliation_candidates": len(self.association_assessments),
|
|
"dossiers": len(self.dossiers),
|
|
"findings": len(self.findings),
|
|
"categories": dict(categories),
|
|
"severity": dict(severities),
|
|
"actionability": dict(actionable),
|
|
"root_causes": dict(root_causes),
|
|
"confirmed_actionable": sum(1 for f in self.findings if f.category == "confirmed" and f.actionable),
|
|
"policy_only": sum(1 for f in self.findings if f.category == "policy"),
|
|
}
|
|
|
|
def _write_csvs(self, output_dir: Path) -> Dict[str, str]:
|
|
findings_path = output_dir / "forensic_findings.csv"
|
|
association_path = output_dir / "association_matrix.csv"
|
|
action_path = output_dir / "action_chain_matrix.csv"
|
|
visibility_path = output_dir / "task_visibility_matrix.csv"
|
|
|
|
fields = list(Finding.__dataclass_fields__.keys())
|
|
with findings_path.open("w", encoding="utf-8", newline="") as handle:
|
|
writer = csv.DictWriter(handle, fieldnames=fields)
|
|
writer.writeheader()
|
|
for finding in self.findings:
|
|
row = asdict(finding)
|
|
row["evidence"] = json.dumps(row["evidence"], ensure_ascii=False, default=_json_default)
|
|
row["proposed_next_step"] = json.dumps(row["proposed_next_step"], ensure_ascii=False, default=_json_default)
|
|
writer.writerow(row)
|
|
|
|
association_fields = [
|
|
"item_id", "source_system", "external_type", "external_ref", "status",
|
|
"assigned_opportunity_id", "exact_link_opportunity_ids", "match_kind",
|
|
"classification", "blocks_opportunity", "confidence", "identity",
|
|
]
|
|
with association_path.open("w", encoding="utf-8", newline="") as handle:
|
|
writer = csv.DictWriter(handle, fieldnames=association_fields)
|
|
writer.writeheader()
|
|
for assessment in self.association_assessments:
|
|
row = asdict(assessment)
|
|
row["exact_link_opportunity_ids"] = ",".join(row["exact_link_opportunity_ids"])
|
|
row["identity"] = json.dumps(row["identity"], ensure_ascii=False, default=_json_default)
|
|
writer.writerow({key: row.get(key) for key in association_fields})
|
|
|
|
action_fields = list(ActionChain.__dataclass_fields__.keys())
|
|
with action_path.open("w", encoding="utf-8", newline="") as handle:
|
|
writer = csv.DictWriter(handle, fieldnames=action_fields)
|
|
writer.writeheader()
|
|
for chain in self.action_chains.values():
|
|
row = asdict(chain)
|
|
for key in ("association_blocker_ids", "visible_actions", "pending_actions"):
|
|
row[key] = ",".join(row[key])
|
|
writer.writerow(row)
|
|
|
|
visibility_fields = [
|
|
"task_id", "opportunity_id", "customer_name", "action_code", "due_at",
|
|
"route", "priority", "visible", "reason", "expected_hidden",
|
|
]
|
|
with visibility_path.open("w", encoding="utf-8", newline="") as handle:
|
|
writer = csv.DictWriter(handle, fieldnames=visibility_fields)
|
|
writer.writeheader()
|
|
for task_id, assessment in self.visibility.items():
|
|
task = assessment.get("task") or {}
|
|
oid = _s(task.get("opportunity_id"))
|
|
row = {
|
|
"task_id": task_id,
|
|
"opportunity_id": oid,
|
|
"customer_name": self._common(oid).get("customer_name"),
|
|
"action_code": _upper(task.get("action_code")),
|
|
"due_at": task.get("due_at"),
|
|
"route": task.get("route"),
|
|
"priority": task.get("priority"),
|
|
"visible": assessment.get("visible"),
|
|
"reason": assessment.get("reason"),
|
|
"expected_hidden": assessment.get("expected_hidden"),
|
|
}
|
|
writer.writerow(row)
|
|
|
|
return {
|
|
"findings_csv": str(findings_path),
|
|
"association_csv": str(association_path),
|
|
"action_chain_csv": str(action_path),
|
|
"visibility_csv": str(visibility_path),
|
|
}
|
|
|
|
def report(self, output_dir: Path) -> Dict[str, str]:
|
|
output_dir.mkdir(parents=True, exist_ok=True)
|
|
json_path = output_dir / "work_center_forensic_audit.json"
|
|
markdown_path = output_dir / "WORK_CENTER_FORENSIC_AUDIT.md"
|
|
dossiers_path = output_dir / "CASE_DOSSIERS.md"
|
|
plan_path = output_dir / "review_plan.json"
|
|
sql_path = output_dir / "diagnostic_queries.sql"
|
|
|
|
payload = {
|
|
"summary": self.summary,
|
|
"findings": [asdict(f) for f in self.findings],
|
|
"association_assessments": [asdict(a) for a in self.association_assessments],
|
|
"reconstructed_assessments": {oid: asdict(value) for oid, value in self.reconstructed.items()},
|
|
"action_chains": {oid: asdict(value) for oid, value in self.action_chains.items()},
|
|
"task_visibility": self.visibility,
|
|
"dossiers": self.dossiers,
|
|
}
|
|
json_path.write_text(json.dumps(payload, ensure_ascii=False, indent=2, default=_json_default), encoding="utf-8")
|
|
|
|
confirmed = [f for f in self.findings if f.category == "confirmed"]
|
|
needs_review = [f for f in self.findings if f.category == "needs_review"]
|
|
policy = [f for f in self.findings if f.category == "policy"]
|
|
lines = [
|
|
"# Auditoria forense do Centro de Trabalho e associações",
|
|
"",
|
|
f"Gerada em: `{self.summary.get('generated_at')}`",
|
|
f"Itens visíveis: **{self.summary.get('work_items_visible', 0)}**",
|
|
f"Oportunidades abertas: **{self.summary.get('open_opportunities_audited', 0)}**",
|
|
f"Tasks pendentes: **{self.summary.get('pending_tasks_audited', 0)}**",
|
|
f"Candidatos abertos: **{self.summary.get('open_reconciliation_candidates', 0)}**",
|
|
f"Achados confirmados: **{len(confirmed)}**",
|
|
f"Casos a rever: **{len(needs_review)}**",
|
|
f"Lacunas/políticas: **{len(policy)}**",
|
|
"",
|
|
"## Como interpretar",
|
|
"",
|
|
"- **Confirmado:** a base de dados contém evidência suficiente de incoerência.",
|
|
"- **A rever:** existe uma diferença, mas falta contexto para chamar bug.",
|
|
"- **Política:** comportamento explicado pelo desenho atual; requer decisão de produto, não correção de dados em massa.",
|
|
"",
|
|
]
|
|
for heading, group in (("Problemas confirmados", confirmed), ("Casos a rever", needs_review), ("Lacunas de política", policy)):
|
|
lines.extend([f"## {heading}", ""])
|
|
if not group:
|
|
lines.append("Nenhum.")
|
|
lines.append("")
|
|
continue
|
|
for index, finding in enumerate(group, 1):
|
|
lines.extend([
|
|
f"### {index}. [{finding.severity.upper()}] {finding.title}",
|
|
"",
|
|
f"- Código: `{finding.code}`",
|
|
f"- Confiança: **{finding.confidence}** · Acionável: **{'sim' if finding.actionable else 'não'}**",
|
|
f"- Cliente: {finding.customer_name or '—'}",
|
|
f"- Oportunidade: `{finding.opportunity_id or '—'}` {finding.opportunity_title}",
|
|
f"- Task: `{finding.task_id or finding.work_item_id or '—'}`",
|
|
f"- Candidato: `{finding.association_item_id or '—'}` · `{finding.external_ref or '—'}`",
|
|
f"- Ação visível: `{finding.displayed_action or '—'}`",
|
|
f"- Ação persistida: `{finding.stored_action or '—'}`",
|
|
f"- Ação central: `{finding.central_action or '—'}`",
|
|
f"- Primeira ação suportada: `{finding.effective_action or '—'}`",
|
|
f"- Causa provável: `{finding.root_cause or '—'}`",
|
|
f"- Explicação: {finding.explanation}",
|
|
f"- Recomendação: {finding.recommendation}",
|
|
f"- Evidência: `{json.dumps(finding.evidence, ensure_ascii=False, default=_json_default)}`",
|
|
"",
|
|
])
|
|
lines.extend([
|
|
"## Garantia de segurança",
|
|
"",
|
|
"O script não possui `--apply`. As consultas próprias usam uma transação PostgreSQL `READ ONLY`; os serviços chamados são de leitura.",
|
|
"",
|
|
])
|
|
markdown_path.write_text("\n".join(lines), encoding="utf-8")
|
|
|
|
dossier_lines = ["# Dossiers por oportunidade", ""]
|
|
for oid, dossier in self.dossiers.items():
|
|
opp = dossier.get("opportunity") or {}
|
|
chain = dossier.get("action_chain") or {}
|
|
dossier_lines.extend([
|
|
f"## {dossier.get('opportunity', {}).get('fiscal_customer_name') or opp.get('customer_name') or oid}",
|
|
"",
|
|
f"- ID: `{oid}`",
|
|
f"- Título: {opp.get('title') or '—'}",
|
|
f"- Fase: `{opp.get('stage') or '—'}`",
|
|
f"- Ação central: `{chain.get('central_action') or '—'}`",
|
|
f"- Bloqueio explícito: `{chain.get('explicit_blocker') or 'nenhum'}`",
|
|
f"- Primeira ação suportada: `{chain.get('effective_action') or '—'}`",
|
|
f"- Estado reconstruído: `{chain.get('reconstructed_state') or '—'}`",
|
|
f"- Tasks pendentes: `{', '.join(chain.get('pending_actions') or []) or 'nenhuma'}`",
|
|
f"- Ações visíveis: `{', '.join(chain.get('visible_actions') or []) or 'nenhuma'}`",
|
|
"",
|
|
"### Associações",
|
|
"",
|
|
])
|
|
associations = dossier.get("association_assessments") or []
|
|
if not associations:
|
|
dossier_lines.append("Sem candidatos relacionados.")
|
|
for assessment in associations:
|
|
dossier_lines.append(
|
|
f"- `{assessment.get('external_ref') or assessment.get('item_id')}`: "
|
|
f"{assessment.get('classification')} · bloqueia={assessment.get('blocks_opportunity')} · "
|
|
f"links={assessment.get('exact_link_opportunity_ids')}"
|
|
)
|
|
dossier_lines.extend(["", "### Tasks pendentes", ""])
|
|
pending = dossier.get("pending_tasks") or []
|
|
if not pending:
|
|
dossier_lines.append("Sem tasks pendentes.")
|
|
for task in pending:
|
|
visibility = self.visibility.get(_s(task.get("id")), {})
|
|
dossier_lines.append(
|
|
f"- `{task.get('id')}` `{task.get('action_code')}` · rota={task.get('route')} · "
|
|
f"due={task.get('due_at')} · visível={visibility.get('visible')} ({visibility.get('reason')})"
|
|
)
|
|
dossier_lines.extend(["", "### Documentos", ""])
|
|
docs = dossier.get("documents") or []
|
|
if not docs:
|
|
dossier_lines.append("Sem documentos ligados.")
|
|
for doc in docs:
|
|
dossier_lines.append(
|
|
f"- `{doc.get('document_kind')}` `{doc.get('document_number') or doc.get('external_id')}` "
|
|
f"· {doc.get('total_amount')} {doc.get('currency')} · role={doc.get('role')} · active={doc.get('is_active')}"
|
|
)
|
|
dossier_lines.extend(["", "### Estado operacional", ""])
|
|
cards = (dossier.get("operation_snapshot") or {}).get("cards") or []
|
|
if not cards:
|
|
dossier_lines.append("Sem cards operacionais.")
|
|
for card in cards:
|
|
dossier_lines.append(
|
|
f"- `{card.get('system')}/{card.get('external_type')}` `{card.get('external_id')}` · "
|
|
f"status={card.get('status')} · {card.get('title') or ''} {card.get('subtitle') or ''}"
|
|
)
|
|
dossier_lines.append("")
|
|
dossiers_path.write_text("\n".join(dossier_lines), encoding="utf-8")
|
|
|
|
review_plan = {
|
|
"generated_at": self.summary.get("generated_at"),
|
|
"read_only": True,
|
|
"warning": "Não aplicar automaticamente. Usar os guardas e dossiers para preparar correções específicas.",
|
|
"confirmed_actionable": [
|
|
{
|
|
"code": f.code,
|
|
"opportunity_id": f.opportunity_id,
|
|
"task_id": f.task_id,
|
|
"association_item_id": f.association_item_id,
|
|
"external_ref": f.external_ref,
|
|
"effective_action": f.effective_action,
|
|
"next_step": f.proposed_next_step,
|
|
"evidence": f.evidence,
|
|
}
|
|
for f in confirmed if f.actionable
|
|
],
|
|
"needs_policy_decision": [
|
|
{
|
|
"code": f.code,
|
|
"opportunity_id": f.opportunity_id,
|
|
"central_action": f.central_action,
|
|
"next_step": f.proposed_next_step,
|
|
}
|
|
for f in policy
|
|
],
|
|
}
|
|
plan_path.write_text(json.dumps(review_plan, ensure_ascii=False, indent=2, default=_json_default), encoding="utf-8")
|
|
|
|
sql_path.write_text("""-- Diagnóstico forense complementar (executar manualmente em READ ONLY)
|
|
|
|
-- 1. Cadeia de tasks por oportunidade, incluindo metadata e eventos
|
|
SELECT o.customer_name, o.id AS opportunity_id, o.stage,
|
|
t.id AS task_id, t.action_code, t.status, t.route, t.priority,
|
|
t.due_at, t.created_at, t.done_at, t.metadata,
|
|
COALESCE(json_agg(json_build_object(
|
|
'event_type', te.event_type, 'created_at', te.created_at,
|
|
'created_by', te.created_by, 'payload', te.payload
|
|
) ORDER BY te.created_at) FILTER (WHERE te.id IS NOT NULL), '[]') AS events
|
|
FROM opportunities o
|
|
LEFT JOIN tasks t ON t.opportunity_id = o.id
|
|
LEFT JOIN task_events te ON te.task_id = t.id
|
|
WHERE o.status = 'open'
|
|
GROUP BY o.customer_name, o.id, o.stage, t.id
|
|
ORDER BY o.customer_name, t.created_at;
|
|
|
|
-- 2. Candidatos abertos e ligações exatas potenciais
|
|
SELECT ri.id AS candidate_id, ri.source_system, ri.external_type,
|
|
ri.external_id, ri.document_number, ri.opportunity_id AS assigned_opportunity,
|
|
ri.customer_name, ri.customer_tax_id, ri.status,
|
|
ol.id AS link_id, ol.opportunity_id AS linked_opportunity,
|
|
ol.external_id AS link_external_id, ol.external_name AS link_external_name,
|
|
ol.status AS link_status
|
|
FROM reconciliation_items ri
|
|
LEFT JOIN operation_links ol
|
|
ON lower(ol.system) = lower(ri.source_system)
|
|
AND COALESCE(ol.status, '') NOT IN ('ignored','deleted')
|
|
AND (
|
|
lower(COALESCE(ol.external_id,'')) IN (lower(COALESCE(ri.external_id,'')), lower(COALESCE(ri.document_number,'')))
|
|
OR lower(COALESCE(ol.external_name,'')) IN (lower(COALESCE(ri.external_id,'')), lower(COALESCE(ri.document_number,'')))
|
|
)
|
|
WHERE ri.status IN ('open','needs_review','conflict')
|
|
ORDER BY ri.updated_at DESC;
|
|
|
|
-- 3. Registos reconstruídos sem estado persistido de revisão
|
|
SELECT id, customer_name, title, stage,
|
|
metadata->>'clientflow_record_mode' AS record_mode,
|
|
metadata->>'historical_evidence_review_required' AS review_required,
|
|
metadata->>'reconstructed_process_validated' AS review_validated
|
|
FROM opportunities
|
|
WHERE status = 'open'
|
|
AND metadata->>'clientflow_record_mode' IN (
|
|
'reconstructed_invoice_review','historical_reconstructed','legacy_review'
|
|
)
|
|
ORDER BY updated_at DESC;
|
|
|
|
-- 4. Tasks pendentes que deveriam ser elegíveis para o Centro de Trabalho
|
|
SELECT t.id, t.opportunity_id, o.customer_name, t.action_code,
|
|
t.route, t.priority, t.due_at, t.source_system, t.metadata
|
|
FROM tasks t
|
|
LEFT JOIN opportunities o ON o.id = t.opportunity_id
|
|
WHERE t.status = 'pending'
|
|
AND NOT (t.action_code LIKE 'FOLLOW_UP_%' AND t.due_at IS NOT NULL AND t.due_at > now())
|
|
ORDER BY
|
|
CASE lower(COALESCE(t.priority,'normal'))
|
|
WHEN 'urgente' THEN 0 WHEN 'alta' THEN 1 WHEN 'high' THEN 1 WHEN 'normal' THEN 2 ELSE 3
|
|
END,
|
|
t.created_at DESC;
|
|
""", encoding="utf-8")
|
|
|
|
paths = {
|
|
"json": str(json_path),
|
|
"markdown": str(markdown_path),
|
|
"dossiers": str(dossiers_path),
|
|
"review_plan": str(plan_path),
|
|
"diagnostic_sql": str(sql_path),
|
|
}
|
|
paths.update(self._write_csvs(output_dir))
|
|
return paths
|
|
|
|
|
|
def _self_test() -> None:
|
|
assert _canonical_action("SHIP_ORDER") == "CREATE_SHIPMENT"
|
|
assert _actions_match("SHIP_ORDER", "CREATE_SHIPMENT")
|
|
assert _actions_match("RECONCILE_DOCUMENTS", "ASSOCIATE_OPPORTUNITY")
|
|
assert not _actions_match("SEND_INVOICE", "CONFIRM_PAYMENT")
|
|
assert _reference_variants("S00325", system="odoo", external_type="sale_order") == {"s00325", "s325"}
|
|
assert "323" not in _reference_variants("S00323", system="odoo", external_type="sale_order")
|
|
assert _norm_tax_id("PT 510 177 441") == "510177441"
|
|
assert _is_legacy_record({"clientflow_record_mode": "historical_reconstructed"})
|
|
assert not _is_legacy_record({"clientflow_record_mode": "normal"})
|
|
assert _truthy_metadata({"historical_evidence_review_required": True}, {"historical_evidence_review_required"})
|
|
assert _falsey_metadata({"historical_evidence_review_required": False}, {"historical_evidence_review_required"})
|
|
print("Self-test OK")
|
|
|
|
|
|
def _parser() -> argparse.ArgumentParser:
|
|
parser = argparse.ArgumentParser(description=__doc__)
|
|
parser.add_argument("--limit", type=int, default=500, help="Máximo de itens normalizados do Centro de Trabalho.")
|
|
parser.add_argument("--opportunity-limit", type=int, default=1000, help="Máximo de oportunidades abertas.")
|
|
parser.add_argument("--focus", action="append", default=[], help="Filtrar relatório por cliente, UUID, referência ou código; repetível.")
|
|
parser.add_argument("--include-clean", action="store_true", help="Incluir dossiers de oportunidades sem achados.")
|
|
parser.add_argument("--output-dir", default="", help="Diretório de saída.")
|
|
parser.add_argument(
|
|
"--fail-on",
|
|
choices=["confirmed_high", "critical", "high", "medium", "low", "info", "never"],
|
|
default="confirmed_high",
|
|
help="Código 2 quando existe achado nesse limiar. confirmed_high ignora policy/needs_review.",
|
|
)
|
|
parser.add_argument("--terminal-limit", type=int, default=80, help="Máximo de achados mostrados no terminal.")
|
|
parser.add_argument("--self-test", action="store_true")
|
|
return parser
|
|
|
|
|
|
def main() -> int:
|
|
args = _parser().parse_args()
|
|
if args.self_test:
|
|
_self_test()
|
|
return 0
|
|
|
|
auditor = ForensicAuditor(
|
|
limit=args.limit,
|
|
opportunity_limit=args.opportunity_limit,
|
|
focus=args.focus,
|
|
include_clean=args.include_clean,
|
|
)
|
|
auditor.load()
|
|
auditor.audit()
|
|
|
|
if args.output_dir:
|
|
output_dir = Path(args.output_dir)
|
|
else:
|
|
stamp = datetime.now().strftime("%Y%m%d_%H%M%S")
|
|
output_dir = PROJECT_ROOT / "audit_reports" / f"work_center_forensic_{stamp}"
|
|
paths = auditor.report(output_dir)
|
|
|
|
summary = auditor.summary
|
|
print("Auditoria forense do Centro de Trabalho em modo READ ONLY")
|
|
print(f"Itens visíveis: {summary.get('work_items_visible', 0)}")
|
|
print(f"Oportunidades abertas: {summary.get('open_opportunities_audited', 0)}")
|
|
print(f"Tasks pendentes: {summary.get('pending_tasks_audited', 0)}")
|
|
print(f"Candidatos abertos: {summary.get('open_reconciliation_candidates', 0)}")
|
|
print(f"Achados confirmados: {summary.get('categories', {}).get('confirmed', 0)}")
|
|
print(f"A rever: {summary.get('categories', {}).get('needs_review', 0)}")
|
|
print(f"Política: {summary.get('categories', {}).get('policy', 0)}")
|
|
|
|
terminal_limit = max(1, int(args.terminal_limit))
|
|
for finding in auditor.findings[:terminal_limit]:
|
|
left = f"{finding.category.upper():12} | {finding.severity.upper():8} | {finding.code:52}"
|
|
action = f"{finding.displayed_action or finding.stored_action or '-'} -> {finding.effective_action or '-'}"
|
|
print(f"{left} | {action:42} | {finding.customer_name or finding.external_ref or '-'}")
|
|
if len(auditor.findings) > terminal_limit:
|
|
print(f"... {len(auditor.findings) - terminal_limit} achado(s) adicionais no relatório.")
|
|
|
|
print("Relatórios:")
|
|
for key, value in paths.items():
|
|
print(f" {key}: {value}")
|
|
|
|
if args.fail_on == "never":
|
|
return 0
|
|
if args.fail_on == "confirmed_high":
|
|
if any(
|
|
finding.category == "confirmed"
|
|
and finding.severity in {"critical", "high"}
|
|
for finding in auditor.findings
|
|
):
|
|
return 2
|
|
return 0
|
|
threshold = args.fail_on
|
|
if any(SEVERITY_ORDER.get(f.severity, 99) <= SEVERITY_ORDER.get(threshold, -1) for f in auditor.findings):
|
|
return 2
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|