From b60217c0dfede95ec0536171ec033183b9251b5e Mon Sep 17 00:00:00 2001 From: Mateo Date: Thu, 24 Sep 2026 12:58:15 +0200 Subject: [PATCH 1/3] feat: add replay compose override and simulated clock helpers --- .gitignore | 4 +++ docker-compose.replay.yml | 23 ++++++++++++++ replay_control/.gitkeep | 0 scripts/replay_clock.py | 64 +++++++++++++++++++++++++++++++++++++++ 4 files changed, 91 insertions(+) create mode 100644 docker-compose.replay.yml create mode 100644 replay_control/.gitkeep create mode 100644 scripts/replay_clock.py diff --git a/.gitignore b/.gitignore index 50c0312..39766fd 100644 --- a/.gitignore +++ b/.gitignore @@ -155,3 +155,7 @@ containers/notebooks/app/.Trash-0 data/alert_samples/ data/triangulated_sequences/ + +# day replay runtime files +replay_control/*.json +.replay_map.jsonl diff --git a/docker-compose.replay.yml b/docker-compose.replay.yml new file mode 100644 index 0000000..12da1b7 --- /dev/null +++ b/docker-compose.replay.yml @@ -0,0 +1,23 @@ +# Override for accelerated day replays (scripts/replay_day.py). +# +# Clock mode (recommended, scripts/replay_day.py --clock): the API runs with a +# simulated accelerated clock (patch on branch replay-clock-20260710 of pyro-api). +# created_at values land directly on the replayed day's hours and every +# utcnow-anchored window behaves exactly as in production — no window scaling. +# The three REPLAY_CLOCK_* variables are exported by replay_day.py --clock when +# it restarts pyro_api. +# +# Usage: +# docker compose -f docker-compose.yml -f docker-compose.replay.yml up -d ... +services: + pyro_api: + environment: + - REPLAY_CLOCK_SPEED=${REPLAY_CLOCK_SPEED:-} + - REPLAY_CLOCK_SIM_ORIGIN=${REPLAY_CLOCK_SIM_ORIGIN:-} + - REPLAY_CLOCK_REAL_ORIGIN=${REPLAY_CLOCK_REAL_ORIGIN:-} + # interactive mode: piecewise clock controlled live via scripts/replay_ctl.py + - REPLAY_CLOCK_FILE=/replay/clock.json + # disable the temporal-model gate for replays (fail-open, era-faithful) + - TEMPORAL_API_URL= + volumes: + - ./replay_control:/replay:ro diff --git a/replay_control/.gitkeep b/replay_control/.gitkeep new file mode 100644 index 0000000..e69de29 diff --git a/scripts/replay_clock.py b/scripts/replay_clock.py new file mode 100644 index 0000000..0d23ee7 --- /dev/null +++ b/scripts/replay_clock.py @@ -0,0 +1,64 @@ +"""Shared helpers for the file-driven replay clock (clock.json segments).""" + +import json +from datetime import datetime, timezone +from pathlib import Path + +ENVDEV_ROOT = Path(__file__).resolve().parent.parent +CONTROL_DIR = ENVDEV_ROOT / "replay_control" +CLOCK_FILE = CONTROL_DIR / "clock.json" +STEPS_FILE = CONTROL_DIR / "steps.json" + + +def real_utcnow(): + return datetime.now(timezone.utc).replace(tzinfo=None) + + +def read_segments(): + try: + data = json.loads(CLOCK_FILE.read_text()) + except (OSError, json.JSONDecodeError): + return [] + segs = [ + (datetime.fromisoformat(s["real0"]), datetime.fromisoformat(s["sim0"]), float(s["speed"])) + for s in data.get("segments", []) + ] + return sorted(segs, key=lambda s: s[0]) + + +def write_segments(segments): + CONTROL_DIR.mkdir(exist_ok=True) + payload = { + "segments": [ + {"real0": r.isoformat(), "sim0": s.isoformat(), "speed": sp} for r, s, sp in segments + ] + } + tmp = CLOCK_FILE.with_suffix(".tmp") + tmp.write_text(json.dumps(payload, indent=1)) + tmp.replace(CLOCK_FILE) + + +def sim_now(segments=None, now=None): + segments = read_segments() if segments is None else segments + now = now or real_utcnow() + if not segments: + return now, 1.0 + active = None + for seg in segments: + if seg[0] <= now: + active = seg + else: + break + if active is None: + return segments[0][1], 0.0 + real0, sim0, speed = active + return sim0 + (now - real0) * speed, speed + + +def append_segment(sim0, speed, real0=None): + """Anchor a new segment at `real0` (default: now) starting from sim time `sim0`.""" + segments = read_segments() + real0 = real0 or real_utcnow() + segments = [s for s in segments if s[0] < real0] + segments.append((real0, sim0, float(speed))) + write_segments(segments) From b8471150a32073b1f7f1215fe5238e733628e89b Mon Sep 17 00:00:00 2001 From: Mateo Date: Thu, 24 Sep 2026 12:58:15 +0200 Subject: [PATCH 2/3] feat: add day replay driver and clock controller --- scripts/live_retimer.py | 115 +++++++++++ scripts/replay_ctl.py | 139 +++++++++++++ scripts/replay_day.py | 437 ++++++++++++++++++++++++++++++++++++++++ 3 files changed, 691 insertions(+) create mode 100644 scripts/live_retimer.py create mode 100644 scripts/replay_ctl.py create mode 100644 scripts/replay_day.py diff --git a/scripts/live_retimer.py b/scripts/live_retimer.py new file mode 100644 index 0000000..9a079f4 --- /dev/null +++ b/scripts/live_retimer.py @@ -0,0 +1,115 @@ +#!/usr/bin/env python3 +"""Companion of replay_day.py: shifts replayed data to the real prod timestamps *while +the replay runs*, so the platform's historical-day page fills in live. + +Rewriting a row that is still inside the API's grouping windows would break sequence +matching and triangulation, so a sequence is only retimed once it is complete (all its +detections posted) AND idle for --age seconds — past every scaled window. Alerts have +their bounds recomputed from their sequences on every pass (idempotent), so an alert's +started_at flips to the historical time as soon as its first sequence is retimed. + +Run in parallel with the replay: + python3 scripts/replay_day.py --day 2026-07-10 --speed 60 --no-restore-times & + python3 scripts/live_retimer.py + +Exits after a final full pass when the replay writes its end-marker. +""" + +import argparse +import json +import subprocess +import time +from pathlib import Path + +ENVDEV_ROOT = Path(__file__).resolve().parent.parent + +CASCADE_SQL = """ +UPDATE sequences s SET started_at = sub.mn, last_seen_at = sub.mx +FROM (SELECT sequence_id, MIN(created_at) mn, MAX(created_at) mx + FROM detections WHERE sequence_id IN ({seq_subquery}) GROUP BY sequence_id) sub +WHERE s.id = sub.sequence_id; +UPDATE alerts a SET started_at = sub.mn, last_seen_at = sub.mx +FROM (SELECT asq.alert_id, MIN(s.started_at) mn, MAX(s.last_seen_at) mx + FROM alerts_sequences asq JOIN sequences s ON s.id = asq.sequence_id + GROUP BY asq.alert_id) sub +WHERE a.id = sub.alert_id; +""" + + +def psql(sql): + return subprocess.run( + ["docker", "compose", "exec", "-T", "db", "psql", + "-U", "dummy_pg_user", "-d", "dummy_pg_db", "-v", "ON_ERROR_STOP=1", "-f", "-"], + input=sql, text=True, cwd=ENVDEV_ROOT, capture_output=True, + ) + + +def retime(entries): + values = ",".join(f"({e['id']},'{e['ts']}')" for e in entries) + ids = ",".join(str(e["id"]) for e in entries) + seq_subquery = f"SELECT DISTINCT sequence_id FROM detections WHERE id IN ({ids}) AND sequence_id IS NOT NULL" + sql = ( + "BEGIN;\n" + f"UPDATE detections d SET created_at = v.ts::timestamp FROM (VALUES {values}) AS v(id, ts) WHERE d.id = v.id;\n" + + CASCADE_SQL.format(seq_subquery=seq_subquery) + + "COMMIT;\n" + ) + return psql(sql) + + +def main(): + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--map-file", type=Path, default=ENVDEV_ROOT / ".replay_map.jsonl") + ap.add_argument("--age", type=float, default=150.0, + help="seconds a complete sequence must stay idle before retiming (default 150)") + ap.add_argument("--interval", type=float, default=15.0) + args = ap.parse_args() + + while not args.map_file.exists(): + time.sleep(1) + + done_marker = False + retimed_seqs = set() + print(f"⏱ live retimer démarré (age {args.age:g}s, passe toutes les {args.interval:g}s)") + while True: + by_seq = {} + for line in args.map_file.read_text().splitlines(): + try: + e = json.loads(line) + except json.JSONDecodeError: + continue + if e.get("done"): + done_marker = True + continue + by_seq.setdefault(e["seq"], []).append(e) + + now = time.time() + batch, batch_seqs = [], [] + for seq, entries in by_seq.items(): + if seq in retimed_seqs: + continue + complete = len(entries) >= entries[0]["total"] + idle = now - max(e["posted"] for e in entries) + if done_marker or (complete and idle >= args.age): + batch.extend(entries) + batch_seqs.append(seq) + + if batch: + proc = retime(batch) + if proc.returncode != 0: + print(f" ✗ retiming KO: {proc.stderr.strip()[:300]}") + else: + retimed_seqs.update(batch_seqs) + first = min(e["ts"] for e in batch)[11:19] + last = max(e["ts"] for e in batch)[11:19] + print(f" ✓ {len(batch_seqs)} séquence(s) re-datée(s) ({len(batch)} détections, {first}->{last} UTC)" + f" — total {len(retimed_seqs)} séquences", flush=True) + + if done_marker: + print("🏁 replay terminé, dernière passe faite — retimer stoppé") + break + time.sleep(args.interval) + + +if __name__ == "__main__": + main() diff --git a/scripts/replay_ctl.py b/scripts/replay_ctl.py new file mode 100644 index 0000000..954b991 --- /dev/null +++ b/scripts/replay_ctl.py @@ -0,0 +1,139 @@ +#!/usr/bin/env python3 +"""Live controller for interactive day replays (replay_day.py --ctl). + +Drives the simulated clock that both the API and the replay driver follow, by +appending segments to replay_control/clock.json. Everything reacts instantly — +no restarts. + +Commands: + status current simulated time, speed, and the next steps + pause freeze the simulation + play [SPEED] resume at SPEED x (default 60) + slow [SPEED] shorthand for a watchable pace (default 10) + next [SPEED] fast-forward to 3 simulated minutes after the next + sequence start (its images are already visible), then + PAUSE (default). Pass SPEED to resume instead, e.g. + `next 10`. + goto HH:MM [SPEED] fast-forward to a simulated time of day, then PAUSE + (default). Pass SPEED to resume playing instead, e.g. + `goto 12h00 10`. HH:MM is PARIS time (what the platform + displays); accepts 12:00, 12h00 or 12h. + +Examples: + python3 scripts/replay_ctl.py status + python3 scripts/replay_ctl.py next # saute au prochain épisode, 10x + python3 scripts/replay_ctl.py play 60 # reprend la journée à 60x +""" + +import json +import sys +from datetime import datetime, timedelta +from zoneinfo import ZoneInfo + +from replay_clock import STEPS_FILE, append_segment, read_segments, real_utcnow, sim_now + +# Bounded so the per-camera posting queue never lags the clock: at higher speeds +# the timestamps of in-flight frames bunch up at the end of the seek and the +# bbox chain can break, splitting a sequence in two. 60x is the fastest pace +# validated against prod (full-day run, zero splits). +SEEK_SPEED = 60.0 +PARIS = ZoneInfo("Europe/Paris") + + +def paris(dt): + return dt.replace(tzinfo=ZoneInfo("UTC")).astimezone(PARIS).strftime("%Hh%Mm%Ss") + + +def load_steps(): + try: + return [ + {"sim_t": datetime.fromisoformat(s["sim_t"]), "label": s["label"]} + for s in json.loads(STEPS_FILE.read_text()) + ] + except (OSError, json.JSONDecodeError): + return [] + + +def cmd_status(): + sim, speed = sim_now() + state = "⏸ pause" if speed == 0 else f"▶ {speed:g}x" + print(f"heure simulée: {paris(sim)} (Paris) | {state}") + upcoming = [s for s in load_steps() if s["sim_t"] > sim][:5] + if upcoming: + print("prochaines étapes:") + for s in upcoming: + print(f" {paris(s['sim_t'])} {s['label']}") + else: + print("plus d'étape à venir") + + +def seek_to(target, cruise): + sim, _ = sim_now() + if target <= sim: + print(f"déjà passé ({paris(target)} <= {paris(sim)})") + return + now = real_utcnow() + seek_real_duration = (target - sim) / SEEK_SPEED + append_segment(sim, SEEK_SPEED, real0=now) + append_segment(target, cruise, real0=now + seek_real_duration) + then = "pause" if cruise == 0 else f"{cruise:g}x" + print( + f"⏩ avance rapide {paris(sim)} -> {paris(target)}" + f" ({seek_real_duration.total_seconds():.0f} s réelles), puis {then}" + ) + + +def main(): + args = sys.argv[1:] + # tolerate `goto status`, `goto pause`...: the command word wins + if len(args) >= 2 and args[0] == "goto" and args[1] in ("status", "pause", "play", "slow", "next"): + args = args[1:] + cmd = args[0] if args else "status" + if cmd == "status": + cmd_status() + elif cmd == "pause": + sim, _ = sim_now() + append_segment(sim, 0) + print(f"⏸ pause à {paris(sim)}") + elif cmd in ("play", "slow"): + default = 60.0 if cmd == "play" else 10.0 + speed = float(args[1]) if len(args) > 1 else default + sim, _ = sim_now() + append_segment(sim, speed) + print(f"▶ {speed:g}x depuis {paris(sim)}") + elif cmd == "next": + cruise = float(args[1]) if len(args) > 1 else 0.0 + sim, _ = sim_now() + upcoming = [s for s in load_steps() if s["sim_t"] > sim] + if not upcoming: + print("plus d'étape à venir") + return + step = upcoming[0] + print(f"étape suivante: {step['label']}") + # land 3 simulated minutes AFTER the sequence start, so its first images + # are already ingested and visible on the platform + seek_to(step["sim_t"] + timedelta(minutes=3), cruise) + elif cmd == "goto": + if len(args) < 2: + sys.exit("usage: goto HH:MM [SPEED] (heure de Paris, ex: goto 12h00)") + m_ = __import__("re").match(r"^(\d{1,2})[:hH](\d{0,2})$", args[1]) + if not m_: + sys.exit( + f"heure invalide: {args[1]} (attendu HH:MM ou HHhMM)\n" + "rappel: status/pause/play/slow/next sont des commandes directes," + " ex: python3 scripts/replay_ctl.py status" + ) + h, m = int(m_.group(1)), int(m_.group(2) or 0) + cruise = float(args[2]) if len(args) > 2 else 0.0 + sim, _ = sim_now() + # HH:MM is Paris local; the simulated clock is naive UTC + paris_now = sim.replace(tzinfo=ZoneInfo("UTC")).astimezone(PARIS) + target_paris = paris_now.replace(hour=h, minute=m, second=0, microsecond=0) + target = target_paris.astimezone(ZoneInfo("UTC")).replace(tzinfo=None) + seek_to(target, cruise) + else: + sys.exit(__doc__) + + +if __name__ == "__main__": + main() diff --git a/scripts/replay_day.py b/scripts/replay_day.py new file mode 100644 index 0000000..9c66d34 --- /dev/null +++ b/scripts/replay_day.py @@ -0,0 +1,437 @@ +#!/usr/bin/env python3 +"""Replay a real production day against the local pyro-envdev API, time-accelerated. + +Reads the day dump produced by the analyse77 download pipeline +(//alert_*/sequence_*/{sequence.json,detections.json,images/*.jpg}), +maps prod cameras/poses onto the local dev environment (creating missing poses +with the prod azimuths), then re-posts every detection (image + bboxes) in +chronological order, with all inter-detection gaps divided by --speed. + +Two modes: + +--clock (recommended): the API image must carry the simulated-clock patch + (pyro-api branch replay-clock-20260710). The script restarts pyro_api with an + accelerated clock anchored on TODAY at the day's real hours: detections land in + the DB with the historical times-of-day, the platform's live view shows the day + unfolding, and every utcnow-anchored window behaves exactly as in production + (no window scaling needed). + +legacy mode (no --clock): posts at "now" and optionally rewrites timestamps to + the historical date at the end (see --no-restore-times / live_retimer.py); + requires the scaled-window overrides. + +Usage: + python3 scripts/replay_day.py --day 2026-07-10 --speed 60 --clock +""" + +import argparse +import ast +import json +import os +import re +import subprocess +import sys +import threading +import time +from concurrent.futures import ThreadPoolExecutor +from datetime import date, datetime, timedelta, timezone +from pathlib import Path + +import requests + +sys.path.insert(0, str(Path(__file__).resolve().parent)) +from replay_clock import CLOCK_FILE, STEPS_FILE, append_segment, sim_now, write_segments # noqa: E402 + +ENVDEV_ROOT = Path(__file__).resolve().parent.parent + +DEFAULT_DATA = Path.home() / "pyronear/test/analyse77" +DEFAULT_API = "http://localhost:5050" + + +def api_login(api, login, pwd): + r = requests.post(f"{api}/api/v1/login/creds", data={"username": login, "password": pwd}, timeout=30) + r.raise_for_status() + return r.json()["access_token"] + + +def auth(token): + return {"Authorization": f"Bearer {token}"} + + +def parse_bboxes(s): + if not s: + return [] + try: + v = ast.literal_eval(s) + boxes = [tuple(b) for b in v] if isinstance(v, (list, tuple)) else [] + except (ValueError, SyntaxError): + return [] + # drop degenerate boxes (e.g. the (0,0,0,0,0) prod placeholder): the API 422s on xmin>=xmax + return [b for b in boxes if round(b[0], 3) < round(b[2], 3) and round(b[1], 3) < round(b[3], 3)] + + +def _fmt_float(x): + # API FLOAT_PATTERN accepts `0`, `1` or `0.xxx` (max 3 decimals) — never `1.0`/`0.0` + s = f"{round(float(x), 3):.3f}".rstrip("0").rstrip(".") + return s or "0" + + +def fmt_bboxes(boxes): + return "[" + ",".join("(" + ",".join(_fmt_float(x) for x in b) + ")" for b in boxes) + "]" + + +def load_day(data_dir, day): + """Collect unique sequences and their detections + image paths.""" + day_dir = data_dir / day + if not day_dir.is_dir(): + sys.exit(f"day folder not found: {day_dir}") + sequences = {} # prod_seq_id -> {"meta":..., "dets":[...]} + for sdir in sorted(day_dir.glob("alert_*/sequence_*")): + meta = json.load(open(sdir / "sequence.json")) + if meta["id"] in sequences: + continue # sequence shared by several alerts: keep first folder + dets = json.load(open(sdir / "detections.json")) + img_by_det = {} + for img in (sdir / "images").glob("*.jpg"): + m = re.search(r"det(\d+)", img.name) + if m: + img_by_det[int(m.group(1))] = img + sequences[meta["id"]] = {"meta": meta, "dets": dets, "imgs": img_by_det} + return sequences + + +def _recreate_api(env): + subprocess.run( + ["docker", "compose", "-f", "docker-compose.yml", "-f", "docker-compose.replay.yml", + "up", "-d", "--force-recreate", "pyro_api"], + cwd=ENVDEV_ROOT, env=env, check=True, capture_output=True, text=True, + ) + deadline = time.time() + 120 + while time.time() < deadline: + try: + if requests.get(f"{DEFAULT_API}/status", timeout=3).ok: + print(" API prête") + return + except requests.RequestException: + pass + time.sleep(2) + sys.exit("l'API n'est pas revenue après le redémarrage") + + +def restart_api_with_clock(speed, sim0, real0): + """Recreate pyro_api with the fixed simulated-clock env, then wait for health.""" + CLOCK_FILE.unlink(missing_ok=True) # the clock file would take priority over the env + env = { + **os.environ, + "REPLAY_CLOCK_SPEED": f"{speed:g}", + "REPLAY_CLOCK_SIM_ORIGIN": sim0.isoformat(), + "REPLAY_CLOCK_REAL_ORIGIN": real0.isoformat(), + } + print(f"🕰 redémarrage de l'API avec l'horloge simulée: {sim0.isoformat()} @ {speed:g}x") + _recreate_api(env) + + +def restart_api_plain(): + """Recreate pyro_api without fixed-clock env: it follows replay_control/clock.json.""" + print("🕰 redémarrage de l'API sur l'horloge pilotable (replay_control/clock.json)") + _recreate_api({**os.environ}) + + +def main(): + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--day", required=True, help="e.g. 2026-07-10") + ap.add_argument("--speed", type=float, default=60.0, help="time compression factor (default 60)") + ap.add_argument("--api", default=DEFAULT_API) + ap.add_argument("--data", type=Path, default=DEFAULT_DATA) + ap.add_argument("--login", default="mateo", help="local superadmin login (.env)") + ap.add_argument("--pwd", default="mateo", help="local superadmin password (.env)") + ap.add_argument("--view-login", default="test77", help="org agent used for the final summary") + ap.add_argument("--view-pwd", default="test") + ap.add_argument("--clock", action="store_true", + help="simulated-clock mode: today's date + the day's real hours (recommended)") + ap.add_argument("--ctl", action="store_true", + help="interactive mode: file-driven clock, starts PAUSED; drive it with" + " scripts/replay_ctl.py (pause/play/slow/next/goto)") + ap.add_argument("--dry-run", action="store_true", help="map + schedule only, no POSTs") + ap.add_argument( + "--no-restore-times", action="store_true", + help="legacy mode only: skip rewriting DB timestamps to the prod times after the replay", + ) + ap.add_argument( + "--map-file", type=Path, default=ENVDEV_ROOT / ".replay_map.jsonl", + help="legacy mode only: mapping file consumed by scripts/live_retimer.py", + ) + args = ap.parse_args() + + sequences = load_day(args.data, args.day) + first_t = min( + datetime.fromisoformat(det["created_at"]) for s in sequences.values() for det in s["dets"] + ) + + real0 = sim0 = None + if args.ctl: + args.clock = True + # today's date, the day's real hours; simulation starts PAUSED just before the + # first detection — advance it with scripts/replay_ctl.py + sim0 = datetime.combine(date.today(), first_t.time()) - timedelta(seconds=30) + if not args.dry_run: + write_segments([(datetime.now(timezone.utc).replace(tzinfo=None), sim0, 0.0)]) + print(f"🕹 mode interactif: horloge en PAUSE à {sim0.isoformat()} — pilotez avec" + f" scripts/replay_ctl.py (status/pause/play/slow/next/goto)") + restart_api_plain() + elif args.clock: + # today's date, the day's real hours + sim0 = datetime.combine(date.today(), first_t.time()) - timedelta(seconds=30) + real0 = datetime.now(timezone.utc).replace(tzinfo=None) + timedelta(seconds=45) + if not args.dry_run: + restart_api_with_clock(args.speed, sim0, real0) + else: + print(f"⚙️ For --speed {args.speed:g}, the API must run with scaled windows, e.g.:") + for var, default in ( + ("SEQUENCE_MIN_INTERVAL_SECONDS", 300), + ("SEQUENCE_RELAXATION_SECONDS", 7200), + ("TRIANGULATION_RELAXATION_SECONDS", 1800), + ): + print(f" {var}={max(1, round(default / args.speed))}") + print(" (see docker-compose.replay.yml)\n") + + admin = api_login(args.api, args.login, args.pwd) + + # Map prod camera names -> local ids + r = requests.get(f"{args.api}/api/v1/cameras/", headers=auth(admin), timeout=30) + r.raise_for_status() + local_cams = {c["name"]: c["id"] for c in r.json()} + + prod_cams = {} # prod camera id -> name (from analyse77 cameras.json) + cam_file = args.data / args.day / "cameras.json" + if not cam_file.exists(): + cam_file = args.data / "cameras.json" + for c in json.load(open(cam_file)): + prod_cams[c["id"]] = c["name"] + + # Camera tokens + cam_tokens = {} + for name, cid in local_cams.items(): + if name not in {prod_cams[s["meta"]["camera_id"]] for s in sequences.values()}: + continue + r = requests.post(f"{args.api}/api/v1/cameras/{cid}/token", headers=auth(admin), timeout=30) + r.raise_for_status() + cam_tokens[name] = r.json()["access_token"] + + # Ensure poses: prod (camera, pose_id) -> local pose id with same azimuth + pose_map = {} # (cam_name, prod_pose_id) -> local_pose_id + local_poses = {} # cam_name -> {azimuth: id} + for name, token in cam_tokens.items(): + r = requests.get(f"{args.api}/api/v1/poses/", headers=auth(token), timeout=30) + r.raise_for_status() + local_poses[name] = {round(float(p["azimuth"]), 1): p["id"] for p in r.json()} + for seq in sequences.values(): + m = seq["meta"] + name = prod_cams[m["camera_id"]] + key = (name, m["pose_id"]) + if key in pose_map: + continue + az = round(float(m["camera_azimuth"]), 1) + if az in local_poses[name]: + pose_map[key] = local_poses[name][az] + continue + payload = {"camera_id": local_cams[name], "azimuth": az, "patrol_id": m["pose_id"]} + r = requests.post(f"{args.api}/api/v1/poses/", headers=auth(admin), json=payload, timeout=30) + r.raise_for_status() + pose_map[key] = r.json()["id"] + local_poses[name][az] = pose_map[key] + print(f"pose créée: {name} az={az}° (prod pose {m['pose_id']}) -> local pose {pose_map[key]}") + + # Build the chronological event list + events = [] + for seq in sequences.values(): + m = seq["meta"] + name = prod_cams[m["camera_id"]] + for det in seq["dets"]: + img = seq["imgs"].get(det["id"]) + if img is None: + continue + boxes = parse_bboxes(det.get("bbox")) + parse_bboxes(det.get("others_bboxes")) + if not boxes: + continue + t = datetime.fromisoformat(det["created_at"]) + events.append({ + "t": t, + "sim_t": datetime.combine(date.today(), t.time()) if args.clock else t, + "cam": name, + "pose": pose_map[(name, m["pose_id"])], + "bboxes": fmt_bboxes(boxes), + "img": img, + "seq": m["id"], + }) + events.sort(key=lambda e: e["t"]) + if not events: + sys.exit("no detections to replay") + if args.ctl and not args.dry_run: + # one step per sequence start, for replay_ctl.py `next` + seq_count = {} + for ev in events: + seq_count[ev["seq"]] = seq_count.get(ev["seq"], 0) + 1 + steps, seen = [], set() + for ev in events: + if ev["seq"] in seen: + continue + seen.add(ev["seq"]) + steps.append({ + "sim_t": ev["sim_t"].isoformat(), + "label": f"seq {ev['seq']} {ev['cam']} ({seq_count[ev['seq']]} dets)", + }) + STEPS_FILE.write_text(json.dumps(steps, indent=1)) + t0, tn = events[0]["t"], events[-1]["t"] + real = (tn - t0).total_seconds() + print(f"\n▶️ {len(events)} détections de {len(sequences)} séquences") + print(f" journée réelle {t0:%H:%M:%S}->{tn:%H:%M:%S} UTC ({real / 3600:.1f} h)" + f" -> replay ~{real / args.speed / 60:.1f} min à {args.speed:g}x") + if args.clock: + print(f" heures simulées sur AUJOURD'HUI ({date.today().isoformat()}) — la vue live" + " de la plateforme suit la journée\n") + if args.dry_run: + return + + stats = {"ok": 0, "err": 0} + id_times = [] # (local_detection_id, prod_created_at_iso) — legacy retiming + id_lock = threading.Lock() + seq_totals = {} + for ev in events: + seq_totals[ev["seq"]] = seq_totals.get(ev["seq"], 0) + 1 + map_fh = None if args.clock else open(args.map_file, "w") # legacy: live_retimer.py + + def post(ev): + try: + with open(ev["img"], "rb") as f: + r = requests.post( + f"{args.api}/api/v1/detections/", + headers=auth(cam_tokens[ev["cam"]]), + data={"bboxes": ev["bboxes"], "pose_id": ev["pose"]}, + files={"file": (ev["img"].name, f, "image/jpeg")}, + timeout=60, + ) + if r.status_code == 201: + stats["ok"] += 1 + with id_lock: + id_times.append((r.json()["id"], ev["t"].isoformat())) + if map_fh is not None: + map_fh.write(json.dumps({ + "id": r.json()["id"], "ts": ev["t"].isoformat(), "seq": ev["seq"], + "posted": time.time(), "total": seq_totals[ev["seq"]], + }) + "\n") + map_fh.flush() + else: + stats["err"] += 1 + print(f" ✗ seq {ev['seq']} {ev['img'].name}: {r.status_code} {r.text[:120]}") + except Exception as e: # noqa: BLE001 + stats["err"] += 1 + print(f" ✗ seq {ev['seq']} {ev['img'].name}: {e}") + + # One single-worker executor per camera: a real camera posts sequentially, and + # concurrent same-camera posts race the API's sequence creation (duplicate and + # orphan sequences). Cross-camera concurrency stays. + cam_pools = {name: ThreadPoolExecutor(max_workers=1) for name in cam_tokens} + + if args.ctl: + # sim-time-driven loop: follows the controllable clock (pause/seek aware) + i = 0 + last_report = 0.0 + while i < len(events): + sim, speed = sim_now() + while i < len(events) and events[i]["sim_t"] <= sim: + cam_pools[events[i]["cam"]].submit(post, events[i]) + i += 1 + if time.time() - last_report > 30: + state = "pause" if speed == 0 else f"{speed:g}x" + print(f" {i}/{len(events)} envoyées | heure simulée {sim:%H:%M:%S} UTC ({state})" + f" | ok={stats['ok']} err={stats['err']}") + last_report = time.time() + time.sleep(0.25) + for pool in cam_pools.values(): + pool.shutdown(wait=True) + else: + if args.clock: + real0_epoch = real0.replace(tzinfo=timezone.utc).timestamp() + + def target_epoch(ev): + return real0_epoch + (ev["sim_t"] - sim0).total_seconds() / args.speed + else: + start = time.time() + + def target_epoch(ev): + return start + (ev["t"] - t0).total_seconds() / args.speed + + for i, ev in enumerate(events): + delay = target_epoch(ev) - time.time() + if delay > 0: + time.sleep(delay) + cam_pools[ev["cam"]].submit(post, ev) + if (i + 1) % 200 == 0: + print(f" {i + 1}/{len(events)} envoyées (heure simulée {ev['t']:%H:%M} UTC," + f" ok={stats['ok']} err={stats['err']})") + for pool in cam_pools.values(): + pool.shutdown(wait=True) + + if args.ctl: + # settle the clock at real-time pace so the day stays visible in the live view + sim, _ = sim_now() + append_segment(sim, 1.0) + print(f"🕰 fin des événements: horloge stabilisée à 1x ({sim:%H:%M:%S} UTC simulé)") + + print(f"\n⏳ envoi terminé ({stats['ok']} ok, {stats['err']} erreurs), attente validation (15 s)...") + time.sleep(15) + if map_fh is not None: + with id_lock: + map_fh.write(json.dumps({"done": True}) + "\n") + map_fh.close() + + result_date = date.today().isoformat() + if not args.clock and not args.no_restore_times and id_times: + print(f"🕐 réécriture des vraies heures prod dans la base locale ({len(id_times)} détections)...") + values = ",".join(f"({i},'{ts}')" for i, ts in id_times) + sql = f""" +BEGIN; +UPDATE detections d SET created_at = v.ts::timestamp +FROM (VALUES {values}) AS v(id, ts) WHERE d.id = v.id; +UPDATE sequences s SET started_at = sub.mn, last_seen_at = sub.mx +FROM (SELECT sequence_id, MIN(created_at) mn, MAX(created_at) mx + FROM detections WHERE sequence_id IS NOT NULL GROUP BY sequence_id) sub +WHERE s.id = sub.sequence_id; +UPDATE alerts a SET started_at = sub.mn, last_seen_at = sub.mx +FROM (SELECT asq.alert_id, MIN(s.started_at) mn, MAX(s.last_seen_at) mx + FROM alerts_sequences asq JOIN sequences s ON s.id = asq.sequence_id + GROUP BY asq.alert_id) sub +WHERE a.id = sub.alert_id; +COMMIT; +""" + proc = subprocess.run( + ["docker", "compose", "exec", "-T", "db", "psql", + "-U", "dummy_pg_user", "-d", "dummy_pg_db", "-v", "ON_ERROR_STOP=1", "-f", "-"], + input=sql, text=True, cwd=ENVDEV_ROOT, capture_output=True, + ) + if proc.returncode != 0: + print(f" ✗ retiming KO: {proc.stderr.strip()[:400]}") + else: + result_date = args.day + print(f" ✓ horodatages restaurés — la plateforme affiche la journée du {args.day}") + + # pre-#633 the admin cannot list other orgs' alerts: use the org agent account + viewer = api_login(args.api, args.view_login, args.view_pwd) + r = requests.get( + f"{args.api}/api/v1/alerts/all/fromdate", params={"from_date": result_date, "limit": 100}, + headers=auth(viewer), timeout=30, + ) + alerts = r.json() if r.ok else [] + print(f"\n📊 RÉSULTAT: {len(alerts)} alertes créées par le replay") + for a in sorted(alerts, key=lambda a: a["started_at"]): + seqs = a.get("sequences") or [] + loc = f"({a['lat']:.5f},{a['lon']:.5f})" if a.get("lat") is not None else "sans loc" + cams = sorted({s.get("camera_id") for s in seqs}) + print(f" alerte {a['id']}: {a['started_at'][11:19]}Z->{a['last_seen_at'][11:19]}Z | {loc}" + f" | {len(seqs)} seq | cams locales {cams}") + print("\n👀 Frontend: http://localhost:8080/ · API: http://localhost:5050/docs") + + +if __name__ == "__main__": + main() From e4d95ec6046eff27e97e64a4d578068f57b089cf Mon Sep 17 00:00:00 2001 From: Mateo Date: Thu, 24 Sep 2026 12:58:15 +0200 Subject: [PATCH 3/3] docs: add day replay guide --- REPLAY.md | 149 ++++++++++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 149 insertions(+) create mode 100644 REPLAY.md diff --git a/REPLAY.md b/REPLAY.md new file mode 100644 index 0000000..3b5e0eb --- /dev/null +++ b/REPLAY.md @@ -0,0 +1,149 @@ +# Rejouer une journée réelle sur la plateforme locale + +Rejoue une journée de prod SDIS-77 (dump `~/pyronear/test/analyse77//`, jours +disponibles : 2026-07-10 à 2026-07-13) contre la stack locale, avec une **horloge +simulée pilotable** : la plateforme affiche la journée aux vraies heures (datée +d'aujourd'hui), et on avance/pause/ralentit à la demande. + +## 1. Prérequis (une fois) + +**Image API avec l'horloge simulée** — deux branches dans le worktree +`~/pyronear/api/pyro-api-replay-20260710`, selon la journée à rejouer (la prod a +déployé #629 entre le 10 au soir et le 12 au matin — vérifié sur les données : +la tolérance bbox de #629 est visible dans les séquences du 12) : + +| Journée rejouée | Branche à builder | +|---|---| +| 2026-07-10 (et 11 matin) | `replay-clock-20260710` (avant #629 — les doublons d'alertes sont fidèles) | +| 2026-07-12, 2026-07-13 | `replay-clock-20260712` (#629+#630+#631 + patch horloge) | + +```bash +cd ~/pyronear/api/pyro-api-replay-20260710 +git checkout replay-clock-20260712 # ou replay-clock-20260710 selon la journée +make build # -> image pyronear/alert-api:latest +``` + +(Pour rejouer avec le code actuel : cherry-pick des 2 commits d'horloge sur main, +puis `make build`.) + +`/etc/hosts` doit contenir `127.0.0.1 minio` (pour voir les images dans la plateforme). + +## 2. Lancer la stack + +```bash +cd ~/pyronear/devops/pyro-envdev +docker compose -f docker-compose.yml -f docker-compose.replay.yml up -d db minio pyro_api init_script frontend +``` + +Attendre que le seed se termine (`docker logs -f init` → "completed successfully"). + +- Plateforme : **http://localhost:8080** — compte **`test77` / `test`** +- API : http://localhost:5050/docs — superadmin `mateo`/`mateo` (cf. `.env`) + +## 3. Lancer le replay interactif + +```bash +python3 scripts/replay_day.py --day 2026-07-10 --ctl +``` + +Le script mappe les caméras/poses de prod sur l'env local, redémarre l'API sur +l'horloge pilotable, et démarre **en pause** juste avant la première détection. +Il tourne en avant-plan jusqu'à la fin de la journée (le laisser ouvert dans son +terminal, piloter depuis un autre). + +## 4. Piloter (`scripts/replay_ctl.py`) + +Effet immédiat, aucun redémarrage : + +```bash +python3 scripts/replay_ctl.py status # heure simulée, vitesse, prochaines étapes +python3 scripts/replay_ctl.py next # ➜ étape suivante : avance rapide, atterrit + # 3 min APRÈS le début de la séquence + # (images déjà visibles), puis PAUSE +python3 scripts/replay_ctl.py next 10 # idem mais reprend la lecture à 10x +python3 scripts/replay_ctl.py goto 12h00 # avance jusqu'à 12h00 HEURE DE PARIS, puis PAUSE +python3 scripts/replay_ctl.py goto 16:50 10 # idem mais reprend à 10x en arrivant +python3 scripts/replay_ctl.py play 60 # lecture 60x (journée en ~11 min) +python3 scripts/replay_ctl.py slow 5 # ralenti 5x +python3 scripts/replay_ctl.py pause # gel de la simulation +``` + +Repères pour le 10/07 (heures affichées, Paris) : **12h00-12h15 feu de Barbizon** +(croix-augas-01 puis triangulation nemours-02, enchevêtrement de doublons pré-#629), +16h00-17h30 série d'alertes de l'après-midi, 18h50-19h40 cluster du soir. + +On ne peut pas revenir en arrière (les données sont en base) — pour revoir un +moment, reset (§6) et `goto`. + +## 5. Rejouer une autre journée (11, 12, 13/07…) + +Toujours **reset d'abord** (l'ancien jour pollue les fenêtres de regroupement et la +vue live). Si on change d'époque (10-11/07 ↔ 12-13/07), rebuilder l'image sur la +bonne branche (§1) **avant** le `up`. Puis relancer le driver avec la date voulue : + +```bash +cd ~/pyronear/devops/pyro-envdev +# arrêter le driver en cours (Ctrl-C dans son terminal, ou pkill) +pkill -f replay_day.py +rm -f replay_control/clock.json replay_control/steps.json +docker compose -f docker-compose.yml -f docker-compose.replay.yml down -v +docker compose -f docker-compose.yml -f docker-compose.replay.yml up -d db minio pyro_api init_script frontend +# attendre le seed: docker logs -f init -> "completed successfully" + +python3 scripts/replay_day.py --day 2026-07-12 --ctl +``` + +Journées disponibles dans le dump : `2026-07-10`, `2026-07-11`, `2026-07-12`, +`2026-07-13`. Le pilotage (`replay_ctl.py`) est identique quel que soit le jour. + +Repères (heures affichées, Paris) : +- **12/07** : 14h41 détection du feu d'Achères/Noisy sur croix-augas-02 + (localisation fausse à 5,5 km via le panache, corrigée à 17h35 — alertes + multiples), après-midi très chargée 14h-17h30. +- **13/07** : 09h24 fumée sur l'axe Faisanderie (croix-augas-02), 14h32 la + séquence « détection SDIS » (azimut 222°), 19h04 confirmation moret — journée + la plus chaotique (le feu de 800 ha brûle toute la journée). + +## 6. Tout remettre à zéro + +```bash +cd ~/pyronear/devops/pyro-envdev +pkill -f replay_day.py +rm -f replay_control/clock.json replay_control/steps.json +docker compose -f docker-compose.yml -f docker-compose.replay.yml down -v +docker compose -f docker-compose.yml -f docker-compose.replay.yml up -d db minio pyro_api init_script frontend +``` + +À faire **entre deux journées rejouées** (sinon l'ancien jour pollue les fenêtres +de regroupement et la vue live). + +## Comment ça marche + +- L'API (patchée) lit son horloge dans `replay_control/clock.json` (monté dans le + conteneur) : segments `(real0, sim0, speed)` relus à chaud — `replay_ctl.py` ne + fait qu'ajouter des segments. Tout le pipeline (fenêtres de séquences, + triangulation, vue live « 24 h ») vit dans ce temps simulé ; les JWT restent sur + l'horloge réelle. +- `replay_day.py --ctl` poste les détections dès que l'horloge les rend « dues », + avec **une file séquentielle par caméra** (comme en prod — la concurrence + intra-caméra crée des séquences dupliquées) ; l'avance rapide est plafonnée à + 60x, sinon la file prend du retard sur l'horloge et une séquence active peut se + scinder à la frontière du saut. +- Fidélité mesurée vs prod (10/07 à 60x) : mêmes regroupements sur les événements + multi-caméras, localisations identiques à ~5 m, horodatages à ±3 s. + +## Modes secondaires + +- `--clock` (sans `--ctl`) : horloge fixe à `--speed`, sans pilotage (le mode du + premier test de bout en bout). +- Sans `--clock` : posts à l'heure réelle + réécriture SQL des horodatages + historiques à la fin (`--no-restore-times` pour désactiver) ; nécessite les + fenêtres réduites — obsolète, préférer `--ctl`. +- `--dry-run` : vérifie le mapping caméras/poses sans rien poster. + +## Limites + +- Seules les séquences validées en prod sont dans le dump (le bruit non validé + n'est pas rejoué) ; les 2 séquences > 100 détections sont tronquées à 100. +- 5 détections à bbox dégénérée `(0,0,0,0,0)` sont filtrées (l'API les rejette). +- Le gate temporel est désactivé (fail-open, fidèle au comportement du 10/07).