soft3/lytics/rs/ingest/src/reports.rs

// ---
// tags: lytics, rust
// crystal-type: source
// crystal-domain: cyber
// ---
//! read-time projections over the append-only event stream.
//!
//! arrivals, passages and windows per the attention model โ€” no sessions,
//! no timeouts. all arithmetic is integer; ratios are rendered at the
//! dashboard edge.
//!
//! `inf_reports` answers every one of these reports through real inf
//! datalog now โ€” `main.rs` calls only `inf_reports::` functions. everything
//! below `Stored`/`in_window` is `#[cfg(test)]`: it exists solely as the
//! differential-testing oracle `inf_reports.rs`'s tests compare against, and
//! is not part of the release binary at all โ€” not dead code kept around out
//! of caution, genuinely absent from a non-test build.

use crate::enrich::{Attribution, Device};
#[cfg(test)]
use lytics_event::Actor;
use lytics_event::{EventBody, Kind, Navigation};
use serde::{Deserialize, Serialize};
use serde_json::json;
use std::collections::BTreeMap;
#[cfg(test)]
use std::collections::BTreeSet;

/// a stored, enriched event โ€” payload-log plaintext.
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Stored {
    pub body: EventBody,
    pub event_hash: String,
    pub attribution: Attribution,
    pub device: Device,
    #[serde(default)]
    pub geo: Option<crate::geo::Geo>,
    pub received_at: u64,
}

impl Stored {
    pub fn is_pageview(&self) -> bool {
        self.body.kind == Kind::Pageview
    }

    /// arrival: a pageview whose navigation is external or direct, or that
    /// carries utm. attribution attaches here.
    pub fn is_arrival(&self) -> bool {
        self.is_pageview()
            && (matches!(
                self.body.navigation,
                Some(Navigation::External) | Some(Navigation::Direct)
            ) || self.body.utm.is_some())
    }

    pub fn attention_ms(&self) -> u64 {
        self.body.attention.map(|a| a.ms).unwrap_or(0)
    }

    /// peak scroll percent (0โ€“100) carried on attention events; `None` if the
    /// event is not an attention sample.
    pub fn scroll_depth(&self) -> Option<u8> {
        self.body.attention.as_ref().map(|a| a.scroll_depth)
    }
}

/// max scroll reached per (neuron, pathname) from attention events โ€” the unit
/// of "how far did this reader get on this page". short pages that never
/// fire attention contribute nothing (honest empty, not a fake 0).
pub fn scroll_reach_samples(events: &[&Stored]) -> Vec<u8> {
    let mut max_by: BTreeMap<(&str, &str), u8> = BTreeMap::new();
    for e in events {
        let Some(d) = e.scroll_depth() else {
            continue;
        };
        let key = (e.body.neuron.as_str(), e.body.pathname.as_str());
        max_by
            .entry(key)
            .and_modify(|m| {
                if d > *m {
                    *m = d;
                }
            })
            .or_insert(d);
    }
    max_by.into_values().collect()
}

/// per-pathname max-scroll samples (one value per neuron that looked).
pub fn scroll_reach_by_path(events: &[&Stored]) -> BTreeMap<String, Vec<u8>> {
    let mut max_by: BTreeMap<(&str, &str), u8> = BTreeMap::new();
    for e in events {
        let Some(d) = e.scroll_depth() else {
            continue;
        };
        let key = (e.body.neuron.as_str(), e.body.pathname.as_str());
        max_by
            .entry(key)
            .and_modify(|m| {
                if d > *m {
                    *m = d;
                }
            })
            .or_insert(d);
    }
    let mut by_path: BTreeMap<String, Vec<u8>> = BTreeMap::new();
    for ((_, path), d) in max_by {
        by_path.entry(path.to_string()).or_default().push(d);
    }
    by_path
}

fn percentile_u8(sorted: &[u8], p: f64) -> u64 {
    if sorted.is_empty() {
        return 0;
    }
    let i = ((sorted.len() as f64 - 1.0) * p).round() as usize;
    sorted[i.min(sorted.len() - 1)] as u64
}

