Skip to content

ontocast.onto.state

Attributes

PEAK_SUFFIX = '_max' module-attribute

UNIT_SUM_SUFFIX = '/unit_sum' module-attribute

Classes

AgentState

Bases: BasePydanticModel

State for the ontology-based knowledge graph agent.

This class maintains the state of the agent during document processing, including input text, content units, ontologies, and workflow status.

Attributes:

Name Type Description
docling_doc Any

Parsed document in native Docling format.

current_domain str

IRI used for forming document namespace.

doc_hid str

An almost unique hash/id for the parent document.

raw_input dict[str, bytes]

Single raw input payload as {filename: bytes}.

failure_stage FailureStage | None

Stage where failure occurred.

failure_reason str | None

Reason for failure.

status Status

Current workflow status.

max_visits int

Maximum render attempts per unit loop.

max_chunks int | None

Maximum number of source content units to split and process.

Source code in ontocast/onto/state.py
 446
 447
 448
 449
 450
 451
 452
 453
 454
 455
 456
 457
 458
 459
 460
 461
 462
 463
 464
 465
 466
 467
 468
 469
 470
 471
 472
 473
 474
 475
 476
 477
 478
 479
 480
 481
 482
 483
 484
 485
 486
 487
 488
 489
 490
 491
 492
 493
 494
 495
 496
 497
 498
 499
 500
 501
 502
 503
 504
 505
 506
 507
 508
 509
 510
 511
 512
 513
 514
 515
 516
 517
 518
 519
 520
 521
 522
 523
 524
 525
 526
 527
 528
 529
 530
 531
 532
 533
 534
 535
 536
 537
 538
 539
 540
 541
 542
 543
 544
 545
 546
 547
 548
 549
 550
 551
 552
 553
 554
 555
 556
 557
 558
 559
 560
 561
 562
 563
 564
 565
 566
 567
 568
 569
 570
 571
 572
 573
 574
 575
 576
 577
 578
 579
 580
 581
 582
 583
 584
 585
 586
 587
 588
 589
 590
 591
 592
 593
 594
 595
 596
 597
 598
 599
 600
 601
 602
 603
 604
 605
 606
 607
 608
 609
 610
 611
 612
 613
 614
 615
 616
 617
 618
 619
 620
 621
 622
 623
 624
 625
 626
 627
 628
 629
 630
 631
 632
 633
 634
 635
 636
 637
 638
 639
 640
 641
 642
 643
 644
 645
 646
 647
 648
 649
 650
 651
 652
 653
 654
 655
 656
 657
 658
 659
 660
 661
 662
 663
 664
 665
 666
 667
 668
 669
 670
 671
 672
 673
 674
 675
 676
 677
 678
 679
 680
 681
 682
 683
 684
 685
 686
 687
 688
 689
 690
 691
 692
 693
 694
 695
 696
 697
 698
 699
 700
 701
 702
 703
 704
 705
 706
 707
 708
 709
 710
 711
 712
 713
 714
 715
 716
 717
 718
 719
 720
 721
 722
 723
 724
 725
 726
 727
 728
 729
 730
 731
 732
 733
 734
 735
 736
 737
 738
 739
 740
 741
 742
 743
 744
 745
 746
 747
 748
 749
 750
 751
 752
 753
 754
 755
 756
 757
 758
 759
 760
 761
 762
 763
 764
 765
 766
 767
 768
 769
 770
 771
 772
 773
 774
 775
 776
 777
 778
 779
 780
 781
 782
 783
 784
 785
 786
 787
 788
 789
 790
 791
 792
 793
 794
 795
 796
 797
 798
 799
 800
 801
 802
 803
 804
 805
 806
 807
 808
 809
 810
 811
 812
 813
 814
 815
 816
 817
 818
 819
 820
 821
 822
 823
 824
 825
 826
 827
 828
 829
 830
 831
 832
 833
 834
 835
 836
 837
 838
 839
 840
 841
 842
 843
 844
 845
 846
 847
 848
 849
 850
 851
 852
 853
 854
 855
 856
 857
 858
 859
 860
 861
 862
 863
 864
 865
 866
 867
 868
 869
 870
 871
 872
 873
 874
 875
 876
 877
 878
 879
 880
 881
 882
 883
 884
 885
 886
 887
 888
 889
 890
 891
 892
 893
 894
 895
 896
 897
 898
 899
 900
 901
 902
 903
 904
 905
 906
 907
 908
 909
 910
 911
 912
 913
 914
 915
 916
 917
 918
 919
 920
 921
 922
 923
 924
 925
 926
 927
 928
 929
 930
 931
 932
 933
 934
 935
 936
 937
 938
 939
 940
 941
 942
 943
 944
 945
 946
 947
 948
 949
 950
 951
 952
 953
 954
 955
 956
 957
 958
 959
 960
 961
 962
 963
 964
 965
 966
 967
 968
 969
 970
 971
 972
 973
 974
 975
 976
 977
 978
 979
 980
 981
 982
 983
 984
 985
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
class AgentState(BasePydanticModel):
    """State for the ontology-based knowledge graph agent.

    This class maintains the state of the agent during document processing,
    including input text, content units, ontologies, and workflow status.

    Attributes:
        docling_doc: Parsed document in native Docling format.
        current_domain: IRI used for forming document namespace.
        doc_hid: An almost unique hash/id for the parent document.
        raw_input: Single raw input payload as {filename: bytes}.
        failure_stage: Stage where failure occurred.
        failure_reason: Reason for failure.
        status: Current workflow status.
        max_visits: Maximum render attempts per unit loop.
        max_chunks: Maximum number of source content units to split and process.
    """

    # Typed `Any` rather than `DoclingDocument | None` so that importing
    # AgentState does not require the `documents` extra; the validator below
    # still enforces the real type whenever docling-core is installed.
    docling_doc: Any = Field(
        default=None,
        description="Parsed document in native Docling format.",
    )
    current_domain: str = Field(
        description=(
            "IRI used for forming the document namespace. Defaults to the "
            "CURRENT_DOMAIN environment variable, then to DEFAULT_DOMAIN; an "
            "explicit constructor argument always wins."
        ),
        default_factory=lambda: os.getenv("CURRENT_DOMAIN") or DEFAULT_DOMAIN,
    )
    doc_hid: str = Field(
        description="An almost unique hash / id for the parent document of the current unit",
        default="default_doc",
    )
    raw_input: dict[str, bytes] = Field(
        default_factory=dict,
        description="Single raw input payload: {filename: bytes}.",
    )
    content_units: list[ContentUnit] = Field(
        default_factory=list,
        description="Pending content units to process.",
    )
    bibliography_units_skipped: int = Field(
        default=0,
        description=(
            "Prepared chunks dropped by reference-list routing "
            "(CHUNK_BIBLIOGRAPHY_MODE=skip)."
        ),
    )
    undersized_units_skipped: int = Field(
        default=0,
        description="Prepared chunks dropped by the CHUNK_MIN_UNIT_CHARS floor.",
    )
    non_content_units_skipped: int = Field(
        default=0,
        description=(
            "Prepared chunks dropped by front/back-matter routing "
            "(CHUNK_NON_CONTENT_MODE=skip)."
        ),
    )
    ontology_artifacts: list[Ontology] = Field(
        default_factory=list,
        description="Final per-anchor ontology artifacts produced for this document.",
    )
    reduced_ontology_artifacts: list[Ontology] = Field(
        default_factory=list,
        description="Reduced ontology artifacts after explicit ontology reduce step.",
    )
    reduced_ontology_by_anchor: dict[str, Ontology] = Field(
        default_factory=dict,
        description="Reduced ontology artifacts indexed by anchor IRI.",
    )
    ontology_reduce_metrics: dict[str, int | float | str | list[dict[str, str]]] = (
        Field(
            default_factory=dict,
            description=(
                "Metrics emitted by ontology reduce stage. The list-valued "
                "entry is minted_duplicate_pairs: the (minted IRI, catalog "
                "IRI, surface, role) records the reconciliation scan found."
            ),
        )
    )
    unit_patch_sources: dict[int, list[str]] = Field(
        default_factory=dict,
        description="Retrieved ontology source IRIs per content unit index.",
    )
    retrieval_metrics: dict[str, int | float | str | dict[str, Any]] = Field(
        default_factory=dict,
        description="Runtime retrieval/evaluation metrics for observability.",
    )
    aggregated_facts: RDFGraph = Field(
        description="RDF triples representing aggregated facts "
        "from the current document",
        default_factory=RDFGraph,
    )
    facts_ontology_context: RDFGraph = Field(
        default_factory=RDFGraph,
        description=(
            "Merged reduced-ontology graph used as read-only schema for facts. "
            "Derived from reduced_ontology_artifacts, which are frozen once the "
            "ontology stage completes, so the facts fan-out builds it once and "
            "merge/validate reuse it rather than each repeating the merge."
        ),
    )
    ontology_user_instruction: str = Field(
        description="Specific user instructions for ontology extraction, e.g. `Focus on extracting places`",
        default="",
    )

    ontology_selection_user_instruction: str = Field(
        description=(
            "Specific user instructions for ontology selection, "
            "e.g. `Prefer ontologies focused on finance`"
        ),
        default="",
    )

    facts_user_instruction: str = Field(
        description="Specific user instructions for facts extraction, e.g. `Focus on extracting places`",
        default="",
    )

    ontology_context_fixed_ontology_id: str = Field(
        description=(
            "Catalog ontology id when ontology_context_mode is fixed_single_ontology "
            "(resolved via OntologyManager)."
        ),
        default="",
    )

    tenant: str | None = Field(
        default=None,
        description="Tenant id when request selected tenancy via query/CLI.",
    )
    project: str | None = Field(
        default=None,
        description="Project id when request selected tenancy via query/CLI.",
    )

    source_url: str | None = Field(
        description="Source URL from JSON input file (for provenance tracking)",
        default=None,
    )

    document_metadata: dict[str, Any] = Field(
        default_factory=dict,
        description=(
            "Caller-asserted document identity metadata attached to the parent "
            "doc_iri as provenance (DOI, ISBN, scheme+value business ids, title, …)."
        ),
    )

    ontology_updates_applied: list[GraphUpdate] = Field(
        default_factory=list,
        description="A list of graph update that improve the current ontology",
    )

    facts_units: list[ContentUnit] = Field(
        default_factory=list,
        description="Successful per-unit facts outputs collected during parallel map phase",
    )

    facts_loop_telemetry: dict[int, list[LoopAttempt]] = Field(
        default_factory=dict,
        description=(
            "Per-unit facts loop attempt records (render/critic/repair) keyed "
            "by content unit index; makes visit efficacy measurable."
        ),
    )

    ontology_loop_telemetry: dict[int, list[LoopAttempt]] = Field(
        default_factory=dict,
        description=(
            "Per-unit ontology loop attempt records keyed by content unit "
            "index -- the ontology critic's ledger, recorded before its gate "
            "is recalibrated so the incumbent's accept rate is measurable."
        ),
    )

    facts_repairs_applied: dict[int, list[GraphRepairRecord]] = Field(
        default_factory=dict,
        description=(
            "Deterministic rewrites applied to rendered facts graphs, keyed by "
            "content unit index — records which triples the machine altered "
            "from what the LLM asserted."
        ),
    )

    unit_failures: list[UnitFailure] = Field(
        default_factory=list,
        description=(
            "Content units that produced no usable output, with the stage and "
            "reason. Without this, a run in which every unit failed was "
            "indistinguishable from one that found nothing to extract."
        ),
    )

    aggregation_clusters: dict[str, list[str]] = Field(
        default_factory=dict,
        description=(
            "Final URI -> source entities merged into it during facts "
            "aggregation (clusters with >= 2 members only); consumed by the "
            "validation gate to turn findings into un-merge pair vetoes."
        ),
    )

    aggregation_key_clusters: list[str] = Field(
        default_factory=list,
        description=(
            "Final URIs of merge clusters supported by natural-key evidence "
            "(a shared identifier value). The gate downgrades string "
            "multi-value findings on these subjects to warnings: two label "
            "variants on a key-evidenced merge are two names for one thing, "
            "not two things."
        ),
    )

    aggregation_cross_unit_pairs: list[tuple[str, str]] = Field(
        default_factory=list,
        description=(
            "Canonical (subject, predicate) pairs whose IRI objects were "
            "contributed by more than one unit. A multi-valued predicate "
            "asserted by a single unit cannot be the product of an identity "
            "merge, so the gate needs this to tell a merge signature from "
            "legitimate multi-value modelling."
        ),
    )

    facts_validation_findings: list[FactsValidationFinding] = Field(
        default_factory=list,
        description=(
            "Post-aggregation invariant findings (functional violations, "
            "suspect multi-values, degenerate coreference, SHACL) remaining "
            "after the un-merge and SHACL-autofix repair stages."
        ),
    )

    facts_gate_repairs: list[GraphRepairRecord] = Field(
        default_factory=list,
        description=(
            "LLM-free repairs the validation gate applied to the merged graph "
            "(shape-driven retyping, code resolution, placeholder pruning). "
            "Document-level counterpart of facts_repairs_applied, which is "
            "per unit."
        ),
    )

    facts_conformance: dict = Field(
        default_factory=dict,
        description=(
            "Rolled-up validation result: whether SHACL ran and the graph "
            "conforms, plus counts by finding kind, SHACL constraint component "
            "and shape. Grouping is what makes the residue diagnosable — a "
            "flat violation list does not distinguish one systematic modelling "
            "gap from many independent defects."
        ),
    )

    ontology_units: list[ContentUnit] = Field(
        default_factory=list,
        description="Successful per-unit ontology outputs collected during parallel map phase",
    )
    ontology_provenance_artifact: RDFGraph = Field(
        default_factory=RDFGraph,
        description="Provenance/reification triples stripped from normalized ontology.",
    )

    failure_stage: FailureStage | None = None
    failure_reason: str | None = None

    improvements_suggestions: list[str] = Field(
        description="Itemized concrete and actionable instructions for improvements of extraction of facts/ontology",
        default_factory=list,
    )

    status: Status = Status.SUCCESS
    max_visits: int = Field(
        default=1,
        description=(
            "Maximum render attempts per unit loop. Mirrors "
            "``ServerConfig.max_visits_per_node``, which every entry path "
            "supplies; at 1 the critic never runs."
        ),
    )
    max_chunks: int | None = None
    target_sections: list[str] | None = Field(
        default=None,
        description="Sections to include when chunking. None = no filter.",
    )
    exclude_sections: list[str] | None = Field(
        default=None,
        description=(
            "Sections to drop when chunking. None = use the resolved section "
            "schema's default_exclude; [] = no exclusion; list = explicit "
            "denylist."
        ),
    )
    summarize_sections: list[str] | None = Field(
        default=None,
        description="Sections to summarize. None = skip summarization node.",
    )
    summary_max_sentences: int = Field(
        default=5,
        description="Max sentences per chunk summary when summarization is enabled.",
    )
    document_type_hint: str | None = Field(
        default=None,
        description=(
            "Optional free-text hint about the source material (e.g. '10-K filing', "
            "'journal article') used to resolve section label schema and LLM tagging."
        ),
    )
    section_schema_id: str | None = Field(
        default=None,
        description=(
            "Section label schema id from ontocast.config.section_labels (e.g. academic, "
            "financial). Overrides document_type_hint when set."
        ),
    )
    model_config = ConfigDict(arbitrary_types_allowed=True, populate_by_name=True)
    render_mode: RenderMode = Field(
        default=RenderMode.ONTOLOGY_AND_FACTS,
        description=("Rendering mode: ontology, facts, or ontology_and_facts."),
    )
    llm_graph_format: LLMGraphFormat = Field(
        default=LLMGraphFormat.TURTLE,
        description=(
            "Format used by the LLM for emitting RDF graph payloads: "
            "'turtle' (default; Turtle strings) or 'jsonld' (compact JSON-LD "
            "objects embedded directly in the structured response)."
        ),
    )
    ontology_context_mode: OntologyContextMode = Field(
        default=OntologyContextMode.SELECTED_SINGLE_ONTOLOGY,
        description=(
            "Per-unit ontology context: selected_single_ontology (LLM-picked catalog), "
            "selected_vector_search_ontology (vector-store ensemble; Qdrant or LanceDB), "
            "or fixed_single_ontology (catalog ontology_id via ontology_context_fixed_ontology_id)."
        ),
    )
    # Budget Tracking
    budget_tracker: BudgetTracker = Field(
        default_factory=BudgetTracker,
        description="Budget statistics tracker (LLM usage and generated triples)",
    )

    @property
    def needs_section_prepare(self) -> bool:
        """Whether the request carries explicit section-dependent options.

        Section tagging itself is default-on in chunk prepare (driven by
        ``CHUNK_SECTION_CLASSIFIER``); schema-default exclusions apply even
        when this is False.
        """
        return (
            self.target_sections is not None
            or self.summarize_sections is not None
            or self.exclude_sections is not None
        )

    @property
    def use_summarization(self) -> bool:
        """Whether per-unit summaries should be produced in the fan-out."""
        return self.summarize_sections is not None

    @property
    def render_ontology(self) -> bool:
        """Whether ontology rendering should run."""
        return self.render_mode in (
            RenderMode.ONTOLOGY,
            RenderMode.ONTOLOGY_AND_FACTS,
        )

    @property
    def render_facts(self) -> bool:
        """Whether facts rendering should run."""
        return self.render_mode in (
            RenderMode.FACTS,
            RenderMode.ONTOLOGY_AND_FACTS,
        )

    def get_content_unit_progress_info(self) -> tuple[int, int]:
        """Get current content unit number and total content units."""
        total_content_units = len(self.content_units)
        current_content_unit_number = 1 if total_content_units > 0 else 0
        return current_content_unit_number, total_content_units

    def get_content_unit_progress_string(self) -> str:
        """Get a formatted string showing content unit progress."""
        current, total = self.get_content_unit_progress_info()
        if total == 0:
            return "no content units"
        return f"content unit {current}/{total}"

    @classmethod
    def render_updated_graph(
        cls, graph: RDFGraph, updates: list[GraphUpdate], max_triples: int | None = None
    ) -> tuple[RDFGraph, bool]:
        """Create a copy of the given graph with all GraphUpdate objects applied.

        This method:
        1. Creates a copy of the input graph
        2. Generates SPARQL queries from all GraphUpdate objects
        3. Executes the queries on the copied graph
        4. Checks if the updated graph exceeds max_triples limit
        5. Returns the updated graph copy, or original if limit exceeded

        Args:
            graph: The RDFGraph to update
            updates: List of GraphUpdate objects to apply
            max_triples: Maximum number of triples allowed. If None, no limit enforced.

        Returns:
            Tuple of (RDFGraph, bool): The updated graph (or original if limit exceeded),
            and a boolean indicating if the update was applied (True) or skipped (False)
        """
        if not updates:
            return graph, True

        # Create a copy of the input graph
        # Use RDFGraph's copy method to preserve type
        updated_graph = RDFGraph()
        for triple in graph:
            updated_graph.add(triple)
        # Copy namespace bindings
        for prefix, namespace in graph.namespaces():
            updated_graph.bind(prefix, namespace)

        all_prefixes = {}
        for graph_update in updates:
            for op in graph_update.triple_operations:
                # Extract prefixes from TripleOp operations
                if isinstance(op, TripleOp) and op.prefixes:
                    all_prefixes.update(op.prefixes)

        # Bind prefixes to the copied graph
        for prefix, uri in all_prefixes.items():
            updated_graph.bind(prefix, uri)

        # Apply each GraphUpdate to the copied graph
        for graph_update in updates:
            # Generate SPARQL queries from the GraphUpdate
            queries = graph_update.generate_sparql_queries()

            # Execute each query on the copied graph
            for query in queries:
                cls._apply_update_query(updated_graph, query)

        # Reject only updates that GROW the graph past the limit. Comparing the
        # absolute post-apply size alone locked the loop out: a seed snapshot
        # already over the cap made every update for the rest of the run fail
        # this check, discarding the LLM's work with only a warning and no way
        # to shrink back under. Deletions and net-shrinking updates now apply.
        if max_triples is not None and len(updated_graph) > max_triples:
            import logging

            logger = logging.getLogger(__name__)
            if len(graph) > max_triples:
                logger.warning(
                    f"Ontology graph is already above the configured limit "
                    f"({len(graph)} > {max_triples} triples) before this update. "
                    f"Applying it anyway because it does not grow the graph "
                    f"({len(updated_graph)} triples); raise or unset "
                    f"ONTOLOGY_MAX_TRIPLES if this is expected."
                )
            if len(updated_graph) > len(graph):
                logger.warning(
                    f"Ontology update skipped: would exceed limit "
                    f"({len(updated_graph)} > {max_triples} triples). "
                    f"Original size: {len(graph)} triples."
                )
                return graph, False  # Return original, unchanged

        return updated_graph, True

    @classmethod
    def _apply_update_query(cls, graph: RDFGraph, query: str) -> None:
        """Apply one SPARQL update query, splitting compound LLM output proactively."""
        parts = cls._split_compound_sparql_query(query)
        for part in parts:
            graph.update(part)

    @staticmethod
    def _split_compound_sparql_query(query: str) -> list[str]:
        """Split a query string containing concatenated top-level UPDATE statements.

        LLMs frequently emit several ``INSERT DATA`` / ``DELETE DATA`` blocks joined
        after a shared ``PREFIX`` block.  Splitting on top-level keyword boundaries
        before calling ``graph.update`` avoids parse errors entirely.

        A single-statement query is returned as a one-element list.
        """
        stripped = query.strip()
        if not stripped:
            return [stripped]

        starts = [m.start() for m in _TOP_LEVEL_UPDATE_START_RE.finditer(stripped)]
        if len(starts) <= 1:
            return [stripped]

        prefix_block = stripped[: starts[0]].strip()
        parts: list[str] = []
        for i, start in enumerate(starts):
            end = starts[i + 1] if i + 1 < len(starts) else len(stripped)
            body = stripped[start:end].strip()
            if body:
                parts.append(f"{prefix_block}\n{body}" if prefix_block else body)
        return parts or [stripped]

    def set_docling_doc(self, doc: "DoclingDocument") -> None:
        """Set the parsed document and generate document hash.

        Args:
            doc: The DoclingDocument to set.
        """
        self.docling_doc = doc
        self.doc_hid = render_text_hash(doc.model_dump_json())

    @field_validator("docling_doc", mode="before")
    @classmethod
    def _coerce_docling_doc(cls, value: object) -> Any:
        if value is None:
            return None
        docling_document = _docling_document_cls()
        if isinstance(value, docling_document):
            return value
        if isinstance(value, dict):
            return docling_document.model_validate(value)
        raise TypeError(f"Expected DoclingDocument or dict, got {type(value).__name__}")

    def set_failure(self, stage: FailureStage, reason: str) -> None:
        """Set failure state with stage and reason.

        Args:
            stage: The stage where the failure occurred.
            reason: The reason for the failure.
        """
        self.failure_stage = stage
        self.failure_reason = reason
        self.status = Status.FAILED

    def clear_failure(self):
        """Clear failure state and set status to success."""
        self.failure_stage = None
        self.failure_reason = None
        self.status = Status.SUCCESS

    @property
    def doc_iri(self) -> URIRef:
        """Get the document IRI.

        Returns:
            str: The document IRI.
        """
        return URIRef(f"{self.current_domain}/doc/{self.doc_hid}")

    @property
    def doc_namespace(self) -> str:
        """Get the document namespace.

        Returns:
            str: The document namespace.
        """
        return normalize_namespace_iri(self.doc_iri, context="facts")

    @property
    def graph_uri(self):
        return self.doc_namespace

    @property
    def ontology_ids(self) -> list[str]:
        """Ontology ids for all current ontology artifacts."""
        artifacts = (
            self.reduced_ontology_artifacts
            if self.reduced_ontology_artifacts
            else self.ontology_artifacts
        )
        return [ontology.ontology_id for ontology in artifacts if ontology.ontology_id]

Attributes

aggregated_facts = Field(description='RDF triples representing aggregated facts from the current document', default_factory=RDFGraph) class-attribute instance-attribute
aggregation_clusters = Field(default_factory=dict, description='Final URI -> source entities merged into it during facts aggregation (clusters with >= 2 members only); consumed by the validation gate to turn findings into un-merge pair vetoes.') class-attribute instance-attribute
aggregation_cross_unit_pairs = Field(default_factory=list, description='Canonical (subject, predicate) pairs whose IRI objects were contributed by more than one unit. A multi-valued predicate asserted by a single unit cannot be the product of an identity merge, so the gate needs this to tell a merge signature from legitimate multi-value modelling.') class-attribute instance-attribute
aggregation_key_clusters = Field(default_factory=list, description='Final URIs of merge clusters supported by natural-key evidence (a shared identifier value). The gate downgrades string multi-value findings on these subjects to warnings: two label variants on a key-evidenced merge are two names for one thing, not two things.') class-attribute instance-attribute
bibliography_units_skipped = Field(default=0, description='Prepared chunks dropped by reference-list routing (CHUNK_BIBLIOGRAPHY_MODE=skip).') class-attribute instance-attribute
budget_tracker = Field(default_factory=BudgetTracker, description='Budget statistics tracker (LLM usage and generated triples)') class-attribute instance-attribute
content_units = Field(default_factory=list, description='Pending content units to process.') class-attribute instance-attribute
current_domain = Field(description='IRI used for forming the document namespace. Defaults to the CURRENT_DOMAIN environment variable, then to DEFAULT_DOMAIN; an explicit constructor argument always wins.', default_factory=lambda: os.getenv('CURRENT_DOMAIN') or DEFAULT_DOMAIN) class-attribute instance-attribute
doc_hid = Field(description='An almost unique hash / id for the parent document of the current unit', default='default_doc') class-attribute instance-attribute
doc_iri property

Get the document IRI.

Returns:

Name Type Description
str URIRef

The document IRI.

doc_namespace property

Get the document namespace.

Returns:

Name Type Description
str str

The document namespace.

docling_doc = Field(default=None, description='Parsed document in native Docling format.') class-attribute instance-attribute
document_metadata = Field(default_factory=dict, description='Caller-asserted document identity metadata attached to the parent doc_iri as provenance (DOI, ISBN, scheme+value business ids, title, …).') class-attribute instance-attribute
document_type_hint = Field(default=None, description="Optional free-text hint about the source material (e.g. '10-K filing', 'journal article') used to resolve section label schema and LLM tagging.") class-attribute instance-attribute
exclude_sections = Field(default=None, description="Sections to drop when chunking. None = use the resolved section schema's default_exclude; [] = no exclusion; list = explicit denylist.") class-attribute instance-attribute
facts_conformance = Field(default_factory=dict, description='Rolled-up validation result: whether SHACL ran and the graph conforms, plus counts by finding kind, SHACL constraint component and shape. Grouping is what makes the residue diagnosable — a flat violation list does not distinguish one systematic modelling gap from many independent defects.') class-attribute instance-attribute
facts_gate_repairs = Field(default_factory=list, description='LLM-free repairs the validation gate applied to the merged graph (shape-driven retyping, code resolution, placeholder pruning). Document-level counterpart of facts_repairs_applied, which is per unit.') class-attribute instance-attribute
facts_loop_telemetry = Field(default_factory=dict, description='Per-unit facts loop attempt records (render/critic/repair) keyed by content unit index; makes visit efficacy measurable.') class-attribute instance-attribute
facts_ontology_context = Field(default_factory=RDFGraph, description='Merged reduced-ontology graph used as read-only schema for facts. Derived from reduced_ontology_artifacts, which are frozen once the ontology stage completes, so the facts fan-out builds it once and merge/validate reuse it rather than each repeating the merge.') class-attribute instance-attribute
facts_repairs_applied = Field(default_factory=dict, description='Deterministic rewrites applied to rendered facts graphs, keyed by content unit index — records which triples the machine altered from what the LLM asserted.') class-attribute instance-attribute
facts_units = Field(default_factory=list, description='Successful per-unit facts outputs collected during parallel map phase') class-attribute instance-attribute
facts_user_instruction = Field(description='Specific user instructions for facts extraction, e.g. `Focus on extracting places`', default='') class-attribute instance-attribute
facts_validation_findings = Field(default_factory=list, description='Post-aggregation invariant findings (functional violations, suspect multi-values, degenerate coreference, SHACL) remaining after the un-merge and SHACL-autofix repair stages.') class-attribute instance-attribute
failure_reason = None class-attribute instance-attribute
failure_stage = None class-attribute instance-attribute
graph_uri property
improvements_suggestions = Field(description='Itemized concrete and actionable instructions for improvements of extraction of facts/ontology', default_factory=list) class-attribute instance-attribute
llm_graph_format = Field(default=LLMGraphFormat.TURTLE, description="Format used by the LLM for emitting RDF graph payloads: 'turtle' (default; Turtle strings) or 'jsonld' (compact JSON-LD objects embedded directly in the structured response).") class-attribute instance-attribute
max_chunks = None class-attribute instance-attribute
max_visits = Field(default=1, description='Maximum render attempts per unit loop. Mirrors ``ServerConfig.max_visits_per_node``, which every entry path supplies; at 1 the critic never runs.') class-attribute instance-attribute
model_config = ConfigDict(arbitrary_types_allowed=True, populate_by_name=True) class-attribute instance-attribute
needs_section_prepare property

Whether the request carries explicit section-dependent options.

Section tagging itself is default-on in chunk prepare (driven by CHUNK_SECTION_CLASSIFIER); schema-default exclusions apply even when this is False.

non_content_units_skipped = Field(default=0, description='Prepared chunks dropped by front/back-matter routing (CHUNK_NON_CONTENT_MODE=skip).') class-attribute instance-attribute
ontology_artifacts = Field(default_factory=list, description='Final per-anchor ontology artifacts produced for this document.') class-attribute instance-attribute
ontology_context_fixed_ontology_id = Field(description='Catalog ontology id when ontology_context_mode is fixed_single_ontology (resolved via OntologyManager).', default='') class-attribute instance-attribute
ontology_context_mode = Field(default=OntologyContextMode.SELECTED_SINGLE_ONTOLOGY, description='Per-unit ontology context: selected_single_ontology (LLM-picked catalog), selected_vector_search_ontology (vector-store ensemble; Qdrant or LanceDB), or fixed_single_ontology (catalog ontology_id via ontology_context_fixed_ontology_id).') class-attribute instance-attribute
ontology_ids property

Ontology ids for all current ontology artifacts.

ontology_loop_telemetry = Field(default_factory=dict, description="Per-unit ontology loop attempt records keyed by content unit index -- the ontology critic's ledger, recorded before its gate is recalibrated so the incumbent's accept rate is measurable.") class-attribute instance-attribute
ontology_provenance_artifact = Field(default_factory=RDFGraph, description='Provenance/reification triples stripped from normalized ontology.') class-attribute instance-attribute
ontology_reduce_metrics = Field(default_factory=dict, description='Metrics emitted by ontology reduce stage. The list-valued entry is minted_duplicate_pairs: the (minted IRI, catalog IRI, surface, role) records the reconciliation scan found.') class-attribute instance-attribute
ontology_selection_user_instruction = Field(description='Specific user instructions for ontology selection, e.g. `Prefer ontologies focused on finance`', default='') class-attribute instance-attribute
ontology_units = Field(default_factory=list, description='Successful per-unit ontology outputs collected during parallel map phase') class-attribute instance-attribute
ontology_updates_applied = Field(default_factory=list, description='A list of graph update that improve the current ontology') class-attribute instance-attribute
ontology_user_instruction = Field(description='Specific user instructions for ontology extraction, e.g. `Focus on extracting places`', default='') class-attribute instance-attribute
project = Field(default=None, description='Project id when request selected tenancy via query/CLI.') class-attribute instance-attribute
raw_input = Field(default_factory=dict, description='Single raw input payload: {filename: bytes}.') class-attribute instance-attribute
reduced_ontology_artifacts = Field(default_factory=list, description='Reduced ontology artifacts after explicit ontology reduce step.') class-attribute instance-attribute
reduced_ontology_by_anchor = Field(default_factory=dict, description='Reduced ontology artifacts indexed by anchor IRI.') class-attribute instance-attribute
render_facts property

Whether facts rendering should run.

render_mode = Field(default=RenderMode.ONTOLOGY_AND_FACTS, description='Rendering mode: ontology, facts, or ontology_and_facts.') class-attribute instance-attribute
render_ontology property

Whether ontology rendering should run.

retrieval_metrics = Field(default_factory=dict, description='Runtime retrieval/evaluation metrics for observability.') class-attribute instance-attribute
section_schema_id = Field(default=None, description='Section label schema id from ontocast.config.section_labels (e.g. academic, financial). Overrides document_type_hint when set.') class-attribute instance-attribute
source_url = Field(description='Source URL from JSON input file (for provenance tracking)', default=None) class-attribute instance-attribute
status = Status.SUCCESS class-attribute instance-attribute
summarize_sections = Field(default=None, description='Sections to summarize. None = skip summarization node.') class-attribute instance-attribute
summary_max_sentences = Field(default=5, description='Max sentences per chunk summary when summarization is enabled.') class-attribute instance-attribute
target_sections = Field(default=None, description='Sections to include when chunking. None = no filter.') class-attribute instance-attribute
tenant = Field(default=None, description='Tenant id when request selected tenancy via query/CLI.') class-attribute instance-attribute
undersized_units_skipped = Field(default=0, description='Prepared chunks dropped by the CHUNK_MIN_UNIT_CHARS floor.') class-attribute instance-attribute
unit_failures = Field(default_factory=list, description='Content units that produced no usable output, with the stage and reason. Without this, a run in which every unit failed was indistinguishable from one that found nothing to extract.') class-attribute instance-attribute
unit_patch_sources = Field(default_factory=dict, description='Retrieved ontology source IRIs per content unit index.') class-attribute instance-attribute
use_summarization property

Whether per-unit summaries should be produced in the fan-out.

Methods:

clear_failure()

Clear failure state and set status to success.

Source code in ontocast/onto/state.py
def clear_failure(self):
    """Clear failure state and set status to success."""
    self.failure_stage = None
    self.failure_reason = None
    self.status = Status.SUCCESS
get_content_unit_progress_info()

Get current content unit number and total content units.

Source code in ontocast/onto/state.py
def get_content_unit_progress_info(self) -> tuple[int, int]:
    """Get current content unit number and total content units."""
    total_content_units = len(self.content_units)
    current_content_unit_number = 1 if total_content_units > 0 else 0
    return current_content_unit_number, total_content_units
get_content_unit_progress_string()

Get a formatted string showing content unit progress.

Source code in ontocast/onto/state.py
def get_content_unit_progress_string(self) -> str:
    """Get a formatted string showing content unit progress."""
    current, total = self.get_content_unit_progress_info()
    if total == 0:
        return "no content units"
    return f"content unit {current}/{total}"
render_updated_graph(graph, updates, max_triples=None) classmethod

Create a copy of the given graph with all GraphUpdate objects applied.

This method: 1. Creates a copy of the input graph 2. Generates SPARQL queries from all GraphUpdate objects 3. Executes the queries on the copied graph 4. Checks if the updated graph exceeds max_triples limit 5. Returns the updated graph copy, or original if limit exceeded

Parameters:

Name Type Description Default
graph RDFGraph

The RDFGraph to update

required
updates list[GraphUpdate]

List of GraphUpdate objects to apply

required
max_triples int | None

Maximum number of triples allowed. If None, no limit enforced.

None

Returns:

Type Description
RDFGraph

Tuple of (RDFGraph, bool): The updated graph (or original if limit exceeded),

bool

and a boolean indicating if the update was applied (True) or skipped (False)

Source code in ontocast/onto/state.py
@classmethod
def render_updated_graph(
    cls, graph: RDFGraph, updates: list[GraphUpdate], max_triples: int | None = None
) -> tuple[RDFGraph, bool]:
    """Create a copy of the given graph with all GraphUpdate objects applied.

    This method:
    1. Creates a copy of the input graph
    2. Generates SPARQL queries from all GraphUpdate objects
    3. Executes the queries on the copied graph
    4. Checks if the updated graph exceeds max_triples limit
    5. Returns the updated graph copy, or original if limit exceeded

    Args:
        graph: The RDFGraph to update
        updates: List of GraphUpdate objects to apply
        max_triples: Maximum number of triples allowed. If None, no limit enforced.

    Returns:
        Tuple of (RDFGraph, bool): The updated graph (or original if limit exceeded),
        and a boolean indicating if the update was applied (True) or skipped (False)
    """
    if not updates:
        return graph, True

    # Create a copy of the input graph
    # Use RDFGraph's copy method to preserve type
    updated_graph = RDFGraph()
    for triple in graph:
        updated_graph.add(triple)
    # Copy namespace bindings
    for prefix, namespace in graph.namespaces():
        updated_graph.bind(prefix, namespace)

    all_prefixes = {}
    for graph_update in updates:
        for op in graph_update.triple_operations:
            # Extract prefixes from TripleOp operations
            if isinstance(op, TripleOp) and op.prefixes:
                all_prefixes.update(op.prefixes)

    # Bind prefixes to the copied graph
    for prefix, uri in all_prefixes.items():
        updated_graph.bind(prefix, uri)

    # Apply each GraphUpdate to the copied graph
    for graph_update in updates:
        # Generate SPARQL queries from the GraphUpdate
        queries = graph_update.generate_sparql_queries()

        # Execute each query on the copied graph
        for query in queries:
            cls._apply_update_query(updated_graph, query)

    # Reject only updates that GROW the graph past the limit. Comparing the
    # absolute post-apply size alone locked the loop out: a seed snapshot
    # already over the cap made every update for the rest of the run fail
    # this check, discarding the LLM's work with only a warning and no way
    # to shrink back under. Deletions and net-shrinking updates now apply.
    if max_triples is not None and len(updated_graph) > max_triples:
        import logging

        logger = logging.getLogger(__name__)
        if len(graph) > max_triples:
            logger.warning(
                f"Ontology graph is already above the configured limit "
                f"({len(graph)} > {max_triples} triples) before this update. "
                f"Applying it anyway because it does not grow the graph "
                f"({len(updated_graph)} triples); raise or unset "
                f"ONTOLOGY_MAX_TRIPLES if this is expected."
            )
        if len(updated_graph) > len(graph):
            logger.warning(
                f"Ontology update skipped: would exceed limit "
                f"({len(updated_graph)} > {max_triples} triples). "
                f"Original size: {len(graph)} triples."
            )
            return graph, False  # Return original, unchanged

    return updated_graph, True
set_docling_doc(doc)

Set the parsed document and generate document hash.

Parameters:

Name Type Description Default
doc 'DoclingDocument'

The DoclingDocument to set.

required
Source code in ontocast/onto/state.py
def set_docling_doc(self, doc: "DoclingDocument") -> None:
    """Set the parsed document and generate document hash.

    Args:
        doc: The DoclingDocument to set.
    """
    self.docling_doc = doc
    self.doc_hid = render_text_hash(doc.model_dump_json())
set_failure(stage, reason)

Set failure state with stage and reason.

Parameters:

Name Type Description Default
stage FailureStage

The stage where the failure occurred.

required
reason str

The reason for the failure.

required
Source code in ontocast/onto/state.py
def set_failure(self, stage: FailureStage, reason: str) -> None:
    """Set failure state with stage and reason.

    Args:
        stage: The stage where the failure occurred.
        reason: The reason for the failure.
    """
    self.failure_stage = stage
    self.failure_reason = reason
    self.status = Status.FAILED

BudgetTracker

Bases: BasePydanticModel

Lightweight tracker for LLM usage statistics and generated triples.

node_durations follows a key convention that distinguishes wall clock from time summed across concurrent workers -- without it, a fan-out node's entry is ambiguous and the two are silently added together:

"<node>" True wall clock for a pipeline node. Written only by the _timed wrapper in :mod:ontocast.stategraph.create. "<node>/unit_sum" Per-unit loop time summed over every parallel worker. Divided by the wall-clock entry this yields effective workers -- see :meth:parallel_efficiency. "<node>/<stage>" Any other sub-stage measurement (worker_wait, loop_lag_total, llm/provider, ...).

Source code in ontocast/onto/state.py
class BudgetTracker(BasePydanticModel):
    """Lightweight tracker for LLM usage statistics and generated triples.

    ``node_durations`` follows a key convention that distinguishes wall clock
    from time summed across concurrent workers -- without it, a fan-out node's
    entry is ambiguous and the two are silently added together:

    ``"<node>"``
        True wall clock for a pipeline node. Written only by the ``_timed``
        wrapper in :mod:`ontocast.stategraph.create`.
    ``"<node>/unit_sum"``
        Per-unit loop time summed over every parallel worker. Divided by the
        wall-clock entry this yields *effective workers* -- see
        :meth:`parallel_efficiency`.
    ``"<node>/<stage>"``
        Any other sub-stage measurement (``worker_wait``, ``loop_lag_total``,
        ``llm/provider``, ...).
    """

    chars_sent: int = Field(default=0, description="Total characters sent to LLM")
    chars_received: int = Field(
        default=0, description="Total characters received from LLM"
    )
    calls_count: int = Field(default=0, description="Total number of LLM API calls")
    cache_hits: int = Field(
        default=0,
        description="LLM calls satisfied from disk cache (no provider tokens)",
    )
    input_tokens: int = Field(
        default=0, description="Billed input tokens (when reported by provider)"
    )
    output_tokens: int = Field(
        default=0, description="Billed output tokens (when reported by provider)"
    )

    # Kept apart from the billed totals rather than folded in: a cache-replayed
    # run costs nothing, so adding these to input_tokens/output_tokens would
    # report spend that never happened. Reported together they answer the other
    # question -- what the workload costs cold.
    cached_input_tokens: int = Field(
        default=0, description="Input tokens replayed from the OntoCast disk cache"
    )
    cached_output_tokens: int = Field(
        default=0, description="Output tokens replayed from the OntoCast disk cache"
    )

    # Detail keys from LangChain's UsageMetadata, summed over billed and
    # replayed calls alike: they describe the shape of the workload, and a
    # reasoning model's thinking tokens matter whether or not this particular
    # run paid for them.
    reasoning_tokens: int = Field(
        default=0, description="Thinking tokens, counted inside the output totals"
    )
    cache_read_input_tokens: int = Field(
        default=0,
        description=(
            "Input tokens served from the provider's own prompt cache, counted "
            "inside the input totals and billed at a reduced rate"
        ),
    )
    cache_creation_input_tokens: int = Field(
        default=0, description="Input tokens written to the provider's prompt cache"
    )

    # Triple generation tracking
    ontology_triples_generated: int = Field(
        default=0, description="Total number of triples generated for ontology updates"
    )
    facts_triples_generated: int = Field(
        default=0, description="Total number of triples generated for facts"
    )
    ontology_operations_count: int = Field(
        default=0, description="Total number of ontology update operations"
    )
    facts_operations_count: int = Field(
        default=0, description="Total number of facts update operations"
    )

    node_durations: dict[str, float] = Field(
        default_factory=dict,
        description="Accumulated seconds per pipeline node/stage (see class docstring)",
    )
    counters: dict[str, int] = Field(
        default_factory=dict,
        description=(
            "Named event counts (e.g. how often a per-document computation ran). "
            "Summed on merge, like node_durations."
        ),
    )

    def add_duration(self, name: str, seconds: float) -> None:
        """Accumulate seconds for a named node or stage.

        Keys ending in :data:`PEAK_SUFFIX` take the maximum instead, since
        adding two peaks would report a stall that never happened.
        """
        if name.endswith(PEAK_SUFFIX):
            self.node_durations[name] = max(self.node_durations.get(name, 0.0), seconds)
            return
        self.node_durations[name] = self.node_durations.get(name, 0.0) + seconds

    def incr(self, name: str, n: int = 1) -> None:
        """Increment a named counter.

        Args:
            name: Counter key, e.g. ``"ctx/merge_document_ontology.calls"``.
            n: Amount to add.
        """
        self.counters[name] = self.counters.get(name, 0) + n

    def parallel_efficiency(self, node: str) -> float | None:
        """Effective worker count for a fan-out node, or ``None`` if unmeasured.

        The ratio of time summed across unit workers to the node's wall clock.
        A value near ``parallel_workers`` means the fan-out is running at full
        width; a value near ``1.0`` means the units are effectively serialised
        -- typically by synchronous CPU work blocking the event loop, which
        ``"<node>/loop_lag_total"`` quantifies.

        Args:
            node: The node key, e.g. ``str(WorkflowNode.RENDER_FACTS)``.

        Returns:
            float | None: Effective workers, or ``None`` when either the wall
            clock or the ``/unit_sum`` entry is missing or zero.
        """
        wall = self.node_durations.get(node)
        unit_sum = self.node_durations.get(f"{node}{UNIT_SUM_SUFFIX}")
        if not wall or unit_sum is None:
            return None
        return unit_sum / wall

    @computed_field  # type: ignore[prop-decorator]
    @property
    def prefix_cache_hit_rate(self) -> float | None:
        """Share of input tokens served from the *provider's* prompt cache.

        This is the provider-side prefix cache, not the OntoCast disk cache --
        the two are independent, and a deployment whose disk cache is cold (a
        fresh container, say) can still hit this one on every unit after the
        first, because the fan-out sends the same document-invariant prefix N
        times.

        Near ``0.0`` with a fan-out wider than one unit means the workers all
        issued before any prefix had been cached: a cache entry only becomes
        readable once the first response has begun. Compare a long, more
        sequential run against a wide one to see the difference.

        Returns:
            float | None: Ratio in ``[0.0, 1.0]``, or ``None`` when no input
            tokens were reported at all.
        """
        # Denominator spans billed *and* replayed input: cache_read_input_tokens
        # is accumulated for both (see _add_usage_detail), so dividing by
        # input_tokens alone would exceed 1.0 on a partially replayed run.
        total_input = self.input_tokens + self.cached_input_tokens
        if total_input <= 0:
            return None
        return self.cache_read_input_tokens / total_input

    @computed_field  # type: ignore[prop-decorator]
    @property
    def reasoning_share_of_output(self) -> float | None:
        """Share of output tokens the model spent on thinking.

        ``reasoning_tokens`` are counted *inside* the output totals, so this is
        a decomposition of what was already paid for, never an addition to it.
        A high share says the cost lever is the model's reasoning budget rather
        than anything about the prompt; a share of ``0.0`` says the model is
        not a reasoning model, and no reasoning-effort setting will change its
        bill.

        Returns:
            float | None: Ratio in ``[0.0, 1.0]``, or ``None`` when no output
            tokens were reported at all.
        """
        total_output = self.output_tokens + self.cached_output_tokens
        if total_output <= 0:
            return None
        return self.reasoning_tokens / total_output

    def _add_usage_detail(self, usage: TokenUsage) -> None:
        """Accumulate the provider-detail keys shared by billed and cached calls."""
        if usage.reasoning_tokens is not None:
            self.reasoning_tokens += usage.reasoning_tokens
        if usage.cache_read_input_tokens is not None:
            self.cache_read_input_tokens += usage.cache_read_input_tokens
        if usage.cache_creation_input_tokens is not None:
            self.cache_creation_input_tokens += usage.cache_creation_input_tokens

    def add_usage(
        self,
        chars_sent: int,
        chars_received: int,
        *,
        usage: TokenUsage | None = None,
    ) -> None:
        """Record a billed provider call.

        Args:
            chars_sent: Prompt length in characters.
            chars_received: Response length in characters.
            usage: Token counts, when the provider reported any.
        """
        self.chars_sent += chars_sent
        self.chars_received += chars_received
        self.calls_count += 1
        if usage is None:
            return
        if usage.input_tokens is not None:
            self.input_tokens += usage.input_tokens
        if usage.output_tokens is not None:
            self.output_tokens += usage.output_tokens
        self._add_usage_detail(usage)

    def add_cache_hit(
        self,
        chars_sent: int,
        chars_received: int,
        *,
        usage: TokenUsage | None = None,
    ) -> None:
        """Record a disk-cache hit (does not increment calls_count).

        Args:
            chars_sent: Prompt length in characters.
            chars_received: Cached response length in characters.
            usage: Token counts stored with the cache entry. ``None`` for
                entries written before usage was persisted, which report as
                unknown rather than as zero.
        """
        self.cache_hits += 1
        self.chars_sent += chars_sent
        self.chars_received += chars_received
        if usage is None:
            return
        if usage.input_tokens is not None:
            self.cached_input_tokens += usage.input_tokens
        if usage.output_tokens is not None:
            self.cached_output_tokens += usage.output_tokens
        self._add_usage_detail(usage)

    def add_ontology_update(self, num_operations: int, num_triples: int) -> None:
        """Add ontology update statistics.

        Args:
            num_operations: Number of update operations generated
            num_triples: Number of triples in these operations
        """
        self.ontology_operations_count += num_operations
        self.ontology_triples_generated += num_triples

    def add_facts_update(self, num_operations: int, num_triples: int) -> None:
        """Add facts update statistics.

        Args:
            num_operations: Number of update operations generated
            num_triples: Number of triples in these operations
        """
        self.facts_operations_count += num_operations
        self.facts_triples_generated += num_triples

    def merge_from(self, other: BudgetTracker) -> None:
        """Accumulate counters from another tracker (e.g. parallel unit workers)."""
        self.chars_sent += other.chars_sent
        self.chars_received += other.chars_received
        self.calls_count += other.calls_count
        self.cache_hits += other.cache_hits
        self.input_tokens += other.input_tokens
        self.output_tokens += other.output_tokens
        self.cached_input_tokens += other.cached_input_tokens
        self.cached_output_tokens += other.cached_output_tokens
        self.reasoning_tokens += other.reasoning_tokens
        self.cache_read_input_tokens += other.cache_read_input_tokens
        self.cache_creation_input_tokens += other.cache_creation_input_tokens
        self.ontology_triples_generated += other.ontology_triples_generated
        self.facts_triples_generated += other.facts_triples_generated
        self.ontology_operations_count += other.ontology_operations_count
        self.facts_operations_count += other.facts_operations_count
        for name, seconds in other.node_durations.items():
            self.add_duration(name, seconds)
        for name, count in other.counters.items():
            self.incr(name, count)

    def get_summary(self) -> str:
        """Get a summary of LLM usage and generated triples."""
        parts = [
            f"LLM: {self.calls_count} calls, "
            f"{self.chars_sent:,} sent, "
            f"{self.chars_received:,} received",
        ]
        if self.cache_hits > 0:
            parts.append(f"{self.cache_hits:,} cache hits")

        if self.input_tokens > 0 or self.output_tokens > 0:
            parts.append(
                f"{self.input_tokens:,} in / {self.output_tokens:,} out tokens"
            )

        # Reported separately so a replayed run reads as "free this time, but
        # this is what it costs", rather than as no token usage at all.
        if self.cached_input_tokens > 0 or self.cached_output_tokens > 0:
            parts.append(
                f"{self.cached_input_tokens:,} in / "
                f"{self.cached_output_tokens:,} out tokens replayed"
            )

        detail = [
            f"{value:,} {label}"
            for label, value in (
                ("reasoning", self.reasoning_tokens),
                ("provider-cache read", self.cache_read_input_tokens),
                ("provider-cache write", self.cache_creation_input_tokens),
            )
            if value > 0
        ]
        if detail:
            parts.append("of which " + ", ".join(detail))

        # The two ratios that decide where a cost change should be aimed: at the
        # prompt (low cache rate) or at the model's thinking budget (high
        # reasoning share). Both are derived, so they are formatted here rather
        # than stored -- and neither is money: pricing lives outside this repo.
        ratios = [
            f"{value:.0%} {label}"
            for label, value in (
                ("prefix-cache hits", self.prefix_cache_hit_rate),
                ("output is reasoning", self.reasoning_share_of_output),
            )
            if value
        ]
        if ratios:
            parts.append(", ".join(ratios))

        if self.ontology_triples_generated > 0 or self.facts_triples_generated > 0:
            parts.append(
                f"Triples: {self.ontology_triples_generated} ontology, "
                f"{self.facts_triples_generated} facts"
            )

        return " | ".join(parts)

    def get_duration_summary(self) -> str:
        """Get per-node wall-clock durations, slowest first."""
        if not self.node_durations:
            return ""
        ranked = sorted(
            self.node_durations.items(), key=lambda item: item[1], reverse=True
        )
        return "Durations: " + ", ".join(
            f"{name} {seconds:.1f}s" for name, seconds in ranked
        )

    def get_parallelism_summary(self) -> str:
        """Effective worker count and event-loop stall per fan-out node.

        Reports only nodes that recorded a ``/unit_sum`` entry, so it is empty
        for pipelines without a fan-out. ``lag`` is the time the event loop was
        blocked by synchronous work while units were meant to be running
        concurrently -- it is the difference between the width configured and
        the width achieved.
        """
        parts: list[str] = []
        for key in sorted(self.node_durations):
            if not key.endswith(UNIT_SUM_SUFFIX):
                continue
            node = key[: -len(UNIT_SUM_SUFFIX)]
            effective = self.parallel_efficiency(node)
            if effective is None:
                continue
            fragment = f"{node} {effective:.1f}x"
            lag = self.node_durations.get(f"{node}/loop_lag_total")
            if lag:
                fragment += f" (loop lag {lag:.1f}s)"
            parts.append(fragment)
        if not parts:
            return ""
        return "Effective workers: " + ", ".join(parts)

Attributes

cache_creation_input_tokens = Field(default=0, description="Input tokens written to the provider's prompt cache") class-attribute instance-attribute
cache_hits = Field(default=0, description='LLM calls satisfied from disk cache (no provider tokens)') class-attribute instance-attribute
cache_read_input_tokens = Field(default=0, description="Input tokens served from the provider's own prompt cache, counted inside the input totals and billed at a reduced rate") class-attribute instance-attribute
cached_input_tokens = Field(default=0, description='Input tokens replayed from the OntoCast disk cache') class-attribute instance-attribute
cached_output_tokens = Field(default=0, description='Output tokens replayed from the OntoCast disk cache') class-attribute instance-attribute
calls_count = Field(default=0, description='Total number of LLM API calls') class-attribute instance-attribute
chars_received = Field(default=0, description='Total characters received from LLM') class-attribute instance-attribute
chars_sent = Field(default=0, description='Total characters sent to LLM') class-attribute instance-attribute
counters = Field(default_factory=dict, description='Named event counts (e.g. how often a per-document computation ran). Summed on merge, like node_durations.') class-attribute instance-attribute
facts_operations_count = Field(default=0, description='Total number of facts update operations') class-attribute instance-attribute
facts_triples_generated = Field(default=0, description='Total number of triples generated for facts') class-attribute instance-attribute
input_tokens = Field(default=0, description='Billed input tokens (when reported by provider)') class-attribute instance-attribute
node_durations = Field(default_factory=dict, description='Accumulated seconds per pipeline node/stage (see class docstring)') class-attribute instance-attribute
ontology_operations_count = Field(default=0, description='Total number of ontology update operations') class-attribute instance-attribute
ontology_triples_generated = Field(default=0, description='Total number of triples generated for ontology updates') class-attribute instance-attribute
output_tokens = Field(default=0, description='Billed output tokens (when reported by provider)') class-attribute instance-attribute
prefix_cache_hit_rate property

Share of input tokens served from the provider's prompt cache.

This is the provider-side prefix cache, not the OntoCast disk cache -- the two are independent, and a deployment whose disk cache is cold (a fresh container, say) can still hit this one on every unit after the first, because the fan-out sends the same document-invariant prefix N times.

Near 0.0 with a fan-out wider than one unit means the workers all issued before any prefix had been cached: a cache entry only becomes readable once the first response has begun. Compare a long, more sequential run against a wide one to see the difference.

Returns:

Type Description
float | None

float | None: Ratio in [0.0, 1.0], or None when no input

float | None

tokens were reported at all.

reasoning_share_of_output property

Share of output tokens the model spent on thinking.

reasoning_tokens are counted inside the output totals, so this is a decomposition of what was already paid for, never an addition to it. A high share says the cost lever is the model's reasoning budget rather than anything about the prompt; a share of 0.0 says the model is not a reasoning model, and no reasoning-effort setting will change its bill.

Returns:

Type Description
float | None

float | None: Ratio in [0.0, 1.0], or None when no output

float | None

tokens were reported at all.

reasoning_tokens = Field(default=0, description='Thinking tokens, counted inside the output totals') class-attribute instance-attribute

Methods:

add_cache_hit(chars_sent, chars_received, *, usage=None)

Record a disk-cache hit (does not increment calls_count).

Parameters:

Name Type Description Default
chars_sent int

Prompt length in characters.

required
chars_received int

Cached response length in characters.

required
usage TokenUsage | None

Token counts stored with the cache entry. None for entries written before usage was persisted, which report as unknown rather than as zero.

None
Source code in ontocast/onto/state.py
def add_cache_hit(
    self,
    chars_sent: int,
    chars_received: int,
    *,
    usage: TokenUsage | None = None,
) -> None:
    """Record a disk-cache hit (does not increment calls_count).

    Args:
        chars_sent: Prompt length in characters.
        chars_received: Cached response length in characters.
        usage: Token counts stored with the cache entry. ``None`` for
            entries written before usage was persisted, which report as
            unknown rather than as zero.
    """
    self.cache_hits += 1
    self.chars_sent += chars_sent
    self.chars_received += chars_received
    if usage is None:
        return
    if usage.input_tokens is not None:
        self.cached_input_tokens += usage.input_tokens
    if usage.output_tokens is not None:
        self.cached_output_tokens += usage.output_tokens
    self._add_usage_detail(usage)
add_duration(name, seconds)

Accumulate seconds for a named node or stage.

Keys ending in :data:PEAK_SUFFIX take the maximum instead, since adding two peaks would report a stall that never happened.

Source code in ontocast/onto/state.py
def add_duration(self, name: str, seconds: float) -> None:
    """Accumulate seconds for a named node or stage.

    Keys ending in :data:`PEAK_SUFFIX` take the maximum instead, since
    adding two peaks would report a stall that never happened.
    """
    if name.endswith(PEAK_SUFFIX):
        self.node_durations[name] = max(self.node_durations.get(name, 0.0), seconds)
        return
    self.node_durations[name] = self.node_durations.get(name, 0.0) + seconds
add_facts_update(num_operations, num_triples)

Add facts update statistics.

Parameters:

Name Type Description Default
num_operations int

Number of update operations generated

required
num_triples int

Number of triples in these operations

required
Source code in ontocast/onto/state.py
def add_facts_update(self, num_operations: int, num_triples: int) -> None:
    """Add facts update statistics.

    Args:
        num_operations: Number of update operations generated
        num_triples: Number of triples in these operations
    """
    self.facts_operations_count += num_operations
    self.facts_triples_generated += num_triples
add_ontology_update(num_operations, num_triples)

Add ontology update statistics.

Parameters:

Name Type Description Default
num_operations int

Number of update operations generated

required
num_triples int

Number of triples in these operations

required
Source code in ontocast/onto/state.py
def add_ontology_update(self, num_operations: int, num_triples: int) -> None:
    """Add ontology update statistics.

    Args:
        num_operations: Number of update operations generated
        num_triples: Number of triples in these operations
    """
    self.ontology_operations_count += num_operations
    self.ontology_triples_generated += num_triples
add_usage(chars_sent, chars_received, *, usage=None)

Record a billed provider call.

Parameters:

Name Type Description Default
chars_sent int

Prompt length in characters.

required
chars_received int

Response length in characters.

required
usage TokenUsage | None

Token counts, when the provider reported any.

None
Source code in ontocast/onto/state.py
def add_usage(
    self,
    chars_sent: int,
    chars_received: int,
    *,
    usage: TokenUsage | None = None,
) -> None:
    """Record a billed provider call.

    Args:
        chars_sent: Prompt length in characters.
        chars_received: Response length in characters.
        usage: Token counts, when the provider reported any.
    """
    self.chars_sent += chars_sent
    self.chars_received += chars_received
    self.calls_count += 1
    if usage is None:
        return
    if usage.input_tokens is not None:
        self.input_tokens += usage.input_tokens
    if usage.output_tokens is not None:
        self.output_tokens += usage.output_tokens
    self._add_usage_detail(usage)
get_duration_summary()

Get per-node wall-clock durations, slowest first.

Source code in ontocast/onto/state.py
def get_duration_summary(self) -> str:
    """Get per-node wall-clock durations, slowest first."""
    if not self.node_durations:
        return ""
    ranked = sorted(
        self.node_durations.items(), key=lambda item: item[1], reverse=True
    )
    return "Durations: " + ", ".join(
        f"{name} {seconds:.1f}s" for name, seconds in ranked
    )
get_parallelism_summary()

Effective worker count and event-loop stall per fan-out node.

Reports only nodes that recorded a /unit_sum entry, so it is empty for pipelines without a fan-out. lag is the time the event loop was blocked by synchronous work while units were meant to be running concurrently -- it is the difference between the width configured and the width achieved.

Source code in ontocast/onto/state.py
def get_parallelism_summary(self) -> str:
    """Effective worker count and event-loop stall per fan-out node.

    Reports only nodes that recorded a ``/unit_sum`` entry, so it is empty
    for pipelines without a fan-out. ``lag`` is the time the event loop was
    blocked by synchronous work while units were meant to be running
    concurrently -- it is the difference between the width configured and
    the width achieved.
    """
    parts: list[str] = []
    for key in sorted(self.node_durations):
        if not key.endswith(UNIT_SUM_SUFFIX):
            continue
        node = key[: -len(UNIT_SUM_SUFFIX)]
        effective = self.parallel_efficiency(node)
        if effective is None:
            continue
        fragment = f"{node} {effective:.1f}x"
        lag = self.node_durations.get(f"{node}/loop_lag_total")
        if lag:
            fragment += f" (loop lag {lag:.1f}s)"
        parts.append(fragment)
    if not parts:
        return ""
    return "Effective workers: " + ", ".join(parts)
get_summary()

Get a summary of LLM usage and generated triples.

Source code in ontocast/onto/state.py
def get_summary(self) -> str:
    """Get a summary of LLM usage and generated triples."""
    parts = [
        f"LLM: {self.calls_count} calls, "
        f"{self.chars_sent:,} sent, "
        f"{self.chars_received:,} received",
    ]
    if self.cache_hits > 0:
        parts.append(f"{self.cache_hits:,} cache hits")

    if self.input_tokens > 0 or self.output_tokens > 0:
        parts.append(
            f"{self.input_tokens:,} in / {self.output_tokens:,} out tokens"
        )

    # Reported separately so a replayed run reads as "free this time, but
    # this is what it costs", rather than as no token usage at all.
    if self.cached_input_tokens > 0 or self.cached_output_tokens > 0:
        parts.append(
            f"{self.cached_input_tokens:,} in / "
            f"{self.cached_output_tokens:,} out tokens replayed"
        )

    detail = [
        f"{value:,} {label}"
        for label, value in (
            ("reasoning", self.reasoning_tokens),
            ("provider-cache read", self.cache_read_input_tokens),
            ("provider-cache write", self.cache_creation_input_tokens),
        )
        if value > 0
    ]
    if detail:
        parts.append("of which " + ", ".join(detail))

    # The two ratios that decide where a cost change should be aimed: at the
    # prompt (low cache rate) or at the model's thinking budget (high
    # reasoning share). Both are derived, so they are formatted here rather
    # than stored -- and neither is money: pricing lives outside this repo.
    ratios = [
        f"{value:.0%} {label}"
        for label, value in (
            ("prefix-cache hits", self.prefix_cache_hit_rate),
            ("output is reasoning", self.reasoning_share_of_output),
        )
        if value
    ]
    if ratios:
        parts.append(", ".join(ratios))

    if self.ontology_triples_generated > 0 or self.facts_triples_generated > 0:
        parts.append(
            f"Triples: {self.ontology_triples_generated} ontology, "
            f"{self.facts_triples_generated} facts"
        )

    return " | ".join(parts)
incr(name, n=1)

Increment a named counter.

Parameters:

Name Type Description Default
name str

Counter key, e.g. "ctx/merge_document_ontology.calls".

required
n int

Amount to add.

1
Source code in ontocast/onto/state.py
def incr(self, name: str, n: int = 1) -> None:
    """Increment a named counter.

    Args:
        name: Counter key, e.g. ``"ctx/merge_document_ontology.calls"``.
        n: Amount to add.
    """
    self.counters[name] = self.counters.get(name, 0) + n
merge_from(other)

Accumulate counters from another tracker (e.g. parallel unit workers).

Source code in ontocast/onto/state.py
def merge_from(self, other: BudgetTracker) -> None:
    """Accumulate counters from another tracker (e.g. parallel unit workers)."""
    self.chars_sent += other.chars_sent
    self.chars_received += other.chars_received
    self.calls_count += other.calls_count
    self.cache_hits += other.cache_hits
    self.input_tokens += other.input_tokens
    self.output_tokens += other.output_tokens
    self.cached_input_tokens += other.cached_input_tokens
    self.cached_output_tokens += other.cached_output_tokens
    self.reasoning_tokens += other.reasoning_tokens
    self.cache_read_input_tokens += other.cache_read_input_tokens
    self.cache_creation_input_tokens += other.cache_creation_input_tokens
    self.ontology_triples_generated += other.ontology_triples_generated
    self.facts_triples_generated += other.facts_triples_generated
    self.ontology_operations_count += other.ontology_operations_count
    self.facts_operations_count += other.facts_operations_count
    for name, seconds in other.node_durations.items():
        self.add_duration(name, seconds)
    for name, count in other.counters.items():
        self.incr(name, count)
parallel_efficiency(node)

Effective worker count for a fan-out node, or None if unmeasured.

The ratio of time summed across unit workers to the node's wall clock. A value near parallel_workers means the fan-out is running at full width; a value near 1.0 means the units are effectively serialised -- typically by synchronous CPU work blocking the event loop, which "<node>/loop_lag_total" quantifies.

Parameters:

Name Type Description Default
node str

The node key, e.g. str(WorkflowNode.RENDER_FACTS).

required

Returns:

Type Description
float | None

float | None: Effective workers, or None when either the wall

float | None

clock or the /unit_sum entry is missing or zero.

Source code in ontocast/onto/state.py
def parallel_efficiency(self, node: str) -> float | None:
    """Effective worker count for a fan-out node, or ``None`` if unmeasured.

    The ratio of time summed across unit workers to the node's wall clock.
    A value near ``parallel_workers`` means the fan-out is running at full
    width; a value near ``1.0`` means the units are effectively serialised
    -- typically by synchronous CPU work blocking the event loop, which
    ``"<node>/loop_lag_total"`` quantifies.

    Args:
        node: The node key, e.g. ``str(WorkflowNode.RENDER_FACTS)``.

    Returns:
        float | None: Effective workers, or ``None`` when either the wall
        clock or the ``/unit_sum`` entry is missing or zero.
    """
    wall = self.node_durations.get(node)
    unit_sum = self.node_durations.get(f"{node}{UNIT_SUM_SUFFIX}")
    if not wall or unit_sum is None:
        return None
    return unit_sum / wall

Functions: