From 35b6cd9e9ebaa877d4387b142cda71816f580173 Mon Sep 17 00:00:00 2001 From: DCCONSTRUCTIONS Date: Wed, 2 Sep 2026 20:37:21 +0300 Subject: [PATCH] feat(perception): gate cross-host graph with acknowledged source clocks --- .../pilot_binary_bridge.py | 8 +- .../pilot_freshness.py | 44 ++-- .../pilot_grpc_graph.py | 86 +++++++- .../pilot_grpc_source.py | 83 ++++++-- .../pilot_lifecycle.py | 2 + .../pilot_source_control.py | 113 ++++++++++ .../run_joint_pilot.py | 18 +- src/k1link/perception/streaming_clock.py | 58 +++++- src/k1link/perception/streaming_continuity.py | 8 +- .../perception/streaming_control_grpc.py | 114 +++++++++- src/k1link/perception/streaming_ingress.py | 33 +-- src/k1link/perception/streaming_lifecycle.py | 17 +- .../perception/streaming_source_clock.py | 94 +++++++++ .../test_perception_network_graph_adapter.py | 194 +++++++++++++++++- tests/test_perception_realtime_scene.py | 42 ++++ tests/test_perception_source_clock.py | 124 +++++++++++ .../test_perception_streaming_control_grpc.py | 50 +++++ 17 files changed, 1026 insertions(+), 62 deletions(-) create mode 100644 experiments/perception/worker/streaming_profile_stage1/pilot_source_control.py create mode 100644 src/k1link/perception/streaming_source_clock.py create mode 100644 tests/test_perception_source_clock.py diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_binary_bridge.py b/experiments/perception/worker/streaming_profile_stage1/pilot_binary_bridge.py index 8e6ec86..9401b2b 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_binary_bridge.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_binary_bridge.py @@ -126,7 +126,11 @@ class BinaryGraphInput: length = len(e.value[0]) rolling_times[offset : offset + length] = e.time_ns offset += length - due = self.wall_zero + stamp - self.source_zero + due = ( + self.source.due(stamp, time.monotonic_ns())[0] + if hasattr(self.source, "observed") + else self.wall_zero + stamp - self.source_zero + ) bundle = { "input_start": self.epoch, "sequence": event.source_sequence, @@ -158,6 +162,8 @@ class BinaryGraphInput: "enqueued_ns": time.monotonic_ns(), "payload_bytes": size, } + if hasattr(self.source, "observed"): + bundle["clock_observer"] = self.source.observed self.bgr_hashes.append(hashlib.sha256(image).hexdigest()) bundle["enqueued_ns"] = time.monotonic_ns() transferred = self.runtime.admit_reserved(self.epoch, bundle, reservation) diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_freshness.py b/experiments/perception/worker/streaming_profile_stage1/pilot_freshness.py index 09c5cfb..0159590 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_freshness.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_freshness.py @@ -89,15 +89,25 @@ def build_freshness(scene, bundle, ddr_layer, epoch_id): return SceneFreshness(epoch_id, CLOCK_DOMAIN, sequence, stamp, tuple(layers.values())) +def source_observation(bundle, now_ns): + if "clock_observer" in bundle: + stamp, uncertainty_ns = bundle["clock_observer"](now_ns) + return stamp, uncertainty_ns / 1e6 + # Legacy same-kernel path, or source/consumer on the originating Mac clock. + return bundle["time_ns"] + now_ns - bundle["due_ns"], 0 + + def assess(freshness, *, bundle, now_ns): - # Both sides of this pilot use the SAME monotonic clock. Map it back to - # original source arrival time; zero mapping uncertainty is local-only and - # is not a claim of hardware camera/LiDAR synchronization or network quality. + observed, uncertainty = source_observation(bundle, now_ns) + return assess_observation(freshness, observed, uncertainty) + + +def assess_observation(freshness, observed, uncertainty): return freshness.assess( epoch_id=freshness.epoch_id, clock_domain_id=CLOCK_DOMAIN, - observed_source_time_ns=bundle["time_ns"] + now_ns - bundle["due_ns"], - clock_uncertainty_ms=0, + observed_source_time_ns=observed, + clock_uncertainty_ms=uncertainty, maximum_clock_uncertainty_ms=5, maximum_layer_age_ms=250, ) @@ -141,19 +151,23 @@ def prepare_publication(scene, bundle, ddr_layer, *, epoch_id, now_ns, mode="who if mode not in ("whole-scene", "per-cell"): raise ValueError("unknown costmap freshness mode") scene["costmap_freshness_mode"] = mode + observed, uncertainty = source_observation(bundle, now_ns) if mode == "per-cell": apply_cell_expiry( scene, source_time_ns=bundle["time_ns"], - observed_source_time_ns=bundle["time_ns"] + now_ns - bundle["due_ns"], + observed_source_time_ns=observed, + uncertainty_ms=uncertainty, ) freshness = build_freshness(scene, bundle, ddr_layer, epoch_id) - checked = assess(freshness, bundle=bundle, now_ns=now_ns) + # One immutable clock observation per boundary, including cell and policy + # reassessment. A concurrent refresh must not reinterpret this same instant. + checked = assess_observation(freshness, observed, uncertainty) if not checked.fresh_complete: suppress_policy(scene) # The envelope hashes the actual guarded policy, not the discarded one. freshness = build_freshness(scene, bundle, ddr_layer, epoch_id) - checked = assess(freshness, bundle=bundle, now_ns=now_ns) + checked = assess_observation(freshness, observed, uncertainty) scene["freshness"] = freshness.to_dict() scene["freshness_at_publication"] = checked.to_dict() scene["stale_at_publication"] = any(item.state == "stale" for item in checked.layers) @@ -202,9 +216,14 @@ def assess_receipt(scene, freshness, *, bundle, now_ns): their original support time. Layer identity/source timestamps do not change; updated costmap/policy hashes belong to the derived view, not the wire bytes. """ - observed = bundle["time_ns"] + now_ns - bundle["due_ns"] + observed, uncertainty = source_observation(bundle, now_ns) previous = scene.get("freshness_at_receipt") or scene["freshness_at_publication"] - if observed < _read_wire_integer(previous["checked_at_source_time_ns"]): + # Independent clocks can have overlapping intervals even though receipt is + # later. Reject a provably reversed interval, not a shifted midpoint alone. + if observed + int(uncertainty * 1e6) < ( + _read_wire_integer(previous["checked_at_source_time_ns"]) + - int(previous["clock_uncertainty_ms"] * 1e6) + ): raise ValueError("consumer clock moved backwards") view = dict(scene) if scene.get("costmap_freshness_mode") == "per-cell": @@ -212,6 +231,7 @@ def assess_receipt(scene, freshness, *, bundle, now_ns): view, source_time_ns=freshness.source_time_ns, observed_source_time_ns=observed, + uncertainty_ms=uncertainty, ) layers = {item.layer: item for item in freshness.layers[:4]} for name in ("costmap", "policy"): @@ -230,7 +250,7 @@ def assess_receipt(scene, freshness, *, bundle, now_ns): oldest_required_input_time_ns=oldest, ) freshness = replace(freshness, layers=tuple(layers.values())) - checked = assess(freshness, bundle=bundle, now_ns=now_ns) + checked = assess_observation(freshness, observed, uncertainty) view = receipt_view(view, checked) # Global suppression also changes the policy payload. A derived view must # remain internally verifiable, including on a second consumer boundary. @@ -238,7 +258,7 @@ def assess_receipt(scene, freshness, *, bundle, now_ns): if policy.payload_sha256 is not None: policy = replace(policy, payload_sha256=payload_digest(view, "policy")) freshness = replace(freshness, layers=(*freshness.layers[:-1], policy)) - checked = assess(freshness, bundle=bundle, now_ns=now_ns) + checked = assess_observation(freshness, observed, uncertainty) view["freshness"] = freshness.to_dict() view["freshness_at_receipt"] = checked.to_dict() return view, checked diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_grpc_graph.py b/experiments/perception/worker/streaming_profile_stage1/pilot_grpc_graph.py index 049c276..9adab5d 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_grpc_graph.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_grpc_graph.py @@ -5,6 +5,7 @@ accepted from the network. Grant rotation stays in this trusted local adapter. """ import asyncio +import secrets import threading import time import traceback @@ -15,7 +16,9 @@ from pilot_binary_bridge import BinaryGraphInput from pilot_network_control import kernel_clock, read_control, write_control from k1link.perception.streaming_continuity import StreamSuspended +from k1link.perception.streaming_control_grpc import ControlTicket, StreamControlEndpoint from k1link.perception.streaming_grpc import GrpcStreamEndpoint +from k1link.perception.streaming_source_clock import SourceClockMonitor class NetworkGraphBridge(BinaryGraphInput): @@ -33,16 +36,40 @@ class NetworkGraphBridge(BinaryGraphInput): address, reset_temporal, start_delay=2, + cross_host=False, ): - self.clock = kernel_clock() - source = SimpleNamespace( - source_zero=source_zero, wall_zero=time.monotonic_ns() + int(start_delay * 1e9) + self.cross_host = cross_host + self.clock_id = "worker-clock-" + secrets.token_hex(12) + self.clock = ( + {"scope": "cross-host-conditional", "clock_id": self.clock_id} + if cross_host + else kernel_clock() + ) + source = ( + SourceClockMonitor(self.clock_id, source_zero) + if cross_host + else SimpleNamespace( + source_zero=source_zero, wall_zero=time.monotonic_ns() + int(start_delay * 1e9) + ) ) super().__init__(runtime, decoder, source, report, reset_temporal=reset_temporal) self.control = Path(control) self.source_status = Path(source_status) self.certificate, self.private_key, self.address = certificate, private_key, address self.endpoint = GrpcStreamEndpoint(runtime, self.consume, self.notice) + self.access = None + self.ticket = ControlTicket(runtime.start, secrets.token_hex(32)) if cross_host else None + self.control_endpoint = ( + StreamControlEndpoint( + self.ticket, + self._pending, + clock_id=self.clock_id, + observe=source.observe, + ) + if cross_host + else None + ) + self.clock_states = [] self.stopping, self.ready = threading.Event(), threading.Event() self.failure = None self.published = [] @@ -64,6 +91,9 @@ class NetworkGraphBridge(BinaryGraphInput): def _grant(self): access = self.endpoint.issue(self.epoch, "recorded-acquisition", 1) + if self.cross_host: + self.access = access + return write_control( self.control / "grant.json", { @@ -76,6 +106,13 @@ class NetworkGraphBridge(BinaryGraphInput): }, ) + def _pending(self): + try: + self.source.observed(time.monotonic_ns()) + except StreamSuspended: + return None + return self.access if self.endpoint.grant is not None else None + def _run(self): try: asyncio.run(self._serve()) @@ -90,21 +127,55 @@ class NetworkGraphBridge(BinaryGraphInput): self.address, certificate=Path(self.certificate).read_bytes(), private_key=Path(self.private_key).read_bytes(), + control_handlers=(self.control_endpoint.handler(),) if self.cross_host else (), ) try: - self._grant() + if self.cross_host: + self.runtime.pause_input(self.epoch, "source-clock") + write_control( + self.control / "bootstrap.json", + { + "activation": self.runtime.start.to_dict(), + "token": self.ticket.token, + "remote_clock_id": self.clock_id, + "source_zero_ns": str(self.source_zero), + }, + ) + else: + self._grant() self.ready.set() while not self.stopping.is_set() and not self.runtime.stop_event.is_set(): + if self.cross_host: + try: + self.source.observed(time.monotonic_ns()) + clock_ready = True + if not self.report["wall_zero_ns"]: + self.report["wall_zero_ns"] = self.source.due( + self.source_zero, time.monotonic_ns() + )[0] + except StreamSuspended: + clock_ready = False + self.runtime.pause_input(self.epoch, "source-clock") + if not self.clock_states or self.clock_states[-1]["ready"] != clock_ready: + self.clock_states.append( + {"at_ns": time.monotonic_ns(), "ready": clock_ready} + ) + if len(self.clock_states) > 128: + raise ValueError("bounded clock transition count exceeded") if self.endpoint.active is None: ended = self.source_status / "source-end.json" - if ended.exists(): + if self.cross_host and self.source.ended: + self.runtime.mailbox.finish() + elif not self.cross_host and ended.exists(): terminal = read_control(ended) if terminal.get("run_id") != self.runtime.start.run_id: raise ValueError("external source completion binding mismatch") if terminal.get("error"): raise ValueError("external source failed; inspect its bounded report") self.runtime.mailbox.finish() - elif self.runtime.continuity.phase == "waiting": + elif self.runtime.continuity.phase == "waiting" and ( + not self.cross_host or clock_ready + ): try: epoch = self.runtime.begin_input(self.runtime.start) except StreamSuspended: @@ -118,6 +189,8 @@ class NetworkGraphBridge(BinaryGraphInput): await asyncio.sleep(0.01) finally: await server.stop(0) + if self.cross_host: + (self.control / "bootstrap.json").unlink(missing_ok=True) deadline = time.monotonic() + 2 while self.endpoint.active is not None and time.monotonic() < deadline: await asyncio.sleep(0.01) @@ -157,6 +230,7 @@ class NetworkGraphBridge(BinaryGraphInput): accepted_camera_sequences=self.accepted, published=self.published, adapter_failure=self.failure, + clock_states=self.clock_states, ) if not self.thread.is_alive() and self.endpoint.active is None: self.window.close() diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_grpc_source.py b/experiments/perception/worker/streaming_profile_stage1/pilot_grpc_source.py index 80af724..60c9b96 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_grpc_source.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_grpc_source.py @@ -18,8 +18,9 @@ from pathlib import Path import grpc from pilot_binary_source import event_bytes, input_gaps, read_member from pilot_freshness import assess_receipt, validate_receipt -from pilot_network_control import read_grant, write_control +from pilot_network_control import read_control, read_grant, write_control from pilot_source import SensorArchive, camera_events, merged_events +from pilot_source_control import SourceControl from k1link.compute.live_perception import LiveIngressEvent from k1link.perception.streaming_grpc import GrpcStreamClient @@ -32,6 +33,16 @@ def connection_delay(wall_zero_ns, now_ns): return max(0, (wall_zero_ns - now_ns - 100_000_000) / 1e9) +async def finish_source(client, reader, clock_control): + # EOF on Exchange can finish the graph and retire its control endpoint. + # Confirm source completion first while that data stream is still alive. + if clock_control: + await clock_control.finish() + if client is not None: + await client.end() + await asyncio.wait_for(reader, timeout=5) + + async def run(args): report = { "started_utc": datetime.now(UTC).isoformat(), @@ -43,7 +54,9 @@ async def run(args): "models": 0, "source_clock_speed": 1.0, "full_source_prepass": False, - "clock_scope": "two-containers-same-unshifted-Linux-kernel", + "clock_scope": "cross-host-conditional" + if args.cross_host + else "two-containers-same-unshifted-Linux-kernel", "error": None, "actuation_allowed": False, "commands_enabled": False, @@ -59,6 +72,7 @@ async def run(args): gaps, next_gap, outage_until = input_gaps(args.input_gap), 0, 0 previous_epoch = None clock = None + clock_control = None async def disconnect(): nonlocal client, reader @@ -71,11 +85,25 @@ async def run(args): try: deadline = time.monotonic() + 120 - while not grant_path.exists(): + startup_path = control / "bootstrap.json" if args.cross_host else grant_path + while not startup_path.exists(): if time.monotonic() > deadline: raise TimeoutError("trusted graph grant was not issued after warmup") await asyncio.sleep(0.02) - clock, _ = read_grant(grant_path) + if args.cross_host: + bootstrap = read_control(startup_path) + clock_control = SourceControl( + args.target, Path(args.certificate).read_bytes(), bootstrap + ) + clock_control.start() + anchor = await clock_control.wait_anchor() + clock = { + "source_zero_ns": anchor.source_zero_ns, + "wall_zero_ns": anchor.local_zero_ns, + "clock": {"scope": "source-and-consumer-same-Mac-clock"}, + } + else: + clock, _ = read_grant(grant_path) first = next(camera_events(args.camera_index, 1)).time_ns if clock["source_zero_ns"] != first - 500_000_000: raise ValueError("recording prefix does not match admitted source clock") @@ -104,6 +132,20 @@ async def run(args): "due_ns": clock["wall_zero_ns"] + source_stamp - clock["source_zero_ns"], } freshness = validate_receipt(scene, bundle, epoch_id=epoch.epoch_id) + if clock_control and not clock_control.ready(): + rejected = report.setdefault("clock_rejected_results", []) + if len(rejected) >= args.frames: + raise ValueError("bounded rejected receipt ledger exceeded") + rejected.append( + { + "sequence": sequence, + "epoch_id": epoch.epoch_id, + "payload_sha256": hashlib.sha256(payload).hexdigest(), + "reason": "source-clock-unavailable", + } + ) + # No fresh scene/policy authority while clock admission is lost. + continue checked_at = time.monotonic_ns() view, checked = assess_receipt( scene, freshness, bundle=bundle, now_ns=checked_at @@ -140,7 +182,18 @@ async def run(args): async def connect(): nonlocal client, reader, access, previous_epoch - value, candidate = read_grant(grant_path) + if clock_control: + if not clock_control.ready() or clock_control.access is None: + return False + candidate = clock_control.access + value = { + **clock, + "cutoff_ns": clock["source_zero_ns"] + + time.monotonic_ns() + - clock["wall_zero_ns"], + } + else: + value, candidate = read_grant(grant_path) if candidate.epoch.epoch_id == previous_epoch: return False if any( @@ -185,6 +238,8 @@ async def run(args): if event.channel == "camera": released.append(event.sequence) try: + if clock_control and not clock_control.ready(): + await disconnect() if reader is not None and reader.done(): reader.result() # Invalid payload is fatal; transport interruption is not. raise ConnectionError("result stream ended before source EOF") @@ -240,20 +295,23 @@ async def run(args): break source_end = True report["end_sent_monotonic_ns"] = time.monotonic_ns() - if client is not None: - await client.end() - await asyncio.wait_for(reader, timeout=5) - write_control( - status / "source-end.json", {"run_id": access.epoch.run_id, "error": False} - ) + await finish_source(client, reader, clock_control) + if not clock_control: + write_control( + status / "source-end.json", {"run_id": access.epoch.run_id, "error": False} + ) except Exception: report["error"] = traceback.format_exc() - if access is not None: + if access is not None and not args.cross_host: write_control( status / "source-end.json", {"run_id": access.epoch.run_id, "error": True} ) finally: await disconnect() + if clock_control: + await clock_control.close() + report["clock_samples"] = clock_control.samples + report["clock_errors"] = clock_control.errors from run_joint_pilot import distribution report.update( @@ -291,6 +349,7 @@ async def run(args): if __name__ == "__main__": parser = argparse.ArgumentParser() parser.add_argument("--target", required=True) + parser.add_argument("--cross-host", action="store_true") parser.add_argument("--control", required=True) parser.add_argument("--source-status", required=True) parser.add_argument("--certificate", required=True) diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_lifecycle.py b/experiments/perception/worker/streaming_profile_stage1/pilot_lifecycle.py index e215023..600d475 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_lifecycle.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_lifecycle.py @@ -166,6 +166,8 @@ class PilotController: def source_now(self): if self.source_clock is None: raise WorkerLeaseError("source clock mapping not configured") + if hasattr(self.source_clock, "source_now"): + return self.source_clock.source_now(time.monotonic_ns()) return self.source_clock.source_zero + time.monotonic_ns() - self.source_clock.wall_zero def _renew(self): diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_source_control.py b/experiments/perception/worker/streaming_profile_stage1/pilot_source_control.py new file mode 100644 index 0000000..b989fbc --- /dev/null +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_source_control.py @@ -0,0 +1,113 @@ +"""Bounded application-side clock loop; no lease or model control.""" + +import asyncio +import secrets +import time + +import grpc + +from k1link.perception.realtime_contract import StreamStart +from k1link.perception.streaming_clock import ClockMappingError, ClockWindow +from k1link.perception.streaming_control_grpc import ControlTicket, StreamControlClient +from k1link.perception.streaming_source_clock import SourceAnchor + + +class SourceControl: + def __init__(self, target, roots, bootstrap): + ticket = ControlTicket(StreamStart.from_dict(bootstrap["activation"]), bootstrap["token"]) + clock_id = "source-clock-" + secrets.token_hex(12) + self.client = StreamControlClient(target, roots, ticket, clock_id=clock_id) + self.window = ClockWindow( + clock_id, bootstrap["remote_clock_id"], rate_ppm=500, timestamp_error_ns=50_000 + ) + self.source_zero_ns = int(bootstrap["source_zero_ns"]) + self.anchor = None # Latch once after the initial clock handshake, never on resume. + self.access = None + self.confirmed = False + self.failure = None + self.stopping = self.ending = False + self.samples, self.errors = [], [] + self.task = None + + async def update(self): + probe, access = await self.client.poll() + bounds = self.window.add(probe) + now = time.monotonic_ns() + if self.anchor is None: + try: + self.window.current(now).require(now) + except ClockMappingError: + pass # Pre-start clock warmup, not a running source pause or retimestamp. + else: + self.anchor = SourceAnchor(self.source_zero_ns, now + 2_000_000_000) + if self.anchor is not None: + await self.client.acknowledge(probe, self.anchor, ended=self.ending) + self.confirmed = True + self.access = access + now = time.monotonic_ns() + self.samples.append( + { + "at_ns": now, + "bounds": bounds.to_dict(), + "uncertainty_ns": bounds.uncertainty_ns(now), + "ready": self.ready(), + "acknowledged": self.anchor is not None, + } + ) + if len(self.samples) > 2048: + raise ValueError("bounded clock diagnostic count exceeded") + + def ready(self): + if self.failure: + raise RuntimeError("source clock task failed") from self.failure + if not self.confirmed: + return False + try: + now = time.monotonic_ns() + self.window.current(now).require(now) + except ClockMappingError: + return False + return True + + async def loop(self): + while not self.stopping: + try: + await self.update() + if self.ending: + return + except (grpc.RpcError, TimeoutError, OSError) as exc: + self.confirmed = False + self.errors.append({"at_ns": time.monotonic_ns(), "type": type(exc).__name__}) + if len(self.errors) > 128: + self.failure = RuntimeError("bounded clock transport error count exceeded") + return + except Exception as exc: + self.failure = exc + return + await asyncio.sleep(0.05) + + def start(self): + self.task = asyncio.create_task(self.loop()) + + async def wait_anchor(self): + deadline = time.monotonic() + 10 + while self.anchor is None or not self.confirmed: + if self.failure: + raise RuntimeError("initial source clock handshake failed") from self.failure + if time.monotonic() > deadline: + raise TimeoutError("bounded initial source clock handshake") + await asyncio.sleep(0.01) + return self.anchor + + async def finish(self): + self.ending = True + await asyncio.wait_for(asyncio.shield(self.task), timeout=2) + if self.failure: + raise RuntimeError("source clock final acknowledgement failed") from self.failure + + async def close(self): + self.stopping = True + if self.task: + self.task.cancel() + await asyncio.gather(self.task, return_exceptions=True) + await self.client.close() diff --git a/experiments/perception/worker/streaming_profile_stage1/run_joint_pilot.py b/experiments/perception/worker/streaming_profile_stage1/run_joint_pilot.py index 1c4f562..572b60b 100644 --- a/experiments/perception/worker/streaming_profile_stage1/run_joint_pilot.py +++ b/experiments/perception/worker/streaming_profile_stage1/run_joint_pilot.py @@ -276,6 +276,17 @@ def run(args): "costmap_freshness_mode": args.costmap_freshness, "ddrnet_min_source_interval_ms": args.ddrnet_min_source_interval_ms, "input_transport": args.input_transport, + "network_clock_mode": "acknowledged-monotonic/v1" + if getattr(args, "network_cross_host", False) + else "local-only", + "network_clock_envelope": { + "relative_rate_ppm": 500, + "timestamp_error_ns": 50000, + "maximum_age_ns": 2000000000, + "maximum_uncertainty_ns": 5000000, + } + if getattr(args, "network_cross_host", False) + else None, "recover_input": getattr(args, "recover_input", False), "input_gap_plan": getattr(args, "input_gap", []), } @@ -549,6 +560,7 @@ def run(args): private_key=args.network_private_key, address=args.network_address, reset_temporal=reset_temporal, + cross_host=args.network_cross_host, ) controller.source_clock = binary_bridge.source report["scope"] = "full-graph-network-candidate; receiver evidence is separate" @@ -710,7 +722,10 @@ def run(args): "scene_sha256": hashlib.sha256(encoded).hexdigest(), } results.append(result) - except StreamSuspended: + except StreamSuspended as exc: + report.setdefault("suspended_results", []).append( + {"sequence": bundle["sequence"], "reason": str(exc)} + ) mailbox.release(bundle, discard_reason="input-gap") cpu_bundle = bundle = computed = mask = proposals = scene = received = view = ( None @@ -947,6 +962,7 @@ if __name__ == "__main__": parser.add_argument("--network-certificate") parser.add_argument("--network-private-key") parser.add_argument("--network-address", default="[::]:50061") + parser.add_argument("--network-cross-host", action="store_true") parser.add_argument( "--ddrnet-layout", choices=("reference", "channels-last"), default="reference" ) diff --git a/src/k1link/perception/streaming_clock.py b/src/k1link/perception/streaming_clock.py index 8fa3ee3..4c937da 100644 --- a/src/k1link/perception/streaming_clock.py +++ b/src/k1link/perception/streaming_clock.py @@ -69,6 +69,60 @@ class ClockProbe: ) +@dataclass(frozen=True) +class ClockReceipt: + """The responder's view AFTER an authenticated echo of its issued probe. + + Responder receive=t2 and send=t3 bracket only processing, NOT network RTT. + t1-t2 is a lower bound and t4-t3 an upper bound for initiator-minus-responder. + Never construct an ordinary reversed ClockProbe: that would invert causality. + """ + + probe: ClockProbe + acknowledged_ns: int + + def __post_init__(self) -> None: + _stamp(self.acknowledged_ns) + if ( + not self.probe.remote_send_ns + <= self.acknowledged_ns + <= (self.probe.remote_receive_ns + 500_000_000) + ): + raise ClockMappingError("clock acknowledgement outside bounded round trip") + + @property + def local_clock_id(self) -> str: + return self.probe.remote_clock_id + + @property + def remote_clock_id(self) -> str: + return self.probe.local_clock_id + + @property + def nonce(self) -> str: + return self.probe.nonce + + @property + def local_send_ns(self) -> int: + return self.probe.remote_receive_ns + + @property + def local_receive_ns(self) -> int: + return self.acknowledged_ns + + def bounds_at( + self, local_ns: int, *, rate_ppm: int, timestamp_error_ns: int + ) -> tuple[int, int]: + _stamp(local_ns) + if local_ns < self.acknowledged_ns: + raise ClockMappingError("acknowledgement is from the future") + widen = _drift(local_ns - self.probe.remote_receive_ns, rate_ppm) + 2 * timestamp_error_ns + return ( + self.probe.local_send_ns - self.probe.remote_receive_ns - widen, + self.probe.local_receive_ns - self.probe.remote_send_ns + widen, + ) + + @dataclass(frozen=True) class ClockBounds: local_clock_id: str @@ -161,11 +215,11 @@ class ClockWindow: timestamp_error_ns, maximum_age_ns, ) - self.samples: deque[ClockProbe] = deque(maxlen=16) + self.samples: deque[ClockProbe | ClockReceipt] = deque(maxlen=16) self.failed = False self.last_receive_ns = -1 - def add(self, sample: ClockProbe) -> ClockBounds: + def add(self, sample: ClockProbe | ClockReceipt) -> ClockBounds: if self.failed: raise ClockMappingError("clock session quarantined; explicit new session required") if ( diff --git a/src/k1link/perception/streaming_continuity.py b/src/k1link/perception/streaming_continuity.py index 94eb413..ff0efa0 100644 --- a/src/k1link/perception/streaming_continuity.py +++ b/src/k1link/perception/streaming_continuity.py @@ -62,7 +62,13 @@ class InputContinuity: raise StreamSuspended("input is waiting for resynchronization") def pause(self, reason: str) -> None: - if reason not in ("input-disconnected", "input-timeout", "source-gap", "worker-telemetry"): + if reason not in ( + "input-disconnected", + "input-timeout", + "source-gap", + "worker-telemetry", + "source-clock", + ): raise ValueError("unknown recoverable pause reason") if self.phase != "waiting": self.pauses += 1 diff --git a/src/k1link/perception/streaming_control_grpc.py b/src/k1link/perception/streaming_control_grpc.py index 3f6b55f..ab6904c 100644 --- a/src/k1link/perception/streaming_control_grpc.py +++ b/src/k1link/perception/streaming_control_grpc.py @@ -23,10 +23,12 @@ from . import streaming_wire as wire from .realtime_contract import StreamStart from .streaming_clock import ClockProbe from .streaming_grpc import OPTIONS, StreamAccess +from .streaming_source_clock import SourceAnchor SERVICE = "missioncore.perception.v1.StreamControl" METHOD = f"/{SERVICE}/Poll" MAX_CONTROL = 8192 +REPORT_METHOD = f"/{SERVICE}/ReportClock" def _document(raw: bytes) -> dict[str, Any]: @@ -77,6 +79,7 @@ class StreamControlEndpoint: *, clock_id: str, clock_ns: Callable[[], int] = time.monotonic_ns, + observe: Callable[[ClockProbe, int, SourceAnchor, bool], None] | None = None, ) -> None: wire.identifier(clock_id) self.activation, self.pending, self.clock_id, self.clock_ns = ( @@ -89,10 +92,16 @@ class StreamControlEndpoint: self._lock = threading.Lock() self._last_poll_ns = -1 self.accepted = self.rejected = 0 + self.observe = observe + self._challenge: dict[str, Any] | None = None def handler(self) -> Any: return grpc.method_handlers_generic_handler( - SERVICE, {"Poll": grpc.unary_unary_rpc_method_handler(self.poll)} + SERVICE, + { + "Poll": grpc.unary_unary_rpc_method_handler(self.poll), + "ReportClock": grpc.unary_unary_rpc_method_handler(self.report_clock), + }, ) def _authenticate(self, context: Any) -> None: @@ -163,6 +172,19 @@ class StreamControlEndpoint: encoded = wire.canonical(response) if len(encoded) > MAX_CONTROL: raise ValueError("trusted control snapshot exceeds bound") + if self.observe is not None: + with self._lock: + self._challenge = { + key: response[key] + for key in ( + "nonce", + "local_clock_id", + "local_send_ns", + "remote_clock_id", + "remote_receive_ns", + "remote_send_ns", + ) + } except (ValueError, RuntimeError): await context.abort( grpc.StatusCode.UNAVAILABLE, "local controller has no current offer" @@ -171,6 +193,52 @@ class StreamControlEndpoint: self.accepted += 1 return encoded + async def report_clock(self, raw: bytes, context: Any) -> bytes: + received = self.clock_ns() + try: + self._authenticate(context) + value = _document(raw) + if set(value) != {"probe", "anchor", "ended"} or self.observe is None: + raise ValueError("clock receipt not admitted") + evidence = value["probe"] + with self._lock: + challenge = self._challenge + if ( + challenge is None + or not isinstance(evidence, dict) + or set(evidence) != {*challenge, "local_receive_ns"} + or any(evidence[k] != v for k, v in challenge.items()) + ): + raise ValueError("clock receipt is not the current issued challenge") + self._challenge = None # Single-use, including malformed/late receipts. + probe = ClockProbe( + evidence["local_clock_id"], + evidence["remote_clock_id"], + evidence["nonce"], + *( + wire.uint64(evidence[k]) + for k in ( + "local_send_ns", + "remote_receive_ns", + "remote_send_ns", + "local_receive_ns", + ) + ), + ) + anchor = SourceAnchor.from_dict(value["anchor"]) + self.observe(probe, received, anchor, value["ended"]) + except (ValueError, TypeError): + self.rejected += 1 + await context.abort(grpc.StatusCode.INVALID_ARGUMENT, "clock receipt rejected") + return b"" + return wire.canonical( + { + "schema_version": "missioncore.stream-clock-receipt/v1", + "activation_binding": wire.binding(self.activation), + "nonce": probe.nonce, + } + ) + class StreamControlClient: def __init__( @@ -196,6 +264,7 @@ class StreamControlClient: ), ) self.call = self.channel.unary_unary(METHOD) + self.report_call = self.channel.unary_unary(REPORT_METHOD) self.polling = False async def poll(self) -> tuple[ClockProbe, StreamAccess | None]: @@ -267,3 +336,46 @@ class StreamControlClient: async def close(self) -> None: await self.channel.close() + + async def acknowledge( + self, probe: ClockProbe, anchor: SourceAnchor, *, ended: bool = False + ) -> None: + if self.polling or probe.local_clock_id != self.clock_id or type(ended) is not bool: + raise ValueError("one scoped control request at a time") + self.polling = True + try: + evidence = { + name: getattr(probe, name) + for name in ( + "nonce", + "local_clock_id", + "remote_clock_id", + ) + } + evidence.update( + { + name: str(getattr(probe, name)) + for name in ( + "local_send_ns", + "remote_receive_ns", + "remote_send_ns", + "local_receive_ns", + ) + } + ) + response = _document( + await self.report_call( + wire.canonical({"probe": evidence, "anchor": anchor.to_dict(), "ended": ended}), + metadata=self.ticket.metadata(), + timeout=0.5, + wait_for_ready=False, + ) + ) + if response != { + "schema_version": "missioncore.stream-clock-receipt/v1", + "activation_binding": wire.binding(self.ticket.activation), + "nonce": probe.nonce, + }: + raise ValueError("clock receipt response binding mismatch") + finally: + self.polling = False diff --git a/src/k1link/perception/streaming_ingress.py b/src/k1link/perception/streaming_ingress.py index c95372d..dbe92f9 100644 --- a/src/k1link/perception/streaming_ingress.py +++ b/src/k1link/perception/streaming_ingress.py @@ -291,22 +291,23 @@ class StreamingIngress: except Exception as exc: self.error = str(exc) self.terminal = "failed" - if self.opened: - if self.runtime.continuity is not None and isinstance( - exc, (InputInterrupted, StreamSuspended) - ): - self.terminal = "paused" - # A late old socket cannot pause a replacement connection. - with suppress(RuntimeError): # Already replaced, stopping, or fenced. - self.runtime.pause_input( - self.input_epoch, - "input-disconnected" - if isinstance(exc, InputInterrupted) - else "source-gap", - ) - else: - self.runtime.mailbox.finish(self.error) - self.runtime.request_stop("failed") + if self.runtime.continuity is not None and isinstance( + exc, (InputInterrupted, StreamSuspended) + ): + # Admission already consumed the single-use grant. A transport + # loss before the application OPEN is still WAIT, not an active + # owner stranded forever. A late old socket remains fenced. + self.terminal = "paused" + with suppress(RuntimeError): # Already replaced, stopping, or fenced. + self.runtime.pause_input( + self.input_epoch, + "input-disconnected" + if isinstance(exc, InputInterrupted) + else "source-gap", + ) + elif self.opened: + self.runtime.mailbox.finish(self.error) + self.runtime.request_stop("failed") finally: self.connection.close() if reservation is not None: diff --git a/src/k1link/perception/streaming_lifecycle.py b/src/k1link/perception/streaming_lifecycle.py index 9c1c4ec..06bacc6 100644 --- a/src/k1link/perception/streaming_lifecycle.py +++ b/src/k1link/perception/streaming_lifecycle.py @@ -223,11 +223,17 @@ class StreamingLifecycle: return self._check(self.start) self.continuity.check(epoch, synchronizing=synchronizing) - if event_ns is not None: - assert self._source_clock_ns is not None + assert self._source_clock_ns is not None + try: now = self._source_clock_ns() - if event_ns < self.continuity.cutoff_ns or not 0 <= now - event_ns <= 250_000_000: - raise StreamSuspended("obsolete or future input, no backlog replay") + except StreamSuspended: + self.continuity.pause("source-clock") + self.mailbox.pause() + raise + if event_ns is not None and ( + event_ns < self.continuity.cutoff_ns or not 0 <= now - event_ns <= 250_000_000 + ): + raise StreamSuspended("obsolete or future input, no backlog replay") def pause_input(self, epoch: StreamStart, reason: str) -> None: with self._lock: @@ -253,11 +259,12 @@ class StreamingLifecycle: raise StreamSuspended("old epoch callbacks or connection still active") if self.continuity.phase != "waiting": raise StreamSuspended("pause required before a new input epoch") + source_now = self._source_clock_ns() # Check BEFORE mutating the paused mailbox. self.mailbox.begin_epoch() if self._ingress_thread is not None: self._threads = [t for t in self._threads if t is not self._ingress_thread] self._ingress_thread = None - return self.continuity.begin(self._source_clock_ns()) + return self.continuity.begin(source_now) def resume_input( self, epoch: StreamStart, evidence: ResumeEvidence, reset_temporal: Callable[[], None] diff --git a/src/k1link/perception/streaming_source_clock.py b/src/k1link/perception/streaming_source_clock.py new file mode 100644 index 0000000..1352d50 --- /dev/null +++ b/src/k1link/perception/streaming_source_clock.py @@ -0,0 +1,94 @@ +"""Responder-owned source timeline and conditional clock readiness, no GPU authority.""" + +from __future__ import annotations + +import threading +from dataclasses import dataclass + +from . import streaming_wire as wire +from .streaming_clock import ClockMappingError, ClockProbe, ClockReceipt, ClockWindow, _stamp +from .streaming_continuity import StreamSuspended + + +@dataclass(frozen=True) +class SourceAnchor: + source_zero_ns: int + local_zero_ns: int + + def __post_init__(self) -> None: + _stamp(self.source_zero_ns) + _stamp(self.local_zero_ns) + + def to_dict(self) -> dict[str, str]: + return { + "source_zero_ns": str(self.source_zero_ns), + "local_zero_ns": str(self.local_zero_ns), + } + + @classmethod + def from_dict(cls, value: object) -> SourceAnchor: + if not isinstance(value, dict) or set(value) != {"source_zero_ns", "local_zero_ns"}: + raise ValueError("invalid immutable source anchor") + return cls(wire.uint64(value["source_zero_ns"]), wire.uint64(value["local_zero_ns"])) + + +class SourceClockMonitor: + def __init__(self, worker_clock_id: str, source_zero_ns: int) -> None: + wire.identifier(worker_clock_id) + _stamp(source_zero_ns) + self.worker_clock_id, self.source_zero = worker_clock_id, source_zero_ns + self.wall_zero = 0 # Not a comparable host clock; use due()/observed(). + self.anchor: SourceAnchor | None = None + self.window: ClockWindow | None = None + self.ended = False + self._lock = threading.RLock() + + def observe( + self, probe: ClockProbe, received_ns: int, anchor: SourceAnchor, ended: bool + ) -> None: + """Called only after the transport verifies the one-use issued challenge.""" + receipt = ClockReceipt(probe, received_ns) + if probe.remote_clock_id != self.worker_clock_id or type(ended) is not bool: + raise ValueError("foreign clock/status") + with self._lock: + if self.ended or anchor.source_zero_ns != self.source_zero: + raise ValueError("source activation ended or source anchor changed") + if self.anchor is None: + if not 0 <= anchor.local_zero_ns - probe.local_receive_ns <= 10_000_000_000: + raise ValueError("initial source start must be bounded and not in the past") + self.anchor = anchor + self.window = ClockWindow( + self.worker_clock_id, + probe.local_clock_id, + rate_ppm=500, + timestamp_error_ns=50_000, + ) + elif anchor != self.anchor: + raise ValueError("source timeline is immutable for the activation") + assert self.window is not None + self.window.add(receipt) + self.ended = ended + + def observed(self, now_ns: int) -> tuple[int, int]: + """Source time midpoint + integer half-width, freshly checked at use.""" + with self._lock: + try: + if self.window is None or self.anchor is None: + raise ClockMappingError("source clock not observed") + bounds = self.window.current(now_ns) + bounds.require(now_ns) + lower, upper = bounds.offset_at(now_ns) + stamp = self.source_zero + now_ns + (lower + upper) // 2 - self.anchor.local_zero_ns + _stamp(stamp) + return stamp, (upper - lower + 1) // 2 + except ClockMappingError as exc: + raise StreamSuspended("source clock unavailable") from exc + + def source_now(self, now_ns: int) -> int: + stamp, uncertainty = self.observed(now_ns) + # Admission/cutoff/resume use the upper age bound, never a younger estimate. + return stamp + uncertainty + + def due(self, stamp: int, now_ns: int) -> tuple[int, int]: + observed, uncertainty = self.observed(now_ns) + return now_ns + stamp - observed, uncertainty diff --git a/tests/test_perception_network_graph_adapter.py b/tests/test_perception_network_graph_adapter.py index 7e2b7af..ef55f1a 100644 --- a/tests/test_perception_network_graph_adapter.py +++ b/tests/test_perception_network_graph_adapter.py @@ -10,11 +10,18 @@ import pytest pytest.importorskip("grpc") +import time + from test_perception_streaming_grpc import eventually, identity, tls # noqa: F401,E402 +from k1link.perception.streaming_control_grpc import ( # noqa: E402 + ControlTicket, + StreamControlClient, +) from k1link.perception.streaming_grpc import GrpcStreamClient # noqa: E402 from k1link.perception.streaming_lifecycle import StreamingLifecycle # noqa: E402 from k1link.perception.streaming_queue import StreamMailbox # noqa: E402 +from k1link.perception.streaming_source_clock import SourceAnchor # noqa: E402 @pytest.fixture @@ -33,7 +40,8 @@ def adapter(monkeypatch): return module, control -def test_adapter_rotates_grant_only_after_old_stream_drains(tmp_path, tls, adapter): # noqa: F811 +@pytest.mark.parametrize("before_open", [False, True]) +def test_adapter_rotates_grant_only_after_old_stream_drains(tmp_path, tls, adapter, before_open): # noqa: F811 module, control = adapter cert, key = tmp_path / "cert", tmp_path / "key" cert.write_bytes(tls[0]) @@ -78,11 +86,21 @@ def test_adapter_rotates_grant_only_after_old_stream_drains(tmp_path, tls, adapt _, first = control.read_grant(tmp_path / "grant.json") client = GrpcStreamClient(f"localhost:{bridge.port}", tls[0], first) try: - await client.open() + if before_open: + await client.call.initial_metadata() # Grant consumed, no application OPEN sent. + else: + await client.open() await client.close() # True network disconnect, not a synthetic runtime pause. - await eventually( - lambda: control.read_grant(tmp_path / "grant.json")[1].epoch != first.epoch - ) + try: + await eventually( + lambda: control.read_grant(tmp_path / "grant.json")[1].epoch != first.epoch + ) + except TimeoutError: + pytest.fail( + f"epoch did not rotate: phase={runtime.continuity.phase}; " + f"stop={runtime.stop_event.is_set()}; adapter={bridge.failure}; " + f"last={bridge.endpoint.last}; active={bridge.endpoint.active is not None}" + ) _, second = control.read_grant(tmp_path / "grant.json") assert runtime.continuity.phase == "synchronizing" assert second.epoch.lease_generation == first.epoch.lease_generation @@ -114,3 +132,169 @@ def test_delayed_source_start_does_not_open_an_idle_connection(adapter): assert source.connection_delay(3_000_000_000, 1_000_000_000) == 1.9 assert source.connection_delay(3_000_000_000, 2_950_000_000) == 0 assert source.connection_delay(3_000_000_000, 3_010_000_000) == 0 + + +@pytest.mark.parametrize("has_stream", [False, True]) +def test_source_eof_ack_precedes_data_channel_retirement(adapter, has_stream): + source = importlib.import_module("pilot_grpc_source") + calls = [] + + async def acknowledge(): + assert "end" not in calls + calls.append("ack") + + async def end(): + assert calls == ["ack"] + calls.append("end") + + async def read(): + calls.append("drain") + + async def check(): + await source.finish_source( + SimpleNamespace(end=end) if has_stream else None, + read() if has_stream else None, + SimpleNamespace(finish=acknowledge), + ) + + asyncio.run(check()) + assert calls == (["ack", "end", "drain"] if has_stream else ["ack"]) + + +def test_source_anchor_latches_only_after_initial_clock_gate_and_never_on_recovery( + adapter, monkeypatch +): + module = importlib.import_module("pilot_source_control") + healthy, acknowledgements = [False], [] + + class Bounds: + def require(self, now): + if not healthy[0]: + raise module.ClockMappingError("synthetic excessive uncertainty") + + def to_dict(self): + return {} + + def uncertainty_ns(self, now): + return 1_000_000 if healthy[0] else 6_000_000 + + class Client: + def __init__(self, *args, **kwargs): + pass + + async def poll(self): + return object(), None + + async def acknowledge(self, probe, anchor, *, ended): + acknowledgements.append(anchor) + + monkeypatch.setattr(module, "StreamControlClient", Client) + monkeypatch.setattr( + module, + "ClockWindow", + lambda *a, **kw: SimpleNamespace(add=lambda probe: Bounds(), current=lambda now: Bounds()), + ) + control = module.SourceControl( + "unused", + b"", + { + "activation": identity().to_dict(), + "token": "b" * 64, + "remote_clock_id": "worker", + "source_zero_ns": str(10**12), + }, + ) + + async def check(): + await control.update() + assert control.anchor is None and not control.ready() and not acknowledgements + healthy[0] = True + await control.update() + anchor = await control.wait_anchor() + assert anchor == control.anchor and control.ready() + healthy[0] = False + await control.update() + assert not control.ready() and control.anchor == anchor + healthy[0] = True + await control.update() + assert control.ready() and acknowledgements == [anchor] * 3 + assert [row["acknowledged"] for row in control.samples] == [False, True, True, True] + + asyncio.run(check()) + + +def test_cross_host_bridge_clock_expiry_retains_owner_and_rotates_epoch(tmp_path, tls, adapter): # noqa: F811 + module, _ = adapter + cert, key = tmp_path / "cert", tmp_path / "key" + cert.write_bytes(tls[0]) + key.write_bytes(tls[1]) + run = StreamingLifecycle( + identity(), + tmp_path / "lease", + StreamMailbox(), + threading.Event(), + clock_ns=lambda: 1_000_000_000, + recover_input=True, + source_clock_ns=lambda: bridge.source.source_now(time.monotonic_ns()), + ) + run.ready() + bridge = module.NetworkGraphBridge( + run, + SimpleNamespace(close=lambda: None), + {}, + control=tmp_path, + source_status=tmp_path / "unused", + source_zero=10**12, + certificate=cert, + private_key=key, + address="localhost:0", + reset_temporal=lambda: None, + cross_host=True, + ) + real_serve = bridge.endpoint.serve + + async def serve(*args, **kwargs): + server, port = await real_serve(*args, **kwargs) + bridge.port = port + return server, port + + bridge.endpoint.serve = serve + + async def check(): + bridge.start() + ticket = ControlTicket(identity(), bridge.ticket.token) + client = StreamControlClient(f"localhost:{bridge.port}", tls[0], ticket, clock_id="source") + try: + probe, missing = await client.poll() + assert missing is None and run.continuity.phase == "waiting" + anchor = SourceAnchor(10**12, probe.local_receive_ns + 1_000_000_000) + await client.acknowledge(probe, anchor) + await eventually(lambda: bridge.access is not None) + first = bridge.access + stream = GrpcStreamClient(f"localhost:{bridge.port}", tls[0], first) + await stream.open() + # The listener keeps observing clock freshness even with no data. + # Shorten the TEST clock evidence lifetime, not any production gate. + bridge.source.window.maximum_age_ns = 50_000_000 + await eventually(lambda: run.continuity.phase == "waiting") + await stream.close() + await eventually(lambda: bridge.endpoint.active is None) + assert not run.stop_event.is_set() and run.lease.start == identity() + run.renew(identity()) + bridge.source.window.maximum_age_ns = 2_000_000_000 + # Old evidence must not itself resume; new controller epoch still required. + await asyncio.sleep(0.03) + probe, _ = await client.poll() + await client.acknowledge(probe, anchor) + await eventually(lambda: bridge.access.epoch != first.epoch) + assert run.continuity.phase == "synchronizing" + assert len(bridge.clock_states) >= 3 + finally: + await client.close() + + try: + asyncio.run(check()) + finally: + run.request_stop("completed") + assert bridge.close() and run.close() + assert not (tmp_path / "bootstrap.json").exists() diff --git a/tests/test_perception_realtime_scene.py b/tests/test_perception_realtime_scene.py index 04048ec..085a667 100644 --- a/tests/test_perception_realtime_scene.py +++ b/tests/test_perception_realtime_scene.py @@ -362,6 +362,48 @@ def test_consumer_expires_only_old_cell_and_rehashes_view_without_mutating_wire( pilot.assess_receipt(view, derived, bundle=bundle, now_ns=bundle["due_ns"]) +def test_cross_host_uncertainty_expires_per_cell_and_survives_publication(pilot): + payload, bundle, ddr = cell_input() + # Without uncertainty old ground would be 249 ms old; upper bound is253 ms. + bundle["clock_observer"] = lambda _: (NOW + 49_000_000, 4_000_000) + pilot.prepare_publication(payload, bundle, ddr, epoch_id="pilot", now_ns=123, mode="per-cell") + assert payload["cell_assessment"]["expired_ground_cells"] == 1 + assert payload["policy_actions"] == [2, 0] + assert payload["freshness_at_publication"]["clock_uncertainty_ms"] == 4 + fresh = pilot.validate_receipt(payload, bundle, epoch_id="pilot") + # An independently measured later receipt can have a slightly smaller + # midpoint while still overlapping the published interval. + local_bundle = {k: v for k, v in bundle.items() if k != "clock_observer"} + view, checked = pilot.assess_receipt( + payload, fresh, bundle=local_bundle, now_ns=bundle["due_ns"] + 48_000_000 + ) + assert checked.clock_uncertainty_ms == 0 and view["policy_actions"] == [2, 0] + with pytest.raises(ValueError, match="backwards"): + pilot.assess_receipt( + payload, fresh, bundle=local_bundle, now_ns=bundle["due_ns"] + 44_000_000 + ) + + +def test_one_clock_snapshot_per_publication_and_receipt_boundary(pilot): + payload, bundle, ddr = cell_input() + calls = [] + + def observe(now): + assert now not in calls # A later refresh may no longer describe this instant. + calls.append(now) + return NOW + now, 4_000_000 + + bundle["clock_observer"] = observe + pilot.prepare_publication( + payload, bundle, ddr, epoch_id="pilot", now_ns=260_000_000, mode="per-cell" + ) + assert payload["policy_actions"] == [2, 2] # Exercises suppression/reassessment too. + fresh = pilot.validate_receipt(payload, bundle, epoch_id="pilot") + view, checked = pilot.assess_receipt(payload, fresh, bundle=bundle, now_ns=270_000_000) + assert calls == [260_000_000, 270_000_000] + assert not checked.fresh_complete and view["policy_actions"] == [2, 2] + + @pytest.mark.parametrize("missing,held_ms", [(True, 0), (False, 220)]) def test_cell_freshness_cannot_override_missing_lidar_or_stale_segmentation( pilot, missing, held_ms diff --git a/tests/test_perception_source_clock.py b/tests/test_perception_source_clock.py new file mode 100644 index 0000000..a81817a --- /dev/null +++ b/tests/test_perception_source_clock.py @@ -0,0 +1,124 @@ +"""Two-sided clock evidence and recoverable owner-preserving admission.""" + +import threading +from dataclasses import replace + +import pytest +from test_perception_streaming_grpc import identity + +from k1link.perception.streaming_clock import ( + ClockMappingError, + ClockProbe, + ClockReceipt, + ClockWindow, +) +from k1link.perception.streaming_continuity import StreamSuspended +from k1link.perception.streaming_lifecycle import StreamingLifecycle +from k1link.perception.streaming_queue import StreamMailbox +from k1link.perception.streaming_source_clock import SourceAnchor, SourceClockMonitor + + +def sample(number=1, outbound=1_000_000, inbound=2_000_000, offset=10**12): + t1 = 10**16 + number * 100_000_000 + return ClockProbe( + "source", + "worker", + str(number), + t1, + t1 + outbound + offset, + t1 + outbound + offset + 10_000, + t1 + outbound + inbound + 10_000, + ) + + +@pytest.mark.parametrize("outbound,inbound", [(1, 9_000_000), (9_000_000, 1), (1, 1)]) +@pytest.mark.parametrize("offset", [-(10**12), 0, 10**12]) +def test_responder_interval_contains_true_reverse_offset(outbound, inbound, offset): + probe = sample(outbound=outbound, inbound=inbound, offset=offset) + receipt = ClockReceipt(probe, probe.local_receive_ns + offset + 1_000_000) + window = ClockWindow("worker", "source", rate_ppm=500, timestamp_error_ns=50_000) + bounds = window.add(receipt) + lower, upper = bounds.offset_at(receipt.acknowledged_ns) + assert lower <= -offset <= upper + assert bounds.uncertainty_ns(receipt.acknowledged_ns + 1_000_000) > bounds.uncertainty_ns( + receipt.acknowledged_ns + ) + + +def test_source_timeline_immutable_expiring_and_observation_bounds_conservative(): + first = sample() + anchor = SourceAnchor(10**12, first.local_receive_ns + 2_000_000_000) + monitor = SourceClockMonitor("worker", anchor.source_zero_ns) + with pytest.raises(StreamSuspended): + monitor.observed(first.remote_send_ns) + now = first.local_receive_ns + 10**12 + 1_000_000 + monitor.observe(first, now, anchor, False) + midpoint, uncertainty = monitor.observed(now) + actual = anchor.source_zero_ns + (now - 10**12) - anchor.local_zero_ns + assert midpoint - uncertainty <= actual <= monitor.source_now(now) + with pytest.raises(StreamSuspended): + monitor.observed(now + 2_000_000_000) + second = sample(30) + later = second.local_receive_ns + 10**12 + 1_000_000 + with pytest.raises(ValueError, match="immutable"): + monitor.observe( + second, later, replace(anchor, local_zero_ns=anchor.local_zero_ns + 1), False + ) + monitor.observe(second, later, anchor, True) + assert monitor.ended + with pytest.raises(ValueError, match="ended"): + monitor.observe(sample(31), later + 100_000_000, anchor, False) + + +def test_bad_ack_order_expired_jump_and_foreign_clock(): + first = sample() + for stamp in (first.remote_send_ns - 1, first.remote_receive_ns + 500_000_001): + with pytest.raises(ClockMappingError): + ClockReceipt(first, stamp) + monitor = SourceClockMonitor("worker", 10**12) + anchor = SourceAnchor(10**12, first.local_receive_ns + 2_000_000_000) + monitor.observe(first, first.local_receive_ns + 10**12 + 1_000_000, anchor, False) + jumped = sample(2, offset=10**12 + 100_000_000) + with pytest.raises(ClockMappingError, match="envelope"): + monitor.observe(jumped, jumped.local_receive_ns + 10**12 + 101_000_000, anchor, False) + with pytest.raises(StreamSuspended): + monitor.observed(jumped.local_receive_ns + 10**12 + 101_000_000) + + +def test_clock_wait_does_not_stop_owner_or_reopen_mailbox_prematurely(tmp_path): + healthy = [True] + + def source_now(): + if not healthy[0]: + raise StreamSuspended("synthetic expired clock mapping") + return 1_000_000_000 + + run = StreamingLifecycle( + identity(), + tmp_path, + StreamMailbox(), + threading.Event(), + clock_ns=lambda: 1_000_000_000, + recover_input=True, + source_clock_ns=source_now, + ) + run.ready() + try: + healthy[0] = False + with pytest.raises(StreamSuspended): + run.check_input(identity()) + assert run.continuity.reason == "source-clock" and not run.stop_event.is_set() + run.renew(identity()) + with pytest.raises(StreamSuspended): + run.begin_input(identity()) + assert run.mailbox.epoch_drained and run.continuity.phase == "waiting" + healthy[0] = True + epoch = run.begin_input(identity()) + assert epoch != identity() and epoch.lease_generation == 1 + assert run.continuity.phase == "synchronizing" + healthy[0] = False + with pytest.raises(StreamSuspended, match="obsolete"): + run.check_input(identity()) + assert run.continuity.phase == "synchronizing" # Stale peer cannot pause new epoch. + finally: + assert run.close() diff --git a/tests/test_perception_streaming_control_grpc.py b/tests/test_perception_streaming_control_grpc.py index 1887a8a..76c9af7 100644 --- a/tests/test_perception_streaming_control_grpc.py +++ b/tests/test_perception_streaming_control_grpc.py @@ -28,6 +28,7 @@ from k1link.perception.streaming_control_grpc import ( from k1link.perception.streaming_grpc import GrpcStreamClient, GrpcStreamEndpoint, StreamAccess from k1link.perception.streaming_lifecycle import StreamingLifecycle from k1link.perception.streaming_queue import StreamMailbox +from k1link.perception.streaming_source_clock import SourceAnchor, SourceClockMonitor @asynccontextmanager @@ -204,3 +205,52 @@ def test_malicious_response_nonce_or_offer_rejected(tls): await client.close() asyncio.run(check()) + + +def test_acknowledged_clock_is_single_use_and_coexists_with_active_stream(tmp_path, tls): + async def check(): + async with controlled(tmp_path, tls) as (_, endpoint, control, client, pending, port): + monitor = SourceClockMonitor("worker", 10**12) + control.observe = monitor.observe + probe, access = await client.poll() + anchor = SourceAnchor(10**12, probe.local_receive_ns + 2_000_000_000) + await client.acknowledge(probe, anchor) + assert monitor.observed(time.monotonic_ns())[1] > 0 + with pytest.raises(grpc.aio.AioRpcError): + await client.acknowledge(probe, anchor) + stream = GrpcStreamClient(f"localhost:{port}", tls[0], access) + try: + await stream.open() + pending[0] = None + for _ in range(12): + await asyncio.sleep(0.03) + probe, missing = await client.poll() + assert missing is None + await client.acknowledge(probe, anchor) + assert monitor.observed(time.monotonic_ns())[1] < 5_000_000 + assert endpoint.active is not None + finally: + await stream.close() + + asyncio.run(check()) + + +@pytest.mark.parametrize("fault", ["nonce", "remote_send", "anchor", "end_type"]) +def test_wrong_clock_receipt_cannot_update_mapping(tmp_path, tls, fault): + async def check(): + async with controlled(tmp_path, tls) as (_, _, control, client, _, _): + monitor = SourceClockMonitor("worker", 10**12) + control.observe = monitor.observe + probe, _ = await client.poll() + anchor = SourceAnchor(10**12, probe.local_receive_ns + 2_000_000_000) + if fault == "nonce": + probe = replace(probe, nonce="foreign") + elif fault == "remote_send": + probe = replace(probe, remote_send_ns=probe.remote_send_ns + 1) + elif fault == "anchor": + anchor = replace(anchor, source_zero_ns=anchor.source_zero_ns + 1) + with pytest.raises((grpc.aio.AioRpcError, ValueError)): + await client.acknowledge(probe, anchor, ended=1 if fault == "end_type" else False) + assert monitor.anchor is None and monitor.window is None + + asyncio.run(check())