Skip to content

ontocast.stategraph.atomic

Reusable per-unit render/critic retry loops.

These loops are designed for map/reduce execution where each content unit is processed independently. They deep-copy the incoming unit state, retry a failed render up to MAX_VISITS times, then run bounded critic passes whose critiques are compiled and applied as patches. A successful render is never repeated.

Ontology context assembly (resolve_unit_ontology_context) runs at the start of both ontology_loop and facts_loop so each unit chooses its own ontology context according to mode/policy.

Attributes

FACTS_PHASE = LoopPhase(name='facts', render_node=WorkflowNode.TEXT_TO_FACTS, critic_node=WorkflowNode.CRITICISE_FACTS, render_stage=FailureStage.GENERATE_GRAPH_UPDATE_FOR_FACTS, critic_stage=FailureStage.FACTS_CRITIQUE, prepare=_prepare_facts, render=_render_facts_phase, criticise=_criticise_facts_phase, collect_findings=_collect_facts_findings, repair_patch=_repair_facts_patch, critic_passes=lambda atomic: atomic.facts_critic_passes, patch_policy=lambda atomic: atomic.facts_patch_policy, acceptance_policy=lambda atomic: atomic.acceptance_policy, critic_skip_reason=_facts_critic_skip_reason) module-attribute

ONTOLOGY_PHASE = LoopPhase(name='ontology', render_node=WorkflowNode.TEXT_TO_ONTOLOGY, critic_node=WorkflowNode.CRITICISE_ONTOLOGY, render_stage=FailureStage.GENERATE_GRAPH_UPDATE_FOR_ONTOLOGY, critic_stage=FailureStage.ONTOLOGY_CRITIQUE, prepare=_prepare_ontology, render=_render_ontology_phase, criticise=_criticise_ontology_phase, collect_findings=_collect_ontology_findings, repair_patch=_repair_ontology_patch, critic_passes=lambda atomic: atomic.ontology_critic_passes, patch_policy=lambda atomic: atomic.ontology_patch_policy, acceptance_policy=lambda atomic: atomic.ontology_acceptance_policy, critic_skip_reason=_ontology_critic_skip_reason) module-attribute

UnitStateT = TypeVar('UnitStateT', UnitFactsState, UnitOntologyState) module-attribute

logger = logging.getLogger(__name__) module-attribute

Classes

LoopPhase dataclass

Everything that differs between the facts and ontology unit loops.

Source code in ontocast/stategraph/atomic.py
@dataclass(frozen=True)
class LoopPhase:
    """Everything that differs between the facts and ontology unit loops."""

    name: Literal["facts", "ontology"]
    render_node: WorkflowNode
    critic_node: WorkflowNode
    render_stage: FailureStage
    critic_stage: FailureStage
    # The two agent hooks are side-effecting: they mutate the state the driver
    # holds. That is what the agents already do, and saying so here is what lets
    # the driver stay generic over the two unit types instead of widening to
    # their union at every call. `_same_state` keeps the assumption honest.
    prepare: Callable[..., None]
    render: Callable[..., Awaitable[None]]
    criticise: Callable[..., Awaitable[None]]
    collect_findings: Callable[..., list]
    #: Deterministic repairs on a patched graph, run before it is judged. Its
    #: records are kept only if the patch is.
    repair_patch: Callable[..., list[GraphRepairRecord]]
    critic_passes: Callable[[AtomicToolBox], int]
    patch_policy: Callable[[AtomicToolBox], CriticPatchPolicy]
    acceptance_policy: Callable[[AtomicToolBox], FactsAcceptancePolicy]
    #: Why a critic pass should not be spent on the unit, or ``None`` to run
    #: it. Decided per pass, on the current graph.
    critic_skip_reason: Callable[..., str | None]

Attributes

acceptance_policy instance-attribute
collect_findings instance-attribute
critic_node instance-attribute
critic_passes instance-attribute
critic_skip_reason instance-attribute
critic_stage instance-attribute
criticise instance-attribute
name instance-attribute
patch_policy instance-attribute
prepare instance-attribute
render instance-attribute
render_node instance-attribute
render_stage instance-attribute
repair_patch instance-attribute

Methods:

