"""Integration test for telemetry degradation and atomic flush.""" import json from pathlib import Path from src.runtime.cli.telemetry_flush import flush_telemetry_queue from src.runtime.core.config import load_runtime_config from src.runtime.observability.langfuse_tracer import LangfuseRuntimeTracer from src.runtime.storage.sqlite_store import SQLiteStore def test_telemetry_degradation_and_flush(tmp_path: Path): db_file = tmp_path / "telemetry.db" store = SQLiteStore(db_file) cfg_data = json.loads(Path("runtime_config.local.json").read_text(encoding="utf-8")) cfg_data["paths"]["sqlite_db"] = str(db_file) cfg_file = tmp_path / "cfg.json" cfg_file.write_text(json.dumps(cfg_data), encoding="utf-8") cfg = load_runtime_config(cfg_file) tracer = LangfuseRuntimeTracer(cfg, store) # Insert 3 degraded events for i in range(3): tracer.record_trace( trace_id=f"tr_00{i}", fingerprint=f"fp_{i}" + "0" * 60, source_url="https://example.com", status="completed_text", spans_data={}, generations=[], metrics={}, ) assert len(store.get_unflushed_telemetry()) == 3 # Run flush CLI exit_code = flush_telemetry_queue(cfg_file, batch_size=10) assert exit_code == 0 assert len(store.get_unflushed_telemetry()) == 0