From d1e6e42cf413f2d22a3b5fbe79079d165b143300 Mon Sep 17 00:00:00 2001 From: DCCONSTRUCTIONS Date: Tue, 25 Aug 2026 18:38:33 +0300 Subject: [PATCH] chore(perception): trace Python GC pauses --- .../run_m48s_reference_graph_shadow_worker.py | 129 ++++++++++++++---- tests/test_m48s_reference_graph_timing.py | 26 +++- 2 files changed, 127 insertions(+), 28 deletions(-) diff --git a/experiments/perception/run_m48s_reference_graph_shadow_worker.py b/experiments/perception/run_m48s_reference_graph_shadow_worker.py index 3678438..8987c59 100644 --- a/experiments/perception/run_m48s_reference_graph_shadow_worker.py +++ b/experiments/perception/run_m48s_reference_graph_shadow_worker.py @@ -4,6 +4,7 @@ from __future__ import annotations import argparse +import gc import hashlib import json import resource @@ -106,6 +107,58 @@ class GpuTelemetry: self._stop.wait(self.interval_seconds) +class GcPauseTelemetry: + """Attribute stop-the-world cyclic-GC pauses to active graph stages.""" + + def __init__(self) -> None: + self.events: list[dict[str, object]] = [] + self._starts: dict[int, tuple[int, int, str | None, int | None]] = {} + self._stage = threading.local() + + def __enter__(self) -> GcPauseTelemetry: + gc.callbacks.append(self._observe) + return self + + def __exit__(self, *_args: object) -> None: + gc.callbacks.remove(self._observe) + + def enter_stage(self, stage_id: str, sequence: int) -> None: + self._stage.value = (stage_id, sequence) + + def exit_stage(self) -> None: + self._stage.value = None + + def _observe(self, phase: str, info: dict[str, int]) -> None: + thread_id = threading.get_ident() + if phase == "start": + stage = getattr(self._stage, "value", None) + stage_id, sequence = stage if stage is not None else (None, None) + self._starts[thread_id] = ( + time.perf_counter_ns(), + info["generation"], + stage_id, + sequence, + ) + return + if phase != "stop": + return + started = self._starts.pop(thread_id, None) + if started is None: + return + started_ns, generation, stage_id, sequence = started + self.events.append( + { + "duration_ns": max(0, time.perf_counter_ns() - started_ns), + "generation": generation, + "collected": info["collected"], + "uncollectable": info["uncollectable"], + "thread_name": threading.current_thread().name, + "stage_id": stage_id, + "sequence": sequence, + } + ) + + class FrameTimingStore: """Join bounded decode, detector and provider timings by source sequence.""" @@ -175,21 +228,31 @@ class FrameTimingStore: class TimedProviderProxy: """Record one provider's actual call duration without another inference pass.""" - def __init__(self, stage_id: str, provider: object, store: FrameTimingStore) -> None: + def __init__( + self, + stage_id: str, + provider: object, + store: FrameTimingStore, + gc_telemetry: GcPauseTelemetry, + ) -> None: self.stage_id = stage_id self.provider = provider self.store = store + self.gc_telemetry = gc_telemetry self.provider_id = cast(Any, provider).provider_id def _call(self, sequence: int, method: str, *args: object) -> object: started_ns = time.perf_counter_ns() + self.gc_telemetry.enter_stage(self.stage_id, sequence) try: return getattr(self.provider, method)(*args) finally: + duration_ns = max(0, time.perf_counter_ns() - started_ns) + self.gc_telemetry.exit_stage() self.store.observe_provider( self.stage_id, sequence, - max(0, time.perf_counter_ns() - started_ns), + duration_ns, ) def associate(self, packet: SourcePacket, proposals: object) -> object: @@ -275,6 +338,7 @@ def main() -> int: progress.open("x", encoding="utf-8") as progress_stream, frame_ledger.open("x", encoding="utf-8") as frame_ledger_stream, GpuTelemetry(arguments.telemetry_interval_seconds) as gpu, + GcPauseTelemetry() as gc_telemetry, ): for loop_index in range(arguments.loops): loop_completion_ages_ns: list[int] = [] @@ -314,6 +378,7 @@ def main() -> int: stage_id, getattr(runtime.graph, attribute), timing_store, + gc_telemetry, ), ) loop_started_ns = time.monotonic_ns() @@ -370,8 +435,7 @@ def main() -> int: loop_documents.append(loop_document) completion_ages_ms.extend(value / 1_000_000.0 for value in loop_completion_ages_ns) map_output_ages_ms.extend( - delivery.obstacle_map.output_age_ns / 1_000_000.0 - for delivery in result.deliveries + delivery.obstacle_map.output_age_ns / 1_000_000.0 for delivery in result.deliveries ) all_deliveries.extend(result.deliveries) all_pipeline_timings.extend(loop_pipeline_timings) @@ -401,20 +465,15 @@ def main() -> int: accounting.update(cast(Mapping[str, int], loop["terminal_outcomes"])) admitted = sum(cast(int, loop["admitted_count"]) for loop in loop_documents) delivered = len(all_deliveries) - processing_wall_seconds = sum( - cast(float, loop["wall_seconds"]) for loop in loop_documents - ) + processing_wall_seconds = sum(cast(float, loop["wall_seconds"]) for loop in loop_documents) queue_high_watermarks = { stage: max( - cast(dict[str, int], loop["queue_high_watermarks"])[stage] - for loop in loop_documents + cast(dict[str, int], loop["queue_high_watermarks"])[stage] for loop in loop_documents ) for stage in ("detector", "geometry", "temporal", "rolling", "threat") } advisories = tuple( - advisory - for delivery in all_deliveries - for advisory in project_m48s_advisories(delivery) + advisory for delivery in all_deliveries for advisory in project_m48s_advisories(delivery) ) identity = _identity_metrics(all_deliveries) semantic = _semantic_metrics(all_deliveries, advisories) @@ -487,6 +546,7 @@ def main() -> int: "identity_continuity": identity, "semantic_advisory": semantic, "pipeline_timing": _pipeline_timing_metrics(all_pipeline_timings), + "python_gc": _gc_telemetry_summary(gc_telemetry.events), "gpu": _telemetry_summary(gpu.samples), "process_peak_rss_before_mib": round(rss_before_kib / 1024.0, 6), "process_peak_rss_after_mib": round(rss_after_kib / 1024.0, 6), @@ -561,8 +621,7 @@ def _loop_document( ) -> dict[str, object]: outcomes = Counter(item.outcome.value for item in result.terminal_outcomes) outcome_stages = Counter( - f"{item.outcome.value}:{item.stage_id}:{item.reason}" - for item in result.terminal_outcomes + f"{item.outcome.value}:{item.stage_id}:{item.reason}" for item in result.terminal_outcomes ) return { "loop_index": loop_index, @@ -627,11 +686,7 @@ def _semantic_metrics( hints[proposal.semantic_hint or "unclassified-camera"] += 1 motions[MotionState.UNKNOWN.value] += 1 families = Counter(item.family.value for item in advisories) - responses = Counter( - response.value - for item in advisories - for response in item.responses - ) + responses = Counter(response.value for item in advisories for response in item.responses) return { "semantic_hint_counts": dict(sorted(hints.items())), "motion_counts": dict(sorted(motions.items())), @@ -719,9 +774,7 @@ def _pipeline_timing_metrics( key.removesuffix("_duration_ns"): _distribution(values) for key, values in detector_values.items() }, - "provider_ms": { - key: _distribution(values) for key, values in provider_values.items() - }, + "provider_ms": {key: _distribution(values) for key, values in provider_values.items()}, "graph_admission_to_delivery_ms": _distribution( top_level_values["graph_admission_to_delivery_ns"] ), @@ -732,9 +785,7 @@ def _pipeline_timing_metrics( "decode_to_delivery_processing_ms": _distribution( top_level_values["decode_to_delivery_processing_ns"] ), - "maximum_graph_sequence": ( - cast(int, maximum["sequence"]) if maximum is not None else None - ), + "maximum_graph_sequence": (cast(int, maximum["sequence"]) if maximum is not None else None), "additional_inference_passes": 0, } @@ -751,6 +802,34 @@ def _telemetry_summary(samples: list[dict[str, float]]) -> dict[str, Any]: return result +def _gc_telemetry_summary(events: list[dict[str, object]]) -> dict[str, object]: + durations_ms = [cast(int, event["duration_ns"]) / 1_000_000.0 for event in events] + generations = Counter(cast(int, event["generation"]) for event in events) + attributed = [event for event in events if event["stage_id"] is not None] + maximum = max(events, key=lambda event: cast(int, event["duration_ns"]), default=None) + return { + "event_count": len(events), + "generation_counts": { + str(generation): count for generation, count in sorted(generations.items()) + }, + "duration_ms": _distribution(durations_ms), + "pipeline_attributed_event_count": len(attributed), + "maximum_event": ( + { + "duration_ms": round(cast(int, maximum["duration_ns"]) / 1_000_000.0, 6), + "generation": maximum["generation"], + "collected": maximum["collected"], + "uncollectable": maximum["uncollectable"], + "thread_name": maximum["thread_name"], + "stage_id": maximum["stage_id"], + "sequence": maximum["sequence"], + } + if maximum is not None + else None + ), + } + + def _input_digests( paths: ReferenceGraphRuntimePaths, detector_profile: Path, diff --git a/tests/test_m48s_reference_graph_timing.py b/tests/test_m48s_reference_graph_timing.py index cafc713..ead3ba7 100644 --- a/tests/test_m48s_reference_graph_timing.py +++ b/tests/test_m48s_reference_graph_timing.py @@ -1,15 +1,14 @@ from __future__ import annotations import importlib.util +import time from pathlib import Path from k1link.perception.detector import DetectorFrameTiming from k1link.perception.recorded_source import DecodedFrameTiming REPOSITORY_ROOT = Path(__file__).resolve().parents[1] -RUNNER_PATH = ( - REPOSITORY_ROOT / "experiments/perception/run_m48s_reference_graph_shadow_worker.py" -) +RUNNER_PATH = REPOSITORY_ROOT / "experiments/perception/run_m48s_reference_graph_shadow_worker.py" SPEC = importlib.util.spec_from_file_location("m48s_timed_shadow_runner", RUNNER_PATH) assert SPEC is not None and SPEC.loader is not None RUNNER = importlib.util.module_from_spec(SPEC) @@ -81,3 +80,24 @@ def test_pipeline_timing_metrics_preserve_single_pass_stage_breakdown() -> None: assert metrics["graph_unattributed_ms"]["maximum"] == 6.0 assert metrics["maximum_graph_sequence"] == 3 assert metrics["additional_inference_passes"] == 0 + + +def test_gc_pause_telemetry_attributes_collection_to_active_stage() -> None: + telemetry = RUNNER.GcPauseTelemetry() + telemetry.enter_stage("rolling", 3928) + telemetry._observe("start", {"generation": 2}) + time.sleep(0.001) + telemetry._observe( + "stop", + {"generation": 2, "collected": 17, "uncollectable": 0}, + ) + telemetry.exit_stage() + + summary = RUNNER._gc_telemetry_summary(telemetry.events) + + assert summary["event_count"] == 1 + assert summary["pipeline_attributed_event_count"] == 1 + assert summary["maximum_event"]["generation"] == 2 + assert summary["maximum_event"]["stage_id"] == "rolling" + assert summary["maximum_event"]["sequence"] == 3928 + assert summary["maximum_event"]["duration_ms"] >= 1.0