/// p50 / p90 / sample count / fixed depth buckets for a set of scroll samples.
pub fn scroll_stats(mut samples: Vec<u8>) -> serde_json::Value {
    samples.sort_unstable();
    let n = samples.len() as u64;
    let buckets = [
        ("0โ€“24", 0u8, 24u8),
        ("25โ€“49", 25, 49),
        ("50โ€“74", 50, 74),
        ("75โ€“100", 75, 100),
    ]
    .iter()
    .map(|(label, lo, hi)| {
        let count = samples.iter().filter(|&&v| v >= *lo && v <= *hi).count() as u64;
        json!({"label": *label, "min": lo, "max": hi, "neurons": count})
    })
    .collect::<Vec<_>>();
    json!({
        "samples": n,
        "p50": percentile_u8(&samples, 0.5),
        "p90": percentile_u8(&samples, 0.9),
        "max": samples.last().copied().unwrap_or(0) as u64,
        "buckets": buckets,
    })
}

/// overview-shaped scroll block: stats over max-reach samples.
pub fn scroll_overview(events: &[&Stored]) -> serde_json::Value {
    scroll_stats(scroll_reach_samples(events))
}

pub fn in_window(events: &[Stored], from: u64, to: u64) -> Vec<&Stored> {
    events
        .iter()
        .filter(|e| e.body.timestamp >= from && e.body.timestamp < to)
        .collect()
}

/// earliest event timestamp per neuron across the full log.
pub fn first_seen(events: &[Stored]) -> BTreeMap<String, u64> {
    let mut m = BTreeMap::new();
    for e in events {
        let ts = e.body.timestamp;
        m.entry(e.body.neuron.clone())
            .and_modify(|t| {
                if ts < *t {
                    *t = ts;
                }
            })
            .or_insert(ts);
    }
    m
}

/// Restrict windowed events by audience relative to window start `from`.
///
/// - `new` โ€” first-seen โ‰ฅ `from` (neuron appeared inside the window)
/// - `returning` โ€” first-seen < `from` (known before the window, active in it)
/// - `all` / anything else โ€” no filter
pub fn filter_audience<'a>(
    windowed: &[&'a Stored],
    first_seen: &BTreeMap<String, u64>,
    from: u64,
    audience: &str,
) -> Vec<&'a Stored> {
    match audience {
        "new" => windowed
            .iter()
            .copied()
            .filter(|e| first_seen.get(&e.body.neuron).is_some_and(|&fs| fs >= from))
            .collect(),
        "returning" => windowed
            .iter()
            .copied()
            .filter(|e| first_seen.get(&e.body.neuron).is_some_and(|&fs| fs < from))
            .collect(),
        _ => windowed.to_vec(),
    }
}

/// events grouped per neuron, ordered by timestamp.
#[cfg(test)]
fn per_neuron<'a>(events: &[&'a Stored]) -> BTreeMap<&'a str, Vec<&'a Stored>> {
    let mut map: BTreeMap<&str, Vec<&Stored>> = BTreeMap::new();
    for e in events {
        map.entry(e.body.neuron.as_str()).or_default().push(e);
    }
    for list in map.values_mut() {
        list.sort_by_key(|e| e.body.timestamp);
    }
    map
}

/// a passage: the run of a neuron's events from one arrival to the next.
/// a stream that opens without an observed arrival still opens a passage โ€”
/// the first event is where the run began, arrival or no.
#[cfg(test)]
pub struct Passage<'a> {
    pub events: Vec<&'a Stored>,
}

#[cfg(test)]
impl Passage<'_> {
    pub fn views(&self) -> usize {
        self.events.iter().filter(|e| e.is_pageview()).count()
    }
    pub fn entry(&self) -> Option<&str> {
        self.events
            .iter()
            .find(|e| e.is_pageview())
            .map(|e| e.body.pathname.as_str())
    }
    pub fn exit(&self) -> Option<&str> {
        self.events
            .iter()
            .rev()
            .find(|e| e.is_pageview())
            .map(|e| e.body.pathname.as_str())
    }
}

#[cfg(test)]
pub fn passages<'a>(events: &[&'a Stored]) -> Vec<Passage<'a>> {
    let mut out = Vec::new();
    for (_neuron, list) in per_neuron(events) {
        let mut current: Vec<&Stored> = Vec::new();
        for e in list {
            if e.is_arrival() && !current.is_empty() {
                out.push(Passage {
                    events: std::mem::take(&mut current),
                });
            }
            current.push(e);
        }
        if !current.is_empty() {
            out.push(Passage { events: current });
        }
    }
    out
}

#[cfg(test)]
fn distinct_neurons(events: &[&Stored]) -> usize {
    events
        .iter()
        .map(|e| e.body.neuron.as_str())
        .collect::<BTreeSet<_>>()
        .len()
}

