feat: add operational eligibility and scheduled follow-up

This commit is contained in:
plx
2026-08-15 07:56:06 +00:00
parent 71452b1e5c
commit ee8e2de606
13 changed files with 812 additions and 26 deletions

View File

@@ -443,6 +443,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
COALESCE(t.metadata, '{}'::jsonb) AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
COALESCE(o.lifecycle_state, '') AS opportunity_lifecycle_state,
o.next_follow_up_at AS opportunity_next_follow_up_at,
COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM (
@@ -492,6 +494,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
COALESCE(io.payload, '{}'::jsonb) AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
COALESCE(o.lifecycle_state, '') AS opportunity_lifecycle_state,
o.next_follow_up_at AS opportunity_next_follow_up_at,
COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM (
@@ -547,6 +551,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
ELSE '{}'::jsonb END AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
COALESCE(o.lifecycle_state, '') AS opportunity_lifecycle_state,
o.next_follow_up_at AS opportunity_next_follow_up_at,
COALESCE(o.value_amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, 'EUR') AS opportunity_currency
FROM (
@@ -592,6 +598,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
COALESCE(ri.payload, '{}'::jsonb) AS item_metadata,
COALESCE(o.metadata, '{}'::jsonb) AS opportunity_metadata,
COALESCE(o.stage, '') AS opportunity_stage,
COALESCE(o.lifecycle_state, '') AS opportunity_lifecycle_state,
o.next_follow_up_at AS opportunity_next_follow_up_at,
COALESCE(o.value_amount, ri.amount, 0) AS opportunity_value_amount,
COALESCE(o.currency, ri.currency, 'EUR') AS opportunity_currency
FROM (
@@ -645,7 +653,29 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
"candidate_limit": candidate_limit,
}).mappings().all()
normalized_rows = _normalise_work_item_intent([dict(r) for r in work_seed_rows])
conversation_ids = sorted({
str(row.get("conversation_id") or "").strip()
for row in work_seed_rows
if str(row.get("conversation_id") or "").strip()
})
conversation_timeline = {}
if conversation_ids:
timeline_rows = conn.execute(text("""
SELECT conversation_id,
max(created_at) FILTER (WHERE direction = 'inbound') AS latest_public_inbound,
max(created_at) FILTER (WHERE direction = 'outbound') AS latest_public_outbound
FROM messages
WHERE source_system = 'chatwoot'
AND conversation_id = ANY(CAST(:conversation_ids AS TEXT[]))
GROUP BY conversation_id
"""), {"conversation_ids": conversation_ids}).mappings().all()
conversation_timeline = {str(row["conversation_id"]): dict(row) for row in timeline_rows}
seed_rows = [dict(r) for r in work_seed_rows]
for row in seed_rows:
timeline = conversation_timeline.get(str(row.get("conversation_id") or ""), {})
row.update({key: timeline.get(key) for key in ("latest_public_inbound", "latest_public_outbound")})
normalized_rows = _normalise_work_item_intent(seed_rows)
opportunity_ids = opportunity_ids_from_work_seeds(normalized_rows)
decisions = get_opportunity_next_actions(opportunity_ids) if opportunity_ids else {}
projection = canonicalize_operations(normalized_rows, decisions, evidence_rows=[dict(row) for row in evidence_rows])
@@ -661,6 +691,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
# noise awaiting cleanup. The cleanup script still fixes the data source.
cleaned_counts["work_queue_total"] = partition["work_queue_total"]
cleaned_counts["waiting_total"] = partition["waiting_total"]
cleaned_counts["backlog_total"] = partition["backlog_total"]
cleaned_counts["not_current_total"] = partition["not_current_total"]
cleaned_counts["raw_source_count"] = projection["raw_source_count"]
cleaned_counts["canonical_total"] = len(canonical_items)
@@ -673,6 +705,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
"canonical_count": len(canonical_items),
"visible_count": len(cleaned_work_items),
"waiting_total": len(all_waiting_items),
"backlog_total": partition["backlog_total"],
"not_current_total": partition["not_current_total"],
}
return {
"counts": cleaned_counts,
@@ -683,6 +717,8 @@ def get_operations_summary(limit: int = 24) -> Dict[str, Any]:
"recent_communications": [dict(r) for r in recent_communications],
"work_items": cleaned_work_items,
"waiting_items": waiting_items,
"backlog_items": partition["backlog_items"][:display_limit],
"not_current_items": partition["not_current_items"][:display_limit],
**diagnostics,
"projection_metrics": diagnostics,
}