Files

42 lines
1.3 KiB
Python

"""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