import importlib import importlib.util import json import multiprocessing import threading from concurrent.futures import ThreadPoolExecutor from datetime import datetime, timezone from pathlib import Path import pytest from pydantic import ValidationError from agent_lab.application.benchmark import ( BenchmarkCaseId, BenchmarkConfig, BenchmarkModelCallTiming, BenchmarkMode, BenchmarkRunResult, ) def _reporting_module(): spec = importlib.util.find_spec( "agent_lab.infrastructure.benchmark_reporting" ) assert spec is not None, "benchmark_reporting module must be implemented" return importlib.import_module( "agent_lab.infrastructure.benchmark_reporting" ) def _config(**overrides: object) -> BenchmarkConfig: payload = { "schema_version": 1, "base_url": "https://provider.example/v1", "model": "benchmark-model", "runs_per_case": 1, "cases": [BenchmarkCaseId.ORDINARY_CHAT], "modes": [BenchmarkMode.DUAL_AGENT], } payload.update(overrides) return BenchmarkConfig.model_validate(payload) def _run( *, case_id: BenchmarkCaseId = BenchmarkCaseId.ORDINARY_CHAT, mode: BenchmarkMode = BenchmarkMode.DUAL_AGENT, iteration: int = 1, status: str = "passed", initial_provider_ttft_ms: int | None = 10, visible_ttft_ms: int | None = 20, turn_wall_time_ms: int | None = 30, prompt_tokens: int | None = 2, completion_tokens: int | None = 3, total_tokens: int | None = 5, cached_tokens: int | None = 1, model_call_count: int | None = 1, fallback_count: int | None = 0, tool_count: int | None = 0, tool_latencies_ms: list[int | None] | None = None, semantic_failures: list[str] | None = None, error: str | None = None, model_calls: list[BenchmarkModelCallTiming] | None = None, ) -> BenchmarkRunResult: return BenchmarkRunResult( case_id=case_id, mode=mode, iteration=iteration, status=status, initial_provider_ttft_ms=initial_provider_ttft_ms, visible_ttft_ms=visible_ttft_ms, turn_wall_time_ms=turn_wall_time_ms, prompt_tokens=prompt_tokens, completion_tokens=completion_tokens, total_tokens=total_tokens, cached_tokens=cached_tokens, model_call_count=model_call_count, fallback_count=fallback_count, tool_count=tool_count, event_names=[], batch_event_names=[], tool_event_names=[], event_sources=[], tool_statuses=[], tool_latencies_ms=tool_latencies_ms or [], semantic_failures=semantic_failures or [], error=error, model_calls=model_calls or [], ) def _model_call( call_index: int, elapsed_ms: int | None, ) -> BenchmarkModelCallTiming: return BenchmarkModelCallTiming( call_index=call_index, call_kind="chat_completion", first_item_kind="raw_chunk", provider_ttft_ms=1, visible_ttft_ms=2, elapsed_ms=elapsed_ms, usage=None, ) def _marker_report(marker: str): reporting = _reporting_module() return reporting.build_benchmark_report( _config(model=marker), [_run()], mock=True, ) def _process_report_writer( output_dir: str, marker: str, timestamp: str, first_temp_barrier, result_queue, ) -> None: reporting = _reporting_module() original_write_text = reporting._write_text first_write = True def synchronized_write(path: Path, content: str) -> tuple[int, int]: nonlocal first_write identity = original_write_text(path, content) if first_write: first_write = False first_temp_barrier.wait() return identity reporting._write_text = synchronized_write try: json_path, markdown_path = reporting.write_benchmark_report( _marker_report(marker), output_dir, timestamp=timestamp, ) except BaseException as exc: result_queue.put(("error", marker, repr(exc))) return result_queue.put( ("ok", marker, json_path.name, markdown_path.name) ) def _assert_concurrent_report_pairs( output_dir: Path, markers: list[str], ) -> None: json_paths = sorted(output_dir.glob("*.json")) markdown_paths = sorted(output_dir.glob("*.md")) assert len(json_paths) == len(markers) assert len(markdown_paths) == len(markers) assert {path.stem for path in json_paths} == { path.stem for path in markdown_paths } actual_markers = set() for json_path in json_paths: marker = json.loads(json_path.read_text())["target"]["model"] markdown = json_path.with_suffix(".md").read_text() actual_markers.add(marker) assert f"- Model: `{marker}`" in markdown assert actual_markers == set(markers) _assert_no_publication_artifacts(output_dir) def _assert_no_publication_artifacts(output_dir: Path) -> None: leftovers = [ path.name for path in output_dir.iterdir() if path.name.endswith((".tmp", ".lock")) or "placeholder" in path.name ] assert leftovers == [] report_paths = [ path for path in output_dir.iterdir() if path.suffix in {".json", ".md"} ] assert all(path.stat().st_size > 0 for path in report_paths) @pytest.mark.parametrize( ("values", "quantile", "expected"), [ ([], 0.50, None), ([7], 0.50, 7.0), ([7], 0.95, 7.0), ([1, 3, 5, 7], 0.50, 4.0), ([1, 3, 5, 7], 0.95, 6.7), ], ) def test_percentile_uses_linear_interpolation( values: list[int], quantile: float, expected: float | None ): reporting = _reporting_module() assert reporting.percentile(values, quantile) == pytest.approx(expected) def test_report_models_are_strict_and_include_required_summary_fields(): reporting = _reporting_module() generated_at = datetime(2026, 7, 14, 1, 2, 3, tzinfo=timezone.utc) report = reporting.build_benchmark_report( _config(), [_run()], mock=True, generated_at=generated_at, ) assert isinstance(report, reporting.BenchmarkReport) assert report.schema_version == 1 assert report.generated_at == generated_at assert report.target.model_dump() == { "base_url": "https://provider.example/v1", "model": "benchmark-model", "mock": True, } assert report.config.model_dump(mode="json") == { "runs_per_case": 1, "cases": ["ordinary_chat"], "modes": ["dual_agent"], } assert report.overall.model_dump() == { "attempted": 1, "passed": 1, "failed": 0, "warnings": report.overall.warnings, "system_errors": [], } with pytest.raises(ValidationError): reporting.BenchmarkReport.model_validate( report.model_dump() | {"schema_version": "1"} ) with pytest.raises(ValidationError): reporting.BenchmarkAggregate.model_validate( report.aggregates[0].model_dump() | {"unexpected": True} ) def test_aggregate_excludes_failed_and_null_samples_from_metrics(): reporting = _reporting_module() runs = [ _run( iteration=1, initial_provider_ttft_ms=10, visible_ttft_ms=20, turn_wall_time_ms=30, prompt_tokens=2, completion_tokens=3, total_tokens=5, cached_tokens=1, model_call_count=1, fallback_count=0, tool_count=2, ), _run( iteration=2, initial_provider_ttft_ms=None, visible_ttft_ms=40, turn_wall_time_ms=None, prompt_tokens=None, completion_tokens=7, total_tokens=7, cached_tokens=None, model_call_count=3, fallback_count=1, tool_count=None, ), _run( iteration=3, status="failed", initial_provider_ttft_ms=999, visible_ttft_ms=999, turn_wall_time_ms=999, prompt_tokens=999, completion_tokens=999, total_tokens=999, cached_tokens=999, model_call_count=999, fallback_count=999, tool_count=999, error="failed", ), ] aggregate = reporting.build_benchmark_report( _config(runs_per_case=3), runs, mock=True ).aggregates[0] assert (aggregate.attempted, aggregate.passed, aggregate.failed) == (3, 2, 1) assert aggregate.initial_provider_ttft_ms.model_dump() == { "count": 1, "p50": 10.0, "p95": 10.0, } assert aggregate.visible_ttft_ms.model_dump() == { "count": 2, "p50": 30.0, "p95": 39.0, } assert aggregate.turn_wall_time_ms.count == 1 assert aggregate.prompt_tokens.model_dump() == { "count": 1, "total": 2, "mean": 2.0, } assert aggregate.completion_tokens.model_dump() == { "count": 2, "total": 10, "mean": 5.0, } assert aggregate.total_tokens.total == 12 assert aggregate.cached_tokens.total == 1 assert aggregate.model_call_count.model_dump() == { "count": 2, "total": 4, "mean": 2.0, } assert aggregate.fallback_count.total == 1 assert aggregate.tool_count.model_dump() == { "count": 1, "total": 2, "mean": 2.0, } def test_aggregate_flattens_passed_model_and_tool_latency_samples(): reporting = _reporting_module() runs = [ _run( iteration=1, model_calls=[ _model_call(1, 10), _model_call(2, None), _model_call(3, 30), ], tool_latencies_ms=[5, None, 15], ), _run( iteration=2, model_calls=[_model_call(1, 20)], tool_latencies_ms=[25], ), _run( iteration=3, status="failed", model_calls=[_model_call(1, 999)], tool_latencies_ms=[999], ), ] report = reporting.build_benchmark_report( _config(runs_per_case=3), runs, mock=True ) aggregate = report.aggregates[0] markdown = reporting.render_markdown(report) assert aggregate.model_elapsed_ms.model_dump() == { "count": 3, "p50": 20.0, "p95": 29.0, } assert aggregate.tool_latency_ms.model_dump() == { "count": 3, "p50": 15.0, "p95": 24.0, } assert "Model elapsed p50/p95 ms" in markdown assert "Tool latency p50/p95 ms" in markdown assert "20.0/29.0" in markdown assert "15.0/24.0" in markdown def test_zero_latency_samples_emit_low_confidence_warnings(): reporting = _reporting_module() aggregate = reporting.build_benchmark_report( _config(), [ _run( initial_provider_ttft_ms=None, visible_ttft_ms=None, turn_wall_time_ms=None, model_calls=[], tool_latencies_ms=[], ) ], mock=True, ).aggregates[0] assert { warning.split(" p95 low confidence:", 1)[0] for warning in aggregate.warnings } == { "ordinary_chat/dual_agent initial_provider_ttft_ms", "ordinary_chat/dual_agent visible_ttft_ms", "ordinary_chat/dual_agent turn_wall_time_ms", "ordinary_chat/dual_agent model_elapsed_ms", "ordinary_chat/dual_agent tool_latency_ms", } assert all("0 samples" in warning for warning in aggregate.warnings) def test_aggregates_are_stably_sorted_and_warn_for_low_confidence_p95(): reporting = _reporting_module() runs = [ _run( case_id=BenchmarkCaseId.WEB_SEARCH_TWO_ANSWERS, mode=BenchmarkMode.DUAL_AGENT, ), _run( case_id=BenchmarkCaseId.ORDINARY_CHAT, mode=BenchmarkMode.CHAT_AGENT_TOOLS, ), _run( case_id=BenchmarkCaseId.ORDINARY_CHAT, mode=BenchmarkMode.DUAL_AGENT, ), ] report = reporting.build_benchmark_report( _config( cases=[ BenchmarkCaseId.WEB_SEARCH_TWO_ANSWERS, BenchmarkCaseId.ORDINARY_CHAT, ], modes=[ BenchmarkMode.CHAT_AGENT_TOOLS, BenchmarkMode.DUAL_AGENT, ], ), runs, mock=True, ) assert [(item.case_id.value, item.mode.value) for item in report.aggregates] == [ ("ordinary_chat", "chat_agent_tools"), ("ordinary_chat", "dual_agent"), ("web_search_two_answers", "dual_agent"), ] assert all( any("p95 low confidence" in warning for warning in item.warnings) for item in report.aggregates ) assert report.overall.warnings == [ warning for aggregate in report.aggregates for warning in aggregate.warnings ] def test_errors_are_redacted_before_truncation_in_json_and_markdown(): reporting = _reporting_module() real_key = "sk-live-secret-value" leaked_error = ( f"Authorization: Bearer {real_key}; " f"API-Key={real_key}; " f"https://example.test/?api_key={real_key}&token={real_key}; " f'{{"access_token":"{real_key}","key":"{real_key}"}}; ' + "x" * 3000 ) report = reporting.build_benchmark_report( _config(), [ _run( status="failed", semantic_failures=[f"token={real_key}"], error=leaked_error, ) ], mock=False, secrets=[real_key], ) payload = report.model_dump(mode="json") json_text = reporting.report_to_json(report) markdown = reporting.render_markdown(report) assert len(report.runs[0].error) == 2048 assert real_key not in json.dumps(payload) assert real_key not in json_text assert real_key not in markdown assert "[REDACTED]" in json_text assert "[REDACTED]" in markdown def test_error_truncation_is_utf8_safe_and_bounded_after_redaction(): reporting = _reporting_module() real_key = "sk-byte-bound-secret" error = f"Bearer {real_key} " + "密" * 1000 sanitized = reporting.sanitize_error(error, secrets=[real_key]) redacted = reporting.redact_text(error, secrets=[real_key]) expected = redacted.encode("utf-8")[:2048].decode("utf-8", errors="ignore") assert sanitized == expected assert real_key not in sanitized assert len(sanitized.encode("utf-8")) <= 2048 assert "\ufffd" not in sanitized def test_configured_secret_is_redacted_from_all_report_string_surfaces(): reporting = _reporting_module() real_key = "sk-everywhere-secret" run = _run().model_copy( update={ "event_names": [real_key], "event_sources": [f"Bearer {real_key}"], "tool_statuses": [f"token={real_key}"], } ) report = reporting.build_benchmark_report( _config( base_url=f"https://provider.example/{real_key}", model=f"model-{real_key}", ), [run], mock=False, secrets=[real_key], ) rendered = reporting.report_to_json(report) + reporting.render_markdown(report) assert real_key not in rendered assert "[REDACTED]" in rendered def test_system_errors_are_typed_redacted_counted_and_rendered(): reporting = _reporting_module() real_key = "sk-runner-system-secret" report = reporting.build_benchmark_report( _config(), [], mock=False, secrets=[real_key], system_errors=[f"Authorization: Bearer {real_key}"], ) payload = json.loads(reporting.report_to_json(report)) markdown = reporting.render_markdown(report) assert report.overall.model_dump() == { "attempted": 1, "passed": 0, "failed": 1, "warnings": [], "system_errors": ["Authorization: [REDACTED] [REDACTED]"], } assert payload["overall"]["system_errors"] == report.overall.system_errors assert "## System Errors" in markdown assert "Authorization: [REDACTED] [REDACTED]" in markdown assert real_key not in reporting.report_to_json(report) assert real_key not in markdown def test_markdown_is_rendered_from_the_same_json_facts(): reporting = _reporting_module() report = reporting.build_benchmark_report( _config(runs_per_case=2), [ _run(iteration=1, visible_ttft_ms=10, total_tokens=5), _run(iteration=2, visible_ttft_ms=30, total_tokens=9), ], mock=True, generated_at=datetime(2026, 7, 14, tzinfo=timezone.utc), ) payload = json.loads(reporting.report_to_json(report)) markdown = reporting.render_markdown(report) aggregate = payload["aggregates"][0] assert str(aggregate["visible_ttft_ms"]["p50"]) in markdown assert str(aggregate["visible_ttft_ms"]["p95"]) in markdown assert str(aggregate["total_tokens"]["total"]) in markdown assert str(aggregate["total_tokens"]["mean"]) in markdown assert payload == report.model_dump(mode="json") def test_atomic_write_creates_matching_unique_json_and_markdown_paths(tmp_path: Path): reporting = _reporting_module() report = reporting.build_benchmark_report( _config(), [_run()], mock=True ) timestamp = "20260714T010203Z" (tmp_path / f"{timestamp}.json").write_text("existing") (tmp_path / f"{timestamp}.md").write_text("existing") json_path, markdown_path = reporting.write_benchmark_report( report, tmp_path, timestamp=timestamp, ) assert json_path.name == f"{timestamp}-1.json" assert markdown_path.name == f"{timestamp}-1.md" assert json.loads(json_path.read_text()) == report.model_dump(mode="json") assert markdown_path.read_text() == reporting.render_markdown(report) _assert_no_publication_artifacts(tmp_path) def test_json_created_immediately_before_link_wins_and_writer_retries( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ): reporting = _reporting_module() timestamp = "20260714T011223Z" external_json = tmp_path / f"{timestamp}.json" real_link = reporting.os.link link_calls = 0 def create_before_link(source: Path, destination: Path) -> None: nonlocal link_calls link_calls += 1 destination = Path(destination) if link_calls == 1: destination.write_text("external json") real_link(source, destination) monkeypatch.setattr(reporting.os, "link", create_before_link) json_path, markdown_path = reporting.write_benchmark_report( _marker_report("json-link-collision"), tmp_path, timestamp=timestamp, ) assert link_calls == 3 assert external_json.read_text() == "external json" assert not (tmp_path / f"{timestamp}.md").exists() assert json_path.name == f"{timestamp}-1.json" assert markdown_path.name == f"{timestamp}-1.md" _assert_no_publication_artifacts(tmp_path) def test_json_replaced_immediately_before_link_wins_and_writer_retries( tmp_path: Path, monkeypatch: pytest.MonkeyPatch ): reporting = _reporting_module() timestamp = "20260714T011224Z" external_json = tmp_path / f"{timestamp}.json" external_json.write_text("old external json") replacement = tmp_path / "competitor-replacement.json" real_link = reporting.os.link real_replace = reporting.os.replace link_calls = 0 def replace_before_link(source: Path, destination: Path) -> None: nonlocal link_calls link_calls += 1 destination = Path(destination) if link_calls == 1: replacement.write_text("new external json") real_replace(replacement, destination) real_link(source, destination) monkeypatch.setattr(reporting.os, "link", replace_before_link) json_path, markdown_path = reporting.write_benchmark_report( _marker_report("json-link-replacement"), tmp_path, timestamp=timestamp, ) assert link_calls == 3 assert external_json.read_text() == "new external json" assert json_path.name == f"{timestamp}-1.json" assert markdown_path.name == f"{timestamp}-1.md" _assert_no_publication_artifacts(tmp_path) @pytest.mark.parametrize("collision_kind", ["create", "replace"]) def test_markdown_link_conflict_removes_owned_json_link_and_retries( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, collision_kind: str, ): reporting = _reporting_module() timestamp = "20260714T011225Z" base_json = tmp_path / f"{timestamp}.json" external_markdown = tmp_path / f"{timestamp}.md" replacement = tmp_path / "competitor-replacement.md" real_link = reporting.os.link real_replace = reporting.os.replace link_calls = 0 if collision_kind == "replace": external_markdown.write_text("old external markdown") def conflict_on_markdown(source: Path, destination: Path) -> None: nonlocal link_calls link_calls += 1 destination = Path(destination) if destination == external_markdown: if collision_kind == "replace": replacement.write_text("external markdown") real_replace(replacement, destination) else: destination.write_text("external markdown") real_link(source, destination) monkeypatch.setattr(reporting.os, "link", conflict_on_markdown) json_path, markdown_path = reporting.write_benchmark_report( _marker_report("markdown-link-collision"), tmp_path, timestamp=timestamp, ) assert link_calls == 4 assert not base_json.exists() assert external_markdown.read_text() == "external markdown" assert json_path.name == f"{timestamp}-1.json" assert markdown_path.name == f"{timestamp}-1.md" _assert_no_publication_artifacts(tmp_path) def test_markdown_conflict_does_not_remove_replaced_json_target( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ): reporting = _reporting_module() timestamp = "20260714T011226Z" base_json = tmp_path / f"{timestamp}.json" base_markdown = tmp_path / f"{timestamp}.md" competitor_json = tmp_path / "competitor.json" real_link = reporting.os.link real_replace = reporting.os.replace def replace_json_before_markdown_conflict( source: Path, destination: Path, ) -> None: destination = Path(destination) if destination == base_markdown: competitor_json.write_text("external json") real_replace(competitor_json, base_json) base_markdown.write_text("external markdown") real_link(source, destination) monkeypatch.setattr( reporting.os, "link", replace_json_before_markdown_conflict, ) json_path, markdown_path = reporting.write_benchmark_report( _marker_report("replaced-json-before-cleanup"), tmp_path, timestamp=timestamp, ) assert base_json.read_text() == "external json" assert base_markdown.read_text() == "external markdown" assert json_path.name == f"{timestamp}-1.json" assert markdown_path.name == f"{timestamp}-1.md" _assert_no_publication_artifacts(tmp_path) def test_link_exception_cleans_only_owned_links_and_temps( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ): reporting = _reporting_module() timestamp = "20260714T011227Z" base_json = tmp_path / f"{timestamp}.json" base_markdown = tmp_path / f"{timestamp}.md" competitor_json = tmp_path / "competitor.json" real_link = reporting.os.link real_replace = reporting.os.replace def fail_markdown_link(source: Path, destination: Path) -> None: destination = Path(destination) if destination == base_markdown: competitor_json.write_text("external json") real_replace(competitor_json, base_json) base_markdown.write_text("external markdown") raise OSError("simulated markdown link failure") real_link(source, destination) monkeypatch.setattr(reporting.os, "link", fail_markdown_link) with pytest.raises( reporting.BenchmarkReportWriteError, match="failed to write benchmark reports", ): reporting.write_benchmark_report( _marker_report("link-exception"), tmp_path, timestamp=timestamp, ) assert base_json.read_text() == "external json" assert base_markdown.read_text() == "external markdown" _assert_no_publication_artifacts(tmp_path) def test_link_interrupt_cleans_owned_links_and_temps_before_propagating( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ): reporting = _reporting_module() real_link = reporting.os.link link_calls = 0 def interrupt_markdown_link(source: Path, destination: Path) -> None: nonlocal link_calls link_calls += 1 if link_calls == 2: raise KeyboardInterrupt real_link(source, destination) monkeypatch.setattr(reporting.os, "link", interrupt_markdown_link) with pytest.raises(KeyboardInterrupt): reporting.write_benchmark_report( _marker_report("link-interrupt"), tmp_path, timestamp="20260714T011228Z", ) assert list(tmp_path.iterdir()) == [] def test_temp_fsync_failure_removes_owned_temps( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ): reporting = _reporting_module() def failing_fsync(file_descriptor: int) -> None: del file_descriptor raise OSError("simulated temp fsync failure") monkeypatch.setattr(reporting.os, "fsync", failing_fsync) with pytest.raises( reporting.BenchmarkReportWriteError, match="failed to write benchmark reports", ): reporting.write_benchmark_report( _marker_report("temp-fsync-failure"), tmp_path, timestamp="20260714T012334Z", ) assert list(tmp_path.iterdir()) == [] def test_both_temps_are_complete_and_fsynced_before_first_link( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, ): reporting = _reporting_module() timestamp = "20260714T040506Z" real_fsync = reporting.os.fsync real_link = reporting.os.link fsync_calls = 0 link_calls = 0 def record_fsync(file_descriptor: int) -> None: nonlocal fsync_calls fsync_calls += 1 real_fsync(file_descriptor) def inspect_before_link(source: Path, destination: Path) -> None: nonlocal link_calls link_calls += 1 temp_paths = sorted(tmp_path.glob("*.tmp")) assert len(temp_paths) == 2 assert all(path.stat().st_size > 0 for path in temp_paths) assert fsync_calls == 2 real_link(source, destination) monkeypatch.setattr(reporting.os, "fsync", record_fsync) monkeypatch.setattr(reporting.os, "link", inspect_before_link) json_path, markdown_path = reporting.write_benchmark_report( _marker_report("complete-temps"), tmp_path, timestamp=timestamp, ) assert link_calls == 2 assert json.loads(json_path.read_text())["target"]["model"] == "complete-temps" assert "- Model: `complete-temps`" in markdown_path.read_text() _assert_no_publication_artifacts(tmp_path) @pytest.mark.parametrize("stress_round", range(3)) def test_concurrent_thread_writers_reserve_distinct_matching_pairs( tmp_path: Path, monkeypatch: pytest.MonkeyPatch, stress_round: int, ): reporting = _reporting_module() writer_count = 8 markers = [ f"thread-{stress_round}-writer-{index}" for index in range(writer_count) ] first_temp_barrier = threading.Barrier(writer_count, timeout=15) original_write_text = reporting._write_text thread_state = threading.local() def synchronized_write(path: Path, content: str) -> tuple[int, int]: identity = original_write_text(path, content) if not getattr(thread_state, "first_write_done", False): thread_state.first_write_done = True first_temp_barrier.wait() return identity monkeypatch.setattr(reporting, "_write_text", synchronized_write) with ThreadPoolExecutor(max_workers=writer_count) as executor: results = list( executor.map( lambda marker: reporting.write_benchmark_report( _marker_report(marker), tmp_path, timestamp=f"20260714T02030{stress_round}Z", ), markers, ) ) assert len({json_path.stem for json_path, _ in results}) == writer_count _assert_concurrent_report_pairs(tmp_path, markers) @pytest.mark.parametrize("stress_round", range(2)) def test_concurrent_process_writers_reserve_distinct_matching_pairs( tmp_path: Path, stress_round: int, ): context = multiprocessing.get_context("spawn") writer_count = 4 markers = [ f"process-{stress_round}-writer-{index}" for index in range(writer_count) ] first_temp_barrier = context.Barrier(writer_count, timeout=20) result_queue = context.Queue() processes = [ context.Process( target=_process_report_writer, args=( str(tmp_path), marker, f"20260714T03040{stress_round}Z", first_temp_barrier, result_queue, ), ) for marker in markers ] for process in processes: process.start() results = [result_queue.get(timeout=30) for _ in processes] for process in processes: process.join(timeout=30) assert all(process.exitcode == 0 for process in processes) assert all(result[0] == "ok" for result in results), results assert len({result[2][:-5] for result in results}) == writer_count _assert_concurrent_report_pairs(tmp_path, markers)