goatifi / backend /tests /test_intelligence.py
KonoDioDaa's picture
Initial FlowTwin deployment
e7a9f02
Raw
History Blame Contribute Delete
13.9 kB
"""Crowd state, bottleneck detection, prediction, routing and the strategy engine."""
from __future__ import annotations
import numpy as np
import pytest
from flowtwin.config import SETTINGS
from flowtwin.crowd.density import classify, density, time_to_threshold
from flowtwin.crowd.flow import build_alerts, detect_bottlenecks, primary_bottleneck
from flowtwin.prediction.features import N_FEATURES, build_feature_matrix
from flowtwin.prediction.inference import DensityPredictor
from flowtwin.routing.graph import RoutingTables
from flowtwin.simulation.agents import POLICY_ADAPTIVE, POLICY_SHORTEST, POLICY_STATIC
from flowtwin.simulation.engine import RunOverrides, Simulator
from flowtwin.strategy.engine import StrategyEngine
from flowtwin.strategy.interventions import generate_candidates
from flowtwin.venue import compile_venue, load_scenario
@pytest.fixture(scope="module")
def venue():
return compile_venue("circuit_alpha")
@pytest.fixture(scope="module")
def scenario():
return load_scenario("circuit_alpha_post_race")
@pytest.fixture(scope="module")
def congested(venue, scenario):
"""A simulation held at the point where Exit B is genuinely failing."""
sim = Simulator(venue, scenario, SETTINGS, seed=42193,
overrides=RunOverrides(crowd_size=40000))
sim.run_for(1000)
return sim
# ── crowd state ──────────────────────────────────────────────────────
def test_density_is_people_over_area():
assert density(np.array([100.0]), np.array([50.0]))[0] == pytest.approx(2.0)
assert density(np.array([100.0]), np.array([0.0]))[0] == 0.0
def test_density_levels_are_ordered():
levels = classify(np.array([0.1, 1.2, 2.2, 3.5]), warning=2.0, critical=3.0)
assert levels.tolist() == [0, 1, 2, 3]
def test_time_to_threshold_interpolates():
t = time_to_threshold(1.0, [(30.0, 1.5), (60.0, 2.5)], threshold=2.0)
assert t == pytest.approx(45.0, abs=1.0)
assert time_to_threshold(3.0, [(30.0, 3.5)], 2.0) == 0.0
assert time_to_threshold(1.0, [(30.0, 1.1)], 2.0) is None
def test_state_engine_tracks_flow_and_growth(congested):
st = congested.state
assert st.edge_density.max() > 0
assert st.edge_inflow_ppm.max() > 0
assert st.edge_velocity.max() <= SETTINGS.movement.free_speed_mps + 1e-9
assert np.all(st.edge_risk >= 0) and np.all(st.edge_risk <= 1)
assert st.samples > 100
def test_risk_contributions_sum_to_the_risk_score(congested):
idx = int(np.argmax(congested.state.edge_risk))
parts = congested.state.risk_contributions(
idx, congested.venue.venue.warning_density, congested.venue.venue.critical_density)
assert sum(parts.values()) == pytest.approx(congested.state.edge_risk[idx], abs=0.02)
# ── bottleneck detection ─────────────────────────────────────────────
def test_detects_the_degraded_exit_as_the_primary_bottleneck(congested):
primary = primary_bottleneck(congested)
assert primary is not None
assert primary.base_id == "X_E_EXITB", f"expected Exit B approach, got {primary.base_id}"
assert primary.risk > 0.5
assert primary.causes, "a bottleneck must explain itself"
def test_bottlenecks_are_reported_once_per_physical_corridor(congested):
found = detect_bottlenecks(congested, limit=8)
ids = [b.base_id for b in found]
assert len(ids) == len(set(ids))
def test_alerts_carry_severity_cause_and_lead_time(congested):
predictor = DensityPredictor(SETTINGS)
preds = predictor.predict(congested)
alerts = build_alerts(congested, detect_bottlenecks(congested), preds)
assert alerts, "no alert raised for a failing exit"
top = alerts[0]
assert top["severity"] in {"critical", "warning", "watch"}
assert top["causes"]
assert "projection" in top
def test_a_quiet_network_raises_no_critical_alert(venue, scenario):
sim = Simulator(venue, scenario, SETTINGS, seed=5,
overrides=RunOverrides(crowd_size=2000))
sim.run_for(200)
predictor = DensityPredictor(SETTINGS)
alerts = build_alerts(sim, detect_bottlenecks(sim), predictor.predict(sim))
assert not any(a["severity"] == "critical" for a in alerts)
# ── prediction ───────────────────────────────────────────────────────
def test_feature_matrix_shape_and_sanity(congested):
X = build_feature_matrix(congested)
assert X.shape == (congested.venue.n_edges, N_FEATURES)
assert np.isfinite(X).all()
def test_prediction_produces_horizons_and_lead_time(congested):
predictor = DensityPredictor(SETTINGS)
preds = predictor.predict(congested)
idx = primary_bottleneck(congested).index
row = preds[idx]
assert set(row["horizons"]) == {str(h) for h in SETTINGS.prediction.horizons_s}
assert all(v >= 0 for v in row["horizons"].values())
assert row["source"] in {"trained_model", "analytic_baseline"}
def test_prediction_responds_to_a_change_in_state(venue, scenario):
"""The projection must track the state, not just the recent trend.
Two branches leave the same instant: one keeps the degraded exit, the other
loses more capacity. The physics must respond (measured throughput falls)
and the projection must respond with it.
"""
sim = Simulator(venue, scenario, SETTINGS, seed=42193,
overrides=RunOverrides(crowd_size=40000))
sim.run_for(900)
idx = primary_bottleneck(sim).index
gate = venue.node_index["EXIT_B"]
unchanged = sim.branch()
worse = sim.branch()
worse.node_budget.multiplier[gate] *= 0.4
unchanged.run_for(240)
worse.run_for(240)
assert worse.state.node_throughput_ppm[gate] < unchanged.state.node_throughput_ppm[gate], \
"cutting the gate did not reduce measured throughput"
assert worse.state.node_queue[gate] > unchanged.state.node_queue[gate]
base = DensityPredictor(SETTINGS).predict(unchanged, [idx])[idx]["peak_projected"]
degraded = DensityPredictor(SETTINGS).predict(worse, [idx])[idx]["peak_projected"]
assert degraded > base, (
f"projection did not rise after the exit was cut further ({degraded} vs {base})")
def test_restoring_capacity_raises_measured_throughput(venue, scenario):
"""Opening capacity is a real change to the network, not a label."""
sim = Simulator(venue, scenario, SETTINGS, seed=42193,
overrides=RunOverrides(crowd_size=40000))
sim.run_for(900)
gate = venue.node_index["EXIT_B"]
degraded = sim.branch()
restored = sim.branch()
assert restored.open_alternate("EXIT_B", 4.0)
degraded.run_for(180)
restored.run_for(180)
assert restored.state.node_throughput_ppm[gate] > degraded.state.node_throughput_ppm[gate]
def test_both_directions_of_a_corridor_share_one_projection(congested):
predictor = DensityPredictor(SETTINGS)
proj = predictor.project(congested)
v = congested.venue
for e in range(v.n_edges):
p = int(v.pair_of[e])
if p >= 0:
assert np.allclose(proj[:, e], proj[:, p])
# ── routing ──────────────────────────────────────────────────────────
def test_every_node_can_reach_every_destination(venue, scenario):
sim = Simulator(venue, scenario, SETTINGS, overrides=RunOverrides(crowd_size=500))
for slot in range(len(sim.dest_indices)):
for node in range(venue.n_nodes):
nodes, _ = sim.tables.path_nodes(POLICY_SHORTEST, slot, node)
assert nodes[-1] == sim.dest_indices[slot], \
f"{venue.node_ids[node]} cannot reach {sim.dest_ids[slot]}"
def test_routing_tables_stay_acyclic_under_hysteresis(congested):
for policy in (POLICY_SHORTEST, POLICY_STATIC, POLICY_ADAPTIVE):
for slot, dest in enumerate(congested.dest_indices):
for node in range(congested.venue.n_nodes):
nodes, _ = congested.tables.path_nodes(policy, slot, node)
assert len(nodes) == len(set(nodes)), \
f"cycle in policy {policy} from {congested.venue.node_ids[node]}"
def test_adaptive_routing_avoids_the_congested_asset(congested):
"""The dynamic plan must not still prefer the failing exit."""
branch = congested.branch()
edge = branch.venue.edge_index["X_E_EXITB"]
node = branch.venue.node_index["EXIT_B"]
branch.divert_flow(0.5, {edge, int(branch.venue.pair_of[edge])}, {node})
slot = branch.dest_indices.index(branch.venue.node_index["TRANSPORT_BUS"])
_, edges = branch.tables.path_nodes(POLICY_ADAPTIVE, slot,
branch.venue.node_index["CON_EAST"])
assert edge not in edges, "adaptive plan still routes through the degraded exit"
def test_hysteresis_limits_route_churn(congested):
"""Repeated refreshes on an unchanged state must not keep flipping routes."""
branch = congested.branch()
branch.refresh_routing()
first = branch.tables.next_hop[POLICY_ADAPTIVE].copy()
for _ in range(6):
branch.refresh_routing()
changed = int(np.sum(branch.tables.next_hop[POLICY_ADAPTIVE] != first))
assert changed == 0, f"{changed} next-hops flapped without any state change"
# ── strategy engine ──────────────────────────────────────────────────
def test_candidates_are_generated_from_topology(congested):
bn = primary_bottleneck(congested)
cands = generate_candidates(congested, bn)
ids = [c.id for c in cands]
assert "no_action" in ids
assert sum(1 for i in ids if i.startswith("reroute_")) >= 3
assert len(ids) >= 5
assert len(ids) == len(set(ids))
for c in cands:
assert c.description and c.instruction
def test_counterfactuals_all_start_from_the_same_state(congested):
"""Two evaluations of the same strategy from the same state must agree."""
predictor = DensityPredictor(SETTINGS)
engine = StrategyEngine(SETTINGS, predictor)
a = engine.evaluate(congested, horizon_s=120, strategy_ids=["reroute_30"])
b = engine.evaluate(congested, horizon_s=120, strategy_ids=["reroute_30"])
ma = {s["id"]: s["metrics"] for s in a["strategies"]}
mb = {s["id"]: s["metrics"] for s in b["strategies"]}
assert ma == mb
def test_evaluation_does_not_advance_the_live_simulation(congested):
predictor = DensityPredictor(SETTINGS)
engine = StrategyEngine(SETTINGS, predictor)
t_before = congested.time
pos_before = congested.pop.pos_m.copy()
engine.evaluate(congested, horizon_s=120, strategy_ids=["reroute_20"])
assert congested.time == t_before
assert np.array_equal(congested.pop.pos_m, pos_before)
def test_recommendation_beats_no_action_on_the_score(congested):
engine = StrategyEngine(SETTINGS, DensityPredictor(SETTINGS))
result = engine.evaluate(congested, horizon_s=240)
assert result["available"]
by_id = {s["id"]: s for s in result["strategies"]}
winner = result["recommendation"]["strategy_id"]
assert by_id[winner]["score"] <= by_id["no_action"]["score"]
assert by_id[winner]["recommended"] is True
assert by_id[winner]["metrics"]["peak_density"] <= by_id["no_action"]["metrics"]["peak_density"]
def test_explanation_uses_measured_values(congested):
engine = StrategyEngine(SETTINGS, DensityPredictor(SETTINGS))
result = engine.evaluate(congested, horizon_s=240)
rec = result["recommendation"]
by_id = {s["id"]: s for s in result["strategies"]}
winner, baseline = by_id[rec["strategy_id"]], by_id["no_action"]
for reason in rec["reasons"]:
key = reason["metric"]
attr = {"peak_density": "peak_density",
"critical_duration": "critical_duration_s",
"avg_travel_time": "avg_travel_time_s",
"aggregate_risk": "aggregate_risk",
"max_queue": "max_queue",
"throughput": "throughput"}[key]
assert reason["value"] == pytest.approx(winner["metrics"][attr], abs=0.02)
assert reason["baseline"] == pytest.approx(baseline["metrics"][attr], abs=0.02)
def test_applying_a_strategy_changes_the_live_simulation(congested):
branch = congested.branch()
engine = StrategyEngine(SETTINGS, DensityPredictor(SETTINGS))
result = engine.apply(branch, "reroute_30")
assert result["applied"]
assert result["agents_affected"] > 0
assert branch.applied_interventions
assert np.sum(branch.pop.policy == POLICY_ADAPTIVE) > 0
def test_unknown_strategy_is_refused(congested):
engine = StrategyEngine(SETTINGS, DensityPredictor(SETTINGS))
result = engine.apply(congested.branch(), "teleport_everyone")
assert result["applied"] is False
def test_recommendation_changes_with_the_scenario(venue):
"""A different failure must not produce the same canned answer."""
engine = StrategyEngine(SETTINGS, DensityPredictor(SETTINGS))
egress = Simulator(venue, load_scenario("circuit_alpha_post_race"), SETTINGS,
seed=42193, overrides=RunOverrides(crowd_size=40000))
egress.run_for(1000)
a = engine.evaluate(egress, horizon_s=180)
arrival = Simulator(venue, load_scenario("circuit_alpha_arrival"), SETTINGS,
seed=7717, overrides=RunOverrides(crowd_size=26000))
arrival.run_for(900)
b = engine.evaluate(arrival, horizon_s=180)
assert a["bottleneck"]["base_id"] != b["bottleneck"]["base_id"], \
"the two scenarios were expected to fail in different places"