chore(perception): trace Python GC pauses
This commit is contained in:
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user