#!/usr/bin/python3
"""ft-gazed: the gaze service. The headset's eye tracking, corrected, for the pointer.

Runs ft-gaze (in the dev container), and for every eye tracker sample (90 Hz):

  1. drops blinks: both eyes' openness under half its running median (each eye its own).
     With the default source, mmap set 1, that's all: set 1 is SteamVR's combined gaze,
     which keeps going when the tracker loses one eye (its two directions stay together).
     Set 2's combined direction is the mean of the eyes' own, so with one eye lost it's
     off by half of whatever that eye reads (9 to 14 degrees apart were seen): with set 2,
     samples with an eye under its floor, or the angle between the eyes jumping more than
     1.5 degrees from its median, are dropped too;
  2. smooths it with a fixation lock (the running mean of the current fixation, 1 degree);
  3. corrects it: the calibration from ft-gazeprobe (calibration.json, reloaded when the
     probe changes it) plus what the pointer's corrections have taught since (LiveCorrection,
     saved in pointer-lessons.json);
  4. sends it to the pointer helper: "gz <yaw> <pitch> <raw yaw> <raw pitch>", head-relative
     degrees (yaw +left, pitch +up). The helper uses it only in gaze mode.

Lessons come back from the helper: when you nudge the gaze-placed pointer with the mouse and
click, it sends "lesson <raw yaw> <raw pitch> <true yaw> <true pitch>": where the raw gaze
was when the mouse took over, and where the pointer was when you clicked (you were looking
there). The gap is the tracker's error there, and it's learned, unless it's more than
LESSON_MAX degrees past the correction (then it wasn't a nudge onto what you looked at).

SteamVR's eye tracking log is followed for the headset going on (its eye model starts over,
and the error moves): lessons from before count less then, so the first few after it
relearn the offset.

Nothing here writes to SteamVR, its eye tracker, or its files: ft-gaze reads the eye
tracker's shared memory read-only.

Control socket: abstract unix datagram "@ft_gazed":
  lesson <rhy> <rhp> <thy> <thp>   from the pointer helper (see above)
  status                           reply: one JSON object
  forget                           drop what the lessons taught (the calibration stays)
  reload                           read calibration.json again

Options: --source action|mmap1|mmap2 (default mmap1; set 2 was a little quieter in the probe,
but loses the pointer whenever the tracker loses an eye), -v (a status line every
5 s on stderr), --to NAME (send the gaze to the abstract socket @NAME instead of the pointer
helper; for testing: a helper without gaze mode forwards what it doesn't know to its driver).
"""

import argparse
import json
import os
import selectors
import signal
import socket
import statistics
import subprocess
import sys
import time
from collections import deque
from pathlib import Path

sys.path.insert(0, str(Path(__file__).resolve().parent))
from gazecal import DEFAULT_MODEL, MODELS, STATE, Correction, Fixation, LiveCorrection, SteamEyeLog  # noqa: E402

REPO = Path(__file__).resolve().parents[1]
HELPER = REPO / "gaze" / "build" / "ft-gaze"
ME = "\0ft_gazed"
POINTER = "\0ft_pointer_helper"
CALIBRATION = STATE / "calibration.json"
LESSONS = STATE / "pointer-lessons.json"
LESSON_LOG = STATE / "pointer-lessons.jsonl"
LESSON_MAX = 8.0  # degrees past the correction
RETRY = 3.0       # seconds before starting ft-gaze again


class PointerLessons(LiveCorrection):
    """LiveCorrection, with the offset held back: one lesson moves the whole correction by a
    third of what it measured, not all of it (two alike, half; three, three fifths). In the
    probe, clicks came thick and fast on a stale calibration, where the error was mostly one
    offset. Here lessons are few and the calibration is often fresh: in the first live test a
    6 degree lesson shifted everything 6 degrees, and the next target, 10 degrees away and
    1.4 off before, was 7.1 off (3.6 with this). Near the lesson, the kernel still takes up
    most of it (5.1 of the 6 degrees)."""

    RIDGE = [2.0] + LiveCorrection.RIDGE[1:]


