Files
clientflow_backend/tests/test_canonical_operations_read_model.py
2026-08-15 00:32:38 +00:00

416 lines
20 KiB
Python

from __future__ import annotations
from pathlib import Path
from app.canonical_operations import (
_merge_exact_work_items,
canonicalize_operations,
opportunity_ids_from_work_seeds,
partition_canonical_items,
)
def row(source="task", id="1", opportunity_id="opp-1", action_code="SEND_QUOTE", **extra):
value = {
"source": source,
"id": id,
"opportunity_id": opportunity_id,
"action_code": action_code,
"status": "pending",
"priority": "normal",
"queue": "vendas",
"title": action_code,
"detail": "Requer intervenção",
"source_system": "chatwoot" if source in {"task", "communication"} else source,
"created_at": "2026-08-15T10:00:00+00:00",
"item_metadata": {},
}
value.update(extra)
return value
def decision(code="SEND_QUOTE", **extra):
value = {
"action_code": code,
"label": code.replace("_", " ").title(),
"description": "Ação canónica calculada.",
"reason": "O processo requer esta ação agora.",
"priority": "normal",
"target_url": "/opportunities/opp-1",
"decision_version": "test-v1",
}
value.update(extra)
return value
def project(rows, decisions=None, limit=None, evidence_rows=None):
return canonicalize_operations(rows, decisions or {}, evidence_rows=evidence_rows, display_limit=limit)
def test_task_and_communication_for_same_opportunity_are_one_card_and_evidence():
result = project([
row(id="task-1"),
row(source="communication", id="comm-1"),
], {"opp-1": decision()})
assert result["canonical_count"] == 1
assert {ref["source"] for ref in result["items"][0]["source_refs"]} == {"task", "communication"}
def test_only_task_matching_canonical_action_is_current_and_stale_task_is_evidence():
result = project([
row(id="stale", action_code="SEND_PROFORMA"),
row(id="current", action_code="CONFIRM_PAYMENT"),
], {"opp-1": decision("CONFIRM_PAYMENT")})
item = result["items"][0]
assert item["id"] == "current"
assert item["current_action_code"] == "CONFIRM_PAYMENT"
assert [ref["task_id"] for ref in item["stale_task_refs"]] == ["stale"]
def test_linked_communication_is_evidence_not_duplicate():
result = project([row(source="communication", id="comm-1")], {"opp-1": decision()})
assert len(result["items"]) == 1
assert result["items"][0]["source_refs"][0]["communication_id"] == "comm-1"
def test_standalone_actionable_and_uncertain_communications_remain_visible():
result = project([
row(source="communication", id="a", opportunity_id="", conversation_id="10", action_code="SUPPORT"),
row(source="communication", id="b", opportunity_id="", conversation_id="11", action_code="REVIEW_MANUALLY"),
])
assert result["canonical_count"] == 2
assert {item["current_action_code"] for item in result["items"]} == {"SUPPORT", "REVIEW_MANUALLY"}
def test_same_reliable_standalone_conversation_is_one_process_but_subject_is_not_used():
result = project([
row(source="task", id="a", opportunity_id="", conversation_id="10", title="same"),
row(source="communication", id="b", opportunity_id="", conversation_id="10", title="same"),
row(source="communication", id="c", opportunity_id="", conversation_id="11", title="same"),
])
assert result["canonical_count"] == 2
assert any(item["raw_source_count"] == 2 for item in result["items"])
def test_same_conversation_with_different_actions_remains_two_work_items():
result = project([
row(source="task", id="quote", opportunity_id="", conversation_id="10", action_code="SEND_QUOTE"),
row(source="communication", id="support", opportunity_id="", conversation_id="10", action_code="SUPPORT"),
])
assert result["canonical_count"] == 2
assert {item["current_action_code"] for item in result["items"]} == {"SEND_QUOTE", "SUPPORT"}
def test_explicit_communication_task_relation_can_merge_different_actions():
result = project([
row(source="task", id="task-1", opportunity_id="", conversation_id="10", action_code="SEND_QUOTE"),
row(
source="communication", id="comm-1", opportunity_id="", conversation_id="10",
action_code="REVIEW_MANUALLY", item_metadata={"task_id": "task-1"},
),
])
assert result["canonical_count"] == 1
assert result["items"][0]["current_action_code"] == "SEND_QUOTE"
assert {ref["source"] for ref in result["items"][0]["source_refs"]} == {"task", "communication"}
def test_final_exact_key_merge_combines_independent_same_action_groups_and_keeps_task_handle():
result = project([
row(
source="task", id="task-1", opportunity_id="", conversation_id="2433",
action_code="SEND_INFO", href="/tasks/task-1", priority="normal",
),
row(
source="communication", id="related", opportunity_id="", conversation_id="2433",
action_code="SEND_INFO", item_metadata={"task_id": "task-1"},
),
row(
source="communication", id="independent", opportunity_id="", conversation_id="2433",
action_code="SEND_INFO", priority="alta", customer_name="Cliente conhecido",
),
])
assert result["canonical_count"] == 1
item = result["items"][0]
assert item["work_item_key"] == "conversation:chatwoot:2433:action:SEND_INFO"
assert item["source"] == "task"
assert item["href"] == "/tasks/task-1"
assert item["priority"] == "alta"
assert item["customer_name"] == "Cliente conhecido"
assert item["created_at"] == "2026-08-15T10:00:00+00:00"
assert item["raw_source_count"] == 3
assert {(ref["source"], ref["id"]) for ref in item["source_refs"]} == {
("task", "task-1"), ("communication", "related"), ("communication", "independent")
}
def test_exact_key_merge_unions_unique_refs_and_stale_task_refs():
duplicate_ref = {"source": "communication", "id": "same"}
item = _merge_exact_work_items([
{
"work_item_key": "conversation:chatwoot:2433:action:SEND_INFO",
"process_key": "conversation:chatwoot:2433", "current_action_code": "SEND_INFO",
"source_refs": [duplicate_ref, {"source": "task", "id": "task-1"}],
"stale_task_refs": [{"source": "task", "id": "stale-1"}],
"priority": "normal",
},
{
"work_item_key": "conversation:chatwoot:2433:action:SEND_INFO",
"process_key": "conversation:chatwoot:2433", "current_action_code": "SEND_INFO",
"source_refs": [duplicate_ref, {"source": "communication", "id": "comm-2"}],
"stale_task_refs": [
{"source": "task", "id": "stale-1"},
{"source": "task", "id": "stale-2"},
],
"priority": "normal",
},
])[0]
assert len(item["source_refs"]) == 3
assert {(ref["source"], ref["id"]) for ref in item["source_refs"]} == {
("communication", "same"), ("task", "task-1"), ("communication", "comm-2")
}
assert {ref["id"] for ref in item["stale_task_refs"]} == {"stale-1", "stale-2"}
def test_deterministic_noise_is_excluded():
result = project([
row(opportunity_id="", action_code="IGNORE_BOUNCE", title="Mailer-Daemon"),
row(opportunity_id="", action_code="IGNORE_SPAM"),
row(opportunity_id="", action_code="NO_ACTION", title="Automated notification"),
])
assert result["canonical_count"] == 0
def test_wait_decision_is_classified_waiting():
item = project([row()], {"opp-1": decision("WAIT_PRODUCTION")})["items"][0]
assert item["operational_queue"] == "waiting"
assert item["status"] == "waiting"
def test_waiting_allowlist_does_not_hide_unknown_wait_action():
for code in ("WAIT_CUSTOMER", "WAIT_PRODUCTION"):
assert project([row()], {"opp-1": decision(code)})["items"][0]["operational_queue"] == "waiting"
unknown = project([row()], {"opp-1": decision("WAIT_MANUAL_REVIEW")})["items"][0]
assert unknown["operational_queue"] == "do_now"
assert unknown["current_action_code"] == "WAIT_MANUAL_REVIEW"
def test_blocker_is_only_current_action_and_downstream_task_is_evidence():
result = project([
row(id="downstream", action_code="SEND_INVOICE"),
row(id="blocker", action_code="VALIDATE_FISCAL_CUSTOMER"),
], {"opp-1": decision("VALIDATE_FISCAL_CUSTOMER", reason="Falta cliente fiscal.")})
assert [item["current_action_code"] for item in result["items"]] == ["VALIDATE_FISCAL_CUSTOMER"]
assert result["items"][0]["stale_task_refs"][0]["task_id"] == "downstream"
def test_linked_reconciliation_is_evidence_and_standalone_decision_remains_visible():
result = project([
row(source="reconciliation", id="loose", opportunity_id="", action_code="RECONCILE_DOCUMENTS"),
row(id="task-seed", action_code="RECONCILE_DOCUMENTS"),
], {"opp-1": decision("RECONCILE_DOCUMENTS")}, evidence_rows=[
row(source="reconciliation", id="linked", status="linked", action_code="RECONCILE_DOCUMENTS"),
])
assert result["canonical_count"] == 2
linked = next(item for item in result["items"] if item.get("opportunity_id"))
assert any(ref.get("reconciliation_item_id") == "linked" for ref in linked["source_refs"])
def test_historical_communication_does_not_seed_but_enriches_a_seeded_opportunity():
historical = row(source="communication", id="history", status="done", action_code="SEND_INFO")
evidence_only = project([], {"opp-1": decision()}, evidence_rows=[historical])
assert evidence_only["canonical_count"] == 0
assert evidence_only["evidence_source_count"] == 0
seeded = project([row(id="seed")], {"opp-1": decision()}, evidence_rows=[historical])
assert seeded["canonical_count"] == 1
assert seeded["evidence_source_count"] == 1
assert any(ref.get("communication_id") == "history" for ref in seeded["items"][0]["source_refs"])
def test_linked_and_resolved_reconciliation_do_not_seed_but_can_be_evidence():
linked = row(source="reconciliation", id="linked", status="linked", action_code="RECONCILE_DOCUMENTS")
resolved = row(source="reconciliation", id="resolved", status="resolved", action_code="RECONCILE_DOCUMENTS")
assert project([], {"opp-1": decision()}, evidence_rows=[linked])["canonical_count"] == 0
assert project([], {"opp-1": decision()}, evidence_rows=[resolved])["canonical_count"] == 0
result = project([row(id="seed")], {"opp-1": decision()}, evidence_rows=[linked, resolved])
assert result["canonical_count"] == 1
refs = result["items"][0]["source_refs"]
assert {ref.get("reconciliation_item_id") for ref in refs} >= {"linked", "resolved"}
def test_current_needs_review_communication_and_open_reconciliation_seed_work():
result = project([
row(source="communication", id="review", opportunity_id="", status="needs_review", action_code="REVIEW_MANUALLY"),
row(source="reconciliation", id="open", opportunity_id="", status="open", action_code="RECONCILE_DOCUMENTS"),
])
assert result["canonical_count"] == 2
assert result["seed_source_count"] == 2
def test_evidence_only_opportunity_ids_are_not_decision_candidates():
seeds = [row(id="seed", opportunity_id="seeded-opportunity")]
evidence = [row(source="communication", id="history", opportunity_id="historical-opportunity", status="done")]
assert opportunity_ids_from_work_seeds(seeds) == ["seeded-opportunity"]
assert "historical-opportunity" not in opportunity_ids_from_work_seeds(seeds)
result = project(seeds, {"seeded-opportunity": decision()}, evidence_rows=evidence)
assert result["canonical_count"] == 1
assert result["items"][0]["opportunity_id"] == "seeded-opportunity"
assert result["evidence_source_count"] == 0
def test_evidence_does_not_change_queue_counts_and_limit_is_post_canonical():
seeds = [row(id=str(index), opportunity_id=f"opp-{index}") for index in range(40)]
decisions = {f"opp-{index}": decision() for index in range(40)}
evidence = [
row(source="communication", id=f"history-{index}", opportunity_id=f"opp-{index}", status="done")
for index in range(40)
]
result = project(seeds, decisions, evidence_rows=evidence)
partition = partition_canonical_items(result["items"], display_limit=30)
assert result["seed_source_count"] == 40
assert result["evidence_source_count"] == 40
assert partition["work_queue_total"] == 40
assert len(partition["visible_actionable_items"]) == 30
def test_canonical_count_and_display_limit_follow_final_exact_key_merge():
rows = []
for index in range(35):
conversation_id = str(3000 + index)
rows.append(row(source="task", id=f"task-{index}", opportunity_id="", conversation_id=conversation_id, action_code="SEND_INFO"))
rows.append(row(source="communication", id=f"comm-{index}", opportunity_id="", conversation_id=conversation_id, action_code="SEND_INFO"))
result = project(rows, limit=30)
assert result["canonical_count"] == 35
assert result["visible_count"] == 30
assert len({item["work_item_key"] for item in result["all_items"]}) == 35
partition = partition_canonical_items(result["all_items"], display_limit=30)
assert partition["work_queue_total"] == 35
def test_every_canonical_output_has_unique_work_item_key():
result = project([
row(source="task", id="task", opportunity_id="", conversation_id="10", action_code="SEND_INFO"),
row(source="communication", id="comm", opportunity_id="", conversation_id="10", action_code="SEND_INFO"),
row(source="communication", id="review", opportunity_id="", conversation_id="10", action_code="REVIEW_MANUALLY"),
])
keys = [item["work_item_key"] for item in result["all_items"]]
assert len(keys) == len(set(keys))
assert len(keys) == 2
def test_pending_outbox_is_not_a_candidate_and_failed_exception_is_visible():
# SQL excludes pending rows; pure projection also preserves only what it is
# given, so assert the service query enforces that source boundary.
source = Path("app/operations_service.py").read_text(encoding="utf-8")
work_query = source.split("work_seed_rows = conn.execute", 1)[1].split("normalized_rows =", 1)[0]
assert "WHERE status IN ('failed','blocked')" in work_query
failed = row(source="outbox", id="failure", opportunity_id="", action_code="JASMIN_CREATE", status="failed")
assert project([failed])["items"][0]["operational_queue"] == "exception"
def test_matching_integration_failure_is_evidence_but_independent_failure_is_second_card():
matching = row(source="outbox", id="match", action_code="SEND_INVOICE", status="failed")
independent = row(source="outbox", id="other", action_code="JASMIN_TAX_FAILURE", status="failed")
result = project([row(action_code="SEND_INVOICE"), matching, independent], {"opp-1": decision("SEND_INVOICE")})
assert result["canonical_count"] == 2
opportunity_item = next(item for item in result["items"] if item["source"] == "opportunity")
assert any(ref.get("outbox_id") == "match" for ref in opportunity_item["source_refs"])
assert any(item["process_key"].startswith("integration-exception:") for item in result["items"])
def test_outbox_error_text_does_not_create_an_unstructured_action_match():
failure = row(
source="outbox", id="failure", action_code="JASMIN_FAILURE", status="failed",
title="SEND_INVOICE failed", detail="Error while processing SEND_INVOICE",
)
result = project([row(action_code="SEND_INVOICE"), failure], {"opp-1": decision("SEND_INVOICE")})
assert result["canonical_count"] == 2
opportunity_item = next(item for item in result["items"] if item["source"] == "opportunity")
assert not any(ref.get("outbox_id") == "failure" for ref in opportunity_item["source_refs"])
assert any(item["process_key"].startswith("integration-exception:") for item in result["items"])
def test_no_action_hides_resolved_evidence_but_preserves_explicit_review():
resolved = row(source="communication", id="resolved", status="done", action_code="SEND_INFO")
assert project([resolved], {"opp-1": decision("NO_ACTION")})["canonical_count"] == 0
needs_review = row(source="communication", id="review", status="needs_review", action_code="SEND_INFO")
review_result = project([needs_review], {"opp-1": decision("NO_ACTION")})
assert review_result["canonical_count"] == 1
assert review_result["items"][0]["current_action_code"] == "REVIEW"
assert "decisão central sem ação" in review_result["items"][0]["why_human_required"].lower()
def test_no_action_preserves_explicit_ambiguous_association_review():
ambiguous = row(
id="ambiguous", action_code="SEND_INVOICE",
opportunity_linking_status="ambiguous",
)
result = project([ambiguous], {"opp-1": decision("NO_ACTION")})
assert result["canonical_count"] == 1
assert result["items"][0]["current_action_code"] == "REVIEW"
def test_opportunity_customer_identity_wins_over_blank_communication_identity():
result = project([
row(source="communication", id="comm", customer_name="", contact_display_name=""),
row(id="task", fiscal_customer_name="CONSTRURECUP", customer_name=""),
], {"opp-1": decision()})
assert result["items"][0]["customer_name"] == "CONSTRURECUP"
assert result["items"][0]["contact_display_name"] == "CONSTRURECUP"
def test_display_limit_is_applied_after_canonicalization():
rows = []
decisions = {}
for index in range(40):
opportunity_id = f"opp-{index}"
decisions[opportunity_id] = decision()
rows.extend(row(id=f"{index}-{duplicate}", opportunity_id=opportunity_id) for duplicate in range(10))
result = project(rows, decisions, limit=30)
assert result["raw_source_count"] == 400
assert result["canonical_count"] == 40
assert result["visible_count"] == 30
def test_canonical_counts_and_keys_are_stable_and_action_sensitive():
rows = [row()]
first = project(rows, {"opp-1": decision("SEND_QUOTE")})
again = project(rows, {"opp-1": decision("SEND_QUOTE")})
changed = project(rows, {"opp-1": decision("CONFIRM_PAYMENT")})
assert first["canonical_count"] == len(first["all_items"]) == 1
assert first["items"][0]["process_key"] == again["items"][0]["process_key"] == "opportunity:opp-1"
assert first["items"][0]["work_item_key"] == again["items"][0]["work_item_key"]
assert first["items"][0]["work_item_key"] != changed["items"][0]["work_item_key"]
def test_projection_and_operations_get_path_contain_no_writes():
projection = Path("app/canonical_operations.py").read_text(encoding="utf-8").upper()
service = Path("app/operations_service.py").read_text(encoding="utf-8")
get_body = service.split("def get_operations_summary", 1)[1].split("def get_system_health_summary", 1)[0].upper()
next_actions = Path("app/opportunity_next_action_service.py").read_text(encoding="utf-8")
bulk_body = next_actions.split("def get_opportunity_next_actions", 1)[1]
for token in ("INSERT INTO", "UPDATE TASKS", "DELETE FROM", "ENSURE_PENDING_TASK"):
assert token not in projection
assert token not in get_body
assert "ensure_operation_schema()" not in bulk_body
def test_operations_query_separates_seed_eligibility_and_exposes_top_level_diagnostics():
source = Path("app/operations_service.py").read_text(encoding="utf-8")
seed_query = source.split("work_seed_rows = conn.execute", 1)[1].split("seed_opportunity_ids =", 1)[0]
evidence_query = source.split("evidence_rows = conn.execute", 1)[1].split("normalized_rows =", 1)[0]
assert "WHERE status IN ('new','classified','needs_review')" in seed_query
assert "OR opportunity_id IS NOT NULL" not in seed_query
assert "WHERE status IN ('open','needs_review','conflict')" in seed_query
assert "active_document_link_exclusion_sql" in source
assert "c.status NOT IN ('new','classified','needs_review')" in evidence_query
assert "ri.status IN ('linked','resolved')" in evidence_query
for key in (
"raw_source_count", "clean_source_count", "seed_source_count",
"evidence_source_count", "canonical_count", "visible_count",
"candidate_limit", "waiting_total",
):
assert f'"{key}"' in source