__init__(name, render_node, critic_node, render_stage, critic_stage, prepare, render, criticise, collect_findings, repair_patch, critic_passes, patch_policy, acceptance_policy, critic_skip_reason)

PatchOutcome dataclass

What one critic pass changed, and whether the loop should run another.

Source code in ontocast/stategraph/atomic.py
@dataclass(frozen=True)
class PatchOutcome:
    """What one critic pass changed, and whether the loop should run another."""

    applied: int = 0
    residual: int = 0
    noop: int = 0
    #: Fixes applied and undone on their own; the rest of the pass stands.
    rolled_back: int = 0
    mandatory_after: int = 0

    @property
    def converged(self) -> bool:
        """True when another pass has nothing left to work with.

        A pass that changed nothing -- every fix rolled back, or nothing to
        apply and nothing mandatory left -- counts as converged: the next pass
        would be handed the same graph and the same findings, so repeating it
        buys a second identical answer at full price.
        """
        return self.applied == 0 and (self.rolled_back > 0 or self.mandatory_after == 0)

Attributes

applied = 0 class-attribute instance-attribute
converged property

True when another pass has nothing left to work with.

A pass that changed nothing -- every fix rolled back, or nothing to apply and nothing mandatory left -- counts as converged: the next pass would be handed the same graph and the same findings, so repeating it buys a second identical answer at full price.

mandatory_after = 0 class-attribute instance-attribute
noop = 0 class-attribute instance-attribute
residual = 0 class-attribute instance-attribute
rolled_back = 0 class-attribute instance-attribute

Methods:

__init__(applied=0, residual=0, noop=0, rolled_back=0, mandatory_after=0)

Functions:

facts_loop(state, tools, document_context, max_visits_per_node=None, pre_resolved_context=None) async

Run the render/critic loop for one content unit's facts.

Source code in ontocast/stategraph/atomic.py
async def facts_loop(
    state: UnitFactsState,
    tools: ToolBox,
    document_context: UnitLoopContext,
    max_visits_per_node: int | None = None,
    pre_resolved_context: UnitOntologyContext | None = None,
) -> UnitFactsState:
    """Run the render/critic loop for one content unit's facts."""
    return await run_unit_loop(
        state,
        tools,
        document_context,
        FACTS_PHASE,
        max_visits_per_node=max_visits_per_node,
        pre_resolved_context=pre_resolved_context,
    )

ontology_loop(state, tools, document_context, max_visits_per_node=None) async

Run the render/critic loop for one content unit's ontology delta.

Source code in ontocast/stategraph/atomic.py
async def ontology_loop(
    state: UnitOntologyState,
    tools: ToolBox,
    document_context: UnitLoopContext,
    max_visits_per_node: int | None = None,
) -> UnitOntologyState:
    """Run the render/critic loop for one content unit's ontology delta."""
    return await run_unit_loop(
        state,
        tools,
        document_context,
        ONTOLOGY_PHASE,
        max_visits_per_node=max_visits_per_node,
    )

run_unit_loop(state, tools, document_context, phase, max_visits_per_node=None, pre_resolved_context=None) async

Extract a unit once, then review and patch it a bounded number of times.

The two budgets mean different things and no longer trade against each other. MAX_VISITS retries a render that failed; a render that succeeded is never repeated, because re-extracting a whole unit is the expensive answer to a local defect and reliably introduces new ones. The critic passes then improve what came back, each one re-running the deterministic checks for free before paying for the critique.

Parameters:

Name Type Description Default
state UnitStateT

Unit state to run the loop over.

required
tools ToolBox

Tool container.

required
document_context UnitLoopContext

Document-level inputs, shared read-only.

required
phase LoopPhase

Which of the two loops this is.

required
max_visits_per_node int | None

Override for the render-failure bound.

None
pre_resolved_context UnitOntologyContext | None

Ontology context resolved once by the caller. The merged document ontology depends only on document-level state, so the fan-out builds it once and hands the same object to every unit; resolving it per unit cost a full rdflib merge and two graph copies each time.

