1
0
Fork 0
Code Issues Pull requests Projects Releases 2 Packages Wiki Activity Actions Pages
DocForge2/tests/test_observability.py

385 lines
16 KiB
Python
Raw Normal View History

2026-07-29 05:07:16 -04:00
from __future__ import annotations
import asyncio
import io
import json
import shutil
import tempfile
import threading
import unittest
from concurrent.futures import ThreadPoolExecutor
from dataclasses import replace
from pathlib import Path
from unittest import mock
import jsonschema
from docforge.cli import main as cli_main
from docforge.context import compile_context
from docforge.errors import DocForgeError
from docforge.index import ProjectIndex
from docforge.mcp_server import DocForgeService
from docforge.project import Project
from docforge.rendering import RenderService
from docforge.telemetry import (
COUNTER_NAMES,
OPERATION_NAMES,
STAGE_NAMES,
increment,
request,
stage,
)
from docforge.viewer_manager import ViewerManagerClient
ROOT = Path(__file__).resolve().parents[1]
FIXTURES = ROOT / "tests" / "fixtures"
RESULT_SCHEMA = json.loads((ROOT / "schemas" / "result.schema.json").read_text())
CLIENT_CONFIGURATION_SCHEMA = json.loads(
(ROOT / "schemas" / "client-configuration.schema.json").read_text()
)
DOCTOR_RESULT_SCHEMA = json.loads((ROOT / "schemas" / "doctor-result.schema.json").read_text())
2026-07-29 05:07:16 -04:00
ZERO_WORK_COUNTERS = (
"project_loads",
"source_files_parsed",
"source_bytes_parsed",
"adapter_projection_loads",
"adapter_source_extractions",
"index_synchronizations",
"index_builds",
"render_prepare_calls",
"render_output_bytes_built",
"render_output_bytes_hashed",
)
class TelemetryContractTests(unittest.TestCase):
def copy_fixture(self, destination: Path) -> Path:
root = destination / "alpha"
shutil.copytree(FIXTURES / "alpha", root)
return root
def test_disabled_collection_reads_no_clock_and_emits_nothing(self) -> None:
with (
mock.patch(
"docforge.telemetry.time.perf_counter_ns",
side_effect=AssertionError("disabled telemetry read the clock"),
),
request("test", enabled=False) as collector,
):
increment("project_loads")
with stage("source.parse"):
pass
self.assertIsNone(collector)
def test_fixed_names_aggregation_and_exception_cleanup(self) -> None:
with (
self.assertRaisesRegex(ValueError, "operation"),
request(
"unknown",
enabled=True,
),
):
pass
with self.assertRaisesRegex(ValueError, "counter"):
increment("unknown")
with self.assertRaisesRegex(ValueError, "stage"), stage("unknown"):
pass
with request("test", enabled=True) as collector:
increment("source_files_parsed", 2)
with stage("source.parse"):
pass
with stage("source.parse"):
pass
assert collector is not None
diagnostics = collector.as_dict(outcome="ok")
self.assertEqual(2, diagnostics["counters"]["source_files_parsed"])
self.assertEqual(2, diagnostics["stages"]["source.parse"]["calls"])
with (
self.assertRaisesRegex(RuntimeError, "failed"),
request(
"test",
enabled=True,
),
stage("index.read"),
):
raise RuntimeError("failed")
with request("test", enabled=True) as next_collector:
increment("project_loads")
assert next_collector is not None
self.assertEqual(1, next_collector.as_dict(outcome="ok")["counters"]["project_loads"])
with request("test", enabled=True) as outer_collector:
increment("project_loads")
with request("test", enabled=True) as inner_collector:
increment("project_loads", 5)
increment("project_loads")
assert outer_collector is not None
assert inner_collector is not None
self.assertEqual(2, outer_collector.as_dict(outcome="ok")["counters"]["project_loads"])
self.assertEqual(5, inner_collector.as_dict(outcome="ok")["counters"]["project_loads"])
def test_schema_fixed_names_match_the_implementation(self) -> None:
properties = RESULT_SCHEMA["$defs"]["diagnostics"]["properties"]
self.assertEqual(
set(OPERATION_NAMES),
set(properties["operation"]["enum"]),
)
self.assertEqual(
set(STAGE_NAMES),
set(properties["stages"]["propertyNames"]["enum"]),
)
self.assertEqual(
set(COUNTER_NAMES),
set(properties["counters"]["required"]),
)
self.assertEqual(
set(COUNTER_NAMES),
set(properties["counters"]["properties"]),
)
for schema in (CLIENT_CONFIGURATION_SCHEMA, DOCTOR_RESULT_SCHEMA):
dedicated = schema["$defs"]["diagnostics"]["properties"]
self.assertEqual(
set(STAGE_NAMES),
set(dedicated["stages"]["propertyNames"]["enum"]),
)
self.assertEqual(
set(COUNTER_NAMES),
set(dedicated["counters"]["required"]),
)
2026-07-29 05:07:16 -04:00
def test_thread_and_async_request_contexts_are_isolated(self) -> None:
barrier = threading.Barrier(2)
def thread_worker(amount: int) -> int:
with request("test", enabled=True) as collector:
increment("project_loads", amount)
barrier.wait()
barrier.wait()
assert collector is not None
return int(collector.as_dict(outcome="ok")["counters"]["project_loads"])
with ThreadPoolExecutor(max_workers=2) as executor:
futures = [executor.submit(thread_worker, amount) for amount in (1, 3)]
self.assertEqual([1, 3], [future.result() for future in futures])
async def verify_async_isolation() -> list[int]:
ready = [asyncio.Event(), asyncio.Event()]
proceed = asyncio.Event()
async def async_worker(position: int, amount: int) -> int:
with request("test", enabled=True) as collector:
increment("project_loads", amount)
ready[position].set()
await proceed.wait()
assert collector is not None
return int(collector.as_dict(outcome="ok")["counters"]["project_loads"])
tasks = [
asyncio.create_task(async_worker(position, amount))
for position, amount in enumerate((2, 5))
]
await asyncio.gather(*(event.wait() for event in ready))
proceed.set()
return list(await asyncio.gather(*tasks))
self.assertEqual([2, 5], asyncio.run(verify_async_isolation()))
def test_generic_positive_controls_and_warm_zero_work_invariants(self) -> None:
with tempfile.TemporaryDirectory() as directory:
root = self.copy_fixture(Path(directory))
project = Project.open(root)
with request("test", enabled=True) as load_collector:
project.load()
assert load_collector is not None
load_counters = load_collector.as_dict(outcome="ok")["counters"]
self.assertEqual(1, load_counters["project_loads"])
self.assertGreater(load_counters["source_files_parsed"], 0)
self.assertGreater(load_counters["source_bytes_parsed"], 0)
index = ProjectIndex(project)
with request("test", enabled=True) as build_collector:
index.build()
assert build_collector is not None
build_counters = build_collector.as_dict(outcome="ok")["counters"]
self.assertEqual(1, build_counters["index_builds"])
self.assertGreater(build_counters["project_loads"], 0)
renderer = RenderService(project)
with request("test", enabled=True) as render_collector:
renderer.render("manual")
assert render_collector is not None
render_counters = render_collector.as_dict(outcome="ok")["counters"]
self.assertGreater(render_counters["render_prepare_calls"], 0)
self.assertGreater(render_counters["render_output_bytes_built"], 0)
with request("test", enabled=True) as deep_collector:
renderer.deep_status("manual")
assert deep_collector is not None
deep_counters = deep_collector.as_dict(outcome="ok")["counters"]
self.assertGreater(deep_counters["render_output_bytes_hashed"], 0)
service = DocForgeService(project, diagnostics=True)
warm_operations = {
"sync": service.synchronize,
"node": lambda: service.invoke(
lambda: service.index.get_node("guide.workflow"),
operation_name="mcp.get_node",
),
"search": lambda: service.invoke(
lambda: service.index.search("workflow", limit=5),
operation_name="mcp.search",
),
"filter": lambda: service.invoke(
lambda: service.index.filter_nodes(family="guide", limit=5),
operation_name="mcp.filter",
),
"backlinks": lambda: service.invoke(
lambda: service.index.backlinks("guide.foundation", limit=5),
operation_name="mcp.backlinks",
),
"dependencies": lambda: service.invoke(
lambda: service.index.dependencies("guide.workflow", depth=2, limit=5),
operation_name="mcp.dependencies",
),
"impact": lambda: service.invoke(
lambda: service.index.impact("guide.foundation", depth=2, limit=5),
operation_name="mcp.impact",
),
"context": lambda: service.invoke(
lambda: compile_context(service.index, "active"),
operation_name="mcp.context",
),
}
warm_results: dict[str, dict[str, object]] = {}
for name, operation in warm_operations.items():
with self.subTest(operation=name):
result = operation()
warm_results[name] = result
diagnostics = result["diagnostics"]
for counter in ZERO_WORK_COUNTERS:
expected = (
1 if name == "sync" and counter == "index_synchronizations" else 0
)
self.assertEqual(
expected,
diagnostics["counters"][counter],
counter,
)
self.assertGreater(diagnostics["counters"]["source_generation_checks"], 0)
node_result = warm_results["node"]
node_diagnostics = node_result["diagnostics"]
self.assertGreater(node_diagnostics["counters"]["index_checks"], 0)
missing_result = service.invoke(
lambda: service.index.get_node("missing.node"),
operation_name="mcp.get_node",
)
self.assertEqual("error", missing_result["status"])
for counter in ZERO_WORK_COUNTERS:
self.assertEqual(
0,
missing_result["diagnostics"]["counters"][counter],
counter,
)
render_result = service.render_status("manual")
render_diagnostics = render_result["diagnostics"]
for counter in ZERO_WORK_COUNTERS:
self.assertEqual(0, render_diagnostics["counters"][counter], counter)
self.assertEqual(1, render_diagnostics["stages"]["render.status"]["calls"])
jsonschema.validate(node_result, RESULT_SCHEMA)
jsonschema.validate(render_result, RESULT_SCHEMA)
def test_structured_errors_include_diagnostics_and_schema_validation(self) -> None:
with tempfile.TemporaryDirectory() as directory:
project = Project.open(self.copy_fixture(Path(directory)))
service = DocForgeService(project, diagnostics=True)
def fail() -> dict[str, object]:
raise DocForgeError("intentional", "Intentional telemetry error")
result = service.invoke(
fail,
synchronize=False,
load_error_identity=False,
operation_name="mcp.invoke",
)
self.assertEqual("error", result["status"])
self.assertEqual("error", result["diagnostics"]["outcome"])
jsonschema.validate(result, RESULT_SCHEMA)
def test_visualization_status_error_has_one_request_and_no_hidden_work(self) -> None:
with tempfile.TemporaryDirectory() as directory:
root = self.copy_fixture(Path(directory))
project = Project.open(root)
service = DocForgeService(project, diagnostics=True)
service.visualization = ViewerManagerClient(
service.index,
state_path=root / ".docforge" / "missing-viewer-manager.json",
)
result = service.visualization_status()
self.assertEqual("error", result["status"])
counters = result["diagnostics"]["counters"]
for counter in ZERO_WORK_COUNTERS:
self.assertEqual(0, counters[counter], counter)
self.assertEqual(1, counters["viewer_manager_requests"])
self.assertEqual(
1,
result["diagnostics"]["stages"]["visualization.status"]["calls"],
)
self.assertEqual(
1,
result["diagnostics"]["stages"]["viewer.manager"]["calls"],
)
def test_diagnostics_are_dropped_before_the_primary_result(self) -> None:
with tempfile.TemporaryDirectory() as directory:
original = Project.open(self.copy_fixture(Path(directory)))
limits = replace(original.descriptor.limits, max_tool_output_chars=600)
project = Project(replace(original.descriptor, limits=limits))
service = DocForgeService(project, diagnostics=True)
result = service.invoke(
lambda: {
"status": "ok",
"project_id": project.descriptor.project_id,
"revision": "unversioned",
"source_hash": "0" * 64,
},
synchronize=False,
operation_name="mcp.invoke",
)
self.assertEqual("ok", result["status"])
self.assertNotIn("diagnostics", result)
def test_cli_diagnostics_flag_is_additive_and_defaults_off(self) -> None:
base_arguments = ["--project-root", "/unused", "info"]
with mock.patch("docforge.cli._run", return_value={"status": "ok"}):
with mock.patch("sys.stdout", new_callable=io.StringIO) as output:
self.assertEqual(0, cli_main(base_arguments))
self.assertNotIn("diagnostics", json.loads(output.getvalue()))
with mock.patch("sys.stdout", new_callable=io.StringIO) as output:
self.assertEqual(
0,
cli_main(
[
"--project-root",
"/unused",
"--diagnostics",
"info",
]
),
)
result = json.loads(output.getvalue())
self.assertEqual("cli.info", result["diagnostics"]["operation"])
if __name__ == "__main__":
unittest.main()