diff --git a/app/authoritative_operational_adapter.py b/app/authoritative_operational_adapter.py index f3438ae..a552bd6 100644 --- a/app/authoritative_operational_adapter.py +++ b/app/authoritative_operational_adapter.py @@ -77,6 +77,12 @@ def _active_obligations(obligations: Iterable[Mapping[str, Any]]) -> list[Mappin and not row.get("resolved_at") and not row.get("superseded_by_task_id")] +def _unanswered_inbound(row: Mapping[str, Any]) -> bool: + inbound = _dt(row.get("latest_public_inbound")) + outbound = _dt(row.get("latest_public_outbound")) + return bool(inbound and (outbound is None or outbound <= inbound)) + + 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")), @@ -149,7 +155,11 @@ 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") - if state == "INQUIRY" or business_action == "SEND_INFO": + valid_send_info = [row for row in by_code.get("SEND_INFO", []) if ( + state != "AWAITING_CUSTOMER" or _unanswered_inbound(row) + )] + if valid_send_info and (state in {"INQUIRY", "AWAITING_CUSTOMER"} or business_action == "SEND_INFO"): + by_code["SEND_INFO"] = valid_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: diff --git a/app/operations_service.py b/app/operations_service.py index 70d1762..5b37e12 100644 --- a/app/operations_service.py +++ b/app/operations_service.py @@ -730,7 +730,7 @@ def get_operations_summary( seeded = {str(row.get("opportunity_id") or "") for row in seed_rows} with (nullcontext(connection) if connection is not None else engine.connect()) as conn: v2_seeds = conn.execute(text(""" - SELECT p.opportunity_id::text, p.business_next_action, p.derived_at, + SELECT p.opportunity_id::text, p.business_state, p.business_next_action, p.derived_at, o.title, o.value_amount, o.currency, o.lifecycle_state, c.name AS customer_name FROM opportunity_flow_state_v2 p @@ -748,7 +748,10 @@ def get_operations_summary( "source": "flow_v2", "id": oid, "opportunity_id": oid, "created_at": projection_seed.get("derived_at"), "due_at": None, "priority": "normal", "queue": "vendas", - "action_code": projection_seed.get("business_next_action"), "status": "projected", + "action_code": (projection_seed.get("business_next_action") + or ("WAIT_PAYMENT" if projection_seed.get("business_state") == "AWAITING_PAYMENT" + else "WAIT_CUSTOMER")), + "status": "projected", "title": projection_seed.get("title") or "Oportunidade", "opportunity_title": projection_seed.get("title") or "", "customer_name": projection_seed.get("customer_name") or "", diff --git a/app/opportunity_next_action_service.py b/app/opportunity_next_action_service.py index 4ad3f07..e1821c1 100644 --- a/app/opportunity_next_action_service.py +++ b/app/opportunity_next_action_service.py @@ -481,9 +481,16 @@ 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, 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' + SELECT t.id::text, t.opportunity_id::text, t.action_code, t.status, t.due_at, + t.source_system, t.conversation_id, t.resolved_at, t.resolution_code, + t.superseded_by_task_id::text, t.created_at, t.metadata, + (SELECT max(m.created_at) FROM messages m + WHERE m.source_system = 'chatwoot' AND m.conversation_id = t.conversation_id + AND m.direction = 'inbound') AS latest_public_inbound, + (SELECT max(m.created_at) FROM messages m + WHERE m.source_system = 'chatwoot' AND m.conversation_id = t.conversation_id + AND m.direction = 'outbound') AS latest_public_outbound + FROM tasks t WHERE t.opportunity_id::text IN :opportunity_ids AND t.status = 'pending' """, ids) customers = _bulk_rows(conn, """ SELECT o.id::text AS opportunity_id, diff --git a/scripts/simulate_blif_flow_v2_authoritative.py b/scripts/simulate_blif_flow_v2_authoritative.py index f0696d4..7e17f12 100644 --- a/scripts/simulate_blif_flow_v2_authoritative.py +++ b/scripts/simulate_blif_flow_v2_authoritative.py @@ -26,6 +26,10 @@ NAMED = { "RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS", "review"), "RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", "REVIEW_REQUIRED", None, "suppressed"), } +DUPLICATE_CANONICAL = { + "1816a06e-9a69-4a9b-9279-1263156892d3": "dc89a466-db24-401b-bfe9-d47644b2d0c8", + "434124fb-ac19-4d78-909a-55761d7e8daa": "fd221608-e007-4043-a23d-07e0c119a345", +} def parse_args(argv: list[str] | None = None) -> argparse.Namespace: @@ -151,11 +155,18 @@ def build_report( actual_queue = "suppressed" if suppressed else decision.get("operational_queue") actual_action = decision.get("effective_action") actual_business = decision.get("business_next_action") - passed = actual_business == expected_business and actual_action == expected_action and actual_queue == expected_queue + duplicate_case = expected_queue == "suppressed" + expected_canonical = DUPLICATE_CANONICAL.get(oid) + canonical_ok = not duplicate_case or decision.get("canonical_opportunity_id") == expected_canonical + passed = ((duplicate_case or actual_business == expected_business) + and actual_action == expected_action and actual_queue == expected_queue and canonical_ok) 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} + if duplicate_case: + named_cases[name].update({"expected_canonical_opportunity_id": expected_canonical, + "actual_canonical_opportunity_id": decision.get("canonical_opportunity_id")}) fiscal = [] for item in v2_current.values(): diff --git a/tests/test_authoritative_operational_adapter.py b/tests/test_authoritative_operational_adapter.py index eb12d9c..b68096f 100644 --- a/tests/test_authoritative_operational_adapter.py +++ b/tests/test_authoritative_operational_adapter.py @@ -50,9 +50,30 @@ def test_due_and_future_customer_followup_semantics(): assert (due.effective_action, due.queue) == ("FOLLOW_UP_CUSTOMER_REVIEW", "do_now") future = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW", NOW + timedelta(days=1))], now=NOW) assert (future.effective_action, future.queue) == (None, "waiting") + future_item = canonicalize_operations([{ + "source": "flow_v2", "id": "opp", "opportunity_id": "opp", + "status": "projected", "action_code": "WAIT_CUSTOMER", "created_at": NOW, + }], {"opp": future.to_dict()})["items"] + assert len(future_item) == 1 and future_item[0]["operational_queue"] == "waiting" assert decide_authoritative_operation(projection(), now=NOW).queue == "waiting" +def test_unanswered_send_info_remains_valid_on_awaiting_customer(): + inbound = NOW - timedelta(hours=2) + obligation = task("SEND_INFO", latest_public_inbound=inbound, latest_public_outbound=None) + got = decide_authoritative_operation(projection(), obligations=[obligation], now=NOW) + assert (got.effective_action, got.queue, got.reason_code) == ( + "SEND_INFO", "do_now", "ACTIVE_RESPONSE_OBLIGATION") + + +def test_send_info_satisfied_by_later_outbound_is_not_current(): + inbound = NOW - timedelta(hours=2) + obligation = task("SEND_INFO", latest_public_inbound=inbound, + latest_public_outbound=inbound + timedelta(minutes=10)) + got = decide_authoritative_operation(projection(), obligations=[obligation], now=NOW) + assert (got.effective_action, got.queue) == (None, "waiting") + + def test_only_active_non_superseded_tasks_override_and_precedence(): stale = task("SUPPORT", resolved_at=NOW) assert decide_authoritative_operation(projection("PROFORMA_REQUIRED", "CREATE_PROFORMA"), obligations=[stale], now=NOW).effective_action == "CREATE_PROFORMA" @@ -81,6 +102,18 @@ def test_duplicate_is_suppressed_even_with_pending_legacy_task(): assert canonicalize_operations([source], {"opp": decision})["items"] == [] +def test_xmat_duplicate_suppresses_diagnostic_review_business_action(): + got = decide_authoritative_operation({ + **projection("REVIEW_REQUIRED", "REVIEW_REQUIRED"), + "opportunity_id": "1816a06e-9a69-4a9b-9279-1263156892d3", + "canonical_opportunity_id": "dc89a466-db24-401b-bfe9-d47644b2d0c8", + "is_duplicate_representation": True, + }, now=NOW).to_dict() + assert got["business_next_action"] == "REVIEW_REQUIRED" + assert got["effective_action"] is None and got["suppress_current_card"] is True + assert got["canonical_opportunity_id"] == "dc89a466-db24-401b-bfe9-d47644b2d0c8" + + def test_authoritative_queue_and_action_survive_canonical_boundary(): decision = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW")], now=NOW).to_dict() source = {"source": "task", "id": "old", "status": "pending", "action_code": "SEND_INVOICE", diff --git a/tests/test_blif_flow_v2_authoritative_simulation.py b/tests/test_blif_flow_v2_authoritative_simulation.py index 38890ff..76a5915 100644 --- a/tests/test_blif_flow_v2_authoritative_simulation.py +++ b/tests/test_blif_flow_v2_authoritative_simulation.py @@ -43,7 +43,8 @@ def named_decisions(fail=False): 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"} + "suppress_current_card": queue == "suppressed", + "canonical_opportunity_id": sim.DUPLICATE_CANONICAL.get(oid, oid)} if fail: result[sim.NAMED["Instalbeira"][0]]["effective_action"] = "SEND_INVOICE" return result