None
Source code in ontocast/stategraph/atomic.py
async def run_unit_loop(
    state: UnitStateT,
    tools: ToolBox,
    document_context: UnitLoopContext,
    phase: LoopPhase,
    max_visits_per_node: int | None = None,
    pre_resolved_context: UnitOntologyContext | None = None,
) -> UnitStateT:
    """Extract a unit once, then review and patch it a bounded number of times.

    The two budgets mean different things and no longer trade against each
    other. ``MAX_VISITS`` retries a render that *failed*; a render that
    succeeded is never repeated, because re-extracting a whole unit is the
    expensive answer to a local defect and reliably introduces new ones. The
    critic passes then improve what came back, each one re-running the
    deterministic checks for free before paying for the critique.

    Args:
        state: Unit state to run the loop over.
        tools: Tool container.
        document_context: Document-level inputs, shared read-only.
        phase: Which of the two loops this is.
        max_visits_per_node: Override for the render-failure bound.
        pre_resolved_context: Ontology context resolved once by the caller. The
            merged document ontology depends only on document-level state, so
            the fan-out builds it once and hands the *same object* to every
            unit; resolving it per unit cost a full rdflib merge and two graph
            copies each time.
    """
    atomic = tools.get_atomic_tools()
    unit_state = state.model_copy(deep=True)
    # Charge resolver LLM calls (e.g. ontology selection) to this unit's
    # tracker -- the copy that survives the loop and is merged by the caller.
    # Shallow copy: retrieval_metrics stays shared with the caller's context.
    document_context = document_context.model_copy(
        update={"budget_tracker": unit_state.budget_tracker}
    )
    # The stage the loop is currently in, so an unhandled exception is
    # attributed to where it happened rather than always naming the critique.
    stage = phase.render_stage
    # Catalog prefix map installed after context resolution (see below); the
    # finally restores whatever was here so a unit cannot leak its map into
    # the next task. ContextVar storage is per-task under the fan-out.
    previous_prefixes: dict[str, str] | None = None
    installed_prefixes = False
    try:
        if pre_resolved_context is not None:
            _apply_unit_ontology_context(unit_state, pre_resolved_context)
        elif isinstance(unit_state, UnitFactsState):
            await _apply_facts_ontology_context(unit_state, document_context, tools)
        else:
            # The ontology loop is the caller that can answer an empty context
            # by inventing vocabulary: ``render_ontology`` branches on an empty
            # seed into ``render_ontology_fresh``.
            resolved = await resolve_unit_ontology_context(
                document_context,
                tools,
                unit_state.content_unit,
                can_create_vocabulary=True,
            )
            _apply_unit_ontology_context(unit_state, resolved)
        if (
            isinstance(unit_state, UnitFactsState)
            and unit_state.conformance_selection_pending
        ):
            _select_conformance_chapter(unit_state, tools)
        phase.prepare(unit_state)

        # Install the catalog prefix map for the whole unit loop so critic and
        # completion payloads reconcile against it, not only the unit graph's
        # bindings. Render agents save/restore rather than clear to None.
        previous_prefixes = RDFGraph.get_known_prefixes()
        catalog_prefixes = build_llm_prefix_map(
            unit_state.ontology_snapshot,
            _supplemental_ontologies_for_unit(document_context, unit_state, tools),
        )
        RDFGraph.set_known_prefixes(catalog_prefixes if catalog_prefixes else None)
        installed_prefixes = True

        max_visits = _resolve_max_visits_limit(
            unit_state.max_visits_per_node, max_visits_per_node
        )
        unit_state.max_visits_per_node = max_visits

        rendered = False
        render_attempt = 0
        supplemental: list[Ontology] = []
        for render_attempt in range(1, max_visits + 1):
            stage = phase.render_stage
            supplemental = _supplemental_ontologies_for_unit(
                document_context, unit_state, tools
            )
            rendered = await _render_with_evidence(
                unit_state,
                atomic,
                phase,
                supplemental,
                render_attempt=render_attempt,
            )
            if rendered:
                break

        if not rendered:
            logger.info("Unit %s loop exhausted render retries", phase.name)
            unit_state.deterministic_findings = phase.collect_findings(
                unit_state, atomic
            )
            return unit_state

        stage = phase.critic_stage
        for pass_index in range(1, phase.critic_passes(atomic) + 1):
            findings = phase.collect_findings(unit_state, atomic)
            unit_state.deterministic_findings = findings
            mandatory_before = sum(1 for finding in findings if finding.mandatory)

            skip_reason = phase.critic_skip_reason(unit_state, atomic)
            if skip_reason is not None:
                # No call is billed, and the record says why: an empty render
                # reviewed by the critic scores perfect for nothing, and a
                # citation-metadata unit has no domain facts to review.
                logger.info(
                    "Unit %s critic pass %s skipped (%s)",
                    phase.name,
                    pass_index,
                    skip_reason,
                )
                if isinstance(unit_state, UnitFactsState):
                    unit_state.critic_outcome = "skipped"
                _record_attempt(
                    unit_state,
                    kind="critic_skipped",
                    render_attempt=render_attempt,
                    critic_attempt=pass_index,
                    n_findings=len(findings),
                    n_mandatory=mandatory_before,
                    accept_reason=skip_reason,
                )
                break

            unit_state.node_visits[phase.critic_node] += 1
            _reset_node_evidence_context(unit_state, phase.critic_node)
            await phase.criticise(unit_state, atomic)
            if unit_state.status != Status.SUCCESS and not _critic_unavailable(
                unit_state
            ):
                request = unit_state.get_external_evidence_request(phase.critic_node)
                if request.initiate_search:
                    await plan_external_evidence_for_node(
                        unit_state, atomic, phase.critic_node
                    )
                    await fetch_external_evidence_for_node(
                        unit_state, atomic, phase.critic_node
                    )
                    await phase.criticise(unit_state, atomic)
            if _critic_unavailable(unit_state):
                # The critic produced no critique, so there is nothing to
                # compile and nothing to re-evaluate acceptance from. The
                # render is kept as it was and the unit leaves FAILED at the
                # critique stage: unreviewed, not accepted. Applying the
                # previous pass's suggestions here would patch against a
                # critique of a graph that has since changed.
                logger.warning(
                    "Unit %s critic unavailable on pass %s; render kept unreviewed",
                    phase.name,
                    pass_index,
                )
                break

            outcome = _apply_critic_patch(
                unit_state,
                atomic,
                phase,
                render_attempt=render_attempt,
                pass_index=pass_index,
                mandatory_before=mandatory_before,
                findings_before=findings,
            )
            logger.info(
                "Unit %s critic pass %s/%s: %d applied, %d residual, %d no-op, "
                "%d rolled back",
                phase.name,
                pass_index,
                phase.critic_passes(atomic),
                outcome.applied,
                outcome.residual,
                outcome.noop,
                outcome.rolled_back,
            )
            if outcome.converged:
                break
        else:
            if phase.critic_passes(atomic) == 0:
                # No pass ran, so nothing collected the residual the document
                # metric sums. Without this the denominator silently counts
                # only units that happened to be criticised.
                unit_state.deterministic_findings = phase.collect_findings(
                    unit_state, atomic
                )

        if (
            phase.name == "facts"
            and atomic.facts_completion_passes > 0
            and isinstance(unit_state, UnitFactsState)
        ):
            await _run_completion_passes(
                unit_state, atomic, phase, render_attempt=render_attempt
            )

        mandatory = sum(
            1 for finding in unit_state.deterministic_findings if finding.mandatory
        )
        if mandatory:
            logger.warning(
                "%d mandatory deterministic finding(s) remain unresolved", mandatory
            )
        return unit_state
    except (OntologyContextConfigError, LLMConfigurationError):
        # Not a unit failure. Both describe the deployment, not the unit, and
        # are identical for every other unit in the document: an unresolvable
        # ontology context (the whole point of ONTOLOGY_CONTEXT_REQUIRED is
        # that the run stops rather than emitting an ungrounded graph), or a
        # request the provider refuses outright. Recording either per unit
        # instead let the fan-out finish, write a zero-triple manifest and exit
        # successfully -- the vacuous pass these errors exist to prevent, now
        # with a traceback per unit to bury the cause.
        logger.error(
            "Deployment fault while running %s units; stopping the run",
            phase.name,
        )
        raise
    except Exception as exc:
        logger.exception("Unhandled exception in %s unit loop", phase.name)
        unit_state.set_failure(stage, str(exc))
        return unit_state
    finally:
        if installed_prefixes:
            RDFGraph.set_known_prefixes(previous_prefixes)