def make_render_facts_node(tools: ToolBox):
async def render_facts(state: AgentState) -> AgentState:
if not state.content_units:
state.facts_units = []
if state.status != Status.FAILED:
state.status = Status.SUCCESS
return state
worker_limit = max(1, tools.config.server.parallel_workers)
semaphore = asyncio.Semaphore(worker_limit)
# Built once for the whole document, not once per unit. It reads only
# reduced_ontology_artifacts, which the ontology stage froze upstream,
# so every unit would otherwise pay an identical full rdflib merge plus
# two graph copies -- synchronously, on the event loop, stalling every
# other unit's in-flight provider call. None means no ontology stage ran
# (facts-only mode); units then resolve their own context as before.
merged_context = build_merged_document_ontology_context(
UnitLoopContext.from_agent_state(state)
)
if (
merged_context is None
and tools.config.server.ontology_context_scope
== OntologyContextScope.DOCUMENT
):
# No ontology stage ran, so there is nothing to merge -- but the
# deployment has asked for one chapter per document rather than one
# per unit. Resolve every unit's context up front and union it, so
# the fan-out shares a prompt prefix instead of N distinct ones.
merged_context = await build_unioned_document_ontology_context(
UnitLoopContext.from_agent_state(state), tools, state.content_units
)
if merged_context is not None:
# Hand the same graph to merge/validate downstream instead of
# letting each rebuild it.
state.facts_ontology_context = merged_context.snapshot.graph
async def process_unit(
unit_index: int,
) -> tuple[int, UnitFactsState, str, list[str], OntologyAssemblyMode]:
async with _unit_slot(semaphore) as worker_wait:
unit_budget = BudgetTracker()
# No-op when the ontology fan-out already summarised this unit;
# does the work when facts run without an ontology stage.
await ensure_unit_summary(state, unit_index, tools, unit_budget)
unit_context = UnitLoopContext.from_agent_state(state, unit_budget)
conformance_chapter, contract_terms, selection_pending = (
tools.shapes_prompt_contract()
)
facts_state = UnitFactsState(
content_unit=state.content_units[unit_index],
ontology_snapshot=_empty_unit_snapshot(),
ontology_patch_sources=[],
facts_user_instruction=state.facts_user_instruction,
conformance_chapter=conformance_chapter,
shapes_contract_terms=contract_terms,
conformance_selection_pending=selection_pending,
budget_tracker=unit_budget,
max_visits_per_node=state.max_visits,
llm_graph_format=state.llm_graph_format,
llm_output_layout=tools.config.server.llm_output_layout,
ontology_context_max_triples=tools.config.server.ontology_context_max_triples,
ontology_chapter_format=tools.config.server.ontology_chapter_format,
ontology_text_caps=tools.config.server.ontology_text_caps,
)
loop_start = time.perf_counter()
result = await facts_loop(
facts_state,
tools,
unit_context,
pre_resolved_context=merged_context,
)
result.budget_tracker.add_duration(
f"{WorkflowNode.RENDER_FACTS}{UNIT_SUM_SUFFIX}",
time.perf_counter() - loop_start,
)
result.budget_tracker.add_duration(
f"{WorkflowNode.RENDER_FACTS}/worker_wait", worker_wait
)
# Per-unit resolver metrics previously landed on a discarded
# deep copy; fold them back (last writer wins on shared keys).
state.retrieval_metrics.update(unit_context.retrieval_metrics)
return (
unit_index,
result,
result.assembly_anchor_iri,
list(result.ontology_patch_sources),
result.assembly_mode_used,
)
# A provider's prefix cache is populated by a request that has already
# completed, so a fan-out that issues every call at once has all of them
# miss a prefix they all share. Running the first unit alone turns the
# rest into hits -- but only where the chapter is shared, which is why
# this is a separate knob rather than an implication of the scope.
warmup = min(
tools.config.server.fanout_warmup_units,
len(state.content_units),
)
raw_results: list = []
unit_errors = 0
if warmup:
warm_results, warm_errors = await _gather_units(
WorkflowNode.RENDER_FACTS,
state,
[process_unit(i) for i in range(warmup)],
)
raw_results.extend(warm_results)
unit_errors += warm_errors
rest = [process_unit(i) for i in range(warmup, len(state.content_units))]
if rest:
more_results, more_errors = await _gather_units(
WorkflowNode.RENDER_FACTS, state, rest, index_offset=warmup
)
raw_results.extend(more_results)
unit_errors += more_errors
ordered_results = sorted(raw_results, key=lambda item: item[0])
facts_units: list[ContentUnit] = []
failed_without_output_count = unit_errors
salvaged_failed_count = 0
# Accumulated over *every* unit, whether or not it ran a repair render,
# so the residual has "units" as its denominator.
findings_residual = 0
mandatory_residual = 0
critic_fixes_applied = 0
critic_fixes_residual = 0
critic_fixes_noop = 0
critic_fixes_rolled_back = 0
critic_fixes_junk_refused = 0
critic_fixes_unresolved_prefix = 0
critic_units_unreviewed = 0
critic_units_skipped = 0
unit_contexts: dict[int, tuple[str, list[str], OntologyAssemblyMode]] = {}
for (
unit_index,
result,
anchor_iri,
patch_sources,
assembly_mode,
) in ordered_results:
state.budget_tracker.merge_from(result.budget_tracker)
findings_residual += len(result.deterministic_findings)
mandatory_residual += sum(
1 for finding in result.deterministic_findings if finding.mandatory
)
critic_fixes_applied += result.critic_fixes_applied
critic_fixes_residual += result.critic_fixes_residual
critic_fixes_noop += result.critic_fixes_noop
critic_fixes_rolled_back += result.critic_fixes_rolled_back
critic_fixes_junk_refused += result.critic_fixes_junk_refused
critic_fixes_unresolved_prefix += result.critic_fixes_unresolved_prefix
# A unit whose critic call failed used to leave the loop SUCCESS,
# indistinguishable from one the critic accepted; a unit the loop
# never sent to the critic looked the same. Both are counted.
critic_units_unreviewed += result.critic_outcome == "unavailable"
critic_units_skipped += result.critic_outcome == "skipped"
if result.attempt_log:
state.facts_loop_telemetry[unit_index] = list(result.attempt_log)
if result.applied_repairs:
state.facts_repairs_applied[unit_index] = list(result.applied_repairs)
unit_contexts[unit_index] = (anchor_iri, patch_sources, assembly_mode)
has_output = len(result.content_unit.graph) > 0
if not has_output:
failed_without_output_count += 1
state.unit_failures.append(
UnitFailure(
unit_index=unit_index,
phase="facts",
stage=(
result.failure_stage.value
if result.failure_stage is not None
else None
),
reason=result.failure_reason,
)
)
continue
facts_units.append(result.content_unit)
if result.status != Status.SUCCESS:
salvaged_failed_count += 1
if failed_without_output_count:
logger.warning(
"Parallel facts map failed without usable output for "
f"{failed_without_output_count}/{len(state.content_units)} unit(s)"
)
if salvaged_failed_count:
logger.warning(
"Parallel facts map salvaged output from non-converged loop(s): "
f"{salvaged_failed_count}/{len(state.content_units)} unit(s)"
)
if merged_context is None and (
tools.config.get_tool_config().facts_validation.context_from_units
):
# No ontology stage ran, so build_merged_document_ontology_context
# had nothing to merge. Without this both MERGE_FACTS and
# VALIDATE_FACTS run against an empty ontology graph although every
# unit rendered against a real one.
unit_ontology_context = _union_unit_ontology_context(ordered_results)
if len(unit_ontology_context) > 0:
state.facts_ontology_context = unit_ontology_context
state.retrieval_metrics[RetrievalMetric.ONTOLOGY_SNAPSHOT_TRIPLES] = (
len(unit_ontology_context)
)
_, state.unit_patch_sources, _, anchor_counts = aggregate_writable_metrics(
unit_contexts
)
state.retrieval_metrics[RetrievalMetric.FACTS_ANCHOR_COUNT] = len(anchor_counts)
state.retrieval_metrics[RetrievalMetric.FACTS_ANCHOR_UNITS] = sum(
anchor_counts.values()
)
all_attempts = [
attempt
for attempts in state.facts_loop_telemetry.values()
for attempt in attempts
]
# Patch passes undone for leaving the unit worse -- deleting without
# writing, shrinking without resolving, or manufacturing new mandatory
# findings. Non-zero means the critique is provoking data-destroying
# responses, the failure mode that has cost runs a large share of their
# value nodes while logging nothing a run could be judged by.
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_PATCHES_ROLLED_BACK] = sum(
1 for attempt in all_attempts if attempt.patch_rolled_back
)
# Residual is read off each unit's final findings, not off the attempt
# log. Summing `attempts[-1]` where `kind == "llm_repair"` silently
# contributed 0 for every unit that never ran a repair render -- the
# clean ones and the ones whose loop exhausted its retries -- so the
# metric's denominator was "units that needed repair", not "units", and
# a change that made *fewer* units enter repair read as a drop in
# residual findings. It also summed total findings, so advisory
# NUMERIC_COVERAGE (which fires on nearly every unit of numeric prose)
# dominated the number that was supposed to track mandatory defects.
state.retrieval_metrics[RetrievalMetric.FACTS_FINDINGS_RESIDUAL] = (
findings_residual
)
state.retrieval_metrics[RetrievalMetric.FACTS_MANDATORY_RESIDUAL] = (
mandatory_residual
)
# The critic's own ledger. `node_visits` counting CRITICISE_FACTS lived
# on the per-unit state copy and died with it, so the number of critic
# calls a run bought was not recoverable from its own artifacts.
critic_attempts = [a for a in all_attempts if a.kind == "critic"]
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_CALLS] = len(
critic_attempts
)
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_ACCEPTED] = sum(
1 for attempt in critic_attempts if attempt.success
)
# What the critique actually bought, split by how it was paid for. The
# loop used to discard every fix on an accepted render, so the honest
# value of these two together was once "far less than the call count
# suggests" -- and nothing recorded it.
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_FIXES_APPLIED] = (
critic_fixes_applied
)
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_FIXES_RESIDUAL] = (
critic_fixes_residual
)
# Fixes that removed exactly what they re-added. A critique made mostly
# of these is a critic producing motion rather than corrections, which
# nothing could see before: they used to be counted as fixes that
# landed.
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_FIXES_NOOP] = (
critic_fixes_noop
)
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_FIXES_ROLLED_BACK] = (
critic_fixes_rolled_back
)
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_FIXES_JUNK_REFUSED] = (
critic_fixes_junk_refused
)
state.retrieval_metrics[
RetrievalMetric.FACTS_CRITIC_FIXES_UNRESOLVED_PREFIX
] = critic_fixes_unresolved_prefix
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_UNITS_UNREVIEWED] = (
critic_units_unreviewed
)
state.retrieval_metrics[RetrievalMetric.FACTS_CRITIC_UNITS_SKIPPED] = (
critic_units_skipped
)
# The insert-only completion pass, read off the attempt log like the
# critic's own ledger: calls billed, inserts that stayed, and missed
# measurements the inventory stopped listing.
completion_attempts = [a for a in all_attempts if a.kind == "completion"]
state.retrieval_metrics[RetrievalMetric.FACTS_COMPLETION_CALLS] = len(
completion_attempts
)
state.retrieval_metrics[RetrievalMetric.FACTS_COMPLETION_TRIPLES_INSERTED] = (
sum(a.n_triples_inserted for a in completion_attempts)
)
state.retrieval_metrics[
RetrievalMetric.FACTS_COMPLETION_MEASUREMENTS_RECOVERED
] = sum(a.n_measurements_recovered for a in completion_attempts)
state.facts_units = facts_units
state.status = _map_stage_status(
failed_without_output_count, len(state.content_units)
)
return state
return render_facts