1
0
Fork 0
Code Issues Pull requests Projects Releases 2 Packages Wiki Activity Actions Pages

Add structured compiler diagnostics

This commit is contained in:
Andraxion 2026-07-29 05:07:16 -04:00
parent 0fe968c475
commit 24bd13f9d9
20 changed files with 1386 additions and 58 deletions

View file

@ -38,6 +38,7 @@ from .models import (
ProposalWriter,
RenderConfig,
)
from .telemetry import increment, stage
@dataclass(frozen=True)
@ -185,6 +186,21 @@ MAX_IMPLEMENTATION_FILES = 4_096
MAX_IMPLEMENTATION_BYTES = 64_000_000
def _load_adapter_projection(loader: AdapterLoader) -> AdapterProjection:
increment("adapter_projection_loads")
with stage("adapter.projection"):
return loader.load_projection()
def _extract_adapter_source(
loader: IncrementalAdapterLoader,
source: AdapterSource,
) -> AdapterSourceProjection:
increment("adapter_source_extractions")
with stage("adapter.extract"):
return loader.extract_source(source)
@dataclass(frozen=True)
class AdapterImplementation:
"""One confined implementation boundary that must remain stable for a process."""
@ -256,7 +272,7 @@ class AdapterProject:
allowed_relations = manifest.allowed_relations
estimated_nodes = manifest.estimated_nodes
else:
initial = loader.load_projection()
initial = _load_adapter_projection(loader)
validate_projection(initial)
root = initial.root
project_id = initial.project_id
@ -370,13 +386,14 @@ class AdapterProject:
self._implementation_snapshot = self._capture_implementation(initial=True)
def load(self) -> ProjectSnapshot:
increment("project_loads")
self.validate_runtime()
canonical_sources = self.canonical_source_paths()
captured = {path: path.read_bytes() for path in canonical_sources}
projection = (
self._load_incremental()
if self._incremental_loader is not None
else self.loader.load_projection()
else _load_adapter_projection(self.loader)
)
validate_projection(projection)
identity = (
@ -411,6 +428,11 @@ class AdapterProject:
def incremental_state(self) -> ProjectState | None:
"""Return current source identity without reconstructing the complete projection."""
increment("source_generation_checks")
with stage("source.generation"):
return self._incremental_state()
def _incremental_state(self) -> ProjectState | None:
self.validate_runtime()
loader = self._incremental_loader
if loader is None:
@ -509,7 +531,7 @@ class AdapterProject:
"incremental_disabled", "Adapter does not implement incremental extraction"
)
incremental = self._load_incremental()
full = self.loader.load_projection()
full = _load_adapter_projection(self.loader)
validate_projection(full)
fields = {
"project_id": incremental.project_id == full.project_id,
@ -578,7 +600,7 @@ class AdapterProject:
hits: list[str] = []
for source in manifest.sources:
if source.source_id in invalidated:
contribution = loader.extract_source(source)
contribution = _extract_adapter_source(loader, source)
reparsed.append(source.source_id)
cache_record = CachedSource(
source_id=source.source_id,

View file

@ -15,12 +15,18 @@ from .index import ProjectIndex
from .onboarding import assess_project, scaffold_project
from .project import Project, project_root_fingerprint
from .rendering import RenderService
from .telemetry import request
from .viewer_manager import ViewerManagerClient
def _parser() -> argparse.ArgumentParser:
parser = argparse.ArgumentParser(prog="docforge")
parser.add_argument("--project-root", type=Path, required=True)
parser.add_argument(
"--diagnostics",
action="store_true",
help="Attach bounded request-local stage timings and counters",
)
commands = parser.add_subparsers(dest="command", required=True)
onboard = commands.add_parser("onboard")
onboard.add_argument("--language", action="append", default=[])
@ -214,12 +220,20 @@ def _run(arguments: argparse.Namespace) -> dict[str, object]:
def main(argv: list[str] | None = None) -> int:
parser = _parser()
arguments = parser.parse_args(argv)
try:
result = _run(arguments)
code = 0
except DocForgeError as error:
result = {"status": "error", "error": error.as_dict()}
code = 2
with request(
f"cli.{arguments.command}",
enabled=arguments.diagnostics,
) as collector:
try:
result = _run(arguments)
code = 0
except DocForgeError as error:
result = {"status": "error", "error": error.as_dict()}
code = 2
if collector is not None:
result["diagnostics"] = collector.as_dict(
outcome="ok" if code == 0 else "error",
)
print(json.dumps(result, sort_keys=True, indent=2))
return code

View file

@ -33,6 +33,7 @@ from .models import (
ProjectState,
)
from .project import project_root_fingerprint
from .telemetry import increment, stage
INDEX_SCHEMA_VERSION = 3
APPLICATION_ID = 1_146_683_778
@ -173,6 +174,11 @@ class ProjectIndex:
def synchronize(self) -> dict[str, object]:
"""Return a current index, rebuilding disposable state when necessary."""
increment("index_synchronizations")
with stage("index.synchronize"):
return self._synchronize()
def _synchronize(self) -> dict[str, object]:
started = time.perf_counter()
try:
checked = self.check(verify_rows=False)
@ -223,6 +229,11 @@ class ProjectIndex:
return {**checked, "synchronization": synchronization}
def _build_locked(self) -> dict[str, object]:
increment("index_builds")
with stage("index.build"):
return self._build_locked_core()
def _build_locked_core(self) -> dict[str, object]:
snapshot = self.project.load()
logic = self._logic_projections()
status = _status(snapshot, logic)
@ -458,6 +469,11 @@ class ProjectIndex:
def _read_snapshot(self) -> Generator[_IndexReadSnapshot, None, None]:
"""Pin one verified index and source generation for a complete read request."""
with stage("index.read"), self._read_snapshot_core() as snapshot:
yield snapshot
@contextmanager
def _read_snapshot_core(self) -> Generator[_IndexReadSnapshot, None, None]:
checked = self.check(verify_rows=False)
signature = self._verified_index_signature
if signature is None or self._index_signature() != signature:
@ -530,6 +546,11 @@ class ProjectIndex:
return snapshot.result(**payload)
def check(self, *, verify_rows: bool = True) -> dict[str, object]:
increment("index_checks")
with stage("index.check"):
return self._check(verify_rows=verify_rows)
def _check(self, *, verify_rows: bool = True) -> dict[str, object]:
if isinstance(self.project, IncrementalStateProject):
state = self.project.incremental_state()
if state is not None:

View file

@ -16,9 +16,10 @@ from .changesets import ChangesetStore
from .context import compile_context
from .errors import DocForgeError
from .index import ProjectIndex
from .models import ProjectService, RuntimeValidatedProject
from .models import IncrementalStateProject, ProjectService, RuntimeValidatedProject
from .project import Project, project_root_fingerprint
from .rendering import RenderService
from .telemetry import request, stage
from .viewer_manager import ViewerManagerClient
SERVER_VERSION = "1.3.0.dev0"
@ -125,6 +126,7 @@ class DocForgeService:
tool_surface: tuple[str, ...] | None = None,
binding_metadata: Mapping[str, object] | None = None,
no_ast: bool = False,
diagnostics: bool = False,
) -> None:
self.project = project
self.index = ProjectIndex(self.project, allow_logic=not no_ast)
@ -140,6 +142,7 @@ class DocForgeService:
self.context_provider = context_provider
self.binding_metadata = dict(binding_metadata or {})
self.no_ast = no_ast
self.diagnostics = diagnostics
self.tool_surface = tool_surface or (
*ALL_TOOLS,
*(APPLICATION_TOOLS if self.application.enabled else ()),
@ -177,6 +180,29 @@ class DocForgeService:
synchronize: bool = True,
mutation: _MutationPolicy | None = None,
load_error_identity: bool = True,
operation_name: str = "mcp.invoke",
) -> dict[str, Any]:
with request(operation_name, enabled=self.diagnostics) as collector:
result = self._invoke_core(
operation,
synchronize=synchronize,
mutation=mutation,
load_error_identity=load_error_identity,
)
if collector is None:
return result
diagnostics = collector.as_dict(outcome="ok" if result.get("status") == "ok" else "error")
with_diagnostics = {**result, "diagnostics": diagnostics}
maximum = self.project.descriptor.limits.max_tool_output_chars
return with_diagnostics if self._encoded_length(with_diagnostics) <= maximum else result
def _invoke_core(
self,
operation: Callable[[], dict[str, object]],
*,
synchronize: bool = True,
mutation: _MutationPolicy | None = None,
load_error_identity: bool = True,
) -> dict[str, Any]:
synchronization: dict[str, object] | None = None
maximum = self.project.descriptor.limits.max_tool_output_chars
@ -192,7 +218,8 @@ class DocForgeService:
try:
try:
if isinstance(self.project, RuntimeValidatedProject):
self.project.validate_runtime()
with stage("mcp.runtime_validation"):
self.project.validate_runtime()
result: dict[str, Any] = operation()
except DocForgeError as error:
if not synchronize or error.code not in RECOVERABLE_INDEX_ERROR_CODES:
@ -214,13 +241,26 @@ class DocForgeService:
}
if load_error_identity:
try:
snapshot = self.project.load()
result.update(
{
"revision": snapshot.revision,
"source_hash": snapshot.source_hash,
}
state = (
self.project.incremental_state()
if isinstance(self.project, IncrementalStateProject)
else None
)
if state is not None:
result.update(
{
"revision": state.revision,
"source_hash": state.source_hash,
}
)
else:
snapshot = self.project.load()
result.update(
{
"revision": snapshot.revision,
"source_hash": snapshot.source_hash,
}
)
except DocForgeError:
result.update({"revision": "unknown", "source_hash": None})
else:
@ -451,7 +491,11 @@ class DocForgeService:
return None
def synchronize(self) -> dict[str, object]:
return self.invoke(self.index.synchronize, synchronize=False)
return self.invoke(
self.index.synchronize,
synchronize=False,
operation_name="mcp.sync",
)
def bootstrap(self) -> dict[str, object]:
def operation() -> dict[str, object]:
@ -506,6 +550,7 @@ class DocForgeService:
operation,
synchronize=False,
load_error_identity=False,
operation_name="mcp.bootstrap",
)
def project_info(self) -> dict[str, object]:
@ -533,7 +578,7 @@ class DocForgeService:
"index_health": index_health,
}
return self.invoke(operation)
return self.invoke(operation, operation_name="mcp.project_info")
def contract(self) -> dict[str, object]:
def operation() -> dict[str, object]:
@ -602,7 +647,7 @@ class DocForgeService:
"project_switching_allowed": False,
}
return self.invoke(operation)
return self.invoke(operation, operation_name="mcp.contract")
def get_logic(self, owner_node_id: str) -> dict[str, object]:
"""Return one Logic projection unless the binding preserves a no-AST adapter."""
@ -618,8 +663,15 @@ class DocForgeService:
),
)
return self.invoke(forbidden, synchronize=False)
return self.invoke(lambda: self.index.get_logic(owner_node_id))
return self.invoke(
forbidden,
synchronize=False,
operation_name="mcp.get_logic",
)
return self.invoke(
lambda: self.index.get_logic(owner_node_id),
operation_name="mcp.get_logic",
)
def validate_project(self) -> dict[str, object]:
def operation() -> dict[str, object]:
@ -635,7 +687,7 @@ class DocForgeService:
"edge_count": len(snapshot.edges),
}
return self.invoke(operation)
return self.invoke(operation, operation_name="mcp.validate_project")
def render_status(
self,
@ -652,10 +704,14 @@ class DocForgeService:
operation,
synchronize=False,
load_error_identity=False,
operation_name="mcp.render_status",
)
def context(self, profile: str, budget: int | None = None) -> dict[str, Any]:
return self.invoke(lambda: self.context_provider(self.index, profile, budget))
return self.invoke(
lambda: self.context_provider(self.index, profile, budget),
operation_name="mcp.context",
)
def visualize(
self,
@ -680,13 +736,21 @@ class DocForgeService:
"visualization": visualization,
}
return self.invoke(operation)
return self.invoke(operation, operation_name="mcp.visualize")
def stop_visualization(self) -> dict[str, object]:
return self.invoke(self.visualization.stop)
return self.invoke(
self.visualization.stop,
operation_name="mcp.stop_visualization",
)
def visualization_status(self) -> dict[str, object]:
return self.invoke(self.visualization.status)
return self.invoke(
self.visualization.status,
synchronize=False,
load_error_identity=False,
operation_name="mcp.visualization_status",
)
def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMCP:
@ -752,7 +816,10 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
def get_node(node_id: str) -> dict[str, Any]:
"""Return one exact stable node from the current validated project index."""
return service.invoke(lambda: service.index.get_node(node_id))
return service.invoke(
lambda: service.index.get_node(node_id),
operation_name="mcp.get_node",
)
@server.tool(name="docforge_get_logic")
def get_logic(owner_node_id: str) -> dict[str, Any]:
@ -764,7 +831,10 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
def search(query: str, limit: int | None = None) -> dict[str, Any]:
"""Run bounded lexical search over the current validated project index."""
return service.invoke(lambda: service.index.search(query, limit=limit))
return service.invoke(
lambda: service.index.search(query, limit=limit),
operation_name="mcp.search",
)
@server.tool(name="docforge_filter_nodes")
def filter_nodes(
@ -783,7 +853,8 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
status=status,
tag=tag,
limit=limit,
)
),
operation_name="mcp.filter",
)
@server.tool(name="docforge_backlinks")
@ -795,7 +866,8 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
"""Return bounded incoming relationships for one exact stable node."""
return service.invoke(
lambda: service.index.backlinks(node_id, relation=relation, limit=limit)
lambda: service.index.backlinks(node_id, relation=relation, limit=limit),
operation_name="mcp.backlinks",
)
@server.tool(name="docforge_dependencies")
@ -806,7 +878,10 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
) -> dict[str, Any]:
"""Traverse declared depends_on relationships within the configured depth limit."""
return service.invoke(lambda: service.index.dependencies(node_id, depth=depth, limit=limit))
return service.invoke(
lambda: service.index.dependencies(node_id, depth=depth, limit=limit),
operation_name="mcp.dependencies",
)
@server.tool(name="docforge_impact")
def impact(
@ -816,7 +891,10 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
) -> dict[str, Any]:
"""Traverse bounded incoming relationships and report exact paths."""
return service.invoke(lambda: service.index.impact(node_id, depth=depth, limit=limit))
return service.invoke(
lambda: service.index.impact(node_id, depth=depth, limit=limit),
operation_name="mcp.impact",
)
@server.tool(name="docforge_get_context")
def get_context(profile: str, budget: int | None = None) -> dict[str, Any]:
@ -890,6 +968,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
return service.invoke(
lambda: service.changesets.create(changeset_id),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.create",
"changeset",
@ -908,6 +987,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
return service.invoke(
lambda: service.changesets.register(changeset_id, operations),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.register",
"changeset",
@ -927,14 +1007,18 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
lambda: service.changesets.list_changesets(
include_history=include_history,
status=status,
)
),
operation_name="mcp.changeset",
)
@server.tool(name="docforge_get_changeset")
def get_changeset(changeset_id: str) -> dict[str, Any]:
"""Inspect a stored proposal even when its canonical base has become stale."""
return service.invoke(lambda: service.changesets.inspect(changeset_id))
return service.invoke(
lambda: service.changesets.inspect(changeset_id),
operation_name="mcp.changeset",
)
@server.tool(name="docforge_rebase_changeset")
def rebase_changeset(
@ -949,6 +1033,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
expected_changeset_hash,
),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.rebase",
"changeset",
@ -972,6 +1057,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
reason,
),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.abandon",
"changeset",
@ -1005,6 +1091,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
rationale=rationale,
),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.append_create",
"changeset",
@ -1038,6 +1125,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
rationale=rationale,
),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.append_update",
"changeset",
@ -1067,6 +1155,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
rationale=rationale,
),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.append_move",
"changeset",
@ -1096,6 +1185,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
rationale=rationale,
),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.append_relationship_update",
"changeset",
@ -1125,6 +1215,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
rationale=rationale,
),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.append_delete",
"changeset",
@ -1137,13 +1228,19 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
def validate_changeset(changeset_id: str) -> dict[str, Any]:
"""Validate a proposal against its exact canonical base and other active proposals."""
return service.invoke(lambda: service.changesets.validate(changeset_id))
return service.invoke(
lambda: service.changesets.validate(changeset_id),
operation_name="mcp.changeset",
)
@server.tool(name="docforge_get_changeset_diff")
def get_changeset_diff(changeset_id: str) -> dict[str, Any]:
"""Return a deterministic structured and textual diff without applying the proposal."""
return service.invoke(lambda: service.changesets.diff(changeset_id))
return service.invoke(
lambda: service.changesets.diff(changeset_id),
operation_name="mcp.changeset",
)
@server.tool(name="docforge_preview_changeset")
def preview_changeset(changeset_id: str, view_id: str) -> dict[str, Any]:
@ -1152,6 +1249,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
return service.invoke(
lambda: service.rendering.preview(changeset_id, view_id),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"render.preview",
"preview",
@ -1188,6 +1286,7 @@ def _create_bound_server(service: DocForgeService, *, read_only: bool) -> FastMC
return service.invoke(
lambda: service.application.apply(changeset_id, expected_changeset_hash),
synchronize=False,
operation_name="mcp.mutation",
mutation=service.mutation(
"changeset.apply",
"application",
@ -1206,6 +1305,7 @@ def create_server(
*,
canonical_applier_id: str | None = None,
no_ast: bool = False,
diagnostics: bool = False,
) -> FastMCP:
project = Project.open(project_root)
return create_project_server(
@ -1220,6 +1320,7 @@ def create_server(
"adapter_mode": "generic",
},
no_ast=no_ast,
diagnostics=diagnostics,
)
@ -1232,6 +1333,7 @@ def create_project_server(
context_provider: ContextProvider = compile_context,
binding_metadata: Mapping[str, object] | None = None,
no_ast: bool = False,
diagnostics: bool = False,
) -> FastMCP:
"""Create the full fixed MCP surface for one explicitly configured project service."""
@ -1243,6 +1345,7 @@ def create_project_server(
context_provider=context_provider,
binding_metadata=binding_metadata,
no_ast=no_ast,
diagnostics=diagnostics,
)
return _create_bound_server(service, read_only=False)
@ -1253,6 +1356,7 @@ def create_read_only_server(
context_provider: ContextProvider = compile_context,
binding_metadata: Mapping[str, object] | None = None,
no_ast: bool = False,
diagnostics: bool = False,
) -> FastMCP:
"""Create an adapter-capable MCP server exposing only the fixed read tool surface."""
@ -1262,6 +1366,7 @@ def create_read_only_server(
tool_surface=READ_TOOLS,
binding_metadata=binding_metadata,
no_ast=no_ast,
diagnostics=diagnostics,
)
return _create_bound_server(service, read_only=True)
@ -1279,12 +1384,18 @@ def main() -> None:
"and function-Logic extraction changes"
),
)
parser.add_argument(
"--diagnostics",
action="store_true",
help="Attach bounded request-local stage timings and counters",
)
arguments = parser.parse_args()
create_server(
arguments.project_root,
arguments.proposal_writer,
canonical_applier_id=arguments.canonical_applier,
no_ast=arguments.no_ast,
diagnostics=arguments.diagnostics,
).run(transport="stdio")

View file

@ -35,6 +35,7 @@ from .models import (
ProposalWriter,
)
from .render_config import load_render_config
from .telemetry import increment, stage
SOURCE_GENERATION_SCHEMA_VERSION = 1
GENERIC_SOURCE_CONTRACT = "docforge-core:0.7.1:index:1"
@ -709,6 +710,7 @@ class Project:
return cls(_load_descriptor(root))
def load(self) -> ProjectSnapshot:
increment("project_loads")
descriptor_bytes = self.descriptor.descriptor_path.read_bytes()
if hashlib.sha256(descriptor_bytes).hexdigest() != self.descriptor.descriptor_hash:
raise DocForgeError(
@ -730,7 +732,15 @@ class Project:
nodes: list[Node] = []
edges: list[Edge] = []
for path in ordered_sources:
source_nodes, source_edges = _load_source_file(self.descriptor, path, captured[path])
raw = captured[path]
increment("source_files_parsed")
increment("source_bytes_parsed", len(raw))
with stage("source.parse"):
source_nodes, source_edges = _load_source_file(
self.descriptor,
path,
raw,
)
nodes.extend(source_nodes)
edges.extend(source_edges)
if len(nodes) > self.descriptor.limits.max_nodes:
@ -804,6 +814,11 @@ class Project:
def incremental_state(self) -> ProjectState | None:
"""Return current source identity without reading or parsing canonical source bytes."""
increment("source_generation_checks")
with stage("source.generation"):
return self._incremental_state()
def _incremental_state(self) -> ProjectState | None:
path = self.generation_path
if not path.is_file() or path.is_symlink():
return None

View file

@ -27,6 +27,7 @@ from .models import (
)
from .project import project_root_fingerprint
from .render_contract import PreparedRender, relative_output, renderer_for
from .telemetry import increment, stage
RENDER_RECEIPT_SCHEMA_VERSION = 1
MAX_RENDER_RECEIPT_BYTES = 64_000
@ -42,6 +43,10 @@ class RenderService:
def status(self, view_id: str | None = None) -> dict[str, object]:
"""Report publication state from bounded receipts without rendering canonical content."""
with stage("render.status"):
return self._status(view_id)
def _status(self, view_id: str | None = None) -> dict[str, object]:
descriptor = self.project.descriptor
config = descriptor.render
current_state = self._current_state()
@ -113,7 +118,9 @@ class RenderService:
state = "oversized"
else:
raw = output.read_bytes()
actual_hash = hashlib.sha256(raw).hexdigest()
increment("render_output_bytes_hashed", len(raw))
with stage("render.output_hash"):
actual_hash = hashlib.sha256(raw).hexdigest()
state = "current" if actual_hash == prepared.output_hash else "stale"
result = self._view_result(
snapshot,
@ -725,13 +732,16 @@ class RenderService:
*,
changeset_hash: str | None,
) -> tuple[PreparedRender, bytes]:
increment("render_prepare_calls")
template = self._template_bytes(snapshot, view)
prepared = renderer_for(view).prepare(
snapshot,
view,
template,
changeset_hash=changeset_hash,
)
with stage("render.prepare"):
prepared = renderer_for(view).prepare(
snapshot,
view,
template,
changeset_hash=changeset_hash,
)
increment("render_output_bytes_built", len(prepared.output))
if len(prepared.output) > snapshot.descriptor.limits.max_render_bytes:
raise DocForgeError("render_too_large", "Rendered output exceeds the configured limit")
return prepared, template

220
src/docforge/telemetry.py Normal file
View file

@ -0,0 +1,220 @@
"""Bounded request-local diagnostics for repository gates and explicit profiling."""
from __future__ import annotations
import time
from collections.abc import Generator
from contextlib import contextmanager
from contextvars import ContextVar
from dataclasses import dataclass, field
from typing import Literal
CounterName = Literal[
"project_loads",
"source_files_parsed",
"source_bytes_parsed",
"adapter_projection_loads",
"adapter_source_extractions",
"source_generation_checks",
"index_checks",
"index_synchronizations",
"index_builds",
"render_prepare_calls",
"render_output_bytes_built",
"render_output_bytes_hashed",
"viewer_manager_requests",
]
StageName = Literal[
"source.generation",
"source.parse",
"adapter.projection",
"adapter.extract",
"index.check",
"index.synchronize",
"index.build",
"index.read",
"render.status",
"render.prepare",
"render.output_hash",
"visualization.status",
"viewer.manager",
"mcp.runtime_validation",
]
COUNTER_NAMES: tuple[CounterName, ...] = (
"project_loads",
"source_files_parsed",
"source_bytes_parsed",
"adapter_projection_loads",
"adapter_source_extractions",
"source_generation_checks",
"index_checks",
"index_synchronizations",
"index_builds",
"render_prepare_calls",
"render_output_bytes_built",
"render_output_bytes_hashed",
"viewer_manager_requests",
)
STAGE_NAMES: frozenset[StageName] = frozenset(
{
"source.generation",
"source.parse",
"adapter.projection",
"adapter.extract",
"index.check",
"index.synchronize",
"index.build",
"index.read",
"render.status",
"render.prepare",
"render.output_hash",
"visualization.status",
"viewer.manager",
"mcp.runtime_validation",
}
)
OPERATION_NAMES = frozenset(
{
"test",
"benchmark.m1",
"mcp.invoke",
"mcp.bootstrap",
"mcp.sync",
"mcp.project_info",
"mcp.contract",
"mcp.get_node",
"mcp.get_logic",
"mcp.search",
"mcp.filter",
"mcp.backlinks",
"mcp.dependencies",
"mcp.impact",
"mcp.context",
"mcp.validate_project",
"mcp.render_status",
"mcp.visualize",
"mcp.visualization_status",
"mcp.stop_visualization",
"mcp.changeset",
"mcp.mutation",
"cli.onboard",
"cli.info",
"cli.validate",
"cli.build",
"cli.reindex",
"cli.sync",
"cli.check",
"cli.validate-index",
"cli.show",
"cli.search",
"cli.filter",
"cli.backlinks",
"cli.dependencies",
"cli.impact",
"cli.context",
"cli.render",
"cli.render-status",
"cli.preview",
"cli.apply",
"cli.visualize",
"cli.visualization-status",
"cli.visualization-stop",
}
)
@dataclass
class _StageAggregate:
calls: int = 0
elapsed_ns: int = 0
@dataclass
class Collector:
"""One bounded aggregate owned by the current request context."""
operation: str
counters: dict[CounterName, int] = field(
default_factory=lambda: {name: 0 for name in COUNTER_NAMES}
)
stages: dict[StageName, _StageAggregate] = field(default_factory=lambda: {})
elapsed_ns: int = 0
def as_dict(self, *, outcome: str) -> dict[str, object]:
if outcome not in {"ok", "error"}:
raise ValueError("Telemetry outcome must be ok or error")
return {
"schema_version": 1,
"operation": self.operation,
"outcome": outcome,
"elapsed_ns": self.elapsed_ns,
"stages": {
name: {
"calls": aggregate.calls,
"elapsed_ns": aggregate.elapsed_ns,
}
for name, aggregate in sorted(self.stages.items())
},
"counters": {name: self.counters[name] for name in COUNTER_NAMES},
}
_CURRENT: ContextVar[Collector | None] = ContextVar(
"docforge_telemetry",
default=None,
)
@contextmanager
def request(
operation: str,
*,
enabled: bool,
) -> Generator[Collector | None, None, None]:
"""Collect one explicit request without affecting the disabled path."""
if operation not in OPERATION_NAMES:
raise ValueError("Unknown telemetry operation")
if not enabled:
yield None
return
collector = Collector(operation=operation)
token = _CURRENT.set(collector)
started = time.perf_counter_ns()
try:
yield collector
finally:
collector.elapsed_ns = time.perf_counter_ns() - started
_CURRENT.reset(token)
def increment(counter: CounterName | str, amount: int = 1) -> None:
"""Increment one fixed counter when a request collector is active."""
if counter not in COUNTER_NAMES:
raise ValueError("Unknown telemetry counter")
if type(amount) is not int or amount < 0:
raise ValueError("Telemetry increments must be nonnegative integers")
collector = _CURRENT.get()
if collector is not None:
collector.counters[counter] += amount
@contextmanager
def stage(name: StageName | str) -> Generator[None, None, None]:
"""Aggregate one fixed stage while avoiding a clock read when disabled."""
if name not in STAGE_NAMES:
raise ValueError("Unknown telemetry stage")
collector = _CURRENT.get()
if collector is None:
yield
return
started = time.perf_counter_ns()
try:
yield
finally:
aggregate = collector.stages.setdefault(name, _StageAggregate())
aggregate.calls += 1
aggregate.elapsed_ns += time.perf_counter_ns() - started

View file

@ -26,6 +26,7 @@ from typing import BinaryIO, cast
from .errors import DocForgeError
from .index import ProjectIndex
from .project import project_root_fingerprint
from .telemetry import increment, stage
from .visualization import VISUALIZATION_TEMPLATE, VisualizationIndexSnapshot
MANAGER_PROTOCOL = "docforge-viewer-manager@1"
@ -635,7 +636,8 @@ class ViewerManagerClient:
return self._lifecycle_request("stop")
def status(self) -> dict[str, object]:
return self._lifecycle_request("status")
with stage("visualization.status"):
return self._lifecycle_request("status")
def _lifecycle_request(self, action: str) -> dict[str, object]:
descriptor = self.index.project.descriptor
@ -654,6 +656,11 @@ class ViewerManagerClient:
}
def _request(self, request: dict[str, object]) -> dict[str, object]:
increment("viewer_manager_requests")
with stage("viewer.manager"):
return self._request_core(request)
def _request_core(self, request: dict[str, object]) -> dict[str, object]:
state = self._read_state()
host = state["host"]
port = state["port"]