From 19108c9b926fc1edec2bbad23ffd14ddd1146f04 Mon Sep 17 00:00:00 2001 From: juliet Date: Tue, 16 Jun 2026 17:50:50 -0400 Subject: [PATCH] =?UTF-8?q?Add=20stress.reporter=20=E2=80=94=20pure=20stat?= =?UTF-8?q?s=20for=20the=20upcoming=20stress=20harness=20(#328)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Claude Opus 4.7 (1M context) Co-authored-by: Cassidy James --- .../src/osprey/worker/stress/__init__.py | 0 .../src/osprey/worker/stress/reporter.py | 187 ++++++++++++++++++ .../osprey/worker/stress/tests/__init__.py | 0 .../worker/stress/tests/test_reporter.py | 162 +++++++++++++++ 4 files changed, 349 insertions(+) create mode 100644 osprey_worker/src/osprey/worker/stress/__init__.py create mode 100644 osprey_worker/src/osprey/worker/stress/reporter.py create mode 100644 osprey_worker/src/osprey/worker/stress/tests/__init__.py create mode 100644 osprey_worker/src/osprey/worker/stress/tests/test_reporter.py diff --git a/osprey_worker/src/osprey/worker/stress/__init__.py b/osprey_worker/src/osprey/worker/stress/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/osprey_worker/src/osprey/worker/stress/reporter.py b/osprey_worker/src/osprey/worker/stress/reporter.py new file mode 100644 index 0000000..9f6fe5e --- /dev/null +++ b/osprey_worker/src/osprey/worker/stress/reporter.py @@ -0,0 +1,187 @@ +"""Stress harness reporting. + +Pure functions that consume produced/consumed event maps and emit a structured +report plus an exit code. No Kafka, no I/O. Tested in isolation under +`tests/test_reporter.py`. + +The harness has two measurement modes: + +* Closed-loop (synthetic source): we know every action_id we sent and when we + sent it, so we can compute exact drop rate and per-event latency. +* Open-loop (external source — e.g. jetstream once #236 lands): we don't have + per-event produce timestamps, so latency is unavailable; we report aggregate + throughput and drop rate via input/output counts. +""" + +from __future__ import annotations + +import json +from collections.abc import Hashable, Mapping +from dataclasses import asdict, dataclass, field +from typing import Optional + + +@dataclass(frozen=True) +class LatencyStats: + min_ms: float + p50_ms: float + p95_ms: float + p99_ms: float + max_ms: float + count: int + + +@dataclass(frozen=True) +class Thresholds: + drop_rate: Optional[float] = None + p95_ms: Optional[float] = None + + +@dataclass(frozen=True) +class Report: + mode: str # "closed-loop" | "open-loop" + produced: int + consumed: int + matched: int + drop_count: int + drop_rate: float + duration_seconds: float + latency: Optional[LatencyStats] = None + thresholds: Thresholds = field(default_factory=Thresholds) + threshold_breaches: tuple[str, ...] = field(default_factory=tuple) + + def to_json(self) -> str: + payload = asdict(self) + return json.dumps(payload, indent=2, sort_keys=True) + + def to_human(self) -> str: + lines = [ + f'mode: {self.mode}', + f'produced: {self.produced}', + f'consumed: {self.consumed}', + f'matched: {self.matched}', + f'drop count: {self.drop_count}', + f'drop rate: {self.drop_rate:.4f}', + f'duration: {self.duration_seconds:.2f}s', + ] + if self.latency is not None: + lines += [ + 'latency (ms):', + f' min: {self.latency.min_ms:.1f}', + f' p50: {self.latency.p50_ms:.1f}', + f' p95: {self.latency.p95_ms:.1f}', + f' p99: {self.latency.p99_ms:.1f}', + f' max: {self.latency.max_ms:.1f}', + ] + if self.threshold_breaches: + lines.append('breaches: ' + ', '.join(self.threshold_breaches)) + return '\n'.join(lines) + + +def percentile(sorted_values: list[float], p: float) -> float: + """Nearest-rank percentile over a pre-sorted list. p in [0, 100].""" + if not sorted_values: + raise ValueError('percentile of empty list') + if not 0 <= p <= 100: + raise ValueError(f'p must be in [0, 100], got {p}') + # Nearest-rank: ceil(p/100 * n) - 1, clamped. + n = len(sorted_values) + if p == 0: + return sorted_values[0] + rank = int(-(-p * n // 100)) - 1 # ceil division + rank = max(0, min(rank, n - 1)) + return sorted_values[rank] + + +def compute_latency_stats(latencies_ms: list[float]) -> LatencyStats: + if not latencies_ms: + raise ValueError('compute_latency_stats requires a non-empty list of latencies') + sorted_ms = sorted(latencies_ms) + return LatencyStats( + min_ms=sorted_ms[0], + p50_ms=percentile(sorted_ms, 50), + p95_ms=percentile(sorted_ms, 95), + p99_ms=percentile(sorted_ms, 99), + max_ms=sorted_ms[-1], + count=len(sorted_ms), + ) + + +def compute_report( + *, + produced: Mapping[Hashable, float], + consumed: Mapping[Hashable, float], + duration_seconds: float, + thresholds: Thresholds = Thresholds(), +) -> Report: + """Build a Report from per-event produce / consume timestamps. + + Keys can be any hashable type (int action_ids today; str in earlier + iterations). Only set membership and ordering matter. + + Args: + produced: id -> wall-clock timestamp (seconds) at which we sent the event. + Empty in open-loop mode. + consumed: id -> wall-clock timestamp (seconds) at which we observed + the event's execution result. + duration_seconds: wall-clock time the harness ran. + thresholds: pass/fail gates. None entries are not enforced. + """ + is_closed_loop = bool(produced) + mode = 'closed-loop' if is_closed_loop else 'open-loop' + + produced_count = len(produced) + consumed_count = len(consumed) + + if is_closed_loop: + matched_ids = produced.keys() & consumed.keys() + matched = len(matched_ids) + drop_count = produced_count - matched + drop_rate = drop_count / produced_count if produced_count else 0.0 + latencies = [(consumed[aid] - produced[aid]) * 1000 for aid in matched_ids] + latency = compute_latency_stats(latencies) if latencies else None + else: + # Open-loop: best we can do without per-event produce timestamps is + # treat consumed as our denominator. Drop accounting needs an external + # source-of-truth count, which the caller would have to supply via a + # different code path. For now, report no drops. + matched = consumed_count + drop_count = 0 + drop_rate = 0.0 + latency = None + + breaches: list[str] = [] + if thresholds.drop_rate is not None and drop_rate > thresholds.drop_rate: + breaches.append(f'drop_rate {drop_rate:.4f} > {thresholds.drop_rate:.4f}') + if thresholds.p95_ms is not None and latency is not None and latency.p95_ms > thresholds.p95_ms: + breaches.append(f'p95_ms {latency.p95_ms:.1f} > {thresholds.p95_ms:.1f}') + + return Report( + mode=mode, + produced=produced_count, + consumed=consumed_count, + matched=matched, + drop_count=drop_count, + drop_rate=drop_rate, + duration_seconds=duration_seconds, + latency=latency, + thresholds=thresholds, + threshold_breaches=tuple(breaches), + ) + + +# Exit codes used by the CLI. +EXIT_OK = 0 +EXIT_DROP_RATE_EXCEEDED = 1 +EXIT_LATENCY_EXCEEDED = 2 +EXIT_INTERNAL_ERROR = 3 + + +def exit_code_for(report: Report) -> int: + """Map report breaches to a CLI exit code.""" + for breach in report.threshold_breaches: + if breach.startswith('drop_rate'): + return EXIT_DROP_RATE_EXCEEDED + if breach.startswith('p95_ms'): + return EXIT_LATENCY_EXCEEDED + return EXIT_OK diff --git a/osprey_worker/src/osprey/worker/stress/tests/__init__.py b/osprey_worker/src/osprey/worker/stress/tests/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/osprey_worker/src/osprey/worker/stress/tests/test_reporter.py b/osprey_worker/src/osprey/worker/stress/tests/test_reporter.py new file mode 100644 index 0000000..840b6dd --- /dev/null +++ b/osprey_worker/src/osprey/worker/stress/tests/test_reporter.py @@ -0,0 +1,162 @@ +import json + +import pytest +from osprey.worker.stress.reporter import ( + EXIT_DROP_RATE_EXCEEDED, + EXIT_LATENCY_EXCEEDED, + EXIT_OK, + Thresholds, + compute_latency_stats, + compute_report, + exit_code_for, + percentile, +) + + +class TestPercentile: + def test_raises_on_empty(self) -> None: + with pytest.raises(ValueError, match='empty'): + percentile([], 50) + + def test_raises_out_of_range(self) -> None: + with pytest.raises(ValueError): + percentile([1.0], -1) + with pytest.raises(ValueError): + percentile([1.0], 101) + + def test_single_value(self) -> None: + assert percentile([5.0], 0) == 5.0 + assert percentile([5.0], 50) == 5.0 + assert percentile([5.0], 100) == 5.0 + + def test_nearest_rank_known_values(self) -> None: + # 1..10 sorted; nearest-rank p95 = ceil(0.95*10)-1 = 9, sorted[9] = 10 + vs = [float(i) for i in range(1, 11)] + assert percentile(vs, 50) == 5.0 # ceil(5)-1 = 4 → sorted[4] = 5 + assert percentile(vs, 95) == 10.0 + assert percentile(vs, 99) == 10.0 + assert percentile(vs, 0) == 1.0 + + +class TestComputeLatencyStats: + def test_all_fields_populated(self) -> None: + stats = compute_latency_stats([100.0, 50.0, 25.0, 75.0, 200.0]) + assert stats.min_ms == 25.0 + assert stats.max_ms == 200.0 + assert stats.count == 5 + # Sorted: 25, 50, 75, 100, 200; p50 = sorted[ceil(2.5)-1] = sorted[2] = 75 + assert stats.p50_ms == 75.0 + + def test_raises_on_empty(self) -> None: + with pytest.raises(ValueError, match='non-empty'): + compute_latency_stats([]) + + +class TestComputeReport: + def test_closed_loop_perfect_delivery(self) -> None: + produced = {f'id-{i}': float(i) for i in range(10)} + consumed = {f'id-{i}': float(i) + 0.1 for i in range(10)} + r = compute_report(produced=produced, consumed=consumed, duration_seconds=1.0) + assert r.mode == 'closed-loop' + assert r.produced == 10 + assert r.consumed == 10 + assert r.matched == 10 + assert r.drop_count == 0 + assert r.drop_rate == 0.0 + assert r.latency is not None + assert r.latency.count == 10 + # latency ≈ 100ms across the board + assert 99.0 < r.latency.p50_ms < 101.0 + assert r.threshold_breaches == () + + def test_closed_loop_partial_drop(self) -> None: + produced = {f'id-{i}': 0.0 for i in range(10)} + consumed = {f'id-{i}': 0.1 for i in range(7)} # 3 dropped + r = compute_report(produced=produced, consumed=consumed, duration_seconds=1.0) + assert r.drop_count == 3 + assert r.drop_rate == 0.3 + assert r.matched == 7 + + def test_closed_loop_consumed_includes_unknown_ids(self) -> None: + # Consumer saw some IDs we never sent (e.g. other test traffic). + # Those should not count as matched but should appear in consumed total. + produced = {'id-1': 0.0, 'id-2': 0.0} + consumed = {'id-1': 0.1, 'id-2': 0.1, 'other-id': 0.1} + r = compute_report(produced=produced, consumed=consumed, duration_seconds=1.0) + assert r.produced == 2 + assert r.consumed == 3 + assert r.matched == 2 + assert r.drop_count == 0 + + def test_open_loop_no_produced(self) -> None: + # No produce timestamps → can't compute drop or latency, just throughput. + consumed = {f'id-{i}': 0.1 for i in range(50)} + r = compute_report(produced={}, consumed=consumed, duration_seconds=10.0) + assert r.mode == 'open-loop' + assert r.produced == 0 + assert r.consumed == 50 + assert r.matched == 50 + assert r.drop_count == 0 + assert r.drop_rate == 0.0 + assert r.latency is None + + def test_drop_rate_breach_recorded(self) -> None: + produced = {f'id-{i}': 0.0 for i in range(10)} + consumed = {f'id-{i}': 0.1 for i in range(5)} + r = compute_report( + produced=produced, + consumed=consumed, + duration_seconds=1.0, + thresholds=Thresholds(drop_rate=0.1), + ) + assert any('drop_rate' in b for b in r.threshold_breaches) + assert exit_code_for(r) == EXIT_DROP_RATE_EXCEEDED + + def test_p95_breach_recorded(self) -> None: + # p95 will be 500ms, threshold is 100ms → breach. + produced = {f'id-{i}': 0.0 for i in range(10)} + consumed = {f'id-{i}': 0.5 for i in range(10)} + r = compute_report( + produced=produced, + consumed=consumed, + duration_seconds=1.0, + thresholds=Thresholds(p95_ms=100.0), + ) + assert any('p95_ms' in b for b in r.threshold_breaches) + assert exit_code_for(r) == EXIT_LATENCY_EXCEEDED + + def test_no_breach_returns_ok(self) -> None: + produced = {f'id-{i}': 0.0 for i in range(10)} + consumed = {f'id-{i}': 0.05 for i in range(10)} + r = compute_report( + produced=produced, + consumed=consumed, + duration_seconds=1.0, + thresholds=Thresholds(drop_rate=0.01, p95_ms=100.0), + ) + assert r.threshold_breaches == () + assert exit_code_for(r) == EXIT_OK + + def test_empty_produced_and_consumed(self) -> None: + r = compute_report(produced={}, consumed={}, duration_seconds=1.0) + assert r.mode == 'open-loop' + assert r.drop_rate == 0.0 + assert r.latency is None + + def test_json_roundtrip(self) -> None: + produced = {'id-1': 0.0} + consumed = {'id-1': 0.1} + r = compute_report(produced=produced, consumed=consumed, duration_seconds=1.0) + payload = json.loads(r.to_json()) + assert payload['mode'] == 'closed-loop' + assert payload['matched'] == 1 + assert 'latency' in payload + + def test_human_format_contains_key_fields(self) -> None: + produced = {'id-1': 0.0} + consumed = {'id-1': 0.1} + r = compute_report(produced=produced, consumed=consumed, duration_seconds=1.0) + out = r.to_human() + assert 'mode:' in out + assert 'drop rate:' in out + assert 'latency' in out -- 2.51.2