Compare commits
2 Commits
ac385342f5
...
fix/blif-f
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f96b1b3d6e | ||
|
|
213ba6b41b |
@@ -48,7 +48,8 @@ class AuthoritativeOperationalDecision:
|
|||||||
"reason_if_blocked": self.reason_code if not self.eligible else None,
|
"reason_if_blocked": self.reason_code if not self.eligible else None,
|
||||||
"operational_queue": self.queue,
|
"operational_queue": self.queue,
|
||||||
"authoritative_v2": True,
|
"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
|
return result
|
||||||
|
|
||||||
@@ -76,9 +77,16 @@ def _active_obligations(obligations: Iterable[Mapping[str, Any]]) -> list[Mappin
|
|||||||
and not row.get("resolved_at") and not row.get("superseded_by_task_id")]
|
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]:
|
def _ref(row: Mapping[str, Any]) -> dict[str, Any]:
|
||||||
return {"source": "task", "id": str(row.get("id") or ""), "status": "pending",
|
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(
|
def decide_authoritative_operation(
|
||||||
@@ -110,6 +118,10 @@ def decide_authoritative_operation(
|
|||||||
oid, canonical, state, business_action, hard_blocker, "blocked", True,
|
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,
|
"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)
|
active = _active_obligations(obligations)
|
||||||
by_code: dict[str, list[Mapping[str, Any]]] = {}
|
by_code: dict[str, list[Mapping[str, Any]]] = {}
|
||||||
@@ -143,7 +155,12 @@ def decide_authoritative_operation(
|
|||||||
|
|
||||||
picked = choose(("REVIEW_MANUALLY", "REVIEW_REQUIRED", "REVIEW", "REVIEW_RECONSTRUCTED_PROCESS"), "review", "ACTIVE_REVIEW_OBLIGATION")
|
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(("SUPPORT",), "do_now", "ACTIVE_SUPPORT_OBLIGATION")
|
||||||
picked = picked or choose(("SEND_INFO",), "do_now", "ACTIVE_RESPONSE_OBLIGATION")
|
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")
|
picked = picked or choose(("CALL_CUSTOMER",), "do_now", "EXPLICIT_CALL_OBLIGATION")
|
||||||
if picked:
|
if picked:
|
||||||
return picked
|
return picked
|
||||||
@@ -168,9 +185,6 @@ def decide_authoritative_operation(
|
|||||||
obligation_source_refs=tuple(_ref(row) for row in rows), confidence=confidence,
|
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:
|
if business_action:
|
||||||
queue = "review" if business_action in REVIEW_ACTIONS else "do_now"
|
queue = "review" if business_action in REVIEW_ACTIONS else "do_now"
|
||||||
return AuthoritativeOperationalDecision(oid, canonical, state, business_action, business_action, queue, True,
|
return AuthoritativeOperationalDecision(oid, canonical, state, business_action, business_action, queue, True,
|
||||||
|
|||||||
@@ -408,6 +408,7 @@ def _opportunity_item(
|
|||||||
"business_next_action": decision.get("business_next_action"),
|
"business_next_action": decision.get("business_next_action"),
|
||||||
"effective_action": decision.get("effective_action"),
|
"effective_action": decision.get("effective_action"),
|
||||||
"canonical_opportunity_id": decision.get("canonical_opportunity_id") or opportunity_id,
|
"canonical_opportunity_id": decision.get("canonical_opportunity_id") or opportunity_id,
|
||||||
|
"material_process_key": decision.get("material_process_key"),
|
||||||
})
|
})
|
||||||
return item, independent_exceptions
|
return item, independent_exceptions
|
||||||
|
|
||||||
|
|||||||
@@ -38,6 +38,53 @@ def _int(value: Any) -> int:
|
|||||||
return 0
|
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:
|
def _is_manually_resolved_error(value: Any) -> bool:
|
||||||
text_value = str(value or "").strip().casefold()
|
text_value = str(value or "").strip().casefold()
|
||||||
return "limpo manualmente" in text_value or "resolvido manualmente" in text_value
|
return "limpo manualmente" in text_value or "resolvido manualmente" in text_value
|
||||||
@@ -683,7 +730,7 @@ def get_operations_summary(
|
|||||||
seeded = {str(row.get("opportunity_id") or "") for row in seed_rows}
|
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:
|
with (nullcontext(connection) if connection is not None else engine.connect()) as conn:
|
||||||
v2_seeds = conn.execute(text("""
|
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,
|
o.title, o.value_amount, o.currency, o.lifecycle_state,
|
||||||
c.name AS customer_name
|
c.name AS customer_name
|
||||||
FROM opportunity_flow_state_v2 p
|
FROM opportunity_flow_state_v2 p
|
||||||
@@ -701,7 +748,10 @@ def get_operations_summary(
|
|||||||
"source": "flow_v2", "id": oid, "opportunity_id": oid,
|
"source": "flow_v2", "id": oid, "opportunity_id": oid,
|
||||||
"created_at": projection_seed.get("derived_at"), "due_at": None,
|
"created_at": projection_seed.get("derived_at"), "due_at": None,
|
||||||
"priority": "normal", "queue": "vendas",
|
"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",
|
"title": projection_seed.get("title") or "Oportunidade",
|
||||||
"opportunity_title": projection_seed.get("title") or "",
|
"opportunity_title": projection_seed.get("title") or "",
|
||||||
"customer_name": projection_seed.get("customer_name") or "",
|
"customer_name": projection_seed.get("customer_name") or "",
|
||||||
@@ -722,6 +772,8 @@ def get_operations_summary(
|
|||||||
) if opportunity_ids else {})
|
) if opportunity_ids else {})
|
||||||
projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows])
|
projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows])
|
||||||
canonical_items = _attach_operation_urls(list(projection["items"]))
|
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)
|
partition = partition_canonical_items(canonical_items, display_limit=display_limit)
|
||||||
actionable_items = partition["actionable_items"]
|
actionable_items = partition["actionable_items"]
|
||||||
all_waiting_items = partition["waiting_items"]
|
all_waiting_items = partition["waiting_items"]
|
||||||
@@ -750,7 +802,7 @@ def get_operations_summary(
|
|||||||
"backlog_total": partition["backlog_total"],
|
"backlog_total": partition["backlog_total"],
|
||||||
"not_current_total": partition["not_current_total"],
|
"not_current_total": partition["not_current_total"],
|
||||||
}
|
}
|
||||||
return {
|
result = {
|
||||||
"counts": cleaned_counts,
|
"counts": cleaned_counts,
|
||||||
"recent_outbox": [dict(r) for r in recent_outbox],
|
"recent_outbox": [dict(r) for r in recent_outbox],
|
||||||
"recent_documents": [dict(r) for r in recent_documents],
|
"recent_documents": [dict(r) for r in recent_documents],
|
||||||
|
|||||||
@@ -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
|
FROM opportunity_flow_state_v2 WHERE opportunity_id::text IN :opportunity_ids
|
||||||
""", ids)
|
""", ids)
|
||||||
tasks = _bulk_rows(conn, """
|
tasks = _bulk_rows(conn, """
|
||||||
SELECT id::text, opportunity_id::text, action_code, status, due_at,
|
SELECT t.id::text, t.opportunity_id::text, t.action_code, t.status, t.due_at,
|
||||||
resolved_at, superseded_by_task_id::text, created_at, metadata
|
t.source_system, t.conversation_id, t.resolved_at, t.resolution_code,
|
||||||
FROM tasks WHERE opportunity_id::text IN :opportunity_ids AND status = 'pending'
|
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)
|
""", ids)
|
||||||
customers = _bulk_rows(conn, """
|
customers = _bulk_rows(conn, """
|
||||||
SELECT o.id::text AS opportunity_id,
|
SELECT o.id::text AS opportunity_id,
|
||||||
@@ -514,10 +521,31 @@ def _get_authoritative_opportunity_decisions(ids: list[str], *, connection: Any
|
|||||||
reconciliation_blocking=oid in reconciliation_ids,
|
reconciliation_blocking=oid in reconciliation_ids,
|
||||||
)
|
)
|
||||||
value = decision.to_dict()
|
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:
|
if projection is None:
|
||||||
value["opportunity_id"] = oid
|
value["opportunity_id"] = oid
|
||||||
value["target_url"] = f"/opportunities/{oid}"
|
value["target_url"] = f"/opportunities/{oid}"
|
||||||
else:
|
else:
|
||||||
|
value["material_process_key"] = projection.get("material_process_key")
|
||||||
value["target_url"] = f"/opportunities/{decision.canonical_opportunity_id or oid}"
|
value["target_url"] = f"/opportunities/{decision.canonical_opportunity_id or oid}"
|
||||||
decisions[oid] = value
|
decisions[oid] = value
|
||||||
return decisions
|
return decisions
|
||||||
|
|||||||
@@ -17,14 +17,18 @@ CLASSIFICATIONS = {
|
|||||||
"REPLACED_BY_CORRECT_ACTION", "LEGACY_ONLY", "UNSAFE_FALSE_NEGATIVE",
|
"REPLACED_BY_CORRECT_ACTION", "LEGACY_ONLY", "UNSAFE_FALSE_NEGATIVE",
|
||||||
}
|
}
|
||||||
NAMED = {
|
NAMED = {
|
||||||
"Instalbeira": ("5c33db95-fab8-477a-bddd-0b9cc8f91302", "CREATE_PROFORMA", "do_now"),
|
"Instalbeira": ("5c33db95-fab8-477a-bddd-0b9cc8f91302", "CREATE_PROFORMA", "CREATE_PROFORMA", "do_now"),
|
||||||
"Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"),
|
"Panoramic": ("fd79b9a1-07e6-4f61-95e8-09eab89c155e", None, "FOLLOW_UP_CUSTOMER_REVIEW", "do_now"),
|
||||||
"ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "do_now"),
|
"ENGEXICON": ("61f1c955-a372-4ea7-b9b0-b8528d74a141", "PREPARE_ORDER", "PREPARE_ORDER", "do_now"),
|
||||||
"CONSTRURECUP": ("e3b23ac5-84db-4763-8a31-a684e873032c", "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, "not_current"),
|
"X MAT canonical": ("dc89a466-db24-401b-bfe9-d47644b2d0c8", None, None, "not_current"),
|
||||||
"X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, "suppressed"),
|
"X MAT duplicate": ("1816a06e-9a69-4a9b-9279-1263156892d3", None, None, "suppressed"),
|
||||||
"RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "review"),
|
"RZSOLAR canonical": ("fd221608-e007-4043-a23d-07e0c119a345", "REVIEW_REQUIRED", "REVIEW_RECONSTRUCTED_PROCESS", "review"),
|
||||||
"RZSOLAR duplicate": ("434124fb-ac19-4d78-909a-55761d7e8daa", None, "suppressed"),
|
"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",
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
@@ -69,6 +73,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"))
|
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]:
|
def _classification(old: dict[str, Any], new: dict[str, Any] | None) -> tuple[str, str]:
|
||||||
if new is not None:
|
if new is not None:
|
||||||
return "REPLACED_BY_CORRECT_ACTION", "Flow v2 selected a different current action for the same work item."
|
return "REPLACED_BY_CORRECT_ACTION", "Flow v2 selected a different current action for the same work item."
|
||||||
@@ -90,23 +99,48 @@ def _classification(old: dict[str, Any], new: dict[str, Any] | None) -> tuple[st
|
|||||||
def build_report(
|
def build_report(
|
||||||
*, identity: dict[str, Any], configured_mode: str, v1: dict[str, Any], v2: dict[str, Any],
|
*, 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,
|
named_decisions: dict[str, dict[str, Any]], missing_projections: int,
|
||||||
|
projection_state_counts: dict[str, int] | None = None,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
v1_current = {_key(item): item for item in _all(v1) 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 = {_key(item): item for item in _all(v2) 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 = [], []
|
removed, changed = [], []
|
||||||
for key, old in v1_current.items():
|
for key, old in v1_current.items():
|
||||||
new = v2_current.get(key)
|
new = v2_current.get(key)
|
||||||
if new is not None and new.get("action_code") == old.get("action_code"):
|
if new is not None and new.get("action_code") == old.get("action_code"):
|
||||||
continue
|
continue
|
||||||
classification, reason = _classification(old, new)
|
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"),
|
"v1_action": old.get("action_code"), "simulated_v2_action": (new or {}).get("action_code"),
|
||||||
"classification": classification, "reason": reason,
|
"classification": classification, "reason": reason,
|
||||||
"source_refs": old.get("source_refs") or []}
|
"source_refs": old.get("source_refs") or []}
|
||||||
(changed if new is not None else removed).append(row)
|
(changed if new is not None else removed).append(row)
|
||||||
added = [{"work_item_key": key, "opportunity_id": item.get("opportunity_id"),
|
added = []
|
||||||
"action": item.get("action_code"), "source_refs": item.get("source_refs") or []}
|
for key, item in v2_current.items():
|
||||||
for key, item in v2_current.items() if key not in v1_current]
|
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] = {}
|
material_counts: dict[str, int] = {}
|
||||||
for item in v2_current.values():
|
for item in v2_current.values():
|
||||||
@@ -115,15 +149,24 @@ def build_report(
|
|||||||
duplicates = [{"material_process": key, "count": count} for key, count in material_counts.items() if count > 1]
|
duplicates = [{"material_process": key, "count": count} for key, count in material_counts.items() if count > 1]
|
||||||
|
|
||||||
named_cases = {}
|
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 {}
|
decision = named_decisions.get(oid) or {}
|
||||||
suppressed = decision.get("suppress_current_card") is True
|
suppressed = decision.get("suppress_current_card") is True
|
||||||
actual_queue = "suppressed" if suppressed else decision.get("operational_queue")
|
actual_queue = "suppressed" if suppressed else decision.get("operational_queue")
|
||||||
actual_action = decision.get("effective_action")
|
actual_action = decision.get("effective_action")
|
||||||
passed = actual_action == expected_action and actual_queue == expected_queue
|
actual_business = decision.get("business_next_action")
|
||||||
named_cases[name] = {"opportunity_id": oid, "expected_action": expected_action,
|
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,
|
"expected_queue": expected_queue, "actual_action": actual_action,
|
||||||
"actual_queue": actual_queue, "passed": passed}
|
"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 = []
|
fiscal = []
|
||||||
for item in v2_current.values():
|
for item in v2_current.values():
|
||||||
@@ -134,6 +177,22 @@ def build_report(
|
|||||||
"reason_code": decision.get("reason_code"),
|
"reason_code": decision.get("reason_code"),
|
||||||
"reason": decision.get("description")})
|
"reason": decision.get("description")})
|
||||||
unsafe = sum(row["classification"] == "UNSAFE_FALSE_NEGATIVE" for row in removed)
|
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 {
|
return {
|
||||||
"identity": identity, "configured_mode": configured_mode,
|
"identity": identity, "configured_mode": configured_mode,
|
||||||
"v1_metrics": _metrics(v1), "authoritative_metrics": _metrics(v2),
|
"v1_metrics": _metrics(v1), "authoritative_metrics": _metrics(v2),
|
||||||
@@ -142,6 +201,15 @@ def build_report(
|
|||||||
"duplicate_cards": duplicates, "duplicate_current_cards": len(duplicates),
|
"duplicate_cards": duplicates, "duplicate_current_cards": len(duplicates),
|
||||||
"missing_projections": missing_projections, "unsafe_false_negatives": unsafe,
|
"missing_projections": missing_projections, "unsafe_false_negatives": unsafe,
|
||||||
"named_cases": named_cases, "fiscal_validation_cases": fiscal,
|
"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 +233,8 @@ def exit_code(report: dict[str, Any]) -> int:
|
|||||||
unsafe = int(report.get("unsafe_false_negatives") or 0)
|
unsafe = int(report.get("unsafe_false_negatives") or 0)
|
||||||
missing = int(report.get("missing_projections") or 0)
|
missing = int(report.get("missing_projections") or 0)
|
||||||
duplicates = int(report.get("duplicate_current_cards") 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]:
|
def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str, Any]:
|
||||||
@@ -199,6 +268,10 @@ def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str,
|
|||||||
SELECT (SELECT count(*) FROM opportunities) AS opportunities,
|
SELECT (SELECT count(*) FROM opportunities) AS opportunities,
|
||||||
(SELECT count(*) FROM opportunity_flow_state_v2) AS projections
|
(SELECT count(*) FROM opportunity_flow_state_v2) AS projections
|
||||||
""").mappings().one()
|
""").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_ids = [value[0] for value in NAMED.values()]
|
||||||
named_decisions = get_decisions(named_ids, flow_mode="authoritative", connection=conn)
|
named_decisions = get_decisions(named_ids, flow_mode="authoritative", connection=conn)
|
||||||
if production and os.environ.get("BLIF_FLOW_V2_MODE") != configured_mode:
|
if production and os.environ.get("BLIF_FLOW_V2_MODE") != configured_mode:
|
||||||
@@ -207,6 +280,7 @@ def run(args: argparse.Namespace, *, runtime_loader=_load_runtime) -> dict[str,
|
|||||||
identity=identity, configured_mode=configured_mode or "unset", v1=v1, v2=v2,
|
identity=identity, configured_mode=configured_mode or "unset", v1=v1, v2=v2,
|
||||||
named_decisions=named_decisions,
|
named_decisions=named_decisions,
|
||||||
missing_projections=max(0, int(projection_counts["opportunities"]) - int(projection_counts["projections"])),
|
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)
|
_write_reports(report)
|
||||||
return report
|
return report
|
||||||
|
|||||||
@@ -27,6 +27,19 @@ def test_business_baselines_and_terminal():
|
|||||||
assert decide_authoritative_operation(projection("COMPLETED"), now=NOW).queue == "not_current"
|
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():
|
def test_missing_projection_fails_closed_without_v1():
|
||||||
got = decide_authoritative_operation(None, now=NOW)
|
got = decide_authoritative_operation(None, now=NOW)
|
||||||
assert (got.effective_action, got.queue, got.reason_code) == ("REVIEW_REQUIRED", "review", "MISSING_V2_PROJECTION")
|
assert (got.effective_action, got.queue, got.reason_code) == ("REVIEW_REQUIRED", "review", "MISSING_V2_PROJECTION")
|
||||||
@@ -37,16 +50,38 @@ def test_due_and_future_customer_followup_semantics():
|
|||||||
assert (due.effective_action, due.queue) == ("FOLLOW_UP_CUSTOMER_REVIEW", "do_now")
|
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)
|
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")
|
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"
|
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():
|
def test_only_active_non_superseded_tasks_override_and_precedence():
|
||||||
stale = task("SUPPORT", resolved_at=NOW)
|
stale = task("SUPPORT", resolved_at=NOW)
|
||||||
assert decide_authoritative_operation(projection("PROFORMA_REQUIRED", "CREATE_PROFORMA"), obligations=[stale], now=NOW).effective_action == "CREATE_PROFORMA"
|
assert decide_authoritative_operation(projection("PROFORMA_REQUIRED", "CREATE_PROFORMA"), obligations=[stale], now=NOW).effective_action == "CREATE_PROFORMA"
|
||||||
obligations = [task("CALL_CUSTOMER"), task("SEND_INFO"), task("SUPPORT"), task("REVIEW_MANUALLY")]
|
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=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("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"
|
assert decide_authoritative_operation(projection(), obligations=[task("CALL_CUSTOMER")], now=NOW).effective_action == "CALL_CUSTOMER"
|
||||||
|
|
||||||
|
|
||||||
@@ -67,6 +102,18 @@ def test_duplicate_is_suppressed_even_with_pending_legacy_task():
|
|||||||
assert canonicalize_operations([source], {"opp": decision})["items"] == []
|
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():
|
def test_authoritative_queue_and_action_survive_canonical_boundary():
|
||||||
decision = decide_authoritative_operation(projection(), obligations=[task("FOLLOW_UP_CUSTOMER_REVIEW")], now=NOW).to_dict()
|
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",
|
source = {"source": "task", "id": "old", "status": "pending", "action_code": "SEND_INVOICE",
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ class Result:
|
|||||||
def __init__(self, row): self.row = row
|
def __init__(self, row): self.row = row
|
||||||
def mappings(self): return self
|
def mappings(self): return self
|
||||||
def one(self): return self.row
|
def one(self): return self.row
|
||||||
|
def all(self): return self.row if isinstance(self.row, list) else [self.row]
|
||||||
|
|
||||||
|
|
||||||
class Connection:
|
class Connection:
|
||||||
@@ -26,6 +27,8 @@ class Connection:
|
|||||||
return Result({"database": self.database, "user": "audit", "transaction_read_only": self.readonly})
|
return Result({"database": self.database, "user": "audit", "transaction_read_only": self.readonly})
|
||||||
if "count(*) FROM opportunities" in sql:
|
if "count(*) FROM opportunities" in sql:
|
||||||
return Result({"opportunities": 328, "projections": 328})
|
return Result({"opportunities": 328, "projections": 328})
|
||||||
|
if "GROUP BY business_state" in sql:
|
||||||
|
return Result([])
|
||||||
return Result({})
|
return Result({})
|
||||||
def close(self): pass
|
def close(self): pass
|
||||||
|
|
||||||
@@ -37,10 +40,11 @@ class Engine:
|
|||||||
|
|
||||||
def named_decisions(fail=False):
|
def named_decisions(fail=False):
|
||||||
result = {}
|
result = {}
|
||||||
for name, (oid, action, queue) in sim.NAMED.items():
|
for name, (oid, business_action, action, queue) in sim.NAMED.items():
|
||||||
result[oid] = {"effective_action": action,
|
result[oid] = {"business_next_action": business_action, "effective_action": action,
|
||||||
"operational_queue": "not_current" if queue == "suppressed" else queue,
|
"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:
|
if fail:
|
||||||
result[sim.NAMED["Instalbeira"][0]]["effective_action"] = "SEND_INVOICE"
|
result[sim.NAMED["Instalbeira"][0]]["effective_action"] = "SEND_INVOICE"
|
||||||
return result
|
return result
|
||||||
@@ -136,3 +140,29 @@ def test_named_case_failure_is_nonzero():
|
|||||||
report = sim.build_report(identity={}, configured_mode="compare", v1={"work_items": []},
|
report = sim.build_report(identity={}, configured_mode="compare", v1={"work_items": []},
|
||||||
v2={"work_items": []}, named_decisions=named_decisions(True), missing_projections=0)
|
v2={"work_items": []}, named_decisions=named_decisions(True), missing_projections=0)
|
||||||
assert sim.exit_code(report) != 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"}
|
||||||
|
|||||||
Reference in New Issue
Block a user