"""Offline transport to the production TypeScript segment functions, not a replica.""" from __future__ import annotations import atexit import json import os import subprocess from pathlib import Path ROOT = Path(__file__).resolve().parents[2] class ProductionSegmentBridge: def __init__(self): node = os.environ.get("VARGA_REPLAY_NODE", "node") self.process = subprocess.Popen( [node, "--import", "tsx", "scripts/rectification-segment-bridge.ts"], cwd=ROOT / "frontend", stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=None, text=True, encoding="utf-8", bufsize=1, ) atexit.register(self.close) def close(self): if self.process.poll() is None: self.process.terminate() self.process.wait(timeout=10) def call(self, payload): self.process.stdin.write(json.dumps(payload) + "\n") self.process.stdin.flush() line = self.process.stdout.readline() if not line: raise RuntimeError("production TS bridge exited") value = json.loads(line) if "error" in value: raise RuntimeError(value["error"]) return value def install(self, vr): self.vr = vr # The caller has already distributed raw/percent/uniform mass. Passing # one row per minute retains that exact mass (lead is deliberately not # applied a second time here); production reducer still owns all math. def summary(weights, segments, true_time): value = self.call({ "operation": "summary_weights", "weights": weights, "segments": segments, "offsets": {stamp: vr.offset_of(stamp, true_time) for stamp in weights}, "truth_offset": 0, }) value.pop("tied", None) return value def gain(probe, weights, segments, true_time): return self.call({ "operation": "gain_weights", "weights": weights, "probes": [probe], "segments": segments, "offsets": {stamp: vr.offset_of(stamp, true_time) for stamp in weights}, })["gains"][0] original_weights = vr.minute_weights def minute_weights(state, mode, *, valid_only=True): if mode != "raw" or not valid_only: return original_weights(state, mode, valid_only=valid_only) return self.call({ "operation": "weights", "rows": state["public"], "eliminated": sorted(state["eliminated"]), }) def replay_scores(public, probes, true_time, *, answers=None, ask_count=vr.ASK_COUNT): asked = list(probes)[:ask_count] given = list(answers) if answers is not None else [vr.optimal_answer(p, true_time) for p in asked] result = self.call({"operation": "replay_raw", "rows": public, "probes": asked, "answers": given}) return vr.finish_state(public, result["scores"], set(result["eliminated"]), sum(answer is not None for answer in given[:len(asked)])) vr.replay_scores = replay_scores vr.minute_weights = minute_weights vr.segment_shares = summary vr.segment_information_gain = gain