def log(msg):
    print(f"ft-gazed: {msg}", file=sys.stderr, flush=True)


class Service:
    def __init__(self, source, verbose, to=POINTER):
        self.source, self.verbose, self.to = source, verbose, to
        STATE.mkdir(parents=True, exist_ok=True)
        self.base = Correction()
        self.mode = DEFAULT_MODEL
        self.cal_mtime = None
        self.live = PointerLessons()
        self.dirty = False
        self.load_calibration()
        self.load_lessons()
        self.steam = SteamEyeLog()
        self.steam.poll()
        self.live.wear_time = self.steam.worn()
        self.live.refit(self.base, self.mode)
        self.fix = Fixation(radius=1.0)
        self.opens = (deque(maxlen=90), deque(maxlen=90))  # left, right
        self.vergence = deque(maxlen=90)
        self.counts = {"samples": 0, "sent": 0, "blinks": 0, "one_eye": 0, "dropped": 0, "lessons_taken": 0, "refused": 0}
        self.last_sample = 0.0
        self.last = None

        self.sock = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM | socket.SOCK_CLOEXEC | socket.SOCK_NONBLOCK)
        self.sock.bind(ME)
        self.out = socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM | socket.SOCK_CLOEXEC | socket.SOCK_NONBLOCK)
        self.sel = selectors.DefaultSelector()
        self.sel.register(self.sock, selectors.EVENT_READ, "control")
        self.proc = None
        self.buf = b""
        self.restart_at = 0.0
        self.running = True

    # --- Calibration and lessons ---

    def load_calibration(self):
        try:
            mtime = CALIBRATION.stat().st_mtime
            d = json.loads(CALIBRATION.read_text())
        except (OSError, ValueError):
            return
        self.cal_mtime = mtime
        if self.source in d:
            self.base.from_json(d[self.source])
        mode = d.get("_meta", {}).get("model")
        self.mode = mode if mode in MODELS else DEFAULT_MODEL
        log(f"calibration: {self.mode}, {self.base.samples} samples")

    def load_lessons(self):
        try:
            d = json.loads(LESSONS.read_text())
            if d.get("source") == self.source:
                self.live.samples = d.get("samples", [])[-PointerLessons.KEEP:]
        except (OSError, ValueError):
            pass
        log(f"{len(self.live.samples)} lessons")

    def save_lessons(self):
        tmp = LESSONS.with_suffix(".tmp")
        tmp.write_text(json.dumps({"source": self.source, "samples": self.live.samples}))
        tmp.replace(LESSONS)
        self.dirty = False

    def correction(self, hy, hp):
        by, bp = self.base.get(hy, hp, self.mode)
        ly, lp = self.live.get(hy, hp)
        return by + ly, bp + lp

    def lesson(self, rhy, rhp, thy, thp):
        dy, dp = thy - rhy, thp - rhp  # the whole error there
        cy, cp = self.correction(rhy, rhp)
        left = ((dy - cy) ** 2 + (dp - cp) ** 2) ** 0.5
        rec = {"time": time.time(), "source": self.source, "model": self.mode, "raw": [rhy, rhp],
               "true": [thy, thp], "correction": [cy, cp], "lesson_deg": left, "wear": self.steam.worn()}
        if left > LESSON_MAX:
            rec["refused"] = f"more than {LESSON_MAX} deg past the correction"
            self.counts["refused"] += 1
        else:
            self.live.add({"time": rec["time"], "hy": rhy, "hp": rhp, "dy": dy, "dp": dp, "wy": 1.0, "wp": 1.0,
                           "how": "pointer"}, self.base, self.mode)
            self.counts["lessons_taken"] += 1
            self.dirty = True
        try:
            with open(LESSON_LOG, "a") as f:
                f.write(json.dumps(rec) + "\n")
        except OSError as e:
            log(f"lesson log: {e}")
        return rec

    # --- ft-gaze ---

    def start_helper(self):
        if not HELPER.exists():
            log(f"ft-gaze isn't built: run {REPO}/gaze/build.sh")
            self.restart_at = time.monotonic() + 30
            return
        env = dict(os.environ)
        env["XDG_RUNTIME_DIR"] = f"/run/user/{os.getuid()}"  # podman needs the real one
        subprocess.run([str(REPO / "scripts" / "container-up.sh")], env=env, check=False)
        distrobox = Path.home() / ".local" / "bin" / "distrobox"
        # ft-gaze quits when its stdin closes: the one thing distrobox passes on.
        self.proc = subprocess.Popen([str(distrobox), "enter", "dev", "--", str(HELPER), "--watch-stdin"], env=env,
                                     stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
                                     start_new_session=True)
        os.set_blocking(self.proc.stdout.fileno(), False)
        os.set_blocking(self.proc.stderr.fileno(), False)
        self.sel.register(self.proc.stdout, selectors.EVENT_READ, "stdout")
        self.sel.register(self.proc.stderr, selectors.EVENT_READ, "stderr")
        self.buf = b""
        log("ft-gaze started")

    def stop_helper(self):
        if not self.proc:
            return
        for f in (self.proc.stdout, self.proc.stderr):
            try:
                self.sel.unregister(f)
            except (KeyError, ValueError):
                pass
        if self.proc.stdin and not self.proc.stdin.closed:
            self.proc.stdin.close()
        try:
            self.proc.wait(timeout=2)
        except subprocess.TimeoutExpired:
            try:
                os.killpg(self.proc.pid, signal.SIGTERM)
            except ProcessLookupError:
                pass
        self.proc = None

    def read_stdout(self):
        try:
            data = os.read(self.proc.stdout.fileno(), 65536)
        except BlockingIOError:
            return
        if not data:
            log(f"ft-gaze stopped (exit {self.proc.poll()}); again in {RETRY:.0f} s")
            self.stop_helper()
            self.restart_at = time.monotonic() + RETRY
            return
        self.buf += data
        *lines, self.buf = self.buf.split(b"\n")
        for line in lines:
            try:
                self.on_sample(json.loads(line))
            except (ValueError, KeyError, TypeError) as e:
                log(f"bad sample: {e}")

    def read_stderr(self):
        try:
            data = os.read(self.proc.stderr.fileno(), 65536)
        except BlockingIOError:
            return
        for line in data.decode("utf-8", "replace").splitlines():
            if line.strip():
                log(line)

    def on_sample(self, s):
        src = s["src"].get(self.source) or {}
        if "hy" not in src:
            return
        self.counts["samples"] += 1
        self.last_sample = time.monotonic()
        m1 = s["src"].get("mmap1") or {}
        o = m1.get("open")
        lr = src.get("lr", m1.get("lr"))
        # Blinks and lost eyes, judged against the last second (see steady_samples: relative,
        # because the lids come down looking down, and the vergence depends on distance).
        # An eye's floor comes from its good readings, so a lost eye doesn't drag it to 0.
        low = [False, False]
        if o and len(o) == 2:
            for k in (0, 1):
                hist = self.opens[k]
                good = [v for v in hist if v >= 0.12]
                floor = max(0.12, 0.5 * statistics.median(good)) if len(good) >= 30 else 0.12
                low[k] = o[k] < floor
                hist.append(o[k])
        if all(low):
            self.counts["blinks"] += 1
            return
        if any(low):
            self.counts["one_eye"] += 1
        if self.source == "mmap2":
            jump = (lr is not None and len(self.vergence) >= 30
                    and abs(lr - statistics.median(self.vergence)) > 1.5)
            if lr is not None and not any(low):
                self.vergence.append(lr)
            if any(low) or jump:
                self.counts["dropped"] += 1
                return
        # The fixation lock works in degrees here (1 degree per "pixel").
        fy, fp = self.fix(src["hy"], src["hp"], s["t"], 1.0)
        cy, cp = self.correction(fy, fp)
        self.last = (fy + cy, fp + cp, fy, fp)
        try:
            self.out.sendto(f"gz {fy + cy:.3f} {fp + cp:.3f} {fy:.3f} {fp:.3f}".encode(), self.to)
            self.counts["sent"] += 1
        except OSError:
            pass  # the pointer helper isn't running

    # --- Control ---

    def on_control(self):
        while True:
            try:
                data, addr = self.sock.recvfrom(512)
            except BlockingIOError:
                return
            words = data.decode("utf-8", "replace").split()
            reply = None
            if words[:1] == ["lesson"] and len(words) == 5:
                try:
                    rec = self.lesson(*map(float, words[1:]))
                    reply = "refused" if "refused" in rec else f"ok {rec['lesson_deg']:.2f}"
                    log(f"lesson {rec['lesson_deg']:.2f} deg at {rec['raw'][0]:+.1f},{rec['raw'][1]:+.1f}"
                        + (f": {rec['refused']}" if "refused" in rec else ""))
                except ValueError:
                    reply = "error bad lesson"
            elif words[:1] == ["status"]:
                reply = json.dumps(self.status())
            elif words[:1] == ["forget"]:
                self.live = PointerLessons()
                self.live.wear_time = self.steam.worn()
                self.save_lessons()
                reply = "ok"
            elif words[:1] == ["reload"]:
                self.load_calibration()
                self.live.refit(self.base, self.mode)
                reply = "ok"
            else:
                reply = "error unknown command"
            if reply and addr:
                try:
                    self.sock.sendto(reply.encode(), addr)
                except OSError:
                    pass

    def status(self):
        ly, lp = self.live.offset()
        return {"source": self.source, "model": self.mode, "calibration_samples": self.base.samples,
                "lessons": len(self.live.samples), "lesson_offset": [round(ly, 3), round(lp, 3)],
                "ft_gaze": self.proc is not None, "sample_age_s": round(time.monotonic() - self.last_sample, 2)
                if self.last_sample else None, "headset_on": self.steam.wearing(), "headset_on_since": self.steam.worn(),
                "last": [round(v, 2) for v in self.last] if self.last else None, **self.counts}

    def periodic(self):
        if self.steam.poll() or self.steam.worn() != self.live.wear_time:
            if self.steam.worn() != self.live.wear_time:
                log("headset on again: older lessons count less until new ones come in")
            self.live.wear_time = self.steam.worn()
            self.live.refit(self.base, self.mode)
        try:
            mtime = CALIBRATION.stat().st_mtime
        except OSError:
            mtime = None
        if mtime != self.cal_mtime:
            self.load_calibration()
            self.live.refit(self.base, self.mode)
        if self.dirty:
            self.save_lessons()

    def run(self):
        next_periodic = time.monotonic()
        next_verbose = time.monotonic() + 5
        while self.running:
            now = time.monotonic()
            if not self.proc and now >= self.restart_at:
                self.start_helper()
            for key, _ in self.sel.select(timeout=0.5):
                if key.data == "control":
                    self.on_control()
                elif key.data == "stdout" and self.proc:
                    self.read_stdout()
                elif key.data == "stderr" and self.proc:
                    self.read_stderr()
            if now >= next_periodic:
                self.periodic()
                next_periodic = now + 1.0
            if self.verbose and now >= next_verbose:
                log(json.dumps(self.status()))
                next_verbose = now + 5
        self.stop_helper()
        if self.dirty:
            self.save_lessons()


def main():
    ap = argparse.ArgumentParser(description="The gaze service: corrected eye tracking for the pointer")
    ap.add_argument("--source", choices=["action", "mmap1", "mmap2"], default="mmap1")
    ap.add_argument("-v", "--verbose", action="store_true")
    ap.add_argument("--to", default="ft_pointer_helper", help="abstract socket to send the gaze to")
    args = ap.parse_args()
    try:
        service = Service(args.source, args.verbose, "\0" + args.to)
    except OSError as e:
        log(f"can't bind @ft_gazed (already running?): {e}")
        sys.exit(1)

    def stop(*_):
        service.running = False
    signal.signal(signal.SIGTERM, stop)
    signal.signal(signal.SIGINT, stop)
    service.run()


if __name__ == "__main__":
    main()