#[cfg(test)]
pub fn overview(events: &[&Stored]) -> serde_json::Value {
    let views = events.iter().filter(|e| e.is_pageview()).count();
    let attention: u64 = events.iter().map(|e| e.attention_ms()).sum();
    let agents = events
        .iter()
        .filter(|e| e.body.actor == Actor::Agent)
        .map(|e| e.body.neuron.as_str())
        .collect::<BTreeSet<_>>()
        .len();
    // visit = passage (stream from one external entry to the next). arrivals
    // stay an internal event label for attribution; the public count is visits.
    let p = passages(events);
    let visits = p.len();
    let neurons = distinct_neurons(events);
    // depth / dwell: pageviews and attention per visit (integer milli-precision)
    let views_per_visit_milli = if visits > 0 {
        (views as u64 * 1000) / visits as u64
    } else {
        0
    };
    let attention_ms_per_visit = if visits > 0 {
        attention / visits as u64
    } else {
        0
    };
    let attention_ms_per_neuron = if neurons > 0 {
        attention / neurons as u64
    } else {
        0
    };
    let mut o = json!({
        "neurons": neurons,
        "views": views,
        "visits": visits,
        "attention_ms": attention,
        "agent_neurons": agents,
        "views_per_visit_milli": views_per_visit_milli,
        "attention_ms_per_visit": attention_ms_per_visit,
        "attention_ms_per_neuron": attention_ms_per_neuron,
    });
    if let Some(obj) = o.as_object_mut() {
        obj.insert("scroll".into(), scroll_overview(events));
    }
    o
}

/// bucket ms: "hour" or "day".
#[cfg(test)]
pub fn timeseries(events: &[&Stored], bucket_ms: u64) -> serde_json::Value {
    // per bucket: neurons, visits(=arrivals), views, attention_ms
    let mut buckets: BTreeMap<u64, (BTreeSet<&str>, u64, u64, u64)> = BTreeMap::new();
    for e in events {
        let b = e.body.timestamp / bucket_ms * bucket_ms;
        let slot = buckets.entry(b).or_default();
        slot.0.insert(e.body.neuron.as_str());
        if e.is_arrival() {
            slot.1 += 1;
        }
        if e.is_pageview() {
            slot.2 += 1;
        }
        slot.3 += e.attention_ms();
    }
    let rows: Vec<_> = buckets
        .into_iter()
        .map(|(t, (n, visits, pv, att))| {
            let neurons = n.len() as u64;
            let depth = views_per_visit_milli(pv, visits);
            let dwell = att.checked_div(visits).unwrap_or(0);
            let att_n = att.checked_div(neurons).unwrap_or(0);
            json!({
                "t": t,
                "neurons": neurons,
                "visits": visits,
                "views": pv,
                "attention_ms": att,
                "views_per_visit_milli": depth,
                "attention_ms_per_visit": dwell,
                "attention_ms_per_neuron": att_n,
            })
        })
        .collect();
    json!(rows)
}

#[cfg(test)]
fn views_per_visit_milli(views: u64, visits: u64) -> u64 {
    views
        .checked_mul(1000)
        .and_then(|x| x.checked_div(visits))
        .unwrap_or(0)
}

