Files
clientflow_backend/scripts/apply_v132_operational_coherence.py

755 lines
32 KiB
Python
Executable File

#!/usr/bin/env python3
"""Apply guarded v132 work-center and Odoo operational coherence repairs.
Dry-run is the default. The script:
- persists reconstructed review state;
- corrects Odoo ``assigned`` stages to ORDER_PREPARATION when no shipment or
physical validation exists;
- converts premature sensitive tasks into one reconstructed-review task;
- resolves only stale Odoo reconciliation candidates with one exact, coherent
existing link;
- materializes the first safe action after the transaction commits.
It never changes commercial values, documents, payments or customer mappings.
"""
from __future__ import annotations
import argparse
import json
import re
import sys
import uuid
from dataclasses import asdict, dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Iterable, Mapping
from sqlalchemy import text
PROJECT_ROOT = Path(__file__).resolve().parent.parent
if str(PROJECT_ROOT) not in sys.path:
sys.path.insert(0, str(PROJECT_ROOT))
engine = None
from app.work_center_action_policy import ( # noqa: E402
RECONSTRUCTED_SENSITIVE_ACTIONS,
canonical_action_code,
is_reconstructed_record,
reconstructed_review_metadata_patch,
reconstructed_review_status,
)
OPEN_RECONCILIATION_STATUSES = {"open", "needs_review", "conflict"}
ADVANCED_ASSIGNED_STAGES = {
"ODOO_ORDER_CREATED",
"IN_PRODUCTION",
"ORDER_PREPARATION",
"READY_TO_SHIP",
"SHIPMENT_CREATED",
}
DONE_TASK_STATUSES = {"done", "completed"}
SENSITIVE_CODES = {canonical_action_code(code) for code in RECONSTRUCTED_SENSITIVE_ACTIONS}
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).casefold()
def _json(value: Any) -> str:
return json.dumps(value or {}, ensure_ascii=False, default=str)
def _as_dict(value: Any) -> dict[str, Any]:
return dict(value) if isinstance(value, Mapping) else {}
def _norm_tax(value: Any) -> str:
return "".join(ch for ch in _s(value) if ch.isdigit())
def _norm_name(value: Any) -> str:
raw = _s(value).casefold()
raw = re.sub(r"[^0-9a-zà-ÿ]+", " ", raw)
stop = {"lda", "ltda", "unipessoal", "sa", "s", "a", "sociedade", "limitada"}
return " ".join(token for token in raw.split() if token not in stop)
def _refs(value: Any) -> set[str]:
raw = _s(value)
if not raw:
return set()
compact = re.sub(r"\s+", "", raw).casefold()
result = {raw.casefold(), compact}
match = re.fullmatch(r"s0*(\d+)", compact)
if match:
number = int(match.group(1))
result.update({f"s{number}", f"s{number:05d}"})
return {item for item in result if item}
def _candidate_refs(item: Mapping[str, Any]) -> set[str]:
refs: set[str] = set()
for value in (item.get("external_id"), item.get("document_number")):
refs.update(_refs(value))
payload = _as_dict(item.get("payload"))
record = _as_dict(payload.get("record")) or payload
for key in ("id", "name", "number", "external_id", "document_number"):
refs.update(_refs(record.get(key)))
return refs
def _link_refs(link: Mapping[str, Any]) -> set[str]:
refs: set[str] = set()
for value in (link.get("external_id"), link.get("external_name")):
refs.update(_refs(value))
return refs
def _physical_assigned(link: Mapping[str, Any]) -> bool:
status = _lower(link.get("status"))
payload = _as_dict(link.get("payload"))
if status == "picking_assigned" or _lower(payload.get("physical_status")) == "picking_assigned":
return True
pickings = payload.get("pickings") or payload.get("outgoing_pickings") or []
return any(_lower(item.get("state")) == "assigned" for item in pickings if isinstance(item, Mapping))
def _task_is_reconstructed_review(task: Mapping[str, Any]) -> bool:
if _upper(task.get("action_code")) == "REVIEW_RECONSTRUCTED_PROCESS":
return True
metadata = _as_dict(task.get("metadata"))
return _upper(task.get("action_code")) == "REVIEW_MANUALLY" and _lower(
metadata.get("review_type") or metadata.get("task_type")
) in {"reconstructed_process", "historical_evidence_review"}
@dataclass
class OpportunityPlan:
opportunity_id: str
customer_name: str
title: str
current_stage: str
target_stage: str
physical_assigned: bool
physical_validated: bool
shipment_exists: bool
reconstructed_before: str
reconstructed_after: str
blocked_action: str
task_action: str
task_operation: str
reasons: list[str]
safe: bool
@dataclass
class CandidatePlan:
item_id: str
reference: str
opportunity_id: str
customer_name: str
identity: str
safe: bool
reason: str
def _load_state(conn: Any, *, focus: list[str]) -> tuple[list[dict[str, Any]], list[dict[str, Any]], list[dict[str, Any]], list[dict[str, Any]]]:
params: dict[str, Any] = {}
where = "WHERE o.status = 'open'"
if focus:
params["focus"] = [f"%{item}%" for item in focus]
where += " AND (o.id::text ILIKE ANY(CAST(:focus AS TEXT[])) OR COALESCE(o.customer_name,'') ILIKE ANY(CAST(:focus AS TEXT[])) OR COALESCE(o.title,'') ILIKE ANY(CAST(:focus AS TEXT[])))"
opportunities = [dict(row) for row in conn.execute(text(f"""
SELECT o.id::text, o.title, o.customer_name, o.stage, o.status,
o.last_action_code, COALESCE(o.metadata, '{{}}'::jsonb) AS metadata,
c.name AS fiscal_name, c.tax_id AS fiscal_tax_id
FROM opportunities o
LEFT JOIN customers c ON c.id = o.local_customer_id
{where}
ORDER BY o.updated_at DESC
"""), params).mappings().all()]
opportunity_ids = [row["id"] for row in opportunities]
if not opportunity_ids:
return opportunities, [], [], []
links = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, system, external_type,
external_id, external_name, status, COALESCE(payload, '{}'::jsonb) AS payload
FROM operation_links
WHERE opportunity_id = ANY(CAST(:ids AS UUID[]))
"""), {"ids": opportunity_ids}).mappings().all()]
tasks = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, action_code, route, action, note,
status, due_at, done_at, COALESCE(metadata, '{}'::jsonb) AS metadata
FROM tasks
WHERE opportunity_id = ANY(CAST(:ids AS UUID[]))
ORDER BY created_at
"""), {"ids": opportunity_ids}).mappings().all()]
shipments = [dict(row) for row in conn.execute(text("""
SELECT id::text, opportunity_id::text, system, external_reference,
status, tracking_code
FROM shipments
WHERE opportunity_id = ANY(CAST(:ids AS UUID[]))
"""), {"ids": opportunity_ids}).mappings().all()]
return opportunities, links, tasks, shipments
def _plan_opportunities(
opportunities: list[dict[str, Any]],
links: list[dict[str, Any]],
tasks: list[dict[str, Any]],
shipments: list[dict[str, Any]],
) -> list[OpportunityPlan]:
links_by: dict[str, list[dict[str, Any]]] = {}
tasks_by: dict[str, list[dict[str, Any]]] = {}
shipments_by: dict[str, list[dict[str, Any]]] = {}
for row in links:
links_by.setdefault(_s(row.get("opportunity_id")), []).append(row)
for row in tasks:
tasks_by.setdefault(_s(row.get("opportunity_id")), []).append(row)
for row in shipments:
shipments_by.setdefault(_s(row.get("opportunity_id")), []).append(row)
plans: list[OpportunityPlan] = []
for opportunity in opportunities:
oid = _s(opportunity.get("id"))
metadata = _as_dict(opportunity.get("metadata"))
opp_links = links_by.get(oid, [])
opp_tasks = tasks_by.get(oid, [])
physical_links = [row for row in opp_links if _lower(row.get("system")) == "odoo" and _lower(row.get("external_type")) == "physical_status"]
assigned = any(_physical_assigned(row) for row in physical_links)
validation_links = [row for row in opp_links if _lower(row.get("system")) == "odoo" and _lower(row.get("external_type")) == "physical_validation"]
validated = any(_lower(row.get("status")) in {"validated", "ready_to_ship"} for row in validation_links)
shipment_links = [row for row in opp_links if _lower(row.get("external_type")) in {"shipment", "tracking"} and _lower(row.get("status")) not in {"not_created", "cancelled", "failed"}]
active_shipments = [
row for row in shipments_by.get(oid, [])
if _s(row.get("external_reference"))
or _s(row.get("tracking_code"))
or _lower(row.get("status")) in {"created", "label_created", "in_transit", "shipped", "delivered"}
]
shipment_exists = bool(shipment_links or active_shipments)
pending = [row for row in opp_tasks if _lower(row.get("status")) == "pending"]
completed_reviews = [row for row in opp_tasks if _task_is_reconstructed_review(row) and _lower(row.get("status")) in DONE_TASK_STATUSES]
pending_review = next((row for row in pending if _task_is_reconstructed_review(row)), None)
pending_sensitive = next((row for row in pending if canonical_action_code(row.get("action_code")) in SENSITIVE_CODES), None)
before = reconstructed_review_status(metadata)
after = before
if completed_reviews:
after = "validated"
elif before == "legacy_unset":
after = "required"
blocked_action = canonical_action_code(
(pending_sensitive or {}).get("action_code")
or opportunity.get("last_action_code")
or "VALIDATE_PHYSICAL_ORDER"
)
target_stage = _upper(opportunity.get("stage"))
reasons: list[str] = []
safe = True
if assigned and not validated and not shipment_exists and target_stage in ADVANCED_ASSIGNED_STAGES:
if target_stage != "ORDER_PREPARATION":
reasons.append("picking_assigned_without_validation_requires_order_preparation")
target_stage = "ORDER_PREPARATION"
elif assigned and (validated or shipment_exists):
safe = False
reasons.append("assigned_but_validation_or_shipment_exists")
task_action = ""
task_operation = "none"
if after == "required" and blocked_action in SENSITIVE_CODES:
task_action = "REVIEW_RECONSTRUCTED_PROCESS"
if pending_review:
task_operation = "keep_review"
elif pending_sensitive:
task_operation = "convert_sensitive_to_review"
else:
task_operation = "materialize_review"
elif assigned and not validated and not shipment_exists:
task_action = "VALIDATE_PHYSICAL_ORDER"
if any(_upper(row.get("action_code")) == "VALIDATE_PHYSICAL_ORDER" for row in pending):
task_operation = "keep_validation"
else:
task_operation = "materialize_validation"
if after != before:
reasons.append(f"reconstructed_review:{before}->{after}")
if task_operation != "none":
reasons.append(task_operation)
if not reasons:
continue
plans.append(OpportunityPlan(
opportunity_id=oid,
customer_name=_s(opportunity.get("fiscal_name") or opportunity.get("customer_name")),
title=_s(opportunity.get("title")),
current_stage=_upper(opportunity.get("stage")),
target_stage=target_stage,
physical_assigned=assigned,
physical_validated=validated,
shipment_exists=shipment_exists,
reconstructed_before=before,
reconstructed_after=after,
blocked_action=blocked_action,
task_action=task_action,
task_operation=task_operation,
reasons=reasons,
safe=safe,
))
return plans
def _load_stale_candidates(conn: Any) -> tuple[list[dict[str, Any]], list[dict[str, Any]], dict[str, dict[str, Any]]]:
items = [dict(row) for row in conn.execute(text("""
SELECT id::text, source_system, external_type, external_id, document_number,
opportunity_id::text, customer_name, customer_tax_id, status,
COALESCE(payload, '{}'::jsonb) AS payload
FROM reconciliation_items
WHERE status IN ('open','needs_review','conflict')
AND source_system = 'odoo'
ORDER BY created_at
""")).mappings().all()]
links = [dict(row) for row in conn.execute(text("""
SELECT ol.id::text, ol.opportunity_id::text, ol.system, ol.external_type,
ol.external_id, ol.external_name, ol.status,
o.title, o.customer_name, o.status AS opportunity_status,
c.name AS fiscal_name, c.tax_id AS fiscal_tax_id
FROM operation_links ol
JOIN opportunities o ON o.id = ol.opportunity_id
LEFT JOIN customers c ON c.id = o.local_customer_id
WHERE ol.system = 'odoo'
AND ol.external_type IN ('sale_order','physical_status')
""")).mappings().all()]
opportunities = {
row["id"]: dict(row)
for row in conn.execute(text("""
SELECT o.id::text, o.title, o.customer_name, o.status,
c.name AS fiscal_name, c.tax_id AS fiscal_tax_id
FROM opportunities o
LEFT JOIN customers c ON c.id = o.local_customer_id
""")).mappings().all()
}
return items, links, opportunities
def _plan_candidates(
items: list[dict[str, Any]],
links: list[dict[str, Any]],
opportunities: dict[str, dict[str, Any]],
*,
allow_unknown_identity: bool,
candidate_refs: set[str] | None = None,
allowed_opportunity_ids: set[str] | None = None,
) -> list[CandidatePlan]:
plans: list[CandidatePlan] = []
requested_refs = {ref for value in (candidate_refs or set()) for ref in _refs(value)}
for item in items:
refs = _candidate_refs(item)
if not refs:
continue
if requested_refs and not (refs & requested_refs):
continue
matching = [row for row in links if refs & _link_refs(row)]
opportunity_ids = sorted({_s(row.get("opportunity_id")) for row in matching if _s(row.get("opportunity_id"))})
if len(opportunity_ids) != 1:
continue
oid = opportunity_ids[0]
if allowed_opportunity_ids is not None and oid not in allowed_opportunity_ids:
continue
opportunity = opportunities.get(oid)
if not opportunity:
plans.append(CandidatePlan(
item_id=_s(item.get("id")), reference=_s(item.get("document_number") or item.get("external_id")),
opportunity_id=oid, customer_name=_s(item.get("customer_name")), identity="missing_opportunity",
safe=False, reason="linked_opportunity_not_found",
))
continue
assigned_oid = _s(item.get("opportunity_id"))
if assigned_oid and assigned_oid != oid:
plans.append(CandidatePlan(
item_id=_s(item.get("id")), reference=_s(item.get("document_number") or item.get("external_id")),
opportunity_id=oid, customer_name=_s(item.get("customer_name")), identity="conflict",
safe=False, reason="candidate_points_to_different_opportunity",
))
continue
candidate_tax = _norm_tax(item.get("customer_tax_id"))
opportunity_tax = _norm_tax(opportunity.get("fiscal_tax_id"))
candidate_name = _norm_name(item.get("customer_name"))
opportunity_names = {_norm_name(opportunity.get("fiscal_name")), _norm_name(opportunity.get("customer_name")), _norm_name(opportunity.get("title"))}
opportunity_names.discard("")
if candidate_tax and opportunity_tax and candidate_tax != opportunity_tax:
identity = "nif_mismatch"
safe = False
reason = "nif_conflict"
elif candidate_tax and opportunity_tax and candidate_tax == opportunity_tax:
identity = "nif_match"
safe = True
reason = "exact_link_and_nif_match"
elif candidate_name and any(candidate_name == name or candidate_name in name or name in candidate_name for name in opportunity_names):
identity = "name_match"
safe = True
reason = "exact_link_and_name_match"
else:
identity = "unknown"
safe = bool(allow_unknown_identity)
reason = "exact_link_identity_unknown_allowed" if safe else "identity_unknown_requires_review"
plans.append(CandidatePlan(
item_id=_s(item.get("id")), reference=_s(item.get("document_number") or item.get("external_id")),
opportunity_id=oid, customer_name=_s(item.get("customer_name")), identity=identity,
safe=safe, reason=reason,
))
return plans
def _convert_task_to_review(conn: Any, opportunity_id: str, blocked_action: str, actor: str) -> str | None:
row = conn.execute(text("""
SELECT id::text, action_code, action, note, COALESCE(metadata, '{}'::jsonb) AS metadata
FROM tasks
WHERE opportunity_id = CAST(:opportunity_id AS UUID)
AND status = 'pending'
AND action_code = ANY(CAST(:codes AS TEXT[]))
ORDER BY due_at NULLS LAST, created_at
LIMIT 1
FOR UPDATE
"""), {"opportunity_id": opportunity_id, "codes": sorted(SENSITIVE_CODES)}).mappings().first()
if not row:
return None
task_id = _s(row.get("id"))
conn.execute(text("""
UPDATE tasks
SET action_code = 'REVIEW_RECONSTRUCTED_PROCESS',
route = 'rever',
action = 'Validar processo reconstruído',
note = 'Confirmar cliente, documento principal, valor e evidências antes de executar a ação sensível seguinte.',
priority = 'alta',
source_system = 'clientflow_next_action',
idempotency_key = 'task:v132:reconstructed_review:' || id::text,
metadata = COALESCE(metadata, '{}'::jsonb) || CAST(:metadata AS JSONB),
updated_at = now()
WHERE id = CAST(:task_id AS UUID)
"""), {
"task_id": task_id,
"metadata": _json({
"review_type": "reconstructed_process",
"blocked_action_code": canonical_action_code(blocked_action or row.get("action_code")),
"converted_from_action_code": _upper(row.get("action_code")),
"converted_by": actor,
"converted_version": "v4928.1.5.132",
}),
})
return task_id
def _apply_opportunity_plan(conn: Any, plan: OpportunityPlan, actor: str) -> dict[str, Any]:
if not plan.safe:
return {"opportunity_id": plan.opportunity_id, "applied": False, "reason": "unsafe_plan"}
current = conn.execute(text("""
SELECT stage, COALESCE(metadata, '{}'::jsonb) AS metadata
FROM opportunities
WHERE id = CAST(:opportunity_id AS UUID)
FOR UPDATE
"""), {"opportunity_id": plan.opportunity_id}).mappings().first()
if not current:
return {"opportunity_id": plan.opportunity_id, "applied": False, "reason": "opportunity_not_found"}
patch: dict[str, Any] = {
"v132_operational_coherence_applied": True,
"v132_operational_coherence_actor": actor,
"v132_operational_coherence_reasons": plan.reasons,
}
if plan.reconstructed_after != plan.reconstructed_before:
patch.update(reconstructed_review_metadata_patch(
plan.reconstructed_after,
actor=actor,
reason="Migração v132: estado explícito da revisão reconstruída.",
blocked_action_code=plan.blocked_action,
))
stage_changed = _upper(current.get("stage")) != plan.target_stage
conn.execute(text("""
UPDATE opportunities
SET stage = CAST(:stage AS TEXT),
status = 'open',
last_action_code = CASE WHEN :next_action <> '' THEN CAST(:next_action AS TEXT) ELSE last_action_code END,
metadata = COALESCE(metadata, '{}'::jsonb) || CAST(:metadata AS JSONB),
updated_at = now()
WHERE id = CAST(:opportunity_id AS UUID)
"""), {
"opportunity_id": plan.opportunity_id,
"stage": plan.target_stage,
"next_action": plan.task_action,
"metadata": _json(patch),
})
removed_validation: list[dict[str, Any]] = []
ignored_shipment_tasks: list[dict[str, Any]] = []
assigned_rewind = (
plan.physical_assigned
and not plan.physical_validated
and not plan.shipment_exists
and plan.target_stage == "ORDER_PREPARATION"
)
if assigned_rewind:
# Remove only invalid synthetic evidence generated from ``assigned``.
# Operator-validated and delivery-done links are never deleted.
removed_validation = [dict(row) for row in conn.execute(text("""
DELETE FROM operation_links
WHERE opportunity_id = CAST(:opportunity_id AS UUID)
AND system = 'odoo'
AND external_type = 'physical_validation'
AND status IN ('ready_to_ship','pending')
AND COALESCE(payload->>'delivery_done','false') <> 'true'
AND COALESCE(payload->>'validated_by_operator','false') <> 'true'
RETURNING id::text
"""), {"opportunity_id": plan.opportunity_id}).mappings().all()]
ignored_shipment_tasks = [dict(row) for row in conn.execute(text("""
UPDATE tasks
SET status = 'ignored', done_at = COALESCE(done_at, now()), done_by = COALESCE(done_by, :actor),
metadata = COALESCE(metadata, '{}'::jsonb) || CAST(:metadata AS JSONB),
note = COALESCE(note, '') || E'\n\nIgnorada pela v132: picking assigned ainda exige validação física.',
updated_at = now()
WHERE opportunity_id = CAST(:opportunity_id AS UUID)
AND status = 'pending'
AND action_code = 'CREATE_SHIPMENT'
RETURNING id::text
"""), {
"opportunity_id": plan.opportunity_id,
"actor": actor,
"metadata": _json({"superseded_by": plan.task_action or "VALIDATE_PHYSICAL_ORDER", "superseded_version": "v4928.1.5.132"}),
}).mappings().all()]
converted_task_id = None
if plan.task_operation == "convert_sensitive_to_review":
converted_task_id = _convert_task_to_review(conn, plan.opportunity_id, plan.blocked_action, actor)
conn.execute(text("""
INSERT INTO opportunity_events (
id, opportunity_id, event_type, action_code, from_stage, to_stage,
note, payload, created_by
) VALUES (
CAST(:id AS UUID), CAST(:opportunity_id AS UUID), 'v132_operational_coherence',
CAST(:action_code AS TEXT), CAST(:from_stage AS TEXT), CAST(:to_stage AS TEXT),
CAST(:note AS TEXT), CAST(:payload AS JSONB), CAST(:created_by AS TEXT)
)
"""), {
"id": str(uuid.uuid4()),
"opportunity_id": plan.opportunity_id,
"action_code": plan.task_action,
"from_stage": _upper(current.get("stage")),
"to_stage": plan.target_stage,
"note": "; ".join(plan.reasons),
"payload": _json({
"stage_changed": stage_changed,
"removed_physical_validation_link_ids": [row["id"] for row in removed_validation],
"ignored_create_shipment_task_ids": [row["id"] for row in ignored_shipment_tasks],
"converted_review_task_id": converted_task_id,
}),
"created_by": actor,
})
return {
"opportunity_id": plan.opportunity_id,
"applied": True,
"stage_changed": stage_changed,
"converted_review_task_id": converted_task_id,
}
def _apply_candidate_plan(conn: Any, plan: CandidatePlan, actor: str) -> dict[str, Any]:
if not plan.safe:
return {"item_id": plan.item_id, "applied": False, "reason": plan.reason}
row = conn.execute(text("""
UPDATE reconciliation_items
SET opportunity_id = CAST(:opportunity_id AS UUID),
status = 'linked',
resolved_at = now(),
updated_at = now(),
resolution_note = 'Resolvido pela v132: evidência já ligada exatamente uma vez.',
payload = COALESCE(payload, '{}'::jsonb) || CAST(:payload AS JSONB)
WHERE id = CAST(:item_id AS UUID)
AND status IN ('open','needs_review','conflict')
RETURNING id::text
"""), {
"item_id": plan.item_id,
"opportunity_id": plan.opportunity_id,
"payload": _json({
"resolved_as_existing_link": True,
"resolved_by": actor,
"resolved_version": "v4928.1.5.132",
"identity_guard": plan.identity,
"external_reference": plan.reference,
}),
}).mappings().first()
return {"item_id": plan.item_id, "applied": bool(row), "reason": plan.reason}
def _materialize(affected_ids: Iterable[str]) -> list[dict[str, Any]]:
from app.opportunity_action_task_materializer import ensure_pending_task_for_next_action
from app.opportunity_next_action_service import get_opportunity_next_action
results = []
for oid in sorted(set(affected_ids)):
try:
next_action = get_opportunity_next_action(oid)
result = ensure_pending_task_for_next_action(
oid, next_action, source="v132_operational_coherence", actor="migration_v132"
)
results.append({"opportunity_id": oid, "next_action": next_action, "materialization": result})
except Exception as exc:
results.append({"opportunity_id": oid, "error": str(exc)})
return results
def _write_report(output_dir: Path, payload: dict[str, Any]) -> None:
output_dir.mkdir(parents=True, exist_ok=True)
(output_dir / "v132_operational_coherence.json").write_text(_json(payload), encoding="utf-8")
lines = [
"# v132 — Correção de coerência operacional",
"",
f"Gerado em: `{payload['generated_at']}`",
f"Modo: **{'APPLY' if payload['apply'] else 'DRY-RUN'}**",
"",
"## Oportunidades",
"",
]
for plan in payload["opportunity_plans"]:
lines.append(
f"- **{plan['customer_name'] or plan['title'] or plan['opportunity_id']}** — "
f"{plan['current_stage']}{plan['target_stage']} · revisão "
f"{plan['reconstructed_before']}{plan['reconstructed_after']} · "
f"task `{plan['task_action'] or ''}`/{plan['task_operation']} · safe={plan['safe']}"
)
lines.extend(["", "## Candidatos já ligados", ""])
for plan in payload["candidate_plans"]:
lines.append(
f"- `{plan['reference']}` → `{plan['opportunity_id']}` · "
f"identity={plan['identity']} · safe={plan['safe']} · {plan['reason']}"
)
lines.extend(["", "## Resultado", "", "```json", json.dumps(payload.get("result") or {}, ensure_ascii=False, indent=2, default=str), "```", ""])
(output_dir / "V132_OPERATIONAL_COHERENCE.md").write_text("\n".join(lines), encoding="utf-8")
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description=__doc__)
parser.add_argument("--apply", action="store_true", help="Aplicar alterações; sem esta opção é dry-run.")
parser.add_argument("--focus", action="append", default=[], help="Filtrar oportunidades por UUID/nome/título. Repetível.")
parser.add_argument("--candidate-ref", action="append", default=[], help="Filtrar candidatos Odoo pela referência, por exemplo S00318. Repetível.")
parser.add_argument("--allow-unknown-identity", action="store_true", help="Permitir resolver candidato exato sem NIF/nome confirmável.")
parser.add_argument("--skip-opportunities", action="store_true", help="Não analisar/corrigir oportunidades; útil para tratar apenas candidatos.")
parser.add_argument("--skip-candidates", action="store_true", help="Não analisar/resolver candidatos obsoletos.")
parser.add_argument("--output-dir", default="", help="Diretório do relatório.")
parser.add_argument("--self-test", action="store_true")
return parser.parse_args()
def self_test() -> None:
assert _refs("S00323") == {"s00323", "s323"}
assert _norm_tax("PT 510 177 441") == "510177441"
assert _norm_name("NOLTIA SYSTEM, LDA") == "noltia system"
assert canonical_action_code("SHIP_ORDER") == "CREATE_SHIPMENT"
print("Self-test OK")
def main() -> int:
args = parse_args()
if args.self_test:
self_test()
return 0
global engine
from app.db import engine as db_engine
engine = db_engine
with engine.begin() as conn:
conn.execute(text("SET TRANSACTION READ ONLY"))
if args.skip_opportunities:
opportunities, links, tasks, shipments = [], [], [], []
opportunity_plans: list[OpportunityPlan] = []
else:
opportunities, links, tasks, shipments = _load_state(conn, focus=args.focus)
opportunity_plans = _plan_opportunities(opportunities, links, tasks, shipments)
if args.skip_candidates:
candidate_plans: list[CandidatePlan] = []
else:
items, candidate_links, candidate_opportunities = _load_stale_candidates(conn)
allowed_opportunity_ids = None
if args.focus and not args.candidate_ref:
allowed_opportunity_ids = {_s(row.get("id")) for row in opportunities}
candidate_plans = _plan_candidates(
items, candidate_links, candidate_opportunities,
allow_unknown_identity=args.allow_unknown_identity,
candidate_refs=set(args.candidate_ref),
allowed_opportunity_ids=allowed_opportunity_ids,
)
print(f"Oportunidades avaliadas: {len(opportunities)}")
print(f"Planos de oportunidade: {len(opportunity_plans)}")
for plan in opportunity_plans:
print(
f"{'SAFE' if plan.safe else 'SKIP':5} | {plan.customer_name or plan.title:48.48} | "
f"{plan.current_stage:18} -> {plan.target_stage:18} | "
f"review={plan.reconstructed_before}->{plan.reconstructed_after} | {plan.task_operation}"
)
print(f"Candidatos exatos encontrados: {len(candidate_plans)}")
for plan in candidate_plans:
print(
f"{'SAFE' if plan.safe else 'REVIEW':6} | {plan.reference:12} | "
f"{plan.identity:12} | {plan.reason}"
)
result: dict[str, Any] = {"applied": False}
affected_ids: list[str] = []
if args.apply:
with engine.begin() as conn:
opportunity_results = []
for plan in opportunity_plans:
outcome = _apply_opportunity_plan(conn, plan, "migration_v132")
opportunity_results.append(outcome)
if outcome.get("applied"):
affected_ids.append(plan.opportunity_id)
candidate_results = [
_apply_candidate_plan(conn, plan, "migration_v132")
for plan in candidate_plans
]
materialization = _materialize(affected_ids)
result = {
"applied": True,
"opportunities": opportunity_results,
"candidates": candidate_results,
"materialization": materialization,
}
print(f"Oportunidades corrigidas: {sum(1 for row in opportunity_results if row.get('applied'))}")
print(f"Candidatos resolvidos: {sum(1 for row in candidate_results if row.get('applied'))}")
timestamp = datetime.now(timezone.utc).strftime("%Y%m%d_%H%M%S")
output_dir = Path(args.output_dir) if args.output_dir else PROJECT_ROOT / "audit_reports" / f"v132_operational_coherence_{timestamp}"
payload = {
"generated_at": datetime.now(timezone.utc).isoformat(),
"apply": bool(args.apply),
"focus": args.focus,
"candidate_refs": args.candidate_ref,
"skip_opportunities": bool(args.skip_opportunities),
"allow_unknown_identity": bool(args.allow_unknown_identity),
"opportunity_plans": [asdict(plan) for plan in opportunity_plans],
"candidate_plans": [asdict(plan) for plan in candidate_plans],
"result": result,
}
_write_report(output_dir, payload)
print(f"Relatório: {output_dir}")
if not args.apply:
print("Dry-run only. Use --apply after reviewing SAFE plans.")
return 0
if __name__ == "__main__":
raise SystemExit(main())