From 213ba6b41b5357a1aabf503189d09f5ff1bc3bf5 Mon Sep 17 00:00:00 2001 From: plx Date: Sun, 16 Aug 2026 01:38:39 +0000 Subject: [PATCH] fix: harden Flow v2 authoritative cutover semantics --- app/authoritative_operational_adapter.py | 16 +-- app/canonical_operations.py | 1 + app/operations_service.py | 51 +++++++++- app/opportunity_next_action_service.py | 25 ++++- .../simulate_blif_flow_v2_authoritative.py | 99 +++++++++++++++---- .../test_authoritative_operational_adapter.py | 16 ++- ...t_blif_flow_v2_authoritative_simulation.py | 33 ++++++- 7 files changed, 211 insertions(+), 30 deletions(-) diff --git a/app/authoritative_operational_adapter.py b/app/authoritative_operational_adapter.py index a0b2b18..f3438ae 100644 --- a/app/authoritative_operational_adapter.py +++ b/app/authoritative_operational_adapter.py @@ -48,7 +48,8 @@ class AuthoritativeOperationalDecision: "reason_if_blocked": self.reason_code if not self.eligible else None, "operational_queue": self.queue, "authoritative_v2": True, - "suppress_current_card": self.queue == "not_current", + "is_duplicate_representation": self.reason_code == "DUPLICATE_SUPPRESSED", + "suppress_current_card": self.reason_code == "DUPLICATE_SUPPRESSED", }) return result @@ -78,7 +79,8 @@ def _active_obligations(obligations: Iterable[Mapping[str, Any]]) -> list[Mappin def _ref(row: Mapping[str, Any]) -> dict[str, Any]: return {"source": "task", "id": str(row.get("id") or ""), "status": "pending", - "action_code": _code(row.get("action_code"))} + "action_code": _code(row.get("action_code")), + "source_system": str(row.get("source_system") or "")} def decide_authoritative_operation( @@ -110,6 +112,10 @@ def decide_authoritative_operation( oid, canonical, state, business_action, hard_blocker, "blocked", True, "HARD_FACTUAL_BLOCKER", "Um bloqueio factual ou de sistema impede a ação atual.", hard_blocker, confidence=confidence, ) + # Terminal factual state cannot be reopened by a leftover operational task. + if state in TERMINAL_STATES: + return AuthoritativeOperationalDecision(oid, canonical, state, business_action, None, "not_current", False, + "TERMINAL_BUSINESS_STATE", "O processo factual está concluído.", confidence=confidence) active = _active_obligations(obligations) by_code: dict[str, list[Mapping[str, Any]]] = {} @@ -143,7 +149,8 @@ def decide_authoritative_operation( picked = choose(("REVIEW_MANUALLY", "REVIEW_REQUIRED", "REVIEW", "REVIEW_RECONSTRUCTED_PROCESS"), "review", "ACTIVE_REVIEW_OBLIGATION") picked = picked or choose(("SUPPORT",), "do_now", "ACTIVE_SUPPORT_OBLIGATION") - picked = picked or choose(("SEND_INFO",), "do_now", "ACTIVE_RESPONSE_OBLIGATION") + if state == "INQUIRY" or business_action == "SEND_INFO": + picked = picked or choose(("SEND_INFO",), "do_now", "ACTIVE_RESPONSE_OBLIGATION") picked = picked or choose(("CALL_CUSTOMER",), "do_now", "EXPLICIT_CALL_OBLIGATION") if picked: return picked @@ -168,9 +175,6 @@ def decide_authoritative_operation( obligation_source_refs=tuple(_ref(row) for row in rows), confidence=confidence, ) - if state in TERMINAL_STATES: - return AuthoritativeOperationalDecision(oid, canonical, state, business_action, None, "not_current", False, - "TERMINAL_BUSINESS_STATE", "O processo factual está concluído.", confidence=confidence) if business_action: queue = "review" if business_action in REVIEW_ACTIONS else "do_now" return AuthoritativeOperationalDecision(oid, canonical, state, business_action, business_action, queue, True, diff --git a/app/canonical_operations.py b/app/canonical_operations.py index c09f655..d2fe99b 100644 --- a/app/canonical_operations.py +++ b/app/canonical_operations.py @@ -408,6 +408,7 @@ def _opportunity_item( "business_next_action": decision.get("business_next_action"), "effective_action": decision.get("effective_action"), "canonical_opportunity_id": decision.get("canonical_opportunity_id") or opportunity_id, + "material_process_key": decision.get("material_process_key"), }) return item, independent_exceptions diff --git a/app/operations_service.py b/app/operations_service.py index d496750..70d1762 100644 --- a/app/operations_service.py +++ b/app/operations_service.py @@ -38,6 +38,53 @@ def _int(value: Any) -> int: return 0 +def _authoritative_material_identity(item: Dict[str, Any]) -> str: + return str(item.get("material_process_key") or item.get("canonical_opportunity_id") + or item.get("process_key") or item.get("work_item_key") or item.get("id")) + + +def _collapse_authoritative_material_processes(items: List[Dict[str, Any]]) -> List[Dict[str, Any]]: + """Enforce one current card per material process without deleting evidence.""" + current_queues = {"do_now", "review", "blocked", "exception"} + action_rank = { + "REVIEW_MANUALLY": 10, "REVIEW_RECONSTRUCTED_PROCESS": 11, "REVIEW_REQUIRED": 12, + "SUPPORT": 20, "SEND_INFO": 30, "CALL_CUSTOMER": 40, + "FOLLOW_UP_CUSTOMER_REVIEW": 50, "FOLLOW_UP_PAYMENT": 51, + } + groups: Dict[str, List[Dict[str, Any]]] = {} + noncurrent: List[Dict[str, Any]] = [] + for item in items: + if item.get("operational_queue") not in current_queues: + noncurrent.append(item) + continue + groups.setdefault(_authoritative_material_identity(item), []).append(item) + selected: List[Dict[str, Any]] = [] + for identity, group in groups.items(): + winner = min(group, key=lambda row: ( + 0 if row.get("operational_queue") in {"blocked", "exception"} else 1, + action_rank.get(str(row.get("action_code") or ""), 100), + str(row.get("work_item_key") or ""), + )) + merged = dict(winner) + if len(group) > 1: + merged["suppressed_competing_actions"] = [ + {"work_item_key": row.get("work_item_key"), "action_code": row.get("action_code"), + "source_refs": row.get("source_refs") or []} + for row in group if row is not winner + ] + refs, seen = list(merged.get("source_refs") or []), set() + for ref in refs: + seen.add((str(ref.get("source")), str(ref.get("id")))) + for row in group: + for ref in row.get("source_refs") or []: + key = (str(ref.get("source")), str(ref.get("id"))) + if key not in seen: + refs.append(ref); seen.add(key) + merged["source_refs"] = refs + selected.append(merged) + return selected + noncurrent + + def _is_manually_resolved_error(value: Any) -> bool: text_value = str(value or "").strip().casefold() return "limpo manualmente" in text_value or "resolvido manualmente" in text_value @@ -722,6 +769,8 @@ def get_operations_summary( ) if opportunity_ids else {}) projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows]) canonical_items = _attach_operation_urls(list(projection["items"])) + if active_flow_mode == "authoritative": + canonical_items = _collapse_authoritative_material_processes(canonical_items) partition = partition_canonical_items(canonical_items, display_limit=display_limit) actionable_items = partition["actionable_items"] all_waiting_items = partition["waiting_items"] @@ -750,7 +799,7 @@ def get_operations_summary( "backlog_total": partition["backlog_total"], "not_current_total": partition["not_current_total"], } - return { + result = { "counts": cleaned_counts, "recent_outbox": [dict(r) for r in recent_outbox], "recent_documents": [dict(r) for r in recent_documents], diff --git a/app/opportunity_next_action_service.py b/app/opportunity_next_action_service.py index 3a6f39b..4ad3f07 100644 --- a/app/opportunity_next_action_service.py +++ b/app/opportunity_next_action_service.py @@ -481,8 +481,8 @@ def _get_authoritative_opportunity_decisions(ids: list[str], *, connection: Any FROM opportunity_flow_state_v2 WHERE opportunity_id::text IN :opportunity_ids """, ids) tasks = _bulk_rows(conn, """ - SELECT id::text, opportunity_id::text, action_code, status, due_at, - resolved_at, superseded_by_task_id::text, created_at, metadata + SELECT id::text, opportunity_id::text, action_code, status, due_at, source_system, + resolved_at, resolution_code, superseded_by_task_id::text, created_at, metadata FROM tasks WHERE opportunity_id::text IN :opportunity_ids AND status = 'pending' """, ids) customers = _bulk_rows(conn, """ @@ -514,10 +514,31 @@ def _get_authoritative_opportunity_decisions(ids: list[str], *, connection: Any reconciliation_blocking=oid in reconciliation_ids, ) value = decision.to_dict() + winning_ids = {str(ref.get("id") or "") for ref in value.get("obligation_source_refs") or []} + obligation_audit = [] + for task in tasks_by_id.get(oid, []): + task_id = str(task.get("id") or "") + if task_id in winning_ids: + classification = "VALID_ACTIVE_OBLIGATION" + elif task.get("superseded_by_task_id"): + classification = "SUPERSEDED_OBLIGATION" + elif task.get("resolved_at") or task.get("resolution_code"): + classification = "SATISFIED_OBLIGATION" + elif ((task_code := str(task.get("action_code") or "").upper()).startswith("SEND_") + and task_code != "SEND_INFO") or task_code == "CREATE_JASMIN_QUOTE": + classification = "PREMATURE_OR_STALE_OBLIGATION" + else: + classification = "STALE_LEGACY_OBLIGATION" + obligation_audit.append({ + "id": task_id, "action_code": task.get("action_code"), "status": task.get("status"), + "source_system": task.get("source_system"), "classification": classification, + }) + value["obligation_audit_refs"] = obligation_audit if projection is None: value["opportunity_id"] = oid value["target_url"] = f"/opportunities/{oid}" else: + value["material_process_key"] = projection.get("material_process_key") value["target_url"] = f"/opportunities/{decision.canonical_opportunity_id or oid}" decisions[oid] = value return decisions diff --git a/scripts/simulate_blif_flow_v2_authoritative.py b/scripts/simulate_blif_flow_v2_authoritative.py index 8562479..f0696d4 100644 --- a/scripts/simulate_blif_flow_v2_authoritative.py +++ b/scripts/simulate_blif_flow_v2_authoritative.py @@ -17,14 +17,14 @@ CLASSIFICATIONS = { "REPLACED_BY_CORRECT_ACTION", "LEGACY_ONLY", "UNSAFE_FALSE_NEGATIVE", } NAMED = { - "Instalbeira": ("5c33db95-fab8-477a-bddd-0b9cc8f91302", "CREATE_PROFORMA", "do_now"), - "Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"), - "ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "do_now"), - "CONSTRURECUP": ("e3b23ac5-84db-4763-8a31-a684e873032c", "PREPARE_ORDER", "do_now"), - "X MAT canonical": ("dc89a466-db24-401b-bfe9-d47644b2d0c8", None, "not_current"), - "X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, "suppressed"), - "RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "review"), - "RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", None, "suppressed"), + "Instalbeira": ("5c33db95-fab8-477a-bddd-0b9cc8f91302", "CREATE_PROFORMA", "CREATE_PROFORMA", "do_now"), + "Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", None, "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"), + "ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "PREPARE_ORDER", "do_now"), + "CONSTRURECUP": ("e3b23ac5-84db-4763-8a31-a684e873032c", "PREPARE_ORDER", "PREPARE_ORDER", "do_now"), + "X MAT canonical": ("dc89a466-db24-401b-bfe9-d47644b2d0c8", None, None, "not_current"), + "X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, None, "suppressed"), + "RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS", "review"), + "RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", "REVIEW_REQUIRED", None, "suppressed"), } @@ -69,6 +69,11 @@ def _key(item: dict[str, Any]) -> str: return str(item.get("work_item_key") or item.get("process_key") or item.get("opportunity_id") or item.get("id")) +def _material_identity(item: dict[str, Any]) -> str: + return str(item.get("material_process_key") or item.get("canonical_opportunity_id") + or item.get("opportunity_id") or item.get("process_key") or _key(item)) + + def _classification(old: dict[str, Any], new: dict[str, Any] | None) -> tuple[str, str]: if new is not None: return "REPLACED_BY_CORRECT_ACTION", "Flow v2 selected a different current action for the same work item." @@ -90,23 +95,48 @@ def _classification(old: dict[str, Any], new: dict[str, Any] | None) -> tuple[st def build_report( *, identity: dict[str, Any], configured_mode: str, v1: dict[str, Any], v2: dict[str, Any], named_decisions: dict[str, dict[str, Any]], missing_projections: int, + projection_state_counts: dict[str, int] | None = None, ) -> dict[str, Any]: - v1_current = {_key(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT} - v2_current = {_key(item): item for item in _all(v2) if item.get("operational_queue") in CURRENT} + v1_current = {_material_identity(item): item for item in _all(v1) if item.get("operational_queue") in CURRENT} + v2_current = {_material_identity(item): item for item in _all(v2) if item.get("operational_queue") in CURRENT} removed, changed = [], [] for key, old in v1_current.items(): new = v2_current.get(key) if new is not None and new.get("action_code") == old.get("action_code"): continue classification, reason = _classification(old, new) - row = {"work_item_key": key, "opportunity_id": old.get("opportunity_id"), + row = {"material_process_identity": key, "work_item_key": _key(old), "opportunity_id": old.get("opportunity_id"), "v1_action": old.get("action_code"), "simulated_v2_action": (new or {}).get("action_code"), "classification": classification, "reason": reason, "source_refs": old.get("source_refs") or []} (changed if new is not None else removed).append(row) - added = [{"work_item_key": key, "opportunity_id": item.get("opportunity_id"), - "action": item.get("action_code"), "source_refs": item.get("source_refs") or []} - for key, item in v2_current.items() if key not in v1_current] + added = [] + for key, item in v2_current.items(): + if key in v1_current: + continue + decision = item.get("decision") if isinstance(item.get("decision"), dict) else {} + winning = list(decision.get("obligation_source_refs") or []) + reason_code = decision.get("reason_code") + obligation_kind = ("VALID_ACTIVE_OBLIGATION" if str(reason_code).startswith(("ACTIVE_", "DUE_", "EXPLICIT_")) + else "FACTUAL_V2_BUSINESS_ACTION") + added.append({ + "material_process_identity": key, "work_item_key": _key(item), + "opportunity_id": item.get("opportunity_id"), + "canonical_opportunity_id": item.get("canonical_opportunity_id"), + "material_process_key": item.get("material_process_key"), + "business_state": decision.get("business_state"), + "business_next_action": decision.get("business_next_action"), + "effective_action": decision.get("effective_action") or item.get("action_code"), + "queue": item.get("operational_queue"), "reason_code": reason_code, + "winning_obligation_task_id": (winning[0].get("id") if winning else None), + "task_status": (winning[0].get("status") if winning else None), + "source_system": (winning[0].get("source_system") if winning else None), + "authoritative_basis": obligation_kind, + "why_absent_from_v1": "NO_V1_CURRENT_CARD_FOR_MATERIAL_PROCESS", + "source_refs": item.get("source_refs") or [], + "non_winning_obligations": [ref for ref in decision.get("obligation_audit_refs") or [] + if ref.get("classification") != "VALID_ACTIVE_OBLIGATION"], + }) material_counts: dict[str, int] = {} for item in v2_current.values(): @@ -115,13 +145,15 @@ def build_report( duplicates = [{"material_process": key, "count": count} for key, count in material_counts.items() if count > 1] named_cases = {} - for name, (oid, expected_action, expected_queue) in NAMED.items(): + for name, (oid, expected_business, expected_action, expected_queue) in NAMED.items(): decision = named_decisions.get(oid) or {} suppressed = decision.get("suppress_current_card") is True actual_queue = "suppressed" if suppressed else decision.get("operational_queue") actual_action = decision.get("effective_action") - passed = actual_action == expected_action and actual_queue == expected_queue - named_cases[name] = {"opportunity_id": oid, "expected_action": expected_action, + actual_business = decision.get("business_next_action") + passed = actual_business == expected_business and actual_action == expected_action and actual_queue == expected_queue + named_cases[name] = {"opportunity_id": oid, "expected_business_next_action": expected_business, + "actual_business_next_action": actual_business, "expected_action": expected_action, "expected_queue": expected_queue, "actual_action": actual_action, "actual_queue": actual_queue, "passed": passed} @@ -134,6 +166,22 @@ def build_report( "reason_code": decision.get("reason_code"), "reason": decision.get("description")}) unsafe = sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in removed) + invalid_obligation_cards = [] + for item in v2_current.values(): + decision = item.get("decision") if isinstance(item.get("decision"), dict) else {} + refs = decision.get("obligation_source_refs") or [] + if refs and any(str(ref.get("status") or "").lower() != "pending" for ref in refs): + invalid_obligation_cards.append({"opportunity_id": item.get("opportunity_id"), + "effective_action": item.get("action_code"), "refs": refs}) + grouped_added = {code: [row for row in added if row["effective_action"] == code] + for code in ("SEND_INFO", "SEND_QUOTE", "SEND_PROFORMA", "REVIEW_REQUIRED")} + state_counts = projection_state_counts or {} + waiting_projection_count = sum(state_counts.get(state, 0) for state in ("AWAITING_CUSTOMER", "AWAITING_PAYMENT")) + rendered_waiting = _metrics(v2)["waiting"] + represented_waiting_states = sum( + 1 for item in _all(v2) + if (item.get("decision") or {}).get("business_state") in {"AWAITING_CUSTOMER", "AWAITING_PAYMENT"} + ) return { "identity": identity, "configured_mode": configured_mode, "v1_metrics": _metrics(v1), "authoritative_metrics": _metrics(v2), @@ -142,6 +190,15 @@ def build_report( "duplicate_cards": duplicates, "duplicate_current_cards": len(duplicates), "missing_projections": missing_projections, "unsafe_false_negatives": unsafe, "named_cases": named_cases, "fiscal_validation_cases": fiscal, + "added_card_diagnostics": added, "grouped_added_diagnostics": grouped_added, + "waiting_state_diagnostics": { + "projection_count": waiting_projection_count, "represented_processes": represented_waiting_states, + "rendered_waiting_cards": rendered_waiting, + "promoted_by_due_obligation": max(0, represented_waiting_states - rendered_waiting), + "rendering_policy": "Waiting states carry no invented internal action; they render in waiting only when selected as Operations processes.", + "lost": max(0, waiting_projection_count - represented_waiting_states), + }, + "invalid_obligation_current_cards": invalid_obligation_cards, } @@ -165,7 +222,8 @@ def exit_code(report: dict[str, Any]) -> int: unsafe = int(report.get("unsafe_false_negatives") or 0) missing = int(report.get("missing_projections") or 0) duplicates = int(report.get("duplicate_current_cards") or 0) - return 2 if unsafe or missing or duplicates or failed_named else 0 + invalid_obligations = len(report.get("invalid_obligation_current_cards") or []) + return 2 if unsafe or missing or duplicates or failed_named or invalid_obligations else 0 def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str, Any]: @@ -199,6 +257,10 @@ def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str, SELECT (SELECT count(*) FROM opportunities) AS opportunities, (SELECT count(*) FROM opportunity_flow_state_v2) AS projections """).mappings().one() + state_rows = conn.exec_driver_sql(""" + SELECT business_state, count(*) AS count + FROM opportunity_flow_state_v2 GROUP BY business_state + """).mappings().all() named_ids = [value[0] for value in NAMED.values()] named_decisions = get_decisions(named_ids, flow_mode="authoritative", connection=conn) if production and os.environ.get("BLIF_FLOW_V2_MODE") != configured_mode: @@ -207,6 +269,7 @@ def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str, identity=identity, configured_mode=configured_mode or "unset", v1=v1, v2=v2, named_decisions=named_decisions, missing_projections=max(0, int(projection_counts["opportunities"]) - int(projection_counts["projections"])), + projection_state_counts={str(row["business_state"]): int(row["count"]) for row in state_rows}, ) _write_reports(report) return report diff --git a/tests/test_authoritative_operational_adapter.py b/tests/test_authoritative_operational_adapter.py index 310080c..eb12d9c 100644 --- a/tests/test_authoritative_operational_adapter.py +++ b/tests/test_authoritative_operational_adapter.py @@ -27,6 +27,19 @@ def test_business_baselines_and_terminal(): assert decide_authoritative_operation(projection("COMPLETED"), now=NOW).queue == "not_current" +def test_terminal_state_ignores_leftover_pending_obligation(): + got = decide_authoritative_operation(projection("COMPLETED"), + obligations=[task("REVIEW_RECONSTRUCTED_PROCESS")], now=NOW) + assert (got.effective_action, got.queue, got.reason_code) == (None, "not_current", "TERMINAL_BUSINESS_STATE") + + +def test_specific_reconstructed_review_overrides_generic_business_review(): + got = decide_authoritative_operation(projection("REVIEW_REQUIRED", "REVIEW_REQUIRED"), + obligations=[task("REVIEW_RECONSTRUCTED_PROCESS")], now=NOW) + assert (got.business_next_action, got.effective_action, got.queue) == ( + "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS", "review") + + def test_missing_projection_fails_closed_without_v1(): got = decide_authoritative_operation(None, now=NOW) assert (got.effective_action, got.queue, got.reason_code) == ("REVIEW_REQUIRED", "review", "MISSING_V2_PROJECTION") @@ -46,7 +59,8 @@ def test_only_active_non_superseded_tasks_override_and_precedence(): obligations = [task("CALL_CUSTOMER"), task("SEND_INFO"), task("SUPPORT"), task("REVIEW_MANUALLY")] assert decide_authoritative_operation(projection(), obligations=obligations, now=NOW).effective_action == "REVIEW_MANUALLY" assert decide_authoritative_operation(projection(), obligations=[task("SUPPORT")], now=NOW).effective_action == "SUPPORT" - assert decide_authoritative_operation(projection(), obligations=[task("SEND_INFO")], now=NOW).effective_action == "SEND_INFO" + assert decide_authoritative_operation(projection("INQUIRY", "SEND_INFO"), obligations=[task("SEND_INFO")], now=NOW).effective_action == "SEND_INFO" + assert decide_authoritative_operation(projection(), obligations=[task("SEND_INFO")], now=NOW).effective_action is None assert decide_authoritative_operation(projection(), obligations=[task("CALL_CUSTOMER")], now=NOW).effective_action == "CALL_CUSTOMER" diff --git a/tests/test_blif_flow_v2_authoritative_simulation.py b/tests/test_blif_flow_v2_authoritative_simulation.py index f4c841e..38890ff 100644 --- a/tests/test_blif_flow_v2_authoritative_simulation.py +++ b/tests/test_blif_flow_v2_authoritative_simulation.py @@ -15,6 +15,7 @@ class Result: def __init__(self, row): self.row = row def mappings(self): return self def one(self): return self.row + def all(self): return self.row if isinstance(self.row, list) else [self.row] class Connection: @@ -26,6 +27,8 @@ class Connection: return Result({"database": self.database, "user": "audit", "transaction_read_only": self.readonly}) if "count(*) FROM opportunities" in sql: return Result({"opportunities": 328, "projections": 328}) + if "GROUP BY business_state" in sql: + return Result([]) return Result({}) def close(self): pass @@ -37,8 +40,8 @@ class Engine: def named_decisions(fail=False): result = {} - for name, (oid, action, queue) in sim.NAMED.items(): - result[oid] = {"effective_action": action, + for name, (oid, business_action, action, queue) in sim.NAMED.items(): + result[oid] = {"business_next_action": business_action, "effective_action": action, "operational_queue": "not_current" if queue == "suppressed" else queue, "suppress_current_card": queue == "suppressed"} if fail: @@ -136,3 +139,29 @@ def test_named_case_failure_is_nonzero(): report = sim.build_report(identity={}, configured_mode="compare", v1={"work_items": []}, v2={"work_items": []}, named_decisions=named_decisions(True), missing_projections=0) assert sim.exit_code(report) != 0 + + +def test_action_changes_correlate_by_material_process_not_action_key(): + v1_item = {"work_item_key": "opportunity:o:action:OLD", "opportunity_id": "o", + "action_code": "OLD", "operational_queue": "do_now"} + v2_item = {"work_item_key": "opportunity:o:action:NEW", "opportunity_id": "o", + "canonical_opportunity_id": "o", "action_code": "NEW", "operational_queue": "do_now"} + report = sim.build_report(identity={}, configured_mode="compare", + v1={"work_items": [v1_item]}, v2={"work_items": [v2_item]}, + named_decisions=named_decisions(), missing_projections=0) + assert len(report["changed_actions"]) == 1 + assert report["removed_cards"] == [] and report["added_cards"] == [] + + +def test_authoritative_material_collapse_selects_one_card_and_retains_evidence(): + from app.operations_service import _collapse_authoritative_material_processes + rows = [ + {"process_key": "conversation:chatwoot:1863", "work_item_key": "x:SEND_INFO", + "action_code": "SEND_INFO", "operational_queue": "do_now", "source_refs": [{"source": "task", "id": "1"}]}, + {"process_key": "conversation:chatwoot:1863", "work_item_key": "x:REVIEW_MANUALLY", + "action_code": "REVIEW_MANUALLY", "operational_queue": "review", "source_refs": [{"source": "communication", "id": "2"}]}, + ] + result = _collapse_authoritative_material_processes(rows) + assert len(result) == 1 + assert result[0]["action_code"] == "REVIEW_MANUALLY" + assert {ref["id"] for ref in result[0]["source_refs"]} == {"1", "2"}