#[cfg(test)]
fn top_counts<'a, F: Fn(&&'a Stored) -> Option<String>>(
    events: &[&'a Stored],
    key: F,
    limit: usize,
) -> Vec<(String, u64)> {
    let mut counts: BTreeMap<String, u64> = BTreeMap::new();
    for e in events {
        if let Some(k) = key(e) {
            *counts.entry(k).or_default() += 1;
        }
    }
    let mut rows: Vec<_> = counts.into_iter().collect();
    rows.sort_by(|a, b| b.1.cmp(&a.1).then(a.0.cmp(&b.0)));
    rows.truncate(limit);
    rows
}

/// top particles โ€” funnel: neurons / visits / views + attention.
#[cfg(test)]
pub fn particles(events: &[&Stored], limit: usize) -> serde_json::Value {
    let mut attention: BTreeMap<String, u64> = BTreeMap::new();
    let mut views: BTreeMap<String, u64> = BTreeMap::new();
    let mut visits: BTreeMap<String, u64> = BTreeMap::new();
    let mut neurons: BTreeMap<String, BTreeSet<&str>> = BTreeMap::new();
    for e in events {
        *attention.entry(e.body.pathname.clone()).or_default() += e.attention_ms();
        if e.is_pageview() {
            *views.entry(e.body.pathname.clone()).or_default() += 1;
            neurons
                .entry(e.body.pathname.clone())
                .or_default()
                .insert(e.body.neuron.as_str());
        }
        if e.is_arrival() {
            *visits.entry(e.body.pathname.clone()).or_default() += 1;
        }
    }
    let mut rows: Vec<_> = views
        .into_iter()
        .map(|(path, pv)| {
            let n = neurons.get(&path).map(|s| s.len() as u64).unwrap_or(0);
            let v = visits.get(&path).copied().unwrap_or(0);
            (path, n, v, pv)
        })
        .collect();
    rows.sort_by(|a, b| {
        b.1.cmp(&a.1)
            .then(b.2.cmp(&a.2))
            .then(b.3.cmp(&a.3))
            .then(a.0.cmp(&b.0))
    });
    rows.truncate(limit);
    let scroll_by = scroll_reach_by_path(events);
    let out: Vec<_> = rows
        .into_iter()
        .map(|(path, neurons, visits, pv)| {
            let att = attention.get(&path).copied().unwrap_or(0);
            let vpv = pv
                .checked_mul(1000)
                .and_then(|x| x.checked_div(visits))
                .unwrap_or(0);
            let att_pv = att.checked_div(visits).unwrap_or(0);
            let att_pn = att.checked_div(neurons).unwrap_or(0);
            let scroll = scroll_stats(scroll_by.get(&path).cloned().unwrap_or_default());
            json!({
                "pathname": path,
                "neurons": neurons,
                "visits": visits,
                "views": pv,
                "attention_ms": att,
                "views_per_visit_milli": vpv,
                "attention_ms_per_visit": att_pv,
                "attention_ms_per_neuron": att_pn,
                "scroll_p50": scroll["p50"],
                "scroll_p90": scroll["p90"],
                "scroll_samples": scroll["samples"],
            })
        })
        .collect();
    json!(out)
}

/// visits by source, with views + attention rolled up from each visit's events.
#[cfg(test)]
pub fn sources(events: &[&Stored], limit: usize) -> serde_json::Value {
    visit_breakdown(events, limit, |e| {
        e.attribution
            .source
            .clone()
            .unwrap_or_else(|| "direct".into())
    })
}

/// visits by channel (same rollup as sources).
#[cfg(test)]
pub fn channels(events: &[&Stored]) -> serde_json::Value {
    visit_breakdown(events, 32, |e| {
        format!("{:?}", e.attribution.channel).to_lowercase()
    })
}

/// group visits by a key taken from the visit's first arrival (or first event).
#[cfg(test)]
fn visit_breakdown<F: Fn(&Stored) -> String>(
    events: &[&Stored],
    limit: usize,
    key_of: F,
) -> serde_json::Value {
    // source key โ†’ (visits, views, attention_ms, neurons)
    let mut map: BTreeMap<String, (u64, u64, u64, BTreeSet<&str>)> = BTreeMap::new();
    for pass in passages(events) {
        let first = pass
            .events
            .iter()
            .find(|e| e.is_arrival())
            .or_else(|| pass.events.first());
        let key = first.map(|e| key_of(e)).unwrap_or_else(|| "unknown".into());
        let neuron = first.map(|e| e.body.neuron.as_str()).unwrap_or("");
        let slot = map.entry(key).or_default();
        slot.0 += 1;
        slot.1 += pass.views() as u64;
        slot.2 += pass.events.iter().map(|e| e.attention_ms()).sum::<u64>();
        if !neuron.is_empty() {
            slot.3.insert(neuron);
        }
    }
    let mut rows: Vec<_> = map
        .into_iter()
        .map(|(k, (visits, views, att, neurons))| (k, neurons.len() as u64, visits, views, att))
        .collect();
    rows.sort_by(|a, b| b.1.cmp(&a.1).then(b.2.cmp(&a.2)).then(a.0.cmp(&b.0)));
    rows.truncate(limit);
    json!(
        rows.into_iter()
            .map(|(k, neurons, visits, views, att)| {
                let vpv_milli = (views * 1000).checked_div(visits).unwrap_or(0);
                let att_pv = att.checked_div(visits).unwrap_or(0);
                let att_pn = att.checked_div(neurons).unwrap_or(0);
                json!({
                    "key": k,
                    "source": k,
                    "channel": k,
                    "neurons": neurons,
                    "visits": visits,
                    "views": views,
                    "attention_ms": att,
                    "views_per_visit_milli": vpv_milli,
                    "attention_ms_per_visit": att_pv,
                    "attention_ms_per_neuron": att_pn,
                })
            })
            .collect::<Vec<_>>()
    )
}

#[cfg(test)]
pub fn actors(events: &[&Stored]) -> serde_json::Value {
    let humans: Vec<&Stored> = events
        .iter()
        .filter(|e| e.body.actor == Actor::Human)
        .copied()
        .collect();
    let agents: Vec<&Stored> = events
        .iter()
        .filter(|e| e.body.actor == Actor::Agent)
        .copied()
        .collect();
    let agent_names = top_counts(
        &agents,
        |e| e.body.agent.as_ref().map(|a| a.name.clone()),
        16,
    );
    json!({
        "human": {"neurons": distinct_neurons(&humans), "views": humans.iter().filter(|e| e.is_pageview()).count(), "attention_ms": humans.iter().map(|e| e.attention_ms()).sum::<u64>()},
        "agent": {"neurons": distinct_neurons(&agents), "views": agents.iter().filter(|e| e.is_pageview()).count(), "declared": agent_names.into_iter().map(|(n, c)| json!({"name": n, "views": c})).collect::<Vec<_>>()},
    })
}

#[cfg(test)]
fn dim_funnel_ref(
    events: &[&Stored],
    key: impl Fn(&Stored) -> Option<String>,
    limit: usize,
) -> Vec<serde_json::Value> {
    let mut neurons: BTreeMap<String, BTreeSet<&str>> = BTreeMap::new();
    let mut views: BTreeMap<String, u64> = BTreeMap::new();
    let mut visits: BTreeMap<String, u64> = BTreeMap::new();
    let mut attn: BTreeMap<String, u64> = BTreeMap::new();
    for e in events {
        let Some(k) = key(e) else { continue };
        neurons
            .entry(k.clone())
            .or_default()
            .insert(e.body.neuron.as_str());
        if e.is_pageview() {
            *views.entry(k.clone()).or_default() += 1;
        }
        if e.is_arrival() {
            *visits.entry(k.clone()).or_default() += 1;
        }
        *attn.entry(k.clone()).or_default() += e.attention_ms();
    }
    let mut rows: Vec<_> = neurons
        .into_iter()
        .map(|(k, set)| {
            let n = set.len() as u64;
            let visits = visits.get(&k).copied().unwrap_or(0);
            let views = views.get(&k).copied().unwrap_or(0);
            let att = attn.get(&k).copied().unwrap_or(0);
            (k, n, visits, views, att)
        })
        .collect();
    rows.sort_by(|a, b| {
        b.1.cmp(&a.1)
            .then(b.2.cmp(&a.2))
            .then(b.3.cmp(&a.3))
            .then(a.0.cmp(&b.0))
    });
    rows.truncate(limit);
    rows.into_iter()
        .map(|(k, neurons, visits, views, att)| {
            let vpv = views
                .checked_mul(1000)
                .and_then(|x| x.checked_div(visits))
                .unwrap_or(0);
            let att_pv = att.checked_div(visits).unwrap_or(0);
            let att_pn = att.checked_div(neurons).unwrap_or(0);
            json!({
                "key": k,
                "neurons": neurons,
                "visits": visits,
                "views": views,
                "attention_ms": att,
                "views_per_visit_milli": vpv,
                "attention_ms_per_visit": att_pv,
                "attention_ms_per_neuron": att_pn,
            })
        })
        .collect()
}

#[cfg(test)]
pub fn countries(events: &[&Stored], limit: usize) -> serde_json::Value {
    let rows = dim_funnel_ref(
        events,
        |e| e.geo.as_ref().and_then(|g| g.country.clone()),
        limit,
    );
    json!(
        rows.into_iter()
            .map(|mut r| {
                // rename key โ†’ country for the public shape
                let k = r["key"].as_str().unwrap_or("").to_string();
                r.as_object_mut().unwrap().remove("key");
                r.as_object_mut()
                    .unwrap()
                    .insert("country".into(), json!(k));
                r
            })
            .collect::<Vec<_>>()
    )
}

#[cfg(test)]
pub fn devices(events: &[&Stored], limit: usize) -> serde_json::Value {
    let map_key = |field: &str, rows: Vec<serde_json::Value>| {
        rows.into_iter()
            .map(|mut r| {
                let k = r["key"].as_str().unwrap_or("").to_string();
                r.as_object_mut().unwrap().remove("key");
                r.as_object_mut().unwrap().insert(field.into(), json!(k));
                r
            })
            .collect::<Vec<_>>()
    };
    json!({
        "browsers": map_key("browser", dim_funnel_ref(events, |e| e.device.browser.clone(), limit)),
        "os": map_key("os", dim_funnel_ref(events, |e| e.device.os.clone(), limit)),
        "classes": map_key("class", dim_funnel_ref(events, |e| e.device.device.clone(), limit)),
    })
}

#[cfg(test)]
const WEEK_MS: u64 = 7 * 24 * 3600 * 1000;

/// retention matrix over ALL events (cohorts need full history):
/// cohort week (first-seen) ร— week offset โ†’ distinct neurons active.
#[cfg(test)]
pub fn retention(events: &[Stored], weeks: usize) -> serde_json::Value {
    let mut first_seen: BTreeMap<&str, u64> = BTreeMap::new();
    for e in events {
        let entry = first_seen
            .entry(e.body.neuron.as_str())
            .or_insert(e.body.timestamp);
        if e.body.timestamp < *entry {
            *entry = e.body.timestamp;
        }
    }
    // cohort week โ†’ offset โ†’ set of neurons
    let mut matrix: BTreeMap<u64, BTreeMap<u64, BTreeSet<&str>>> = BTreeMap::new();
    for e in events {
        let first = first_seen[e.body.neuron.as_str()];
        let cohort = first / WEEK_MS;
        let offset = e.body.timestamp / WEEK_MS - cohort;
        if (offset as usize) < weeks {
            matrix
                .entry(cohort)
                .or_default()
                .entry(offset)
                .or_default()
                .insert(&e.body.neuron);
        }
    }
    let rows: Vec<_> = matrix
        .into_iter()
        .map(|(cohort, offsets)| {
            let size = offsets.get(&0).map(|s| s.len()).unwrap_or(0);
            let cells: Vec<usize> =
                (0..weeks as u64).map(|o| offsets.get(&o).map(|s| s.len()).unwrap_or(0)).collect();
            json!({"cohort_week": cohort, "cohort_start_ms": cohort * WEEK_MS, "size": size, "weeks": cells})
        })
        .collect();
    json!(rows)
}

/// ordered funnel: how many neurons complete each prefix of `steps`
/// (pathnames) within the window.
#[cfg(test)]
pub fn funnel(events: &[&Stored], steps: &[String]) -> serde_json::Value {
    if steps.is_empty() {
        return json!([]);
    }
    let mut reached = vec![0usize; steps.len()];
    for (_neuron, list) in per_neuron(events) {
        let mut at = 0usize;
        for e in list {
            if at < steps.len() && e.is_pageview() && e.body.pathname == steps[at] {
                at += 1;
            }
        }
        for slot in reached.iter_mut().take(at) {
            *slot += 1;
        }
    }
    json!(
        steps
            .iter()
            .zip(reached)
            .map(|(s, n)| json!({"step": s, "neurons": n}))
            .collect::<Vec<_>>()
    )
}

/// return probability: of neurons first seen in [from, to), how many came
/// back within `horizon_ms` after their first event. integers only โ€”
/// (returned, total) โ€” the ratio renders at the edge.
#[cfg(test)]
pub fn returns(events: &[Stored], from: u64, to: u64, horizon_ms: u64) -> serde_json::Value {
    let mut first_seen: BTreeMap<&str, u64> = BTreeMap::new();
    for e in events {
        let entry = first_seen
            .entry(e.body.neuron.as_str())
            .or_insert(e.body.timestamp);
        if e.body.timestamp < *entry {
            *entry = e.body.timestamp;
        }
    }
    let cohort: BTreeMap<&str, u64> = first_seen
        .into_iter()
        .filter(|(_, t)| *t >= from && *t < to)
        .collect();
    let mut returned: BTreeSet<&str> = BTreeSet::new();
    for e in events {
        if let Some(first) = cohort.get(e.body.neuron.as_str())
            && e.body.timestamp > *first
            && e.body.timestamp <= first + horizon_ms
        {
            returned.insert(&e.body.neuron);
        }
    }
    json!({"cohort": cohort.len(), "returned": returned.len(), "horizon_ms": horizon_ms})
}

#[cfg(test)]
pub fn passages_report(events: &[&Stored], limit: usize) -> serde_json::Value {
    let p = passages(events);
    let total = p.len();
    let views_total: usize = p.iter().map(|x| x.views()).sum();
    // (passages count, set of neurons)
    let mut entries: BTreeMap<String, (u64, BTreeSet<&str>)> = BTreeMap::new();
    let mut exits: BTreeMap<String, (u64, BTreeSet<&str>)> = BTreeMap::new();
    for x in &p {
        let neuron = x
            .events
            .first()
            .map(|e| e.body.neuron.as_str())
            .unwrap_or("");
        if let Some(e) = x.entry() {
            let slot = entries.entry(e.to_string()).or_default();
            slot.0 += 1;
            if !neuron.is_empty() {
                slot.1.insert(neuron);
            }
        }
        if let Some(e) = x.exit() {
            let slot = exits.entry(e.to_string()).or_default();
            slot.0 += 1;
            if !neuron.is_empty() {
                slot.1.insert(neuron);
            }
        }
    }
    let top = |m: BTreeMap<String, (u64, BTreeSet<&str>)>| {
        let mut rows: Vec<_> = m
            .into_iter()
            .map(|(k, (passages, neurons))| (k, passages, neurons.len() as u64))
            .collect();
        rows.sort_by(|a, b| b.2.cmp(&a.2).then(b.1.cmp(&a.1)).then(a.0.cmp(&b.0)));
        rows.truncate(limit);
        rows.into_iter()
            .map(|(k, passages, neurons)| {
                json!({"pathname": k, "passages": passages, "neurons": neurons})
            })
            .collect::<Vec<_>>()
    };
    json!({
        "passages": total,
        "views_total": views_total,
        "entries": top(entries),
        "exits": top(exits),
    })
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::enrich::Channel;
    use lytics_event::event::Utm;

    fn ev(neuron: &str, path: &str, ts: u64, nav: Option<Navigation>, att: u64) -> Stored {
        Stored {
            body: EventBody {
                neuron: neuron.into(),
                actor: Actor::Human,
                agent: None,
                kind: if att > 0 {
                    Kind::Attention
                } else {
                    Kind::Pageview
                },
                navigation: nav,
                hostname: "cyber.page".into(),
                pathname: path.into(),
                referrer: None,
                utm: None,
                attention: (att > 0).then_some(lytics_event::Attention {
                    ms: att,
                    scroll_depth: 0,
                }),
                props: None,
                revenue: None,
                timestamp: ts,
            },
            event_hash: format!("{neuron}-{path}-{ts}"),
            attribution: Attribution {
                source: None,
                channel: Channel::Direct,
            },
            device: Device {
                browser: None,
                browser_version: None,
                os: None,
                device: None,
            },
            geo: None,
            received_at: ts,
        }
    }

    #[test]
    fn passages_split_on_arrivals_never_on_pauses() {
        let events = [
            ev("n1", "/a", 1000, Some(Navigation::External), 0),
            // a six-hour pause โ€” same passage, no timeout exists
            ev(
                "n1",
                "/b",
                1000 + 6 * 3600 * 1000,
                Some(Navigation::Internal),
                0,
            ),
            // a new arrival โ€” new passage
            ev(
                "n1",
                "/c",
                1000 + 7 * 3600 * 1000,
                Some(Navigation::External),
                0,
            ),
        ];
        let refs: Vec<&Stored> = events.iter().collect();
        let p = passages(&refs);
        assert_eq!(p.len(), 2);
        assert_eq!(p[0].views(), 2);
        assert_eq!(p[0].entry(), Some("/a"));
        assert_eq!(p[0].exit(), Some("/b"));
        assert_eq!(p[1].entry(), Some("/c"));
    }

    #[test]
    fn utm_pageview_is_arrival() {
        let mut e = ev("n1", "/a", 1, Some(Navigation::Internal), 0);
        e.body.utm = Some(Utm {
            source: Some("nl".into()),
            medium: None,
            campaign: None,
            term: None,
            content: None,
        });
        assert!(e.is_arrival());
    }

    #[test]
    fn scroll_stats_use_max_reach_per_neuron_path() {
        // n1 looks at /a twice โ€” max 80 wins; n2 reaches 20 on /a
        let mut a1 = ev("n1", "/a", 1000, Some(Navigation::External), 5000);
        a1.body.kind = Kind::Attention;
        a1.body.attention = Some(lytics_event::Attention {
            ms: 5000,
            scroll_depth: 40,
        });
        let mut a2 = ev("n1", "/a", 2000, Some(Navigation::Internal), 3000);
        a2.body.kind = Kind::Attention;
        a2.body.attention = Some(lytics_event::Attention {
            ms: 3000,
            scroll_depth: 80,
        });
        let mut a3 = ev("n2", "/a", 3000, Some(Navigation::External), 2000);
        a3.body.kind = Kind::Attention;
        a3.body.attention = Some(lytics_event::Attention {
            ms: 2000,
            scroll_depth: 20,
        });
        let events = [a1, a2, a3];
        let r: Vec<&Stored> = events.iter().collect();
        let samples = scroll_reach_samples(&r);
        assert_eq!(samples.len(), 2); // two (neuron, path) pairs
        let mut s = samples;
        s.sort_unstable();
        assert_eq!(s, vec![20, 80]);
        let stats = scroll_stats(s);
        assert_eq!(stats["samples"], 2);
        assert_eq!(stats["max"], 80);
        let p50 = stats["p50"].as_u64().unwrap();
        assert!(p50 == 20 || p50 == 80);
        assert_eq!(stats["p90"], 80);
    }

    #[test]
    fn audience_splits_new_vs_returning_by_first_seen() {
        // n1 first-seen before window, active in window โ†’ returning
        // n2 first-seen inside window โ†’ new
        let events = [
            ev("n1", "/a", 100, Some(Navigation::External), 0),
            ev("n1", "/b", 1100, Some(Navigation::Internal), 0),
            ev("n2", "/a", 1200, Some(Navigation::External), 0),
        ];
        let from = 1000u64;
        let to = 2000u64;
        let windowed = in_window(&events, from, to);
        let fs = first_seen(&events);
        assert_eq!(fs.get("n1"), Some(&100));
        assert_eq!(fs.get("n2"), Some(&1200));

        let neu = filter_audience(&windowed, &fs, from, "new");
        let ret = filter_audience(&windowed, &fs, from, "returning");
        assert_eq!(neu.len(), 1);
        assert_eq!(neu[0].body.neuron, "n2");
        assert_eq!(ret.len(), 1);
        assert_eq!(ret[0].body.neuron, "n1");
        assert_eq!(
            filter_audience(&windowed, &fs, from, "all").len(),
            windowed.len()
        );
    }

    #[test]
    fn funnel_counts_ordered_prefixes() {
        let events = [
            ev("n1", "/a", 1, Some(Navigation::External), 0),
            ev("n1", "/b", 2, Some(Navigation::Internal), 0),
            ev("n2", "/a", 1, Some(Navigation::External), 0),
            ev("n3", "/b", 1, Some(Navigation::External), 0), // out of order: /b before /a
        ];
        let refs: Vec<&Stored> = events.iter().collect();
        let f = funnel(&refs, &["/a".into(), "/b".into()]);
        let rows = f.as_array().unwrap();
        assert_eq!(rows[0]["neurons"], 2); // n1, n2 reached /a
        assert_eq!(rows[1]["neurons"], 1); // only n1 continued to /b
    }

    #[test]
    fn retention_places_neurons_in_cohorts() {
        let w = WEEK_MS;
        let events = vec![
            ev("n1", "/", 0, Some(Navigation::External), 0),
            ev("n1", "/", w + 10, Some(Navigation::External), 0), // returns week 1
            ev("n2", "/", 5, Some(Navigation::External), 0),      // never returns
        ];
        let r = retention(&events, 4);
        let rows = r.as_array().unwrap();
        assert_eq!(rows.len(), 1);
        assert_eq!(rows[0]["size"], 2);
        assert_eq!(rows[0]["weeks"][0], 2);
        assert_eq!(rows[0]["weeks"][1], 1);
    }

    #[test]
    fn returns_counts_cohort_and_returned() {
        let events = vec![
            ev("n1", "/", 100, Some(Navigation::External), 0),
            ev("n1", "/", 200, Some(Navigation::External), 0),
            ev("n2", "/", 150, Some(Navigation::External), 0),
        ];
        let r = returns(&events, 0, 1000, 1000);
        assert_eq!(r["cohort"], 2);
        assert_eq!(r["returned"], 1);
    }

    #[test]
    fn attention_sums_into_overview() {
        let events = [
            ev("n1", "/a", 1, Some(Navigation::External), 0),
            ev("n1", "/a", 2, None, 42_000),
        ];
        let refs: Vec<&Stored> = events.iter().collect();
        let o = overview(&refs);
        assert_eq!(o["attention_ms"], 42_000);
        assert_eq!(o["views"], 1);
        assert_eq!(o["neurons"], 1);
    }
}

Graph