diff --git a/app/external_reconciliation_sync.py b/app/external_reconciliation_sync.py index ac03331..6401eae 100644 --- a/app/external_reconciliation_sync.py +++ b/app/external_reconciliation_sync.py @@ -221,7 +221,21 @@ def _looks_like_company_name(value: Any) -> bool: def _external_customer_key(record: Dict[str, Any], *, source_system: str) -> str: if source_system == "jasmin": - return _clean(_first(record, "partyKey", "customerPartyKey", "naturalKey", "key", "id")) + # A document naturalKey identifies the commercial document + # (for example ORC.ORC2026.136), not the customer. Prefer the + # customer party code exposed by Jasmin and never fall back to the + # document naturalKey when seeding/updating a fiscal customer. + return _clean( + _first( + record, + "partyKey", + "customerPartyKey", + "buyerCustomerParty", + "accountingParty", + "buyerCustomerPartyKey", + "accountingPartyKey", + ) + ) if source_system == "odoo": return _clean(_first(record, "partner_external_id", "id")) return _clean(_first(record, "id", "key", "externalId")) @@ -355,7 +369,36 @@ def _jasmin_external_type(record: Dict[str, Any], default_type: str) -> str: def _jasmin_amount(record: Dict[str, Any]) -> Optional[str]: - return _decimal_or_none(_first(record, "payableAmount", "totalAmount", "total", "grossAmount", "amount")) + # Recent Jasmin payloads expose both flattened numeric fields and nested + # money objects. Prefer the payable total including tax. + direct = _first( + record, + "payableAmountAmount", + "totalAmount", + "grossValueAmount", + "taxExclusiveAmountAmount", + "total", + "grossAmount", + "amount", + ) + if direct not in (None, ""): + parsed = _decimal_or_none(direct) + if parsed is not None: + return parsed + + for key in ("payableAmount", "grossValue", "taxExclusiveAmount"): + money = record.get(key) + if isinstance(money, dict): + parsed = _decimal_or_none( + _first(money, "amount", "baseAmount", "reportingAmount") + ) + if parsed is not None: + return parsed + elif money not in (None, ""): + parsed = _decimal_or_none(money) + if parsed is not None: + return parsed + return None def _jasmin_candidate_from_record(record: Dict[str, Any], *, default_type: str) -> Optional[Dict[str, Any]]: @@ -699,7 +742,6 @@ def _existing_odoo_sale_links(external_id: Any, order_name: Any) -> List[Dict[st JOIN opportunities o ON o.id = ol.opportunity_id WHERE ol.system = 'odoo' AND ol.external_type = 'sale_order' - AND o.status = 'open' AND ( (NULLIF(:external_id, '') IS NOT NULL AND ol.external_id = :external_id) OR (NULLIF(:order_name, '') IS NOT NULL AND UPPER(COALESCE(ol.external_name, '')) = UPPER(:order_name)) diff --git a/scripts/audit_fix_odoo_reconciliation_links.py b/scripts/audit_fix_odoo_reconciliation_links.py new file mode 100755 index 0000000..1ea8190 --- /dev/null +++ b/scripts/audit_fix_odoo_reconciliation_links.py @@ -0,0 +1,418 @@ +#!/usr/bin/env python3 +"""Audit and repair stale Odoo reconciliation candidates. + +Problem addressed +----------------- +A reconciliation item may remain ``open`` even though an exact Odoo sale-order +link already exists in ``operation_links``. The current synchronizer can ignore +links whose opportunity is closed because ``_existing_odoo_sale_links`` filters +with ``o.status = 'open'``. + +Default behaviour is read-only. Use ``--apply`` to: + 1. remove that exact source-code filter, creating a timestamped backup; and + 2. mark unambiguous stale reconciliation items as ``linked``. + +Ambiguous cases with more than one linked opportunity are never changed. +Documents supplied through ``--exclude`` are also never changed. + +Run from the ClientFlow repository root, for example: + + PYTHONPATH="$PWD" .venv/bin/python scripts/audit_fix_odoo_reconciliation_links.py + PYTHONPATH="$PWD" .venv/bin/python scripts/audit_fix_odoo_reconciliation_links.py --apply + +Exit codes: + 0: audit/apply completed + 2: configuration/source validation error + 3: database operation error +""" + +from __future__ import annotations + +import argparse +import json +import shutil +import sys +from collections import Counter +from datetime import datetime, timezone +from pathlib import Path +from typing import Any, Iterable + +from sqlalchemy import text + +from app.db import engine + +SOURCE_FILE = Path("app/external_reconciliation_sync.py") +BUGGY_FILTER = " AND o.status = 'open'\n" + +AUDIT_SQL = text( + r""" + WITH candidate_links AS ( + SELECT + ri.id AS reconciliation_item_id, + ri.document_number, + ri.external_id AS reconciliation_external_id, + ri.customer_name, + ri.customer_tax_id, + ri.status AS reconciliation_status, + ri.opportunity_id AS current_opportunity_id, + ri.resolution_note, + ri.resolved_at, + ol.opportunity_id, + ol.external_id AS link_external_id, + ol.external_name, + o.title AS opportunity_title, + o.status AS opportunity_status, + o.stage AS opportunity_stage, + o.closed_at + FROM reconciliation_items ri + JOIN operation_links ol + ON ol.system = 'odoo' + AND ol.external_type = 'sale_order' + AND ( + ( + NULLIF(BTRIM(COALESCE(ri.external_id, '')), '') IS NOT NULL + AND ol.external_id = ri.external_id + ) + OR ( + NULLIF(BTRIM(COALESCE(ri.document_number, '')), '') IS NOT NULL + AND UPPER(BTRIM(COALESCE(ol.external_name, ''))) + = UPPER(BTRIM(ri.document_number)) + ) + ) + JOIN opportunities o ON o.id = ol.opportunity_id + WHERE ri.source_system = 'odoo' + AND ri.external_type = 'odoo_sale_order' + AND ri.status IN ('open', 'needs_review', 'conflict') + ), grouped AS ( + SELECT + reconciliation_item_id, + document_number, + reconciliation_external_id, + customer_name, + customer_tax_id, + reconciliation_status, + current_opportunity_id, + resolution_note, + resolved_at, + COUNT(DISTINCT opportunity_id) AS opportunity_count, + MIN(opportunity_id::text) AS single_opportunity_id + FROM candidate_links + GROUP BY + reconciliation_item_id, + document_number, + reconciliation_external_id, + customer_name, + customer_tax_id, + reconciliation_status, + current_opportunity_id, + resolution_note, + resolved_at + ) + SELECT + g.reconciliation_item_id::text, + g.document_number, + g.reconciliation_external_id, + g.customer_name, + g.customer_tax_id, + g.reconciliation_status, + g.current_opportunity_id::text, + g.opportunity_count, + CASE WHEN g.opportunity_count = 1 THEN g.single_opportunity_id ELSE NULL END + AS opportunity_id, + CASE WHEN g.opportunity_count = 1 THEN o.title ELSE NULL END + AS opportunity_title, + CASE WHEN g.opportunity_count = 1 THEN o.status ELSE NULL END + AS opportunity_status, + CASE WHEN g.opportunity_count = 1 THEN o.stage ELSE NULL END + AS opportunity_stage, + CASE WHEN g.opportunity_count = 1 THEN o.closed_at ELSE NULL END + AS closed_at, + g.resolution_note, + g.resolved_at + FROM grouped g + LEFT JOIN opportunities o + ON g.opportunity_count = 1 + AND o.id = CAST(g.single_opportunity_id AS UUID) + ORDER BY + CASE WHEN g.opportunity_count = 1 THEN 0 ELSE 1 END, + g.document_number + """ +) + +UPDATE_ONE_SQL = text( + r""" + UPDATE reconciliation_items + SET opportunity_id = CAST(:opportunity_id AS UUID), + status = 'linked', + resolution_note = CASE + WHEN COALESCE(BTRIM(resolution_note), '') = '' + THEN 'Resolvido por auditoria: ligação Odoo exata já existente em operation_links.' + WHEN POSITION('Resolvido por auditoria: ligação Odoo exata' IN resolution_note) > 0 + THEN resolution_note + ELSE resolution_note || E'\nResolvido por auditoria: ligação Odoo exata já existente em operation_links.' + END, + resolved_at = COALESCE(resolved_at, now()), + updated_at = now(), + payload = COALESCE(payload, '{}'::jsonb) || jsonb_build_object( + 'resolved_as_existing_operation_link', TRUE, + 'resolved_by', 'audit_fix_odoo_reconciliation_links', + 'resolved_at_audit', now() + ) + WHERE id = CAST(:reconciliation_item_id AS UUID) + AND status IN ('open', 'needs_review', 'conflict') + RETURNING + id::text AS reconciliation_item_id, + document_number, + opportunity_id::text, + status, + resolved_at + """ +) + + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser( + description="Audita e corrige candidatos Odoo já ligados a oportunidades." + ) + parser.add_argument( + "--apply", + action="store_true", + help="Aplica a correção no código e na base de dados. Sem esta opção é dry-run.", + ) + parser.add_argument( + "--data-only", + action="store_true", + help="Com --apply, corrige apenas a base de dados, sem alterar o código.", + ) + parser.add_argument( + "--code-only", + action="store_true", + help="Com --apply, corrige apenas o código, sem alterar a base de dados.", + ) + parser.add_argument( + "--exclude", + action="append", + default=[], + metavar="DOCUMENTO", + help="Não altera este documento. Pode repetir, por exemplo --exclude S00330.", + ) + parser.add_argument( + "--json", + action="store_true", + help="Emite o relatório de auditoria em JSON.", + ) + parser.add_argument( + "--source-file", + default=str(SOURCE_FILE), + help="Caminho do ficheiro external_reconciliation_sync.py.", + ) + return parser.parse_args() + + +def audit_code(path: Path) -> dict[str, Any]: + if not path.exists(): + return { + "path": str(path), + "exists": False, + "buggy_filter_count": None, + "needs_fix": None, + } + source = path.read_text(encoding="utf-8") + count = source.count(BUGGY_FILTER) + return { + "path": str(path), + "exists": True, + "buggy_filter_count": count, + "needs_fix": count > 0, + } + + +def apply_code_fix(path: Path) -> dict[str, Any]: + if not path.exists(): + raise RuntimeError(f"Ficheiro não encontrado: {path}") + + source = path.read_text(encoding="utf-8") + count = source.count(BUGGY_FILTER) + if count == 0: + return {"changed": False, "reason": "Filtro já não existe."} + if count != 1: + raise RuntimeError( + f"Esperado exatamente 1 filtro {BUGGY_FILTER!r}; encontrados {count}." + ) + + timestamp = datetime.now(timezone.utc).strftime("%Y%m%dT%H%M%SZ") + backup = path.with_suffix(path.suffix + f".bak.{timestamp}") + shutil.copy2(path, backup) + + updated = source.replace(BUGGY_FILTER, "", 1) + path.write_text(updated, encoding="utf-8") + return {"changed": True, "backup": str(backup)} + + +def audit_database() -> list[dict[str, Any]]: + with engine.connect() as conn: + rows = conn.execute(AUDIT_SQL).mappings().all() + return [dict(row) for row in rows] + + +def apply_database_fix(excluded: Iterable[str]) -> list[dict[str, Any]]: + excluded_set = {str(value).strip() for value in excluded if str(value).strip()} + audit_rows = audit_database() + eligible = [ + row for row in audit_rows + if int(row.get("opportunity_count") or 0) == 1 + and row.get("opportunity_id") + and row.get("document_number") not in excluded_set + ] + + changed: list[dict[str, Any]] = [] + with engine.begin() as conn: + # Prevent two operators/jobs from applying the same repair concurrently. + conn.execute(text("SELECT pg_advisory_xact_lock(hashtext(:lock_name))"), { + "lock_name": "clientflow.audit_fix_odoo_reconciliation_links", + }) + for row in eligible: + updated = conn.execute(UPDATE_ONE_SQL, { + "reconciliation_item_id": row["reconciliation_item_id"], + "opportunity_id": row["opportunity_id"], + }).mappings().first() + if updated: + changed.append(dict(updated)) + return changed + + +def serialize(value: Any) -> Any: + if isinstance(value, datetime): + return value.isoformat() + return value + + +def print_audit(rows: list[dict[str, Any]], excluded: set[str]) -> None: + counts = Counter() + for row in rows: + if int(row["opportunity_count"] or 0) == 1: + counts[f"single:{row.get('opportunity_status') or 'unknown'}"] += 1 + else: + counts["ambiguous"] += 1 + if row.get("document_number") in excluded: + counts["excluded"] += 1 + + print("\nAUDITORIA — candidatos Odoo com ligação existente") + print("=" * 72) + print(f"Total encontrado: {len(rows)}") + print(f"Ligações únicas / oportunidade fechada: {counts['single:closed']}") + print(f"Ligações únicas / oportunidade aberta: {counts['single:open']}") + other_single = sum( + count for key, count in counts.items() + if key.startswith("single:") and key not in {"single:closed", "single:open"} + ) + print(f"Ligações únicas / outros estados: {other_single}") + print(f"Ambíguos (mais de uma oportunidade): {counts['ambiguous']}") + print(f"Excluídos por opção: {counts['excluded']}") + + if not rows: + print("Nenhuma inconsistência encontrada.") + return + + print("\nDetalhe:") + for row in rows: + document = row.get("document_number") or "(sem número)" + count = int(row.get("opportunity_count") or 0) + excluded_marker = " [EXCLUÍDO]" if document in excluded else "" + if count == 1: + print( + f"- {document}{excluded_marker}: {row.get('reconciliation_status')} -> " + f"{row.get('opportunity_id')} | " + f"{row.get('opportunity_status')}/{row.get('opportunity_stage')} | " + f"{row.get('opportunity_title')}" + ) + else: + print(f"- {document}{excluded_marker}: AMBÍGUO ({count} oportunidades)") + + +def main() -> int: + args = parse_args() + if args.data_only and args.code_only: + print("ERRO: --data-only e --code-only não podem ser usados em conjunto.", file=sys.stderr) + return 2 + + source_path = Path(args.source_file) + excluded = {value.strip() for value in args.exclude if value.strip()} + + code_report = audit_code(source_path) + try: + rows_before = audit_database() + except Exception as exc: + print(f"ERRO ao auditar a base de dados: {exc}", file=sys.stderr) + return 3 + + if args.json: + report = { + "mode": "apply" if args.apply else "dry-run", + "code": code_report, + "excluded": sorted(excluded), + "database": [{k: serialize(v) for k, v in row.items()} for row in rows_before], + } + print(json.dumps(report, ensure_ascii=False, indent=2)) + else: + print("AUDITORIA DO CÓDIGO") + print("=" * 72) + print(f"Ficheiro: {code_report['path']}") + print(f"Existe: {code_report['exists']}") + print(f"Filtro incorreto encontrado: {code_report['buggy_filter_count']}") + print_audit(rows_before, excluded) + + if not args.apply: + if not args.json: + print("\nDRY-RUN: nenhuma alteração aplicada.") + print("Use --apply para corrigir código e dados.") + return 0 + + code_result: dict[str, Any] | None = None + changed_rows: list[dict[str, Any]] = [] + + try: + if not args.data_only: + code_result = apply_code_fix(source_path) + if not args.code_only: + changed_rows = apply_database_fix(excluded) + except Exception as exc: + print(f"ERRO durante a aplicação: {exc}", file=sys.stderr) + return 3 + + try: + rows_after = audit_database() + except Exception as exc: + print(f"ERRO na auditoria posterior: {exc}", file=sys.stderr) + return 3 + + if args.json: + result = { + "code_result": code_result, + "database_rows_changed": [ + {k: serialize(v) for k, v in row.items()} for row in changed_rows + ], + "remaining_inconsistencies": [ + {k: serialize(v) for k, v in row.items()} for row in rows_after + ], + } + print(json.dumps(result, ensure_ascii=False, indent=2)) + else: + print("\nAPLICAÇÃO") + print("=" * 72) + if code_result is not None: + print(f"Código: {code_result}") + print(f"Itens corrigidos na base de dados: {len(changed_rows)}") + for row in changed_rows: + print( + f"- {row['document_number']} -> {row['opportunity_id']} " + f"({row['status']})" + ) + print_audit(rows_after, excluded) + + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/repair_jasmin_reconciliation_item.py b/scripts/repair_jasmin_reconciliation_item.py new file mode 100644 index 0000000..c553e68 --- /dev/null +++ b/scripts/repair_jasmin_reconciliation_item.py @@ -0,0 +1,157 @@ +#!/usr/bin/env python3 +"""Repair a Jasmin reconciliation item affected by incorrect customer mapping. + +Dry-run is the default. Use --apply only after reviewing the printed plan. +This script does not create/link opportunities and does not resolve the item. +""" +from __future__ import annotations + +import argparse +import json +import sys +from decimal import Decimal, InvalidOperation +from pathlib import Path + +from sqlalchemy import text + +ROOT = Path(__file__).resolve().parents[1] +if str(ROOT) not in sys.path: + sys.path.insert(0, str(ROOT)) + +from app.db import engine + + +def _clean(value: object) -> str: + return str(value or "").strip() + + +def _money(record: dict) -> Decimal | None: + values = [ + record.get("payableAmountAmount"), + (record.get("payableAmount") or {}).get("amount") + if isinstance(record.get("payableAmount"), dict) + else record.get("payableAmount"), + record.get("grossValueAmount"), + ] + for value in values: + if value in (None, ""): + continue + try: + return Decimal(str(value)).quantize(Decimal("0.01")) + except (InvalidOperation, ValueError): + continue + return None + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("document_number", help="Ex.: ORC.ORC2026.136") + parser.add_argument("--apply", action="store_true", help="Aplicar a reparação; sem esta flag faz dry-run") + args = parser.parse_args() + + with engine.begin() as conn: + item = conn.execute( + text(""" + SELECT ri.id::text, ri.status, ri.opportunity_id::text, + ri.customer_id::text, ri.customer_tax_id, + ri.document_number, ri.amount, ri.payload, + c.name AS customer_name, c.tax_id, + c.jasmin_customer_party_key + FROM reconciliation_items ri + LEFT JOIN customers c ON c.id = ri.customer_id + WHERE ri.document_number = :document_number + FOR UPDATE OF ri + """), + {"document_number": args.document_number}, + ).mappings().first() + + if not item: + raise SystemExit(f"Documento não encontrado: {args.document_number}") + if not item["customer_id"]: + raise SystemExit("O item não tem customer_id; reparação automática recusada") + + payload = item["payload"] or {} + record = payload.get("record") if isinstance(payload, dict) else None + if not isinstance(record, dict): + raise SystemExit("payload.record não existe ou não é um objeto") + + party_key = _clean(record.get("buyerCustomerParty") or record.get("accountingParty")) + payload_tax_id = _clean(record.get("buyerCustomerPartyTaxId") or record.get("accountingPartyTaxId")) + amount = _money(record) + + if not party_key: + raise SystemExit("Não foi possível obter buyerCustomerParty/accountingParty") + if payload_tax_id and _clean(item["tax_id"]) and payload_tax_id != _clean(item["tax_id"]): + raise SystemExit( + f"NIF divergente: cliente={item['tax_id']} payload={payload_tax_id}; reparação recusada" + ) + + conflict = conn.execute( + text(""" + SELECT id::text, name, tax_id + FROM customers + WHERE jasmin_customer_party_key = :party_key + AND id <> CAST(:customer_id AS UUID) + LIMIT 1 + """), + {"party_key": party_key, "customer_id": item["customer_id"]}, + ).mappings().first() + if conflict: + raise SystemExit( + "Party key já pertence a outro cliente: " + + json.dumps(dict(conflict), ensure_ascii=False, default=str) + ) + + plan = { + "mode": "apply" if args.apply else "dry-run", + "document_number": item["document_number"], + "reconciliation_item_id": item["id"], + "customer_id": item["customer_id"], + "customer_name": item["customer_name"], + "tax_id": item["tax_id"], + "party_key_before": item["jasmin_customer_party_key"], + "party_key_after": party_key, + "amount_before": str(item["amount"]) if item["amount"] is not None else None, + "amount_after": str(amount) if amount is not None else None, + "status_unchanged": item["status"], + "opportunity_id_unchanged": item["opportunity_id"], + } + print(json.dumps(plan, ensure_ascii=False, indent=2, default=str)) + + if not args.apply: + conn.rollback() + print("DRY-RUN: nenhuma alteração aplicada.") + return 0 + + conn.execute( + text(""" + UPDATE customers + SET jasmin_customer_party_key = :party_key, + metadata = jsonb_set( + COALESCE(metadata, '{}'::jsonb), + '{external_customer_key}', + to_jsonb(CAST(:party_key AS text)), + true + ), + updated_at = now() + WHERE id = CAST(:customer_id AS UUID) + """), + {"party_key": party_key, "customer_id": item["customer_id"]}, + ) + if amount is not None: + conn.execute( + text(""" + UPDATE reconciliation_items + SET amount = :amount, + updated_at = now() + WHERE id = CAST(:item_id AS UUID) + """), + {"amount": amount, "item_id": item["id"]}, + ) + + print("Reparação aplicada. O item permanece aberto e sem opportunity_id.") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/test_external_reconciliation_jasmin_mapping.py b/tests/test_external_reconciliation_jasmin_mapping.py new file mode 100644 index 0000000..4a1eaf9 --- /dev/null +++ b/tests/test_external_reconciliation_jasmin_mapping.py @@ -0,0 +1,38 @@ +from app.external_reconciliation_sync import _external_customer_key, _jasmin_amount + + +def test_jasmin_customer_key_prefers_party_code_over_document_natural_key(): + record = { + "naturalKey": "ORC.ORC2026.136", + "buyerCustomerParty": "0569", + "accountingParty": "0569", + } + + assert _external_customer_key(record, source_system="jasmin") == "0569" + + +def test_jasmin_customer_key_does_not_use_document_natural_key(): + record = {"naturalKey": "ORC.ORC2026.136", "id": "document-uuid"} + + assert _external_customer_key(record, source_system="jasmin") == "" + + +def test_jasmin_amount_reads_flattened_payable_total(): + record = { + "payableAmountAmount": 441.57, + "grossValueAmount": 359.00, + } + + assert _jasmin_amount(record) == "441.57" + + +def test_jasmin_amount_reads_nested_money_object(): + record = { + "payableAmount": { + "amount": 441.57, + "baseAmount": 441.57, + "reportingAmount": 441.57, + } + } + + assert _jasmin_amount(record) == "441.57"