From 42ec68ed51a91b4e9ecfad9a43abf1246f9352ba Mon Sep 17 00:00:00 2001 From: saphid <4596216+saphid@users.noreply.github.com> Date: Tue, 29 Sep 2026 12:44:51 +1000 Subject: [PATCH] fix(tracking): reduce eye images while capturing, with a backlog cap The capture tool's stdout is block-buffered, so waiting for its directory name left every image on disk until the capture ended (3,554 PNGs in a 20 s run on the Frame). Detect the new capture directory instead, decode in three worker processes, and stop the capture if more than 900 images wait. On the Frame a 30 s run now peaks at 7 images on disk. Part of #27 Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/tracking.md | 24 ++++++++++++---- frame/tracking/pulse.py | 64 ++++++++++++++++++++++++++++++++--------- tests/test_pulse.py | 30 ++++++++++++++++--- 3 files changed, 95 insertions(+), 23 deletions(-) diff --git a/docs/tracking.md b/docs/tracking.md index 1bb7dfe..7f71070 100644 --- a/docs/tracking.md +++ b/docs/tracking.md @@ -193,13 +193,16 @@ How it works: - **Capture (verified).** SteamVR ships `eyetracking --calib N`, which saves both eye cameras for N seconds as 400×400 8-bit IR PNGs with a monotonic - timestamp per frame. A 2-second test captured 177 stereo pairs, about 90 fps. - SteamVR's live eye tracker, part of `steamvr.service`, gets its frames from - the DSP and stops its cameras when the headset is off; the capture ran - alongside it without errors in its log. **Untested:** whether the capture - and the live tracker coexist while the headset is worn and tracking. + timestamp per frame, at about 90 fps per eye. SteamVR's live eye tracker, + part of `steamvr.service`, gets its frames from the DSP and stops its + cameras when the headset is off. Unworn captures ran alongside it: its PID + and log were unchanged and our OpenXR gaze session still started + afterwards. **Untested:** whether the capture and the live tracker coexist + while the headset is worn and tracking. - **Privacy.** Each image is reduced to a 16×16 grid of patch averages as - soon as it is complete, then deleted. The capture directory is removed on + soon as it is complete, then deleted. Three worker processes do this beside + the capture. If more than 900 images (about five seconds) ever wait, the + capture stops rather than letting eye images accumulate. The capture directory is removed on exit, even after errors. No image is kept or leaves the Frame. The estimate is printed only with `--show`, and sent or saved only with `--osc` or `--log`, as for the strap. @@ -213,6 +216,15 @@ How it works: stands out from the noise. Otherwise the command exits 3 and sends nothing. The thresholds are provisional until checked on real wearers. +**Verified on the Frame, unworn, 2026-09-29:** a 30-second run captured +5,362 eye frames, never had more than 7 images on disk, finished 3 s after +the capture ended and left no capture directory. It reported no clear pulse +(exit 3), as it should with nobody wearing it. Worth knowing: the unworn +patches agreed on a steady rhythm near 129 BPM (2.15 Hz) with low +signal/noise (0.19). That is a camera or illumination artifact, not a pulse, +and the signal/noise gate kept it from being reported. A worn test should +also record an unworn baseline, to rule out the same artifact. + **Verified on synthetic data** (unit tests): a 0.3% brightness pulse in a third of the patches, with noise, drift, blinks and eye movement, is recovered within 1.5 BPM at 58, 72 and 115 BPM; noise and blinks alone are diff --git a/frame/tracking/pulse.py b/frame/tracking/pulse.py index 9af7642..2b582cb 100644 --- a/frame/tracking/pulse.py +++ b/frame/tracking/pulse.py @@ -56,22 +56,49 @@ def load_grid(path): image.get_rowstride(), image.get_n_channels()) +def reduce_file(loader, path): + """Reduce one image to its patch grid and delete it, whatever happens.""" + try: + return loader(path) + finally: + os.unlink(path) + + class Capture: """Run SteamVR's eye-camera capture and reduce frames as they arrive.""" NAME = re.compile(r"^(left|right)_(\d+)\.png$") # Only SteamVR's own capture directories are read and removed. - WRITING = re.compile(r"Writing capture to: (/tmp/etcalib_[\w-]+)") PREFIX = "/tmp/etcalib_" + # About five seconds of frames (~65 MB in RAM-backed /tmp). If reduction + # falls further behind than this, the capture stops rather than letting + # eye images pile up. + MAX_BACKLOG = 900 - def __init__(self, seconds, runner=subprocess.Popen, loader=load_grid): - self.seconds, self.runner, self.loader = seconds, runner, loader + def __init__(self, seconds, runner=subprocess.Popen, loader=load_grid, workers=3): + self.seconds, self.runner, self.loader, self.workers = seconds, runner, loader, workers self.grids = {"left": {}, "right": {}} self.directory = None + self.pool = None + self.submitted = set() + self.futures = {} + + def candidates(self): + parent, stem = os.path.split(self.PREFIX) + return {Path(entry.path) for entry in os.scandir(parent) + if entry.name.startswith(stem) and entry.is_dir(follow_symlinks=False)} def run(self): - # The capture tool's output goes to a private temporary file, read back - # to find the capture directory and failure messages. + # The capture tool's output goes to a private temporary file, read for + # failure messages. It is block-buffered, so the capture directory is + # found by watching for a new one rather than waiting for its name. import tempfile + before = self.candidates() + if self.workers: + # SteamVR pins its eye tracker to cores 0-1; decoding runs beside it. + import concurrent.futures + import multiprocessing + self.pool = concurrent.futures.ProcessPoolExecutor( + self.workers, mp_context=multiprocessing.get_context("fork")) with tempfile.TemporaryFile("w+") as log: process = self.runner([str(ET_BIN), "-b", "CDSP", "-w", str(ET_WEIGHTS), "--calib", str(self.seconds)], cwd=str(ET_BIN.parent), stdout=log, stderr=subprocess.STDOUT) @@ -79,10 +106,10 @@ class Capture: deadline = time.monotonic() + self.seconds + 30 while True: if not self.directory: - log.seek(0) - match = self.WRITING.search(log.read()) - if match: - self.directory = Path(match.group(1)) + new = self.candidates() - before + if len(new) > 1: + raise RuntimeError("another eye-camera capture is running") + self.directory = new.pop() if new else None finished = process.poll() is not None self.reduce(final=finished) if finished: @@ -104,6 +131,8 @@ class Capture: except subprocess.TimeoutExpired: process.kill() process.wait() + if self.pool: + self.pool.shutdown(wait=True, cancel_futures=True) self.remove() def reduce(self, final=False): @@ -116,14 +145,23 @@ class Capture: match = self.NAME.match(entry.name) if match: pending[match.group(1)].append((int(match.group(2)), entry.path)) + if sum(len(files) for files in pending.values()) > self.MAX_BACKLOG: + raise RuntimeError("eye-image processing fell behind; capture stopped") for eye, files in pending.items(): files.sort() ready = files if final else files[:-1] for index, path in ready: - try: - self.grids[eye][index] = self.loader(path) - finally: - os.unlink(path) + if path in self.submitted: + continue + self.submitted.add(path) + if self.pool: + self.futures[(eye, index)] = self.pool.submit(reduce_file, self.loader, path) + else: + self.grids[eye][index] = reduce_file(self.loader, path) + for key, future in list(self.futures.items()): + if final or future.done(): + self.grids[key[0]][key[1]] = future.result() + del self.futures[key] def frames(self): """[(monotonic seconds, eye, grid)] joined with the capture metadata.""" diff --git a/tests/test_pulse.py b/tests/test_pulse.py index 7b2a5b7..bc2d92f 100644 --- a/tests/test_pulse.py +++ b/tests/test_pulse.py @@ -125,7 +125,7 @@ class FakeCaptureTool: stdout.flush() self.returncode = 1 return self - stdout.write(f"Writing capture to: {self.directory}\nCapturing images...\n") + # Real output is block-buffered until exit, so nothing is printed here. stdout.flush() self.directory.mkdir() return self @@ -156,6 +156,7 @@ class CaptureLifecycle(unittest.TestCase): def setUp(self): self.root = Path(tempfile.mkdtemp()) self.directory = self.root / "etcalib_test" + (self.root / "etcalib_older").mkdir() # an earlier capture is not ours self.loaded = [] def tearDown(self): @@ -170,9 +171,8 @@ class CaptureLifecycle(unittest.TestCase): def capture(self, tool): class TestCapture(pulse.Capture): # A temporary directory stands in for /tmp/etcalib_*. - WRITING = __import__("re").compile(r"Writing capture to: (\S+)") - PREFIX = str(self.root) - capture = TestCapture(20, runner=tool, loader=self.loader) + PREFIX = str(self.root / "etcalib_") + capture = TestCapture(20, runner=tool, loader=self.loader, workers=0) with patch.object(pulse.time, "sleep", lambda s: None): return capture, capture.run() @@ -181,6 +181,7 @@ class CaptureLifecycle(unittest.TestCase): capture, frames = self.capture(tool) self.assertEqual(sorted(self.loaded), sorted(f"{e}_{i}.png" for e in ("left", "right") for i in range(6))) self.assertFalse(self.directory.exists()) # images and metadata removed + self.assertTrue((self.root / "etcalib_older").exists()) self.assertEqual(len(frames), 10) # frame 2 invalid in both eyes self.assertEqual(tool.command[-2:], ["--calib", "20"]) self.assertTrue(all(isinstance(t, float) and eye in ("left", "right") for t, eye, _ in frames)) @@ -191,6 +192,27 @@ class CaptureLifecycle(unittest.TestCase): self.capture(tool) self.assertIn("cameras unavailable", str(raised.exception)) + def test_backlog_stops_capture_and_removes_images(self): + class Stalled(FakeCaptureTool): + def poll(self): + for i in range(20): + for eye in ("left", "right"): + (self.directory / f"{eye}_{self.polls * 20 + i}.png").write_bytes(b"png") + self.polls += 1 + return None + tool = Stalled(self.directory) + class TestCapture(pulse.Capture): + PREFIX = str(self.root / "etcalib_") + MAX_BACKLOG = 50 + capture = TestCapture(20, runner=tool, loader=self.loader, workers=0) + capture.reduce = lambda final=False, original=capture.reduce: ( + original(final) if tool.polls > 3 else None) # reduction stalls + with patch.object(pulse.time, "sleep", lambda s: None), self.assertRaises(RuntimeError) as raised: + capture.run() + self.assertIn("fell behind", str(raised.exception)) + self.assertEqual(tool.returncode, -15) # capture tool stopped + self.assertFalse(self.directory.exists()) # no image left behind + def test_only_removes_capture_directories(self): capture = pulse.Capture(20) capture.directory = self.root