#!/usr/bin/python3
"""Live overlay for the Data Flow page (/status/dataflow.py).

Aggregates a few cheap localhost endpoints into one small payload:
  - service_control_api  : which APIs/services are actually running (screen -list)
  - agent_monitor_api    : per-pod field agents (which pod is scanning right now)
  - agent_monitor_api    : Foreman Pickaxe watchers (which CIDRs they're busy on)
  - gen_ingest / flow    : feed freshness for the generator + gas ingest edges

The result is cached on disk for CACHE_TTL seconds, so ten open browsers cost
the same as one. Nothing here writes anything or touches the control path.
"""

import json
import os
import sys
import time

import requests

CACHE_FILE = '/var/www/html/ngon/status/caches/dataflow_live.json'
CACHE_TTL = 5          # seconds — matches the page's poll interval
SERVICES_TTL = 20      # `screen -list` is the priciest call; refresh it less often


def _scan_age(hhmmss):
    """Seconds since an agent's last completed scan.

    The field server and the cloud box run the same local clock, so the agents'
    HH:MM:SS stamps compare directly. Returns None if unparseable or if the
    delta looks like a midnight wrap.
    """
    try:
        h, m, sec = [int(x) for x in str(hhmmss).split(':')]
    except Exception:
        return None
    now = time.localtime()
    delta = (now.tm_hour * 3600 + now.tm_min * 60 + now.tm_sec) - (h * 3600 + m * 60 + sec)
    if delta < 0 or delta > 86000:
        return None
    return delta


def _get(url, timeout=3):
    try:
        r = requests.get(url, timeout=timeout)
        if r.ok:
            return r.json()
    except Exception:
        pass
    return None


def _read_cache():
    try:
        with open(CACHE_FILE) as f:
            return json.load(f)
    except Exception:
        return None


def _write_cache(payload):
    try:
        tmp = CACHE_FILE + '.tmp'
        with open(tmp, 'w') as f:
            json.dump(payload, f)
        os.replace(tmp, CACHE_FILE)
    except Exception:
        pass


def _services(prev):
    """Service up/down map. Reuses the previous value until SERVICES_TTL lapses."""
    now = time.time()
    if prev and now - prev.get('services_ts', 0) < SERVICES_TTL:
        return prev.get('services', {}), prev.get('services_ts', 0)

    data = _get('http://localhost:5015/api/services/status', timeout=6)
    if not data:
        # Keep the last known map rather than flashing everything red
        return (prev or {}).get('services', {}), (prev or {}).get('services_ts', 0)
    return {s['name']: bool(s.get('running')) for s in data.get('services', [])}, now


def build(prev):
    out = {'ts': time.time()}

    services, services_ts = _services(prev)
    out['services'] = services
    out['services_ts'] = services_ts

    # --- field scanner: one agent per pod/CIDR -------------------------------
    agent = _get('http://localhost:5013/agent_state') or {}
    pods = {}
    for a in agent.get('agents', []):
        status = (a.get('status') or '').lower()
        mode = (a.get('mode') or '').lower()
        pods[a.get('pod')] = {
            'status': status,                 # idle | scanning | verify | sleep-rN | wait-net | offline
            'mode': mode,                     # scan | emergency-sleep
            # normalized for display: 'verify' is still a live conversation with the
            # miners, and the EMS hammer rounds report as sleep-rN
            'active': status in ('scanning', 'verify') or status.startswith('sleep-r'),
            'ems': mode.startswith('emergency'),
            'down': status in ('offline', 'wait-net', 'error'),
            'miners': a.get('miners'),
            'hashing': a.get('hashing'),
            'sleeping': a.get('sleeping'),
            'last_scan': a.get('last_scan'),
            'scan_age': _scan_age(a.get('last_scan')),
        }
    out['pods'] = pods
    out['scanner'] = {
        'online': bool(agent.get('online')),
        'age': agent.get('age'),
        'summary': agent.get('summary', {}),
        'peplink': agent.get('peplink', {}),
    }

    # --- Foreman Pickaxes: the other thing scanning our pods -----------------
    pick = _get('http://localhost:5013/pickaxe_state') or {}
    watchers, busy_pods = [], {}
    for w in pick.get('watchers', []):
        watchers.append({
            'name': w.get('name'),
            'ip': w.get('ip'),
            'online': bool(w.get('online')),
            'status': w.get('status'),
            'busy': len(w.get('busy_cidrs') or []),
        })
        for b in (w.get('busy_cidrs') or []):
            for p in (b.get('pods') or []):
                busy_pods[p] = max(busy_pods.get(p, 0), b.get('conns') or 0)
    out['watchers'] = watchers
    out['foreman_pods'] = busy_pods

    # --- ingest feed freshness ----------------------------------------------
    gen = _get('http://localhost:5018/health') or {}
    ages = [v for v in (gen.get('ages_sec') or {}).values() if isinstance(v, (int, float))]
    out['gen_feed'] = {
        'fresh': sum(1 for a in ages if a < 900),
        'total': len(ages),
        'min_age': round(min(ages), 1) if ages else None,
    }

    # wired PTZ cameras: are they resolving to an IP, and how fresh is the frame?
    ptz = _get('http://localhost:5024/api/ptz/health') or {}
    snaps = [v for v in (ptz.get('snaps') or {}).values() if isinstance(v, (int, float))]
    out['ptz'] = {
        'cameras': ptz.get('cameras'),
        'resolved': ptz.get('resolved'),
        'newest_snap': round(min(snaps), 1) if snaps else None,
    }

    # group-text bridge: is the phone still checking in?
    tb = _get('http://localhost:5023/health') or {}
    dev = (tb.get('devices') or [{}])[0]
    out['textbridge'] = {
        'ok': tb.get('status') == 'healthy',
        'device': dev.get('device_id'),
        'state': dev.get('state'),
        'age': round(dev['age_seconds'], 1) if isinstance(dev.get('age_seconds'), (int, float)) else None,
        'last_send': dev.get('last_send_ts'),
    }

    flow = _get('http://localhost:5017/health') or {}
    fages = flow.get('ages_sec') or {}
    out['gas_feed'] = {
        'sites': {k: round(v, 1) for k, v in fages.items()},
        'min_age': round(min(fages.values()), 1) if fages else None,
    }

    return out


def main():
    print("Content-Type: application/json")
    print("Cache-Control: no-store")
    print("")

    prev = _read_cache()
    if prev and time.time() - prev.get('ts', 0) < CACHE_TTL:
        prev['cached'] = True
        print(json.dumps(prev))
        return

    try:
        payload = build(prev)
    except Exception as e:
        if prev:
            prev['cached'] = True
            prev['error'] = str(e)
            print(json.dumps(prev))
            return
        print(json.dumps({'error': str(e), 'ts': time.time()}))
        return

    _write_cache(payload)
    payload['cached'] = False
    print(json.dumps(payload))


main()
