From 1fba818de1e5b206bc0703aedf11ba39a99f708d Mon Sep 17 00:00:00 2001 From: Manas Narra Date: Sat, 26 Sep 2026 14:55:36 +0530 Subject: [PATCH] feat(outreach): the report says for which event we called before the goal MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The Performance tab shows lift as goals met after a conversation over goals met without one. Flipkart asked for the proof behind the "after" half: for each of those customers, which event did we call them for, after which the goal letter arrived. Until now that table was a hand-run Metabase query (25 Sep 2026); this puts the same cut on the report itself. - get_call_facts_by_runs_query returns two more facts per (run, template): last_answered_at, the latest answered call before the run's exited_at (or so far, while open), as a plain max(...) FILTER beside first_answered_at; and last_answered_event, the stage that call was placed for — the call square's stage label (current_stage) when it has one, else the event_name fact — picked in its own DISTINCT ON CTE (newest call first, id as the tie-break) and joined back, so the facts aggregate stays a hash aggregate and never sorts or unpacks `payload`; the payload is read once per answered lead. Both words were written into the payload when the call square QUEUED the lead, and a retry carries its parent's, so the word can trail where the customer stood by the time they answered. - build_report folds `goal_met_after_by_event` on ReportCustomers: every goal counted into goal_met_after_reach adds one under its run's last answered event, the latest across the agents that rang it; "(no event)" when the payload carries neither a stage label nor an event_name. Busiest first, the no-event bucket always last. Its values always sum to goal_met_after_reach, so the split and the lift are the same set of runs. Plans without stage labels or event_name facts get the single "(no event)" row; the console hides the block in that case. - Tests: the fold (two agents, the later call wins; a conversation after the end counts as before; withdrawn adds nothing; no-event last even when it is the biggest bucket), the SQL (last is judged by the one answered rule and stops at exited_at; no array_agg; facts never reads payload; the DISTINCT ON CTE and its join), and the two rule-count assertions. Verified locally against real lead rows and through GET /workflows/{id}/report on a seeded plan: the by-event rows sum to the after count on every window tried; a stage-labelled call lands under its label. On 20,800 runs with 2 KB payloads the read costs the same as the release version and returns rows identical to the ordered-array_agg form it replaces. --- app/crm/outreach/analytics.py | 39 +++++++- app/crm/outreach/schemas.py | 16 +++- .../breeze_buddy/lead_call_tracker.py | 6 +- .../queries/breeze_buddy/lead_call_tracker.py | 39 +++++++- tests/crm/test_console_reads.py | 89 ++++++++++++++++++- 5 files changed, 178 insertions(+), 11 deletions(-) diff --git a/app/crm/outreach/analytics.py b/app/crm/outreach/analytics.py index f07747e93..baed4268e 100644 --- a/app/crm/outreach/analytics.py +++ b/app/crm/outreach/analytics.py @@ -252,7 +252,11 @@ def build_report( ``by_reach`` cuts the same runs by stage (never dialled / dialled, no answer / spoke) with the same time rule, so its "spoke" column is exactly the after-we-spoke pair above; ``open_by_square`` is where the - open ones stand.""" + open ones stand. ``goal_met_after_by_event`` splits the after-we-spoke + goals by the stage the run's LAST answered call before it ended was + placed for (the lead's ``current_stage``, else its ``event_name``, both + written when the call was queued; ``_NO_EVENT`` when it carries + neither), so it always sums to ``goal_met_after_reach``.""" by_reach = { stage: { "runs": 0, @@ -265,6 +269,7 @@ def build_report( for stage in REACH_STAGES } open_by_square: Dict[str, int] = {} + after_by_event: Dict[str, int] = {} customers = { "runs": len(endings), "unique_customers": len({e.enrollment_key for e in endings}), @@ -330,6 +335,9 @@ def build_report( if ending.exit_reason == "goal_met": customers[f"goal_met_{when}_reach"] += 1 by_reach[stage]["goal_met"] += 1 + if spoke_first: + event = _last_answered_event(rows) + after_by_event[event] = after_by_event.get(event, 0) + 1 elif ending.exit_reason == "withdrawn": customers[f"withdrawn_{when}_reach"] += 1 by_reach[stage]["withdrawn"] += 1 @@ -348,6 +356,15 @@ def build_report( call_histogram=_bars(histogram), by_reach={k: ReportReach(**v) for k, v in by_reach.items()}, open_by_square=dict(sorted(open_by_square.items())), + # busiest first; the no-event bucket always last, whatever its + # count — it is the row a plan without stages or event facts + # returns, and it says nothing about a stage + goal_met_after_by_event=dict( + sorted( + after_by_event.items(), + key=lambda kv: (kv[0] == _NO_EVENT, -kv[1], kv[0]), + ) + ), ), calls=_report_calls(plan_calls, dialled, customers["reached"]), by_template=[ @@ -415,6 +432,26 @@ def _bars(histogram: Dict[int, _Bar]) -> List[ReportCallBar]: ] +# The bucket for an after-we-spoke goal whose last answered call carried +# neither a stage label nor an event_name fact — a plan without stage labels +# whose letters hold no event_name, or a lead pushed outside a plan. For such +# a plan the whole split is this one row; the split only says something for +# plans that label their call squares or whose letters carry event_name. +_NO_EVENT = "(no event)" + + +def _last_answered_event(rows: List[Dict[str, Any]]) -> str: + """PURE: the stage the run's last answered call before it ended was + placed for, across every template that rang it (the read is per + template, the latest wins).""" + latest = max( + (r for r in rows if r.get("last_answered_at")), + key=lambda r: r["last_answered_at"], + default=None, + ) + return str((latest or {}).get("last_answered_event") or _NO_EVENT) + + def _stage(placed: int, spoke: bool) -> str: """PURE: which of the three reach stages a run is in.""" if spoke: diff --git a/app/crm/outreach/schemas.py b/app/crm/outreach/schemas.py index da4bfaf58..1cda0fb92 100644 --- a/app/crm/outreach/schemas.py +++ b/app/crm/outreach/schemas.py @@ -886,7 +886,20 @@ class ReportCustomers(BaseModel): .goal_met and goal_met_before_reach is the other two stages' sum. ``open_by_square`` says where the still-open runs stand right now (current_node → runs), so "still open" splits into waiting-for-a-call - and waiting-for-an-event by the plan's own squares.""" + and waiting-for-an-event by the plan's own squares. + + ``goal_met_after_by_event`` is goal_met_after_reach cut by the stage + the run's last answered call before it ended was PLACED FOR: the call + square's stage label (``current_stage``) when it has one, else the + ``event_name`` fact — both written into the lead's payload when the + call was queued, and a retry carries its parent's, so the word can + trail where the customer stood by the time they answered. It answers + "for which event did we call, and then they converted", the proof a + merchant asks for beside the lift; not "where was the customer when we + spoke". Only meaningful for plans whose call squares carry stage labels + or whose letters carry ``event_name``; any other plan gets the single + "(no event)" row. Values sum to goal_met_after_reach, busiest first, + "(no event)" always last.""" runs: int unique_customers: int @@ -904,6 +917,7 @@ class ReportCustomers(BaseModel): open: int by_reach: Dict[str, ReportReach] = Field(default_factory=dict) open_by_square: Dict[str, int] = Field(default_factory=dict) + goal_met_after_by_event: Dict[str, int] = Field(default_factory=dict) # The calls-per-customer chart: every FINISHED call lead, placed or not, # so a call the dialler refused (CALL_LIMIT_REACHED) is on it; one still # queued or on the line is not. calls_per_customer above stays diff --git a/app/database/accessor/breeze_buddy/lead_call_tracker.py b/app/database/accessor/breeze_buddy/lead_call_tracker.py index d605d3034..2f3804406 100644 --- a/app/database/accessor/breeze_buddy/lead_call_tracker.py +++ b/app/database/accessor/breeze_buddy/lead_call_tracker.py @@ -432,8 +432,10 @@ async def get_call_facts_by_runs( runs: Sequence[Tuple[str, datetime, Optional[datetime]]], ) -> Dict[str, List[Dict[str, Any]]]: """Per workflow run (keyed by its id): one row per template that rang - it — template, leads, placed, answered, no_answer, busy, in_progress, - first_answered_at. ``runs`` is (id, entered_at, exited_at) per run: the + it — template, leads, finished, placed, answered, no_answer, busy, + in_progress, outcomes, first_answered_at, last_answered_at and + last_answered_event (the payload's event_name on the last answered call + before the run ended). ``runs`` is (id, entered_at, exited_at) per run: the id finds the leads the run stamped, the lifetime bounds its retries. A run with no lead is absent. Plain dicts: the data layer knows no CRM shape.""" diff --git a/app/database/queries/breeze_buddy/lead_call_tracker.py b/app/database/queries/breeze_buddy/lead_call_tracker.py index 2a027e787..0ae192724 100644 --- a/app/database/queries/breeze_buddy/lead_call_tracker.py +++ b/app/database/queries/breeze_buddy/lead_call_tracker.py @@ -412,7 +412,27 @@ def get_call_facts_by_runs_query( template), placed or not, the latter as {outcome: count}. A lead the dialler ended without ringing (CALL_LIMIT_REACHED) is in both, never in ``placed``; one still queued or on the line is in neither. The - calls-per-customer chart counts these.""" + calls-per-customer chart counts these. + + ``last_answered_at`` / ``last_answered_event`` are the run's LAST + answered call before its exited_at (or so far, while open) and the + stage that call was placed for: ``current_stage`` when the call square + carries a stage label (the walker stamps it into the payload), else the + ``event_name`` fact. Both were written into the payload when the call + square QUEUED the lead, and a retry carries its parent's — so the word + may trail where the customer stood by the time they answered. The + report reads "for which event did we call, and then they converted" + off it. Per (run, template) like the rest; the caller takes the latest + across templates. NULL when the payload carries neither. + + The event is picked in its own CTE (``last_ev``, DISTINCT ON the group, + newest call first) rather than an ordered ``array_agg`` in ``facts``: + an ordered aggregate turns the whole facts read into a sort, and + ``payload`` is a jsonb column that can run to kilobytes on a call-square + lead — sorting it for every lead is what the plain ``max(...) FILTER`` + beside it avoids. ``facts`` stays a hash aggregate that never touches + ``payload``; ``last_ev`` unpacks it once per answered lead.""" + last = f'{_ANSWERED} AND "call_initiated_time" < COALESCE(exited_at, now())' text = f""" {_run_leads_cte()}, by_outcome AS ( @@ -437,15 +457,26 @@ def get_call_facts_by_runs_query( count(*) FILTER (WHERE "status" = 'FINISHED' AND "outcome" = 'NO_ANSWER')::int AS no_answer, count(*) FILTER (WHERE "status" = 'FINISHED' AND "outcome" = 'BUSY')::int AS busy, count(*) FILTER (WHERE "call_initiated_time" IS NOT NULL AND "status" <> 'FINISHED')::int AS in_progress, - min("call_initiated_time") FILTER (WHERE {_ANSWERED}) AS first_answered_at + min("call_initiated_time") FILTER (WHERE {_ANSWERED}) AS first_answered_at, + max("call_initiated_time") FILTER (WHERE {last}) AS last_answered_at FROM mine GROUP BY 1, 2 + ), last_ev AS ( + SELECT DISTINCT ON (run_id, COALESCE("template", '')) + run_id, + COALESCE("template", '') AS template, + COALESCE("payload" ->> 'current_stage', "payload" ->> 'event_name') AS last_answered_event + FROM mine + WHERE {last} + ORDER BY run_id, COALESCE("template", ''), "call_initiated_time" DESC, "id" ) SELECT f.run_id AS enrollment_id, f.template, f.leads, f.finished, f.placed, f.answered, f.no_answer, f.busy, f.in_progress, - f.first_answered_at, COALESCE(o.outcomes, '{{}}'::jsonb) AS outcomes + f.first_answered_at, f.last_answered_at, e.last_answered_event, + COALESCE(o.outcomes, '{{}}'::jsonb) AS outcomes FROM facts f - LEFT JOIN outcomes o ON o.run_id = f.run_id AND o.template = f.template; + LEFT JOIN outcomes o ON o.run_id = f.run_id AND o.template = f.template + LEFT JOIN last_ev e ON e.run_id = f.run_id AND e.template = f.template; """ return text, [merchant_id, enrollment_ids, entered_ats, exited_ats] diff --git a/tests/crm/test_console_reads.py b/tests/crm/test_console_reads.py index dad67a569..204a2b676 100644 --- a/tests/crm/test_console_reads.py +++ b/tests/crm/test_console_reads.py @@ -209,7 +209,8 @@ def test_the_report_and_the_calls_summary_fold_the_same_lead_set() -> None: # the one answered definition: the count, the moment, and the stats # rows' own `spoke` column answered = "'NO_ANSWER', 'NUMBER_UNAVAILABLE', 'FAILED'" - assert facts_sql.count(answered) == 2 and stats_sql.count(answered) == 1 + # facts: answered, first_answered_at, last_answered_at, last_answered_event + assert facts_sql.count(answered) == 4 and stats_sql.count(answered) == 1 assert ") AS spoke" in stats_sql and "GROUP BY 1, 2, 3" in stats_sql @@ -489,6 +490,8 @@ def _fact(placed: int, answered: int = 0, first=None, template="nudge", **more): "busy": more.pop("busy", 0), "in_progress": 0, "first_answered_at": first, + "last_answered_at": more.pop("last", first), + "last_answered_event": more.pop("event", None), } ] @@ -557,6 +560,85 @@ def test_the_report_tells_before_from_after_by_time_alone() -> None: assert sum(stage.runs for stage in by.values()) == c.runs +def test_after_we_spoke_goals_are_split_by_the_event_we_last_spoke_on() -> None: + """The proof beside the lift: for each goal that a conversation + preceded, the event the customer stood on at our LAST answered call + before the run ended — read off the lead's frozen event_name, the + latest across the agents that rang.""" + h = timedelta(hours=1) + endings = [ + _ending(1, "exited", "goal_met", T0), # spoke on OFFERED → after + _ending(2, "exited", "goal_met", T0), # two agents; the later one wins + _ending(3, "exited", "goal_met", T0), # spoke, payload had no event + _ending(4, "exited", "goal_met", T0), # spoke only after it ended → before + _ending(5, "exited", "withdrawn", T0), # spoke, but no goal + _ending(6, "exited", "goal_met", T0), # never dialled + ] + facts = { + "r1": _fact(1, 1, T0 - h, event="OFFERED"), + "r2": _fact(2, 1, T0 - 3 * h, last=T0 - 3 * h, event="OFFERED") + + _fact(1, 1, T0 - h, template="kyc", event="KYC_COMPLETED"), + "r3": _fact(1, 1, T0 - h), + "r4": _fact(1, 1, T0 + h, last=None), + "r5": _fact(1, 1, T0 - h, event="OFFERED"), + } + c = analytics.build_report(endings, facts).customers + assert (c.goal_met_after_reach, c.goal_met_before_reach) == (3, 2) + assert c.goal_met_after_by_event == { + "KYC_COMPLETED": 1, + "OFFERED": 1, + "(no event)": 1, + } + assert sum(c.goal_met_after_by_event.values()) == c.goal_met_after_reach + # busiest first, ties by name, the no-event bucket last + assert list(c.goal_met_after_by_event) == ["KYC_COMPLETED", "OFFERED", "(no event)"] + # ... and last even when it is the biggest bucket — the shape a plan + # without stage labels or event_name facts returns + endings = [_ending(i, "exited", "goal_met", T0) for i in range(1, 4)] + facts = { + "r1": _fact(1, 1, T0 - h), + "r2": _fact(1, 1, T0 - h), + "r3": _fact(1, 1, T0 - h, event="OFFERED"), + } + by = analytics.build_report(endings, facts).customers.goal_met_after_by_event + assert list(by.items()) == [("OFFERED", 1), ("(no event)", 2)] + + +def test_the_last_answered_call_is_the_last_one_before_the_run_ended() -> None: + """The read decides "last" with the same answered rule as "first", and + stops at the run's exited_at — a call picked up after the goal letter + is not where we last spoke BEFORE it. The event rides in its own + DISTINCT ON CTE, newest call first, so the facts aggregate stays a hash + aggregate and never sorts or unpacks ``payload``; the stage label wins + over the merchant's event_name when the call square has one.""" + sql, _ = get_call_facts_by_runs_query("m1", ["r1"], [T0], [None]) + last = ( + '"call_initiated_time" IS NOT NULL AND "status" = \'FINISHED\' ' + "AND COALESCE(\"outcome\", '') NOT IN " + "('NO_ANSWER', 'NUMBER_UNAVAILABLE', 'FAILED') " + 'AND "call_initiated_time" < COALESCE(exited_at, now())' + ) + assert ( + f'max("call_initiated_time") FILTER (WHERE {last}) AS last_answered_at' in sql + ) + assert "array_agg" not in sql + facts = sql[sql.index("facts AS (") : sql.index("last_ev AS (")] + assert "payload" not in facts + assert "SELECT DISTINCT ON (run_id, COALESCE(\"template\", ''))" in sql + assert ( + "COALESCE(\"payload\" ->> 'current_stage', \"payload\" ->> 'event_name') " + "AS last_answered_event" in sql + ) + assert f"WHERE {last}\n" in sql + assert ( + 'ORDER BY run_id, COALESCE("template", \'\'), "call_initiated_time" DESC, "id"' + in sql + ) + assert ( + "LEFT JOIN last_ev e ON e.run_id = f.run_id AND e.template = f.template" in sql + ) + + def test_two_agents_fold_into_the_plan_wide_table() -> None: endings = [_ending(1), _ending(2)] facts = { @@ -671,8 +753,9 @@ def test_the_report_reads_are_windowed_on_entered_at_and_tenant_first() -> None: assert 'l."merchant_id" = $1' in text assert values == ["m1", ["a", "b"], [T0, T0], [None, T0]] assert 'min("call_initiated_time")' in text and "GROUP BY 1, 2" in text - # one definition of answered, used for the count and the moment alike - assert text.count("'NO_ANSWER', 'NUMBER_UNAVAILABLE', 'FAILED'") == 2 + # one definition of answered: the count, the first moment, the last moment + # before the run ended and the event on it + assert text.count("'NO_ANSWER', 'NUMBER_UNAVAILABLE', 'FAILED'") == 4 def test_a_run_is_staged_by_whether_anyone_spoke_before_it_ended() -> None: