diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_binary_bridge_probe.py b/experiments/perception/worker/streaming_profile_stage1/pilot_binary_bridge_probe.py index ec261e6..334f0e7 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_binary_bridge_probe.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_binary_bridge_probe.py @@ -85,6 +85,7 @@ def run(args): renewer.start() client = bridge = None decoded = [] + timings = [] with (root / "decoder.log").open("wb") as log: try: command = ["python3", "-B", "/probe/pilot_fragment_decoder.py"] @@ -115,6 +116,14 @@ def run(args): bridge = BinaryGraphBridge(runtime, client, source, source_report) bridge.start() while (bundle := runtime.mailbox.take()) is not None: + timings.append( + { + "sequence": bundle["sequence"], + "decode_ms": bundle["decode_ms"], + "decode_rpc_ms": bundle["decode_rpc_ms"], + "source_due_to_dequeue_ms": (time.monotonic_ns() - bundle["due_ns"]) / 1e6, + } + ) decoded.append(digest(bundle)) runtime.mailbox.release(bundle) if runtime.mailbox.error and not args.fault_decoder: @@ -132,6 +141,7 @@ def run(args): report["source"] = source_report report["peak_input_bytes"] = runtime.mailbox.peak_bytes report["residual_input_bytes"] = runtime.mailbox.bytes + report["timings"] = timings (root / "report.json").write_text(json.dumps(report, indent=2) + "\n") if args.fault_decoder: assert not decoded and report["released"] and not runtime.mailbox.bytes diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_fragment_decoder.py b/experiments/perception/worker/streaming_profile_stage1/pilot_fragment_decoder.py index 8b74ff9..01ec72b 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_fragment_decoder.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_fragment_decoder.py @@ -28,7 +28,7 @@ def main(): "decode_ms": (time.monotonic_ns() - begin) / 1e6, "frame_index": decoder.frames - 1, }, - image.tobytes(), + memoryview(image), ) del image else: diff --git a/experiments/perception/worker/streaming_profile_stage1/pilot_ipc.py b/experiments/perception/worker/streaming_profile_stage1/pilot_ipc.py index b217df0..d1065ed 100644 --- a/experiments/perception/worker/streaming_profile_stage1/pilot_ipc.py +++ b/experiments/perception/worker/streaming_profile_stage1/pilot_ipc.py @@ -32,10 +32,18 @@ def receive(stream, max_payload=MAX_PAYLOAD): def send(stream, header, payload=b""): - if len(payload) > MAX_PAYLOAD: + # Borrow a contiguous buffer until this synchronous write completes. Never + # concatenate the whole BGR payload with framing or make an implicit copy. + view = memoryview(payload).cast("B") + if len(view) > MAX_PAYLOAD: raise ValueError("IPC payload exceeds budget") - raw = json.dumps({**header, "payload_bytes": len(payload)}, allow_nan=False).encode() + raw = json.dumps({**header, "payload_bytes": len(view)}, allow_nan=False).encode() if len(raw) > 65536: raise ValueError("IPC header exceeds budget") - stream.write(struct.pack("