Skip to content

ontocast.toolbox

Attributes

logger = logging.getLogger(__name__) module-attribute

Classes

ToolBox

A container class for all tools used in the ontology processing workflow.

This class initializes and manages various tools needed for document processing, ontology management, and LLM interactions.

Parameters:

Name Type Description Default
config Config

Configuration object containing all necessary settings.

required
Source code in ontocast/toolbox.py
 145
 146
 147
 148
 149
 150
 151
 152
 153
 154
 155
 156
 157
 158
 159
 160
 161
 162
 163
 164
 165
 166
 167
 168
 169
 170
 171
 172
 173
 174
 175
 176
 177
 178
 179
 180
 181
 182
 183
 184
 185
 186
 187
 188
 189
 190
 191
 192
 193
 194
 195
 196
 197
 198
 199
 200
 201
 202
 203
 204
 205
 206
 207
 208
 209
 210
 211
 212
 213
 214
 215
 216
 217
 218
 219
 220
 221
 222
 223
 224
 225
 226
 227
 228
 229
 230
 231
 232
 233
 234
 235
 236
 237
 238
 239
 240
 241
 242
 243
 244
 245
 246
 247
 248
 249
 250
 251
 252
 253
 254
 255
 256
 257
 258
 259
 260
 261
 262
 263
 264
 265
 266
 267
 268
 269
 270
 271
 272
 273
 274
 275
 276
 277
 278
 279
 280
 281
 282
 283
 284
 285
 286
 287
 288
 289
 290
 291
 292
 293
 294
 295
 296
 297
 298
 299
 300
 301
 302
 303
 304
 305
 306
 307
 308
 309
 310
 311
 312
 313
 314
 315
 316
 317
 318
 319
 320
 321
 322
 323
 324
 325
 326
 327
 328
 329
 330
 331
 332
 333
 334
 335
 336
 337
 338
 339
 340
 341
 342
 343
 344
 345
 346
 347
 348
 349
 350
 351
 352
 353
 354
 355
 356
 357
 358
 359
 360
 361
 362
 363
 364
 365
 366
 367
 368
 369
 370
 371
 372
 373
 374
 375
 376
 377
 378
 379
 380
 381
 382
 383
 384
 385
 386
 387
 388
 389
 390
 391
 392
 393
 394
 395
 396
 397
 398
 399
 400
 401
 402
 403
 404
 405
 406
 407
 408
 409
 410
 411
 412
 413
 414
 415
 416
 417
 418
 419
 420
 421
 422
 423
 424
 425
 426
 427
 428
 429
 430
 431
 432
 433
 434
 435
 436
 437
 438
 439
 440
 441
 442
 443
 444
 445
 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
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
1281
1282
1283
1284
1285
1286
1287
1288
1289
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
1315
1316
1317
1318
1319
1320
1321
1322
1323
1324
1325
1326
1327
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
1353
1354
1355
1356
1357
1358
1359
1360
1361
1362
1363
1364
1365
1366
1367
1368
1369
1370
1371
1372
1373
1374
1375
1376
1377
1378
1379
1380
1381
1382
1383
1384
1385
1386
1387
1388
1389
1390
1391
1392
1393
1394
1395
1396
1397
1398
1399
1400
1401
1402
1403
1404
1405
1406
1407
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
1426
1427
1428
1429
1430
1431
1432
1433
1434
1435
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1451
1452
1453
1454
1455
1456
1457
1458
1459
1460
1461
1462
1463
1464
1465
1466
1467
1468
1469
1470
1471
1472
1473
1474
1475
1476
1477
1478
1479
1480
1481
1482
1483
1484
1485
1486
1487
1488
1489
1490
1491
1492
1493
1494
1495
1496
1497
1498
1499
1500
class ToolBox:
    """A container class for all tools used in the ontology processing workflow.

    This class initializes and manages various tools needed for document processing,
    ontology management, and LLM interactions.

    Args:
        config: Configuration object containing all necessary settings.
    """

    @classmethod
    async def acreate(cls, config: Config) -> "ToolBox":
        """Construct a ToolBox from inside a running event loop.

        Equivalent to ``ToolBox(config)``, except that LLM provider setup is
        awaited rather than driven through :func:`asyncio.run` -- which is
        illegal in a loop and is what makes the plain constructor unusable from
        async code. Embedders should prefer this, and pair it with
        ``async with`` so backend connections are released:

        ```python
        async with await ToolBox.acreate(config) as tools:
            await tools.initialize()
        ```

        Args:
            config: Fully resolved configuration.

        Returns:
            A ready ToolBox. Call :meth:`initialize` to sync ontologies and
            prepare backend schema.
        """
        warn_if_workers_exceed_inflight(config)
        runtime = await ToolBoxRuntime.acreate(config)
        return cls(config, runtime=runtime)

    def __init__(
        self,
        config: Config,
        *,
        llm: LLMTool | None = None,
        runtime: "ToolBoxRuntime | None" = None,
    ):
        """Build a ToolBox bound to whatever partition ``config`` names.

        Args:
            config: Fully resolved configuration.
            llm: Pre-built LLM tool, used when no ``runtime`` is supplied.
                ``LLMTool.create`` otherwise runs, which cannot be called from
                inside a running event loop -- prefer :meth:`acreate` there.
            runtime: Shared tenancy-independent tools. Supplied by
                :class:`~ontocast.registry.ToolBoxRegistry` so scoped ToolBoxes
                do not each load an embedding model; built fresh when omitted.
        """
        # Store the config for later use
        self.config = config

        # Get tool configuration
        tool_config = config.get_tool_config()

        # Tools that do not vary by tenant live on the runtime, so a registry of
        # scoped ToolBoxes shares one LLM client, converter and embedding model.
        if runtime is None:
            # A process-level fact, said where the process builds its LLM
            # client -- not once per tenancy scope over a shared runtime.
            warn_if_workers_exceed_inflight(config)
        self.runtime = runtime or ToolBoxRuntime(config, llm=llm)

        # Fuseki when a URI is configured (auth is optional), otherwise in-memory.
        if tool_config.fuseki.uri:
            self.triple_store_manager: TripleStoreManager = FusekiTripleStoreManager(
                uri=tool_config.fuseki.uri,
                auth=tool_config.fuseki.auth,
                dataset=tool_config.fuseki.dataset,
                ontologies_dataset=tool_config.fuseki.ontologies_dataset,
                shapes_dataset=tool_config.fuseki.shapes_dataset,
            )
        else:
            self.triple_store_manager = InMemoryTripleStoreManager()

        self.ontology_manager: OntologyManager = OntologyManager()
        self.ontology_manager.register_triple_store(self.triple_store_manager)

        # The shapes partition. Partition-scoped like the ontology catalog, and
        # reset on a tenancy switch for the same reason.
        self.shapes_catalog: ShapesCatalog = ShapesCatalog()
        self.shapes_catalog.register_triple_store(self.triple_store_manager)
        # Tenancy the in-memory catalog currently reflects; None until first set.
        self._active_tenancy: tuple[str, str] | None = None
        # Guards the tenancy retarget, which mutates ToolBox-wide state (dataset
        # names, the ontology catalog, vector-store table names) across awaits.
        # It is driven by a per-request query parameter with no concurrency cap,
        # so without this two requests for different tenants can interleave and
        # read or write each other's partition. Created lazily: __init__ may run
        # outside an event loop.
        self._tenancy_lock: asyncio.Lock | None = None
        self._tenancy_lock_loop: asyncio.AbstractEventLoop | None = None
        # Set by attach_registry() when this ToolBox fronts a multi-tenant host.
        self._registry: "ToolBoxRegistry | None" = None

        # Graph algorithms over graphs it is handed; it does not fetch. Built
        # only when something consumes it: its sole consumer is
        # OntologyPatchRetriever, which exists only alongside a vector store,
        # so constructing it unconditionally made the dependency graph claim
        # "always needed" where the truth is "needed with vector retrieval".
        self.sparql_tool: SPARQLTool | None = None

        self.vector_store: VectorStoreManager | None = None
        self.patch_retriever: OntologyPatchRetriever | None = None
        self.vector_store_ready: bool = False
        self.vector_store_last_error: Exception | None = None
        # Seed ontologies the last sync had to supply because the partition
        # served no graph for them. Read by the tenancy switch, which must
        # materialize exactly those and nothing else -- rewriting a whole
        # catalog on every switch is what scoped ToolBoxes exist to avoid.
        self._last_seed_repairs: list[Ontology] = []
        self._vector_store_init_lock = asyncio.Lock()

        # Both supported backends need the BM25 tool, so the only case that
        # skips it is having no vector store at all -- decided by the same
        # resolution the factory uses.
        needs_sparse = resolve_backend(tool_config) is not VectorStoreBackend.NONE
        vector_store = create_vector_store_manager(
            tool_config,
            embedding=self.embedding_tool,
            sparse_embedding=(
                self.runtime.sparse_embedding_tool(tool_config.embedding)
                if needs_sparse
                else None
            ),
        )
        if vector_store is not None:
            self.vector_store = vector_store
            self.sparql_tool = SPARQLTool(
                triple_store_manager=self.triple_store_manager
            )
            self.patch_retriever = OntologyPatchRetriever(
                vector_store=vector_store,
                sparql_tool=self.sparql_tool,
                patch=tool_config.patch_retrieval,
                ontology_manager=self.ontology_manager,
            )
            self.ontology_manager.register_vector_store(self.patch_retriever)

    # -- shared runtime delegates -----------------------------------------
    #
    # These tools do not vary by tenant and live on the runtime, but they were
    # ToolBox attributes for the whole life of the project and are read -- and
    # substituted -- that way across the pipeline, the CLI, the worker and the
    # tests. Each delegates in both directions, so a caller replacing
    # `tools.converter` replaces the shared one, which is what it always meant.

    @property
    def runtime(self) -> ToolBoxRuntime:
        """Shared tenancy-independent tools.

        Materialized empty on first access when it was never assigned. Several
        unit tests build a ToolBox with ``ToolBox.__new__(ToolBox)`` to exercise
        one tool without standing up the whole container, and used to set that
        tool as a plain attribute; this keeps that working, and reading a tool
        that was never set still raises ``AttributeError`` for its own name.
        """
        runtime = self.__dict__.get("_runtime")
        if runtime is None:
            runtime = ToolBoxRuntime.__new__(ToolBoxRuntime)
            self.__dict__["_runtime"] = runtime
        return runtime

    @runtime.setter
    def runtime(self, value: ToolBoxRuntime) -> None:
        self.__dict__["_runtime"] = value

    @property
    def shared_cache(self) -> Cacher:
        """Shared on-disk cache backing the LLM and converter tools."""
        return self.runtime.shared_cache

    @shared_cache.setter
    def shared_cache(self, value: Cacher) -> None:
        self.runtime.shared_cache = value

    @property
    def llm(self) -> LLMTool:
        """Shared LLM tool."""
        return self.runtime.llm

    @llm.setter
    def llm(self, value: LLMTool) -> None:
        self.runtime.llm = value

    @property
    def llm_provider(self):
        """Configured LLM provider."""
        return self.runtime.llm_provider

    @llm_provider.setter
    def llm_provider(self, value) -> None:
        self.runtime.llm_provider = value

    @property
    def search_provider(self):
        """Configured web-search provider, or None when disabled."""
        return self.runtime.search_provider

    @search_provider.setter
    def search_provider(self, value) -> None:
        self.runtime.search_provider = value

    @property
    def atomic_tools(self) -> AtomicToolBox:
        """Per-unit tool surface used by the render/critic loops."""
        return self.runtime.atomic_tools

    @atomic_tools.setter
    def atomic_tools(self, value: AtomicToolBox) -> None:
        self.runtime.atomic_tools = value

    @property
    def converter(self) -> ConverterTool:
        """Document converter."""
        return self.runtime.converter

    @converter.setter
    def converter(self, value: ConverterTool) -> None:
        self.runtime.converter = value

    @property
    def chunker(self) -> ChunkerTool:
        """Text chunker."""
        return self.runtime.chunker

    @chunker.setter
    def chunker(self, value: ChunkerTool) -> None:
        self.runtime.chunker = value

    @property
    def aggregator(self) -> EmbeddingBasedAggregator:
        """Facts aggregator."""
        return self.runtime.aggregator

    @aggregator.setter
    def aggregator(self, value: EmbeddingBasedAggregator) -> None:
        self.runtime.aggregator = value

    @property
    def embedding_tool(self) -> EmbeddingTool:
        """Dense embedding provider."""
        return self.runtime.embedding_tool

    @embedding_tool.setter
    def embedding_tool(self, value: EmbeddingTool) -> None:
        self.runtime.embedding_tool = value

    def get_entity_aligner(
        self,
        embedding_model: str | None = None,
        similarity_threshold: float | None = None,
    ) -> EntityAligner:
        """Return a cached entity aligner for the given embedding settings."""
        tool_config = self.config.get_tool_config()
        return self.runtime.get_entity_aligner(
            embedding_model or tool_config.aggregation.embedding_model,
            similarity_threshold
            if similarity_threshold is not None
            else tool_config.aggregation.similarity_threshold,
        )

    async def get_llm_tool(self, budget_tracker: BudgetTracker) -> LLMTool:
        """Return the shared LLM tool, charging usage to ``budget_tracker``.

        Args:
            budget_tracker: The budget tracker to charge for this task's calls.

        Returns:
            LLMTool: The shared LLM tool.
        """
        return await self.runtime.get_llm_tool(budget_tracker)

    def require_triple_store_manager(self) -> TripleStoreManager:
        """Return the configured triple store manager or raise a clear error."""
        manager = self.triple_store_manager
        if manager is None:
            raise RuntimeError("Triple store backend is not configured")
        return manager

    def require_vector_store(self) -> VectorStoreManager:
        """Return the configured vector store or raise a directive error."""
        if self.vector_store is None:
            raise RuntimeError(
                "No vector store is configured. Set QDRANT_URI (Qdrant server) "
                "or LANCEDB_ENABLED=true (embedded LanceDB); each needs its "
                "matching extra, ontocast[qdrant] or ontocast[lancedb]."
            )
        return self.vector_store

    def require_patch_retriever(self) -> OntologyPatchRetriever:
        """Return the ontology patch retriever or raise a directive error."""
        if self.patch_retriever is None:
            raise RuntimeError(
                "Ontology patch retrieval needs a vector store. Set QDRANT_URI "
                "or LANCEDB_ENABLED=true."
            )
        return self.patch_retriever

    async def aclose(self) -> None:
        """Release every backend connection this ToolBox opened.

        The ToolBox owns an httpx client (Fuseki) and a Qdrant client, neither
        of which was previously closed anywhere -- ``FusekiTripleStoreManager``
        even defined ``close()`` that nothing called. Long-lived hosts that
        build a ToolBox per tenant, and tests that build many, leaked sockets.

        Also closes any scoped ToolBoxes this one spawned through
        :meth:`for_scope`, so shutting down the ToolBox an application holds
        releases every tenant's connections too.

        Safe to call more than once, and never raises: teardown failures are
        logged, since a caller shutting down cannot act on them.
        """
        registry = self._registry
        if registry is not None:
            self._registry = None
            try:
                await registry.aclose()
            except Exception as exc:
                logger.warning("Error closing tenancy registry: %s", exc)

        if self.triple_store_manager is not None:
            try:
                await self.triple_store_manager.close()
            except Exception as exc:
                logger.warning("Error closing triple store manager: %s", exc)

        if self.vector_store is not None:
            try:
                await asyncio.to_thread(self.vector_store.close)
            except Exception as exc:
                logger.warning("Error closing vector store: %s", exc)

    async def __aenter__(self) -> "ToolBox":
        return self

    async def __aexit__(self, *exc_info: object) -> None:
        await self.aclose()

    @property
    def scope(self) -> TenancyScope | None:
        """The partition this ToolBox is bound to, once tenancy is assigned."""
        if self._active_tenancy is None:
            return None
        return TenancyScope.build(*self._active_tenancy)

    def bind_scope(self, scope: TenancyScope) -> None:
        """Record the partition a ToolBox was built for.

        For a ToolBox whose stores were configured for ``scope`` at
        construction (``Config.for_tenancy``): stores and catalogs are left
        alone, unlike ``update_tenancy``, which retargets them.
        """
        self._active_tenancy = scope.key

    def attach_registry(self, registry: "ToolBoxRegistry") -> None:
        """Resolve other partitions through ``registry`` rather than a private one.

        Only needed to share one registry across several ToolBoxes;
        :meth:`for_scope` builds its own on first use otherwise.
        """
        self._registry = registry

    def ensure_tenancy_registry(self) -> "ToolBoxRegistry":
        """Return this ToolBox's registry, creating it on first use.

        Built lazily so a single-tenant embedder never allocates one, and so
        nothing has to be wired at construction time.
        """
        if self._registry is None:
            from ontocast.registry import ToolBoxRegistry

            self._registry = ToolBoxRegistry(
                self.config,
                self.runtime,
                max_scopes=self.config.server.max_tenancy_scopes,
            )
        return self._registry

    async def for_scope(
        self,
        tenant: str,
        project: str,
        *,
        ontology_context_mode: OntologyContextMode | None = None,
        fail_on_vector_store_error: bool = False,
    ) -> "ToolBox":
        """Return a ToolBox bound to ``tenant`` / ``project``.

        Returns ``self`` when the scope already matches. Otherwise resolves
        through the attached registry, which shares this ToolBox's runtime, so
        the new scope costs a triple store and an ontology catalog rather than
        another embedding model.

        Isolation is by construction: each scope owns a deep copy of ``Config``.
        That copy matters -- vector store managers hold their config sections by
        reference and rewrite collection names when tenancy is applied, so
        scopes sharing a ``Config`` would alias each other.

        Args:
            tenant: Tenant identifier.
            project: Project identifier within the tenant.
            ontology_context_mode: Mode to initialize a newly built scope for.
            fail_on_vector_store_error: Raise rather than log when vector store
                preparation fails for a newly built scope.

        Returns:
            A ToolBox bound to the requested partition.
        """
        requested = TenancyScope.build(tenant, project)
        if self._active_tenancy == requested.key:
            return self
        return await self.ensure_tenancy_registry().get(
            requested,
            ontology_context_mode=ontology_context_mode,
            fail_on_vector_store_error=fail_on_vector_store_error,
        )

    async def update_tenancy(self, tenant: str, project: str) -> None:
        """Retarget Fuseki datasets and Qdrant collections for ``tenant`` / ``project``."""
        await self.update_tenancy_with_vector_mode(
            tenant,
            project,
            initialize_vector_store=True,
            fail_on_vector_store_error=True,
        )

    def _get_tenancy_lock(self) -> asyncio.Lock:
        """Return the tenancy lock bound to the running loop.

        Rebuilt when the loop changes: the CLI bootstrap runs several
        ``asyncio.run`` calls, and an ``asyncio.Lock`` created on a closed loop
        cannot be awaited on a later one.
        """
        loop = asyncio.get_running_loop()
        if self._tenancy_lock is None or self._tenancy_lock_loop is not loop:
            self._tenancy_lock = asyncio.Lock()
            self._tenancy_lock_loop = loop
        return self._tenancy_lock

    async def update_tenancy_with_vector_mode(
        self,
        tenant: str,
        project: str,
        *,
        initialize_vector_store: bool,
        fail_on_vector_store_error: bool,
    ) -> None:
        """Retarget tenancy and optionally initialize vector store collections.

        Serialized: the body mutates ToolBox-wide state across awaits, and the
        HTTP layer calls this per request from a ``?tenant=`` query parameter
        with no concurrency cap. Interleaving two switches leaves the catalog
        and the store handles describing different tenants.
        """
        async with self._get_tenancy_lock():
            await self._update_tenancy_with_vector_mode_locked(
                tenant,
                project,
                initialize_vector_store=initialize_vector_store,
                fail_on_vector_store_error=fail_on_vector_store_error,
            )

    async def _update_tenancy_with_vector_mode_locked(
        self,
        tenant: str,
        project: str,
        *,
        initialize_vector_store: bool,
        fail_on_vector_store_error: bool,
    ) -> None:
        t, p = tenant.strip(), project.strip()
        if not t or not p:
            raise ValueError("tenant and project must be non-empty")

        tenancy_changed = (t, p) != self._active_tenancy

        triple = self.triple_store_manager
        if triple is not None and triple.supports_tenancy_partition():
            await triple.update_tenancy(t, p)
            if isinstance(triple, FusekiTripleStoreManager):
                fuseki_cfg = self.config.tool_config.fuseki
                fuseki_cfg.dataset = triple.dataset
                fuseki_cfg.ontologies_dataset = triple.ontologies_dataset
                fuseki_cfg.shapes_dataset = triple.shapes_dataset

        if tenancy_changed:
            # The catalog, its alias-collision ledger, and the graph caches are all
            # partition-scoped. Carrying them across a switch leaks one tenant's
            # ontologies into another's requests -- and its alias ledger can reject
            # a legitimately distinct ontology that reuses an ontology_id.
            # ``None`` means this is the first assignment, which happens at startup
            # before ``initialize()``; leave the population to it rather than
            # fetching twice. Any later switch must repopulate -- including when the
            # partition we are leaving was empty.
            is_first_assignment = self._active_tenancy is None
            self.ontology_manager.reset_catalog()
            self.shapes_catalog.reset()
            self._active_tenancy = (t, p)
            if not is_first_assignment and triple is not None:
                # Seed TTLs are replayed into the new partition, but only where
                # it serves no graph of its own: a scope whose catalog is empty
                # for want of a bootstrap is the same fault as a startup one,
                # and answering it with "the seeds are startup-only" left the
                # tenant extracting against no vocabulary at all. A partition
                # that already serves an ontology keeps it -- the seed never
                # overwrites a terminal a previous run evolved.
                for ontology in await self._synchronize_ontologies():
                    self.ontology_manager.add_ontology(ontology, skip_vector_index=True)
                for seed in self._last_seed_repairs:
                    await self._materialize_ontology(seed)
                # Shapes are still whatever the partition holds plus the
                # configured directory, and may be nothing -- which correctly
                # reads as "SHACL never checked" rather than "conforms".
                await self.shapes_catalog.sync(
                    self.config.tool_config.facts_validation.shapes_dir
                )

        if self.vector_store is not None:
            # No config copy-back: `create_vector_store_manager` passes
            # `tool_config.vector_store` by reference, so `apply_tenancy` has
            # already rewritten the very object Config holds.
            self.vector_store.apply_tenancy(t, p)
            if initialize_vector_store:
                try:
                    await self.vector_store.initialize()
                    self.vector_store_ready = True
                    self.vector_store_last_error = None
                except Exception as exc:
                    self.vector_store_ready = False
                    self.vector_store_last_error = exc
                    if fail_on_vector_store_error:
                        raise
                    logger.warning(
                        "Vector store tenancy initialization failed; continuing without vector retrieval: %s",
                        exc,
                    )

    async def clean_tenancy_data(
        self, tenant: str, project: str, *, include_shapes: bool = False
    ) -> None:
        """Flush triple-store and vector-store partitions for ``tenant`` / ``project``.

        Shapes are retained unless ``include_shapes`` is set: facts and
        ontologies come back from a rerun, but shapes are the deployment's
        validation contract, and dropping them turns the SHACL gate off without
        an error -- the next run reports ``shacl_evaluated: null`` instead of
        failing.

        Takes the tenancy lock: this is destructive, and a concurrent retarget
        would let it resolve partition names against a scope that changed
        mid-flight.

        Args:
            tenant: Tenant identifier.
            project: Project identifier within the tenant.
            include_shapes: Also drop the shapes partition.
        """
        t, p = tenant.strip(), project.strip()
        if not t or not p:
            raise ValueError("tenant and project must be non-empty")

        async with self._get_tenancy_lock():
            triple = self.triple_store_manager
            if triple is not None:
                if not triple.supports_tenancy_partition():
                    raise NotImplementedError(
                        f"Triple store {type(triple).__name__} has no tenant/project partitions"
                    )
                await triple.clean_tenancy(t, p, include_shapes=include_shapes)

            vector = self.vector_store
            if vector is not None and vector.supports_tenancy_partition():
                was_ready = self.vector_store_ready
                # `clean_tenancy` drops the collections and their embedding
                # metadata. Leaving the store marked ready would point every
                # later search in this process at a collection that no longer
                # exists, so readiness is dropped first and only restored by a
                # successful recreate -- the same wipe/recreate order the
                # startup path uses. Nothing is reindexed: the flush emptied
                # the catalog this index would have been built from.
                self.vector_store_ready = False
                await vector.clean_tenancy(t, p)
                if was_ready:
                    try:
                        await vector.initialize()
                        self.vector_store_ready = True
                        self.vector_store_last_error = None
                    except Exception as exc:
                        self.vector_store_last_error = exc
                        logger.warning(
                            "Vector store could not be recreated after flushing "
                            "%s/%s; vector retrieval is off until it is "
                            "reinitialized: %s",
                            t,
                            p,
                            exc,
                        )

        if include_shapes and (t, p) == self._active_tenancy:
            self.shapes_catalog.reset()

    def get_atomic_tools(self) -> AtomicToolBox:
        """Return the minimal toolbox used by atomic render/critic paths.

        The runtime's :class:`AtomicToolBox` is shared by every scope; what is
        handed out here is a shallow copy bound to *this* scope's ontology
        catalog, so the per-unit repairs can ask the whole catalog whether a
        term exists rather than only the unit's retrieved snapshot. The copy
        is per call and carries no state of its own -- the catalog-term memo
        lives on the catalog -- so a tool replaced on the shared instance is
        seen by the next unit.
        """
        return self.atomic_tools.scoped_to_catalog(self.ontology_manager)

    def shapes_prompt_contract(self) -> tuple[str, tuple[str, ...], bool]:
        """Conformance chapter, term exemptions, and the selection flag.

        ``("", (), False)`` when the contract is off or the shapes partition
        is empty -- the prompt is then byte-identical to a shape-less
        deployment's. Full-catalog modes return the memoized whole chapter
        with the flag False. When per-unit selection is in force (mode
        ``context``, or ``auto`` with a catalog that outgrew the line cap)
        the chapter comes back empty with the flag True, and the unit loop
        fills it via :meth:`shapes_chapter_for_context` once the unit's
        ontology snapshot exists. Exemption terms are the full catalog's in
        every mode -- the gate validates against every shape, and the terms
        are catalog IRIs either way.
        """
        cfg = self.config.tool_config.facts_validation
        mode = cfg.shapes_prompt_contract
        if mode == "off":
            return "", (), False
        max_lines = cfg.shapes_prompt_max_lines
        terms = self.shapes_catalog.prompt_contract_terms(max_lines=max_lines)
        select = mode == "context" or (
            mode == "auto" and self.shapes_catalog.needs_selection(max_lines=max_lines)
        )
        if select:
            return "", terms, True
        return (
            self.shapes_catalog.conformance_chapter(max_lines=max_lines),
            terms,
            False,
        )

    def shapes_chapter_for_context(self, context_terms: set[str]) -> str:
        """Per-unit conformance chapter, joined on the unit's context IRIs."""
        cfg = self.config.tool_config.facts_validation
        return self.shapes_catalog.selected_chapter(
            context_terms, max_lines=cfg.shapes_prompt_max_lines
        )

    def serialize(self, state: AgentState) -> None:
        """Persist the document's ontologies and facts.

        Drives :meth:`aserialize` under a single :func:`asyncio.run`, so a
        document with N ontologies costs one event loop and one backend
        connection rather than N of each -- the per-call sync entry points open
        (and tear down) a fresh HTTP client every time.

        Raises:
            RuntimeError: If called from inside a running event loop; await
                :meth:`aserialize` there instead.
        """
        require_no_running_loop("ToolBox.serialize", "ToolBox.aserialize")
        asyncio.run(self.aserialize(state))

    async def aserialize(self, state: AgentState) -> None:
        """Persist the document's ontologies and facts (async form)."""
        ontologies_to_serialize = document_ontology_access(
            state
        ).serialization_targets()
        for ontology in ontologies_to_serialize:
            if ontology and ontology.hash:
                # The async variant is not optional here: the sync
                # add_ontology() refuses to reindex the vector store inside a
                # running event loop, and aserialize is by definition inside
                # one — so with vector retrieval registered, the sync call
                # raised RuntimeError on the first document that actually
                # produced an ontology version to serialize.
                await self.ontology_manager.aadd_ontology(ontology)

        if self.triple_store_manager is not None:
            for ontology in ontologies_to_serialize:
                await self.triple_store_manager.aserialize(ontology)
            if state.render_facts:
                await self.triple_store_manager.aserialize(
                    state.aggregated_facts,
                    graph_uri=state.graph_uri,
                )

    def should_initialize_vector_store(
        self, ontology_context_mode: OntologyContextMode | None
    ) -> bool:
        return (
            self.vector_store is not None
            and ontology_context_mode
            == OntologyContextMode.SELECTED_VECTOR_SEARCH_ONTOLOGY
        )

    def is_vector_store_ready(self) -> bool:
        return self.vector_store is not None and self.vector_store_ready

    async def ensure_vector_store(
        self,
        ontology_context_mode: OntologyContextMode | None,
        *,
        fail_on_vector_store_error: bool,
    ) -> None:
        """Prepare the vector store when ``ontology_context_mode`` needs it.

        A ToolBox is initialized for the mode it started in; a later request in
        vector mode would otherwise find no index and be refused.
        """
        if not self.should_initialize_vector_store(ontology_context_mode):
            return
        async with self._vector_store_init_lock:
            if self.is_vector_store_ready():
                return
            logger.info("Preparing the vector store for a vector-mode request")
            await self.initialize(
                ontology_context_mode=ontology_context_mode,
                fail_on_vector_store_error=fail_on_vector_store_error,
                wipe_vector_store=False,
            )

    async def _check_catalog_index_agreement(
        self, synchronized_ontologies: list[Ontology]
    ) -> None:
        """Refuse to start when the vector index knows ontologies the store lost.

        The two halves of retrieval can disagree, and when they do the pipeline
        does not notice: vector search selects atoms from the index, then the
        induced subgraph is built over the *catalog*, so an index that still
        holds a vocabulary the triple store no longer serves yields perfectly
        healthy retrieval metrics -- the expected seed IRIs, the expected atom
        count -- and an empty graph. Every unit then renders against an empty
        ontology chapter, falls back on generic vocabulary, and passes a
        conformance check that has no node left to constrain.

        Only the unambiguous case is fatal here: an index with content beside
        a catalog with none. Extra indexed IRIs alongside a populated catalog
        are ordinary staleness (that is what orphan pruning is for) and only
        warn. The mirror case -- an empty index, which a wipe guarantees and
        which this check therefore cannot see -- belongs to
        :meth:`_check_catalog_ready`, after materialization has had its chance
        to refill it.

        This does not consult ``ONTOLOGY_CONTEXT_REQUIRED``. That setting says
        whether a run wants a catalog; this says the two halves of retrieval
        disagree about which ontologies exist, which no configuration asks for.
        A run that deliberately extracts without a catalog does not thereby ask
        for a stale index to select atoms nothing can expand. Nor does it fire
        during a legitimate bootstrap, where index and catalog are both empty.

        Args:
            synchronized_ontologies: What the catalog actually served.

        Raises:
            EmptyOntologyContextError: The catalog is empty and the index is
                not.
        """
        # Deferred: retrieval_capabilities imports ToolBox, so a module-level
        # import here would be circular.
        from ontocast.onto.retrieval_capabilities import EmptyOntologyContextError

        if not self.is_vector_store_ready() or self.vector_store is None:
            return
        try:
            indexed = await asyncio.to_thread(
                self.vector_store.list_indexed_ontology_iris
            )
        except Exception as exc:  # noqa: BLE001 - diagnostics must not break startup
            logger.warning("Could not inspect the vector index: %s", exc)
            return
        if not indexed:
            return
        catalog = {o.iri for o in synchronized_ontologies if o.iri}
        if not catalog:
            message = (
                "The ontology catalog is empty but the vector index holds "
                f"{len(indexed)} ontology IRI(s): {sorted(indexed)}. Retrieval "
                "would select atoms the triple store cannot expand, so every "
                "unit would render against an empty ontology and the "
                "conformance gate would report a vacuous pass. Load the "
                "ontologies into the triple store, or wipe the "
                "index (VECTOR_STORE_WIPE_ON_INIT / --wipe-vector-store)."
            )
            raise EmptyOntologyContextError(message)
        orphans = indexed - catalog
        if orphans:
            logger.warning(
                "Vector index holds %d ontology IRI(s) absent from the "
                "catalog: %s. Retrieval can select atoms the induced subgraph "
                "cannot expand. Enable orphan pruning or reindex.",
                len(orphans),
                sorted(orphans),
            )

    def _catalog_sources_description(self) -> str:
        """Name the places a catalog could have come from, for an error message."""
        directory = self.config.tool_config.path_config.ontology_directory
        sources = [
            f"ontology_directory={directory}"
            if directory is not None
            else "ontology_directory is unset"
        ]
        triple = self.triple_store_manager
        if isinstance(triple, FusekiTripleStoreManager):
            sources.append(f"triple-store dataset={triple.ontologies_dataset!r}")
        elif triple is not None:
            sources.append(f"triple store={type(triple).__name__}")
        return "; ".join(sources)

    async def _check_catalog_ready(
        self, synchronized_ontologies: list[Ontology], *, required: bool
    ) -> None:
        """Refuse to run against a catalog nothing can be extracted from.

        Two failures reach this point looking like success, because the wipe in
        :meth:`initialize` is unconditional while the refill that follows it is
        not:

        1. The sync found no ontologies at all, so materialization had nothing
           to write and the index a wipe just emptied stays empty.
        2. The sync found ontologies but every one of them indexed to nothing,
           so retrieval has a catalog it can expand and no atoms to select.

        Either way every content unit resolves an empty ontology context. Under
        ``ONTOLOGY_CONTEXT_REQUIRED`` that is fatal per unit anyway -- failing
        here instead names the subsystem, costs no provider calls, and happens
        before a document is converted.

        The two are gated differently. An empty catalog may be exactly what
        the operator meant -- a server filled later over HTTP, or a render mode
        that builds its own vocabulary -- so it is fatal only where the caller
        says it cannot be. An empty index over a *populated* catalog is never
        meant: it says materialization ran and produced nothing, which no
        configuration asks for, so it does not consult
        ``ONTOLOGY_CONTEXT_REQUIRED`` either.

        Args:
            synchronized_ontologies: What the sync resolved.
            required: Whether an empty catalog is fatal. False for a server that
                legitimately starts empty and is filled through ``POST
                /ontologies``, and for any run whose render mode creates
                ontologies; True for a facts-only batch run, which can only
                produce an ungrounded graph without one.

        Raises:
            EmptyOntologyContextError: The catalog is empty and this entry
                point requires one, or the index is empty over a populated
                catalog.
        """
        # Deferred: retrieval_capabilities imports ToolBox, so a module-level
        # import here would be circular.
        from ontocast.onto.retrieval_capabilities import EmptyOntologyContextError

        fatal = required and self.config.server.ontology_context_required

        def _report(message: str) -> None:
            if fatal:
                raise EmptyOntologyContextError(message)
            logger.warning(message)

        if not synchronized_ontologies:
            if self.config.server.render_mode != RenderMode.FACTS:
                # The starting point of an ontology-rendering run, not a fault:
                # each unit with no context is sent to render_ontology_fresh,
                # which mints a catalog ontology from the text.
                logger.info(
                    "The ontology catalog is empty (%s); render_mode=%s will "
                    "create ontologies from the corpus.",
                    self._catalog_sources_description(),
                    self.config.server.render_mode.value,
                )
                return
            _report(
                "The ontology catalog resolved to zero ontologies "
                f"({self._catalog_sources_description()}) and render_mode="
                f"{self.config.server.render_mode.value} creates none. Every "
                "content unit would render against an empty ontology, fall "
                "back on generic vocabulary, and pass a conformance check with "
                "no node left to constrain. Seed the catalog, or render "
                "ontologies as well as facts."
            )
            return

        if not self.is_vector_store_ready() or self.vector_store is None:
            return
        try:
            indexed = await asyncio.to_thread(
                self.vector_store.list_indexed_ontology_iris
            )
        except Exception as exc:  # noqa: BLE001 - diagnostics must not break startup
            logger.warning("Could not inspect the vector index: %s", exc)
            return
        if not indexed:
            raise EmptyOntologyContextError(
                f"The ontology catalog holds {len(synchronized_ontologies)} "
                "ontolog(ies) but the vector index is empty after "
                "materialization, so vector retrieval has nothing to select. "
                "Check for indexing errors above, or switch "
                "ONTOLOGY_CONTEXT_MODE away from selected_vector_search_ontology."
            )

    async def initialize(
        self,
        *,
        ontology_context_mode: OntologyContextMode | None = None,
        fail_on_vector_store_error: bool = True,
        wipe_vector_store: bool | None = None,
        prune_orphan_iris: bool | None = None,
        require_populated_catalog: bool = False,
    ) -> None:
        """Initialize the toolbox with ontologies and their properties.

        This method synchronizes ontologies between filesystem and triple store,
        then fetches ontologies from the triple store and updates their properties
        using the LLM tool.

        Args:
            ontology_context_mode: When vector search mode, ensure the vector store
                is ready before materializing atoms. ``None`` uses
                ``ONTOLOGY_CONTEXT_MODE``.
            fail_on_vector_store_error: Raise on vector init failure when True.
            wipe_vector_store: Drop the current vector partition before init.
                ``None`` uses ``VECTOR_STORE_WIPE_ON_INIT`` (default False).
            prune_orphan_iris: Delete indexed IRIs absent from the sync catalog.
                ``None`` uses ``VECTOR_STORE_PRUNE_ORPHAN_IRIS_ON_INIT`` (default True).
            require_populated_catalog: Fail rather than warn when the catalog or
                its index comes out empty. A batch run sets this -- it has no
                later chance to be given ontologies, so an empty catalog can
                only yield an ungrounded graph. A server leaves it False:
                starting empty and being filled through ``POST /ontologies`` is
                a supported way to run.

        Raises:
            EmptyOntologyContextError: ``require_populated_catalog`` and
                ``ONTOLOGY_CONTEXT_REQUIRED`` are both set and the catalog or
                its index is empty once synchronization has finished.
        """
        import asyncio
        import time

        init_started = time.perf_counter()
        if ontology_context_mode is None:
            ontology_context_mode = self.config.server.ontology_context_mode
        vsc = self.config.tool_config.vector_store
        do_wipe = vsc.wipe_on_init if wipe_vector_store is None else wipe_vector_store
        do_prune = (
            vsc.prune_orphan_iris_on_init
            if prune_orphan_iris is None
            else prune_orphan_iris
        )

        if self.triple_store_manager is not None:
            await self.triple_store_manager.async_init()

        # The wipe is honoured whatever the context mode is. It used to sit
        # inside the mode-gated branch below, so under any non-vector mode a
        # destructive flag was accepted, did nothing, and said nothing -- and a
        # request that meant to clear a partition before reindexing left the
        # old vectors in place.
        vector_store = self.vector_store
        if do_wipe:
            if vector_store is None:
                logger.warning(
                    "A vector store wipe was requested but no vector store is "
                    "configured; nothing to wipe"
                )
            else:
                logger.warning(
                    "Wiping vector store partition before initialize "
                    "(wipe_vector_store=True)"
                )
                try:
                    await vector_store.wipe_store()
                except Exception as exc:
                    self.vector_store_last_error = exc
                    if fail_on_vector_store_error:
                        raise
                    logger.warning("Vector store wipe failed: %s", exc)
                finally:
                    self.vector_store_ready = False

        if self.should_initialize_vector_store(ontology_context_mode):
            if vector_store is None:
                self.vector_store_ready = False
                self.vector_store_last_error = RuntimeError(
                    "Vector store is not configured"
                )
                if fail_on_vector_store_error:
                    raise self.vector_store_last_error
                logger.warning(
                    "Vector store was requested for initialization but is not configured"
                )
            else:
                try:
                    await vector_store.initialize()
                    self.vector_store_ready = True
                    self.vector_store_last_error = None
                except Exception as exc:
                    self.vector_store_ready = False
                    self.vector_store_last_error = exc
                    if fail_on_vector_store_error:
                        raise
                    logger.warning(
                        "Vector store initialization failed; continuing without vector retrieval: %s",
                        exc,
                    )

        shapes_started = time.perf_counter()
        await self.shapes_catalog.sync(
            self.config.tool_config.facts_validation.shapes_dir
        )
        logger.info(
            "Shapes sync finished in %.2fs", time.perf_counter() - shapes_started
        )

        sync_started = time.perf_counter()
        synchronized_ontologies = await self._synchronize_ontologies()
        logger.info(
            "Ontology sync finished: %d ontolog(ies) in %.2fs",
            len(synchronized_ontologies),
            time.perf_counter() - sync_started,
        )

        if do_prune and self.is_vector_store_ready() and self.vector_store is not None:
            triple = self.triple_store_manager
            catalog_is_authoritative = (
                triple is None or triple.last_catalog_was_complete()
            )
            if not catalog_is_authoritative:
                # Pruning deletes indexed ontologies that the catalog no longer
                # mentions. A catalog that only partly loaded mentions fewer
                # ontologies than exist, so pruning against it deletes live
                # data on the strength of a network error.
                logger.warning(
                    "Skipping vector-store orphan prune: the ontology catalog "
                    "loaded incompletely, so absent IRIs are not evidence of "
                    "deletion."
                )
            else:
                keep_iris = {o.iri for o in synchronized_ontologies if o.iri}
                orphans = await asyncio.to_thread(
                    self.vector_store.prune_orphan_ontology_iris, keep_iris
                )
                if orphans:
                    logger.info(
                        "Pruned %d orphan ontology IRI(s) from vector store: %s",
                        len(orphans),
                        orphans,
                    )

        await self._check_catalog_index_agreement(synchronized_ontologies)

        for ontology in synchronized_ontologies:
            self.ontology_manager.add_ontology(ontology, skip_vector_index=True)

        concurrency = max(1, vsc.reindex_concurrency)
        semaphore = asyncio.Semaphore(concurrency)

        async def _materialize_one(ontology: Ontology) -> None:
            async with semaphore:
                onto_started = time.perf_counter()
                indexed = await self._materialize_ontology(ontology)
                logger.info(
                    "Materialized ontology %s in %.2fs (%d atom(s) indexed)",
                    ontology.iri,
                    time.perf_counter() - onto_started,
                    indexed,
                )

        materialize_started = time.perf_counter()
        await asyncio.gather(
            asyncio.gather(*[_materialize_one(o) for o in synchronized_ontologies]),
            update_ontology_manager(om=self.ontology_manager, llm_tool=self.llm),
        )
        logger.info(
            "Ontology materialize + enrich finished for %d ontolog(ies) in %.2fs "
            "(reindex_concurrency=%d); initialize total %.2fs",
            len(synchronized_ontologies),
            time.perf_counter() - materialize_started,
            concurrency,
            time.perf_counter() - init_started,
        )

        # After materialization, not before: the wipe above is unconditional and
        # the refill is not, so an empty index is only evidence of a fault once
        # the reindex has had its turn.
        await self._check_catalog_ready(
            synchronized_ontologies, required=require_populated_catalog
        )

    def _load_seed_ontologies_from_directory(self) -> list[Ontology]:
        """Load seed ontologies from ``ontology_directory`` (*.ttl).

        Every way this returns nothing says so at INFO or louder. A silent
        empty list here is indistinguishable downstream from a healthy run
        against a catalog that happens to live in the triple store, and it is
        the difference between a five-minute diagnosis and a long one.
        """
        ontology_dir = self.config.tool_config.path_config.ontology_directory
        if ontology_dir is None:
            logger.info(
                "No ontology_directory configured; seed ontologies come from "
                "the triple store only"
            )
            return []
        directory = pathlib.Path(ontology_dir).expanduser()
        if not directory.is_dir():
            logger.warning(
                "Configured ontology_directory %s is not a directory (resolved "
                "from %s); no seed ontologies will be loaded",
                directory.absolute(),
                ontology_dir,
            )
            return []
        ontologies: list[Ontology] = []
        paths = sorted(directory.glob("*.ttl"))
        if not paths:
            logger.warning(
                "No *.ttl files in ontology_directory %s", directory.absolute()
            )
        for path in paths:
            try:
                ontologies.append(Ontology.from_file(path))
                logger.debug("Loaded seed ontology from %s", path)
            except Exception as exc:
                logger.error("Failed to load seed ontology %s: %s", path, exc)
        return ontologies

    @staticmethod
    def _defines_terms(ontology: Ontology) -> bool:
        """Whether an ontology carries terms rather than only its own header.

        The header is not the ontology. A catalog read builds one of these from
        the ``owl:Ontology`` subject and fills in the graph separately, so an
        ontology whose graph never arrived still answers with three triples
        about itself -- a non-empty graph that defines nothing, and that no
        amount of retrieval can expand into context.
        """
        header_subjects = {
            subject
            for subject, _, _ in ontology.graph.triples((None, RDF.type, OWL.Ontology))
        }
        return any(subject not in header_subjects for subject, _, _ in ontology.graph)

    @classmethod
    def _seed_ontologies_missing_from(
        cls, served: list[Ontology], seeds: list[Ontology]
    ) -> list[Ontology]:
        """Seeds the partition serves no usable graph for.

        The test is the graph, not the IRI. A store that *lists* an ontology but
        serves no terms for it -- a lost dataset, a partition that was never
        populated, a graph dropped out from under a still-registered header --
        used to be indistinguishable from a healthy one, so the seed that could
        have repaired it was skipped and the catalog stayed empty.

        A served ontology that does define terms always wins. Ontology-mode runs
        write evolved terminals back to the store, and re-materializing an older
        seed over one of those would silently revert the catalog to whatever is
        on disk.
        """
        served_by_iri = {ontology.iri: ontology for ontology in served if ontology.iri}
        missing: list[Ontology] = []
        for seed in seeds:
            existing = served_by_iri.get(seed.iri)
            if existing is not None and cls._defines_terms(existing):
                continue
            missing.append(seed)
        return missing

    async def _synchronize_ontologies(self) -> list[Ontology]:
        """Resolve the catalog from the triple store, repaired from disk.

        Records the seeds it had to supply on ``_last_seed_repairs`` so a caller
        that is not :meth:`initialize` -- which materializes everything anyway --
        can write back exactly those.
        """
        import asyncio

        seed_ontologies = await asyncio.to_thread(
            self._load_seed_ontologies_from_directory
        )
        logger.info(
            "Found %d seed ontolog(ies) in ontology_directory",
            len(seed_ontologies),
        )

        triple_store_ontologies: list[Ontology] = []
        if self.triple_store_manager is not None:
            triple_store_ontologies = (
                await self.triple_store_manager.afetch_ontologies()
            )
            logger.info(
                "Found %d ontologies in triple store", len(triple_store_ontologies)
            )

        repairs = self._seed_ontologies_missing_from(
            triple_store_ontologies, seed_ontologies
        )
        self._last_seed_repairs = repairs
        repaired_iris = {seed.iri for seed in repairs}
        resolved = [
            ontology
            for ontology in triple_store_ontologies
            if ontology.iri not in repaired_iris
        ]
        for seed_onto in repairs:
            logger.info(
                "Syncing seed ontology to triple store: %s (version: %s)",
                seed_onto.iri,
                seed_onto.version,
            )
            resolved.append(seed_onto)

        return resolved

    async def _materialize_ontology(self, ontology: Ontology) -> int:
        """Write ontology to the triple store and rebuild vector atoms.

        Returns:
            int: Atoms indexed. Zero when no vector store is ready, and zero
            when the ontology atomized to nothing -- which the reindex cannot
            distinguish, having already deleted the previous atoms, so it is
            reported here rather than left to look like a successful pass.
        """
        import asyncio

        if self.triple_store_manager is not None:
            await self.triple_store_manager.aserialize(ontology)

        if not (self.is_vector_store_ready() and self.vector_store is not None):
            return 0
        indexed = await asyncio.to_thread(self.vector_store.reindex_ontology, ontology)
        if not indexed:
            logger.warning(
                "Ontology %s indexed 0 atoms (%d triples). Retrieval cannot "
                "select any of its terms.",
                ontology.iri,
                len(ontology.graph),
            )
        return indexed

    async def ingest_ontology_ttl(
        self, ttl: bytes, *, filename: str | None = None
    ) -> Ontology:
        """Register Turtle in the triple store and the vector index.

        ``ontology_directory`` is a read-only seed fixture consulted once at
        init, so nothing is written there: an ingested ontology lives in the
        triple store and vector index only and does not survive a rebuild from
        seeds.
        """
        import asyncio

        graph = RDFGraph()

        def _parse() -> None:
            graph.parse(BytesIO(ttl), format="turtle")

        await asyncio.to_thread(_parse)
        o = Ontology(graph=graph)
        if not o.iri or o.iri == ONTOLOGY_NULL_IRI:
            raise ValueError("Loaded turtle does not define a valid ontology IRI")
        if not o.hash:
            raise ValueError("Ontology hash could not be computed")
        self.ontology_manager.validate_identity_uniqueness(o)

        await self._materialize_ontology(o)
        self.ontology_manager.add_ontology(o, skip_vector_index=True)
        return o

    async def delete_ontology_by_iri(self, ontology_iri: str) -> None:
        """Remove ontology from manager, vector store, and triple store.

        ``ontology_directory`` is deliberately untouched: unlinking a seed TTL
        that declares this IRI would destroy curated input the next init
        reloads from -- an irreversible edit to the user's files in response
        to a store-level delete.
        """
        import asyncio

        self.ontology_manager.remove_ontology_by_iri(ontology_iri)
        if self.vector_store is not None:
            await asyncio.to_thread(self.vector_store.delete_ontology, ontology_iri)

        if self.triple_store_manager is not None:
            await self.triple_store_manager.drop_all_ontology_graphs_for_iri(
                ontology_iri
            )

    async def ingest_shapes_ttl(
        self, ttl: bytes, *, filename: str | None = None
    ) -> str:
        """Register a SHACL shapes document in the shapes partition.

        Unlike :meth:`ingest_ontology_ttl` this neither requires an ontology
        IRI nor touches the vector index: a shapes document may be a bare set
        of ``sh:NodeShape`` declarations, and shapes are never retrieved by
        similarity. A document that *does* declare an ``owl:Ontology`` header
        is stored under that IRI, so uploading it again replaces it.

        ``shapes_dir`` is a read-only seed fixture consulted at init, so
        nothing is written there.

        Args:
            ttl: Turtle bytes of the shapes document.
            filename: Original filename, used only to name a headerless
                document when the bytes carry no ontology IRI.

        Returns:
            str: The named graph the document was stored under.

        Raises:
            ValueError: If the Turtle does not parse, or holds no triples.
        """
        import asyncio

        graph = RDFGraph()

        def _parse() -> None:
            graph.parse(BytesIO(ttl), format="turtle")

        try:
            await asyncio.to_thread(_parse)
        except Exception as error:
            raise ValueError(f"Invalid Turtle: {error}") from error
        if not len(graph):
            raise ValueError("Shapes document holds no triples")

        graph_uri = shapes_graph_uri(
            graph,
            fallback=(f"urn:shapes:{filename}" if filename else content_graph_uri(ttl)),
        )
        return await self.shapes_catalog.ingest(graph, graph_uri=graph_uri)

    async def delete_shapes_by_uri(self, graph_uri: str) -> None:
        """Remove one shapes document from the shapes partition.

        ``shapes_dir`` is deliberately untouched, for the same reason
        :meth:`delete_ontology_by_iri` leaves ``ontology_directory`` alone: a
        store-level delete must not edit the operator's files. A document that
        came from the seed directory therefore returns on the next init.
        """
        await self.shapes_catalog.delete(graph_uri)

Attributes

aggregator property writable

Facts aggregator.

atomic_tools property writable

Per-unit tool surface used by the render/critic loops.

chunker property writable

Text chunker.

config = config instance-attribute
converter property writable

Document converter.

embedding_tool property writable

Dense embedding provider.

llm property writable

Shared LLM tool.

llm_provider property writable

Configured LLM provider.

ontology_manager = OntologyManager() instance-attribute
patch_retriever = None instance-attribute
runtime property writable

Shared tenancy-independent tools.

Materialized empty on first access when it was never assigned. Several unit tests build a ToolBox with ToolBox.__new__(ToolBox) to exercise one tool without standing up the whole container, and used to set that tool as a plain attribute; this keeps that working, and reading a tool that was never set still raises AttributeError for its own name.

scope property

The partition this ToolBox is bound to, once tenancy is assigned.

search_provider property writable

Configured web-search provider, or None when disabled.

shapes_catalog = ShapesCatalog() instance-attribute
shared_cache property writable

Shared on-disk cache backing the LLM and converter tools.

sparql_tool = None instance-attribute
triple_store_manager = FusekiTripleStoreManager(uri=tool_config.fuseki.uri, auth=tool_config.fuseki.auth, dataset=tool_config.fuseki.dataset, ontologies_dataset=tool_config.fuseki.ontologies_dataset, shapes_dataset=tool_config.fuseki.shapes_dataset) instance-attribute
vector_store = None instance-attribute
vector_store_last_error = None instance-attribute
vector_store_ready = False instance-attribute

Methods:

__aenter__() async
Source code in ontocast/toolbox.py
async def __aenter__(self) -> "ToolBox":
    return self
__aexit__(*exc_info) async
Source code in ontocast/toolbox.py
async def __aexit__(self, *exc_info: object) -> None:
    await self.aclose()
__init__(config, *, llm=None, runtime=None)

Build a ToolBox bound to whatever partition config names.

Parameters:

Name Type Description Default
config Config

Fully resolved configuration.

required
llm LLMTool | None

Pre-built LLM tool, used when no runtime is supplied. LLMTool.create otherwise runs, which cannot be called from inside a running event loop -- prefer :meth:acreate there.

None
runtime ToolBoxRuntime | None

Shared tenancy-independent tools. Supplied by :class:~ontocast.registry.ToolBoxRegistry so scoped ToolBoxes do not each load an embedding model; built fresh when omitted.

None
Source code in ontocast/toolbox.py
def __init__(
    self,
    config: Config,
    *,
    llm: LLMTool | None = None,
    runtime: "ToolBoxRuntime | None" = None,
):
    """Build a ToolBox bound to whatever partition ``config`` names.

    Args:
        config: Fully resolved configuration.
        llm: Pre-built LLM tool, used when no ``runtime`` is supplied.
            ``LLMTool.create`` otherwise runs, which cannot be called from
            inside a running event loop -- prefer :meth:`acreate` there.
        runtime: Shared tenancy-independent tools. Supplied by
            :class:`~ontocast.registry.ToolBoxRegistry` so scoped ToolBoxes
            do not each load an embedding model; built fresh when omitted.
    """
    # Store the config for later use
    self.config = config

    # Get tool configuration
    tool_config = config.get_tool_config()

    # Tools that do not vary by tenant live on the runtime, so a registry of
    # scoped ToolBoxes shares one LLM client, converter and embedding model.
    if runtime is None:
        # A process-level fact, said where the process builds its LLM
        # client -- not once per tenancy scope over a shared runtime.
        warn_if_workers_exceed_inflight(config)
    self.runtime = runtime or ToolBoxRuntime(config, llm=llm)

    # Fuseki when a URI is configured (auth is optional), otherwise in-memory.
    if tool_config.fuseki.uri:
        self.triple_store_manager: TripleStoreManager = FusekiTripleStoreManager(
            uri=tool_config.fuseki.uri,
            auth=tool_config.fuseki.auth,
            dataset=tool_config.fuseki.dataset,
            ontologies_dataset=tool_config.fuseki.ontologies_dataset,
            shapes_dataset=tool_config.fuseki.shapes_dataset,
        )
    else:
        self.triple_store_manager = InMemoryTripleStoreManager()

    self.ontology_manager: OntologyManager = OntologyManager()
    self.ontology_manager.register_triple_store(self.triple_store_manager)

    # The shapes partition. Partition-scoped like the ontology catalog, and
    # reset on a tenancy switch for the same reason.
    self.shapes_catalog: ShapesCatalog = ShapesCatalog()
    self.shapes_catalog.register_triple_store(self.triple_store_manager)
    # Tenancy the in-memory catalog currently reflects; None until first set.
    self._active_tenancy: tuple[str, str] | None = None
    # Guards the tenancy retarget, which mutates ToolBox-wide state (dataset
    # names, the ontology catalog, vector-store table names) across awaits.
    # It is driven by a per-request query parameter with no concurrency cap,
    # so without this two requests for different tenants can interleave and
    # read or write each other's partition. Created lazily: __init__ may run
    # outside an event loop.
    self._tenancy_lock: asyncio.Lock | None = None
    self._tenancy_lock_loop: asyncio.AbstractEventLoop | None = None
    # Set by attach_registry() when this ToolBox fronts a multi-tenant host.
    self._registry: "ToolBoxRegistry | None" = None

    # Graph algorithms over graphs it is handed; it does not fetch. Built
    # only when something consumes it: its sole consumer is
    # OntologyPatchRetriever, which exists only alongside a vector store,
    # so constructing it unconditionally made the dependency graph claim
    # "always needed" where the truth is "needed with vector retrieval".
    self.sparql_tool: SPARQLTool | None = None

    self.vector_store: VectorStoreManager | None = None
    self.patch_retriever: OntologyPatchRetriever | None = None
    self.vector_store_ready: bool = False
    self.vector_store_last_error: Exception | None = None
    # Seed ontologies the last sync had to supply because the partition
    # served no graph for them. Read by the tenancy switch, which must
    # materialize exactly those and nothing else -- rewriting a whole
    # catalog on every switch is what scoped ToolBoxes exist to avoid.
    self._last_seed_repairs: list[Ontology] = []
    self._vector_store_init_lock = asyncio.Lock()

    # Both supported backends need the BM25 tool, so the only case that
    # skips it is having no vector store at all -- decided by the same
    # resolution the factory uses.
    needs_sparse = resolve_backend(tool_config) is not VectorStoreBackend.NONE
    vector_store = create_vector_store_manager(
        tool_config,
        embedding=self.embedding_tool,
        sparse_embedding=(
            self.runtime.sparse_embedding_tool(tool_config.embedding)
            if needs_sparse
            else None
        ),
    )
    if vector_store is not None:
        self.vector_store = vector_store
        self.sparql_tool = SPARQLTool(
            triple_store_manager=self.triple_store_manager
        )
        self.patch_retriever = OntologyPatchRetriever(
            vector_store=vector_store,
            sparql_tool=self.sparql_tool,
            patch=tool_config.patch_retrieval,
            ontology_manager=self.ontology_manager,
        )
        self.ontology_manager.register_vector_store(self.patch_retriever)
aclose() async

Release every backend connection this ToolBox opened.

The ToolBox owns an httpx client (Fuseki) and a Qdrant client, neither of which was previously closed anywhere -- FusekiTripleStoreManager even defined close() that nothing called. Long-lived hosts that build a ToolBox per tenant, and tests that build many, leaked sockets.

Also closes any scoped ToolBoxes this one spawned through :meth:for_scope, so shutting down the ToolBox an application holds releases every tenant's connections too.

Safe to call more than once, and never raises: teardown failures are logged, since a caller shutting down cannot act on them.

Source code in ontocast/toolbox.py
async def aclose(self) -> None:
    """Release every backend connection this ToolBox opened.

    The ToolBox owns an httpx client (Fuseki) and a Qdrant client, neither
    of which was previously closed anywhere -- ``FusekiTripleStoreManager``
    even defined ``close()`` that nothing called. Long-lived hosts that
    build a ToolBox per tenant, and tests that build many, leaked sockets.

    Also closes any scoped ToolBoxes this one spawned through
    :meth:`for_scope`, so shutting down the ToolBox an application holds
    releases every tenant's connections too.

    Safe to call more than once, and never raises: teardown failures are
    logged, since a caller shutting down cannot act on them.
    """
    registry = self._registry
    if registry is not None:
        self._registry = None
        try:
            await registry.aclose()
        except Exception as exc:
            logger.warning("Error closing tenancy registry: %s", exc)

    if self.triple_store_manager is not None:
        try:
            await self.triple_store_manager.close()
        except Exception as exc:
            logger.warning("Error closing triple store manager: %s", exc)

    if self.vector_store is not None:
        try:
            await asyncio.to_thread(self.vector_store.close)
        except Exception as exc:
            logger.warning("Error closing vector store: %s", exc)
acreate(config) async classmethod

Construct a ToolBox from inside a running event loop.

Equivalent to ToolBox(config), except that LLM provider setup is awaited rather than driven through :func:asyncio.run -- which is illegal in a loop and is what makes the plain constructor unusable from async code. Embedders should prefer this, and pair it with async with so backend connections are released:

async with await ToolBox.acreate(config) as tools:
    await tools.initialize()

Parameters:

Name Type Description Default
config Config

Fully resolved configuration.

required

Returns:

Type Description
ToolBox

A ready ToolBox. Call :meth:initialize to sync ontologies and

ToolBox

prepare backend schema.

Source code in ontocast/toolbox.py
@classmethod
async def acreate(cls, config: Config) -> "ToolBox":
    """Construct a ToolBox from inside a running event loop.

    Equivalent to ``ToolBox(config)``, except that LLM provider setup is
    awaited rather than driven through :func:`asyncio.run` -- which is
    illegal in a loop and is what makes the plain constructor unusable from
    async code. Embedders should prefer this, and pair it with
    ``async with`` so backend connections are released:

    ```python
    async with await ToolBox.acreate(config) as tools:
        await tools.initialize()
    ```

    Args:
        config: Fully resolved configuration.

    Returns:
        A ready ToolBox. Call :meth:`initialize` to sync ontologies and
        prepare backend schema.
    """
    warn_if_workers_exceed_inflight(config)
    runtime = await ToolBoxRuntime.acreate(config)
    return cls(config, runtime=runtime)
aserialize(state) async

Persist the document's ontologies and facts (async form).

Source code in ontocast/toolbox.py
async def aserialize(self, state: AgentState) -> None:
    """Persist the document's ontologies and facts (async form)."""
    ontologies_to_serialize = document_ontology_access(
        state
    ).serialization_targets()
    for ontology in ontologies_to_serialize:
        if ontology and ontology.hash:
            # The async variant is not optional here: the sync
            # add_ontology() refuses to reindex the vector store inside a
            # running event loop, and aserialize is by definition inside
            # one — so with vector retrieval registered, the sync call
            # raised RuntimeError on the first document that actually
            # produced an ontology version to serialize.
            await self.ontology_manager.aadd_ontology(ontology)

    if self.triple_store_manager is not None:
        for ontology in ontologies_to_serialize:
            await self.triple_store_manager.aserialize(ontology)
        if state.render_facts:
            await self.triple_store_manager.aserialize(
                state.aggregated_facts,
                graph_uri=state.graph_uri,
            )
attach_registry(registry)

Resolve other partitions through registry rather than a private one.

Only needed to share one registry across several ToolBoxes; :meth:for_scope builds its own on first use otherwise.

Source code in ontocast/toolbox.py
def attach_registry(self, registry: "ToolBoxRegistry") -> None:
    """Resolve other partitions through ``registry`` rather than a private one.

    Only needed to share one registry across several ToolBoxes;
    :meth:`for_scope` builds its own on first use otherwise.
    """
    self._registry = registry
bind_scope(scope)

Record the partition a ToolBox was built for.

For a ToolBox whose stores were configured for scope at construction (Config.for_tenancy): stores and catalogs are left alone, unlike update_tenancy, which retargets them.

Source code in ontocast/toolbox.py
def bind_scope(self, scope: TenancyScope) -> None:
    """Record the partition a ToolBox was built for.

    For a ToolBox whose stores were configured for ``scope`` at
    construction (``Config.for_tenancy``): stores and catalogs are left
    alone, unlike ``update_tenancy``, which retargets them.
    """
    self._active_tenancy = scope.key
clean_tenancy_data(tenant, project, *, include_shapes=False) async

Flush triple-store and vector-store partitions for tenant / project.

Shapes are retained unless include_shapes is set: facts and ontologies come back from a rerun, but shapes are the deployment's validation contract, and dropping them turns the SHACL gate off without an error -- the next run reports shacl_evaluated: null instead of failing.

Takes the tenancy lock: this is destructive, and a concurrent retarget would let it resolve partition names against a scope that changed mid-flight.

Parameters:

Name Type Description Default
tenant str

Tenant identifier.

required
project str

Project identifier within the tenant.

required
include_shapes bool

Also drop the shapes partition.

False
Source code in ontocast/toolbox.py
async def clean_tenancy_data(
    self, tenant: str, project: str, *, include_shapes: bool = False
) -> None:
    """Flush triple-store and vector-store partitions for ``tenant`` / ``project``.

    Shapes are retained unless ``include_shapes`` is set: facts and
    ontologies come back from a rerun, but shapes are the deployment's
    validation contract, and dropping them turns the SHACL gate off without
    an error -- the next run reports ``shacl_evaluated: null`` instead of
    failing.

    Takes the tenancy lock: this is destructive, and a concurrent retarget
    would let it resolve partition names against a scope that changed
    mid-flight.

    Args:
        tenant: Tenant identifier.
        project: Project identifier within the tenant.
        include_shapes: Also drop the shapes partition.
    """
    t, p = tenant.strip(), project.strip()
    if not t or not p:
        raise ValueError("tenant and project must be non-empty")

    async with self._get_tenancy_lock():
        triple = self.triple_store_manager
        if triple is not None:
            if not triple.supports_tenancy_partition():
                raise NotImplementedError(
                    f"Triple store {type(triple).__name__} has no tenant/project partitions"
                )
            await triple.clean_tenancy(t, p, include_shapes=include_shapes)

        vector = self.vector_store
        if vector is not None and vector.supports_tenancy_partition():
            was_ready = self.vector_store_ready
            # `clean_tenancy` drops the collections and their embedding
            # metadata. Leaving the store marked ready would point every
            # later search in this process at a collection that no longer
            # exists, so readiness is dropped first and only restored by a
            # successful recreate -- the same wipe/recreate order the
            # startup path uses. Nothing is reindexed: the flush emptied
            # the catalog this index would have been built from.
            self.vector_store_ready = False
            await vector.clean_tenancy(t, p)
            if was_ready:
                try:
                    await vector.initialize()
                    self.vector_store_ready = True
                    self.vector_store_last_error = None
                except Exception as exc:
                    self.vector_store_last_error = exc
                    logger.warning(
                        "Vector store could not be recreated after flushing "
                        "%s/%s; vector retrieval is off until it is "
                        "reinitialized: %s",
                        t,
                        p,
                        exc,
                    )

    if include_shapes and (t, p) == self._active_tenancy:
        self.shapes_catalog.reset()
delete_ontology_by_iri(ontology_iri) async

Remove ontology from manager, vector store, and triple store.

ontology_directory is deliberately untouched: unlinking a seed TTL that declares this IRI would destroy curated input the next init reloads from -- an irreversible edit to the user's files in response to a store-level delete.

Source code in ontocast/toolbox.py
async def delete_ontology_by_iri(self, ontology_iri: str) -> None:
    """Remove ontology from manager, vector store, and triple store.

    ``ontology_directory`` is deliberately untouched: unlinking a seed TTL
    that declares this IRI would destroy curated input the next init
    reloads from -- an irreversible edit to the user's files in response
    to a store-level delete.
    """
    import asyncio

    self.ontology_manager.remove_ontology_by_iri(ontology_iri)
    if self.vector_store is not None:
        await asyncio.to_thread(self.vector_store.delete_ontology, ontology_iri)

    if self.triple_store_manager is not None:
        await self.triple_store_manager.drop_all_ontology_graphs_for_iri(
            ontology_iri
        )
delete_shapes_by_uri(graph_uri) async

Remove one shapes document from the shapes partition.

shapes_dir is deliberately untouched, for the same reason :meth:delete_ontology_by_iri leaves ontology_directory alone: a store-level delete must not edit the operator's files. A document that came from the seed directory therefore returns on the next init.

Source code in ontocast/toolbox.py
async def delete_shapes_by_uri(self, graph_uri: str) -> None:
    """Remove one shapes document from the shapes partition.

    ``shapes_dir`` is deliberately untouched, for the same reason
    :meth:`delete_ontology_by_iri` leaves ``ontology_directory`` alone: a
    store-level delete must not edit the operator's files. A document that
    came from the seed directory therefore returns on the next init.
    """
    await self.shapes_catalog.delete(graph_uri)
ensure_tenancy_registry()

Return this ToolBox's registry, creating it on first use.

Built lazily so a single-tenant embedder never allocates one, and so nothing has to be wired at construction time.

Source code in ontocast/toolbox.py
def ensure_tenancy_registry(self) -> "ToolBoxRegistry":
    """Return this ToolBox's registry, creating it on first use.

    Built lazily so a single-tenant embedder never allocates one, and so
    nothing has to be wired at construction time.
    """
    if self._registry is None:
        from ontocast.registry import ToolBoxRegistry

        self._registry = ToolBoxRegistry(
            self.config,
            self.runtime,
            max_scopes=self.config.server.max_tenancy_scopes,
        )
    return self._registry
ensure_vector_store(ontology_context_mode, *, fail_on_vector_store_error) async

Prepare the vector store when ontology_context_mode needs it.

A ToolBox is initialized for the mode it started in; a later request in vector mode would otherwise find no index and be refused.

Source code in ontocast/toolbox.py
async def ensure_vector_store(
    self,
    ontology_context_mode: OntologyContextMode | None,
    *,
    fail_on_vector_store_error: bool,
) -> None:
    """Prepare the vector store when ``ontology_context_mode`` needs it.

    A ToolBox is initialized for the mode it started in; a later request in
    vector mode would otherwise find no index and be refused.
    """
    if not self.should_initialize_vector_store(ontology_context_mode):
        return
    async with self._vector_store_init_lock:
        if self.is_vector_store_ready():
            return
        logger.info("Preparing the vector store for a vector-mode request")
        await self.initialize(
            ontology_context_mode=ontology_context_mode,
            fail_on_vector_store_error=fail_on_vector_store_error,
            wipe_vector_store=False,
        )
for_scope(tenant, project, *, ontology_context_mode=None, fail_on_vector_store_error=False) async

Return a ToolBox bound to tenant / project.

Returns self when the scope already matches. Otherwise resolves through the attached registry, which shares this ToolBox's runtime, so the new scope costs a triple store and an ontology catalog rather than another embedding model.

Isolation is by construction: each scope owns a deep copy of Config. That copy matters -- vector store managers hold their config sections by reference and rewrite collection names when tenancy is applied, so scopes sharing a Config would alias each other.

Parameters:

Name Type Description Default
tenant str

Tenant identifier.

required
project str

Project identifier within the tenant.

required
ontology_context_mode OntologyContextMode | None

Mode to initialize a newly built scope for.

None
fail_on_vector_store_error bool

Raise rather than log when vector store preparation fails for a newly built scope.

False

Returns:

Type Description
ToolBox

A ToolBox bound to the requested partition.

Source code in ontocast/toolbox.py
async def for_scope(
    self,
    tenant: str,
    project: str,
    *,
    ontology_context_mode: OntologyContextMode | None = None,
    fail_on_vector_store_error: bool = False,
) -> "ToolBox":
    """Return a ToolBox bound to ``tenant`` / ``project``.

    Returns ``self`` when the scope already matches. Otherwise resolves
    through the attached registry, which shares this ToolBox's runtime, so
    the new scope costs a triple store and an ontology catalog rather than
    another embedding model.

    Isolation is by construction: each scope owns a deep copy of ``Config``.
    That copy matters -- vector store managers hold their config sections by
    reference and rewrite collection names when tenancy is applied, so
    scopes sharing a ``Config`` would alias each other.

    Args:
        tenant: Tenant identifier.
        project: Project identifier within the tenant.
        ontology_context_mode: Mode to initialize a newly built scope for.
        fail_on_vector_store_error: Raise rather than log when vector store
            preparation fails for a newly built scope.

    Returns:
        A ToolBox bound to the requested partition.
    """
    requested = TenancyScope.build(tenant, project)
    if self._active_tenancy == requested.key:
        return self
    return await self.ensure_tenancy_registry().get(
        requested,
        ontology_context_mode=ontology_context_mode,
        fail_on_vector_store_error=fail_on_vector_store_error,
    )
get_atomic_tools()

Return the minimal toolbox used by atomic render/critic paths.

The runtime's :class:AtomicToolBox is shared by every scope; what is handed out here is a shallow copy bound to this scope's ontology catalog, so the per-unit repairs can ask the whole catalog whether a term exists rather than only the unit's retrieved snapshot. The copy is per call and carries no state of its own -- the catalog-term memo lives on the catalog -- so a tool replaced on the shared instance is seen by the next unit.

Source code in ontocast/toolbox.py
def get_atomic_tools(self) -> AtomicToolBox:
    """Return the minimal toolbox used by atomic render/critic paths.

    The runtime's :class:`AtomicToolBox` is shared by every scope; what is
    handed out here is a shallow copy bound to *this* scope's ontology
    catalog, so the per-unit repairs can ask the whole catalog whether a
    term exists rather than only the unit's retrieved snapshot. The copy
    is per call and carries no state of its own -- the catalog-term memo
    lives on the catalog -- so a tool replaced on the shared instance is
    seen by the next unit.
    """
    return self.atomic_tools.scoped_to_catalog(self.ontology_manager)
get_entity_aligner(embedding_model=None, similarity_threshold=None)

Return a cached entity aligner for the given embedding settings.

Source code in ontocast/toolbox.py
def get_entity_aligner(
    self,
    embedding_model: str | None = None,
    similarity_threshold: float | None = None,
) -> EntityAligner:
    """Return a cached entity aligner for the given embedding settings."""
    tool_config = self.config.get_tool_config()
    return self.runtime.get_entity_aligner(
        embedding_model or tool_config.aggregation.embedding_model,
        similarity_threshold
        if similarity_threshold is not None
        else tool_config.aggregation.similarity_threshold,
    )
get_llm_tool(budget_tracker) async

Return the shared LLM tool, charging usage to budget_tracker.

Parameters:

Name Type Description Default
budget_tracker BudgetTracker

The budget tracker to charge for this task's calls.

required

Returns:

Name Type Description
LLMTool LLMTool

The shared LLM tool.

Source code in ontocast/toolbox.py
async def get_llm_tool(self, budget_tracker: BudgetTracker) -> LLMTool:
    """Return the shared LLM tool, charging usage to ``budget_tracker``.

    Args:
        budget_tracker: The budget tracker to charge for this task's calls.

    Returns:
        LLMTool: The shared LLM tool.
    """
    return await self.runtime.get_llm_tool(budget_tracker)
ingest_ontology_ttl(ttl, *, filename=None) async

Register Turtle in the triple store and the vector index.

ontology_directory is a read-only seed fixture consulted once at init, so nothing is written there: an ingested ontology lives in the triple store and vector index only and does not survive a rebuild from seeds.

Source code in ontocast/toolbox.py
async def ingest_ontology_ttl(
    self, ttl: bytes, *, filename: str | None = None
) -> Ontology:
    """Register Turtle in the triple store and the vector index.

    ``ontology_directory`` is a read-only seed fixture consulted once at
    init, so nothing is written there: an ingested ontology lives in the
    triple store and vector index only and does not survive a rebuild from
    seeds.
    """
    import asyncio

    graph = RDFGraph()

    def _parse() -> None:
        graph.parse(BytesIO(ttl), format="turtle")

    await asyncio.to_thread(_parse)
    o = Ontology(graph=graph)
    if not o.iri or o.iri == ONTOLOGY_NULL_IRI:
        raise ValueError("Loaded turtle does not define a valid ontology IRI")
    if not o.hash:
        raise ValueError("Ontology hash could not be computed")
    self.ontology_manager.validate_identity_uniqueness(o)

    await self._materialize_ontology(o)
    self.ontology_manager.add_ontology(o, skip_vector_index=True)
    return o
ingest_shapes_ttl(ttl, *, filename=None) async

Register a SHACL shapes document in the shapes partition.

Unlike :meth:ingest_ontology_ttl this neither requires an ontology IRI nor touches the vector index: a shapes document may be a bare set of sh:NodeShape declarations, and shapes are never retrieved by similarity. A document that does declare an owl:Ontology header is stored under that IRI, so uploading it again replaces it.

shapes_dir is a read-only seed fixture consulted at init, so nothing is written there.

Parameters:

Name Type Description Default
ttl bytes

Turtle bytes of the shapes document.

required
filename str | None

Original filename, used only to name a headerless document when the bytes carry no ontology IRI.

None

Returns:

Name Type Description
str str

The named graph the document was stored under.

Raises:

Type Description
ValueError

If the Turtle does not parse, or holds no triples.

Source code in ontocast/toolbox.py
async def ingest_shapes_ttl(
    self, ttl: bytes, *, filename: str | None = None
) -> str:
    """Register a SHACL shapes document in the shapes partition.

    Unlike :meth:`ingest_ontology_ttl` this neither requires an ontology
    IRI nor touches the vector index: a shapes document may be a bare set
    of ``sh:NodeShape`` declarations, and shapes are never retrieved by
    similarity. A document that *does* declare an ``owl:Ontology`` header
    is stored under that IRI, so uploading it again replaces it.

    ``shapes_dir`` is a read-only seed fixture consulted at init, so
    nothing is written there.

    Args:
        ttl: Turtle bytes of the shapes document.
        filename: Original filename, used only to name a headerless
            document when the bytes carry no ontology IRI.

    Returns:
        str: The named graph the document was stored under.

    Raises:
        ValueError: If the Turtle does not parse, or holds no triples.
    """
    import asyncio

    graph = RDFGraph()

    def _parse() -> None:
        graph.parse(BytesIO(ttl), format="turtle")

    try:
        await asyncio.to_thread(_parse)
    except Exception as error:
        raise ValueError(f"Invalid Turtle: {error}") from error
    if not len(graph):
        raise ValueError("Shapes document holds no triples")

    graph_uri = shapes_graph_uri(
        graph,
        fallback=(f"urn:shapes:{filename}" if filename else content_graph_uri(ttl)),
    )
    return await self.shapes_catalog.ingest(graph, graph_uri=graph_uri)
initialize(*, ontology_context_mode=None, fail_on_vector_store_error=True, wipe_vector_store=None, prune_orphan_iris=None, require_populated_catalog=False) async

Initialize the toolbox with ontologies and their properties.

This method synchronizes ontologies between filesystem and triple store, then fetches ontologies from the triple store and updates their properties using the LLM tool.

Parameters:

Name Type Description Default
ontology_context_mode OntologyContextMode | None

When vector search mode, ensure the vector store is ready before materializing atoms. None uses ONTOLOGY_CONTEXT_MODE.

None
fail_on_vector_store_error bool

Raise on vector init failure when True.

True
wipe_vector_store bool | None

Drop the current vector partition before init. None uses VECTOR_STORE_WIPE_ON_INIT (default False).

None
prune_orphan_iris bool | None

Delete indexed IRIs absent from the sync catalog. None uses VECTOR_STORE_PRUNE_ORPHAN_IRIS_ON_INIT (default True).

None
require_populated_catalog bool

Fail rather than warn when the catalog or its index comes out empty. A batch run sets this -- it has no later chance to be given ontologies, so an empty catalog can only yield an ungrounded graph. A server leaves it False: starting empty and being filled through POST /ontologies is a supported way to run.

False

Raises:

Type Description
EmptyOntologyContextError

require_populated_catalog and ONTOLOGY_CONTEXT_REQUIRED are both set and the catalog or its index is empty once synchronization has finished.

Source code in ontocast/toolbox.py
async def initialize(
    self,
    *,
    ontology_context_mode: OntologyContextMode | None = None,
    fail_on_vector_store_error: bool = True,
    wipe_vector_store: bool | None = None,
    prune_orphan_iris: bool | None = None,
    require_populated_catalog: bool = False,
) -> None:
    """Initialize the toolbox with ontologies and their properties.

    This method synchronizes ontologies between filesystem and triple store,
    then fetches ontologies from the triple store and updates their properties
    using the LLM tool.

    Args:
        ontology_context_mode: When vector search mode, ensure the vector store
            is ready before materializing atoms. ``None`` uses
            ``ONTOLOGY_CONTEXT_MODE``.
        fail_on_vector_store_error: Raise on vector init failure when True.
        wipe_vector_store: Drop the current vector partition before init.
            ``None`` uses ``VECTOR_STORE_WIPE_ON_INIT`` (default False).
        prune_orphan_iris: Delete indexed IRIs absent from the sync catalog.
            ``None`` uses ``VECTOR_STORE_PRUNE_ORPHAN_IRIS_ON_INIT`` (default True).
        require_populated_catalog: Fail rather than warn when the catalog or
            its index comes out empty. A batch run sets this -- it has no
            later chance to be given ontologies, so an empty catalog can
            only yield an ungrounded graph. A server leaves it False:
            starting empty and being filled through ``POST /ontologies`` is
            a supported way to run.

    Raises:
        EmptyOntologyContextError: ``require_populated_catalog`` and
            ``ONTOLOGY_CONTEXT_REQUIRED`` are both set and the catalog or
            its index is empty once synchronization has finished.
    """
    import asyncio
    import time

    init_started = time.perf_counter()
    if ontology_context_mode is None:
        ontology_context_mode = self.config.server.ontology_context_mode
    vsc = self.config.tool_config.vector_store
    do_wipe = vsc.wipe_on_init if wipe_vector_store is None else wipe_vector_store
    do_prune = (
        vsc.prune_orphan_iris_on_init
        if prune_orphan_iris is None
        else prune_orphan_iris
    )

    if self.triple_store_manager is not None:
        await self.triple_store_manager.async_init()

    # The wipe is honoured whatever the context mode is. It used to sit
    # inside the mode-gated branch below, so under any non-vector mode a
    # destructive flag was accepted, did nothing, and said nothing -- and a
    # request that meant to clear a partition before reindexing left the
    # old vectors in place.
    vector_store = self.vector_store
    if do_wipe:
        if vector_store is None:
            logger.warning(
                "A vector store wipe was requested but no vector store is "
                "configured; nothing to wipe"
            )
        else:
            logger.warning(
                "Wiping vector store partition before initialize "
                "(wipe_vector_store=True)"
            )
            try:
                await vector_store.wipe_store()
            except Exception as exc:
                self.vector_store_last_error = exc
                if fail_on_vector_store_error:
                    raise
                logger.warning("Vector store wipe failed: %s", exc)
            finally:
                self.vector_store_ready = False

    if self.should_initialize_vector_store(ontology_context_mode):
        if vector_store is None:
            self.vector_store_ready = False
            self.vector_store_last_error = RuntimeError(
                "Vector store is not configured"
            )
            if fail_on_vector_store_error:
                raise self.vector_store_last_error
            logger.warning(
                "Vector store was requested for initialization but is not configured"
            )
        else:
            try:
                await vector_store.initialize()
                self.vector_store_ready = True
                self.vector_store_last_error = None
            except Exception as exc:
                self.vector_store_ready = False
                self.vector_store_last_error = exc
                if fail_on_vector_store_error:
                    raise
                logger.warning(
                    "Vector store initialization failed; continuing without vector retrieval: %s",
                    exc,
                )

    shapes_started = time.perf_counter()
    await self.shapes_catalog.sync(
        self.config.tool_config.facts_validation.shapes_dir
    )
    logger.info(
        "Shapes sync finished in %.2fs", time.perf_counter() - shapes_started
    )

    sync_started = time.perf_counter()
    synchronized_ontologies = await self._synchronize_ontologies()
    logger.info(
        "Ontology sync finished: %d ontolog(ies) in %.2fs",
        len(synchronized_ontologies),
        time.perf_counter() - sync_started,
    )

    if do_prune and self.is_vector_store_ready() and self.vector_store is not None:
        triple = self.triple_store_manager
        catalog_is_authoritative = (
            triple is None or triple.last_catalog_was_complete()
        )
        if not catalog_is_authoritative:
            # Pruning deletes indexed ontologies that the catalog no longer
            # mentions. A catalog that only partly loaded mentions fewer
            # ontologies than exist, so pruning against it deletes live
            # data on the strength of a network error.
            logger.warning(
                "Skipping vector-store orphan prune: the ontology catalog "
                "loaded incompletely, so absent IRIs are not evidence of "
                "deletion."
            )
        else:
            keep_iris = {o.iri for o in synchronized_ontologies if o.iri}
            orphans = await asyncio.to_thread(
                self.vector_store.prune_orphan_ontology_iris, keep_iris
            )
            if orphans:
                logger.info(
                    "Pruned %d orphan ontology IRI(s) from vector store: %s",
                    len(orphans),
                    orphans,
                )

    await self._check_catalog_index_agreement(synchronized_ontologies)

    for ontology in synchronized_ontologies:
        self.ontology_manager.add_ontology(ontology, skip_vector_index=True)

    concurrency = max(1, vsc.reindex_concurrency)
    semaphore = asyncio.Semaphore(concurrency)

    async def _materialize_one(ontology: Ontology) -> None:
        async with semaphore:
            onto_started = time.perf_counter()
            indexed = await self._materialize_ontology(ontology)
            logger.info(
                "Materialized ontology %s in %.2fs (%d atom(s) indexed)",
                ontology.iri,
                time.perf_counter() - onto_started,
                indexed,
            )

    materialize_started = time.perf_counter()
    await asyncio.gather(
        asyncio.gather(*[_materialize_one(o) for o in synchronized_ontologies]),
        update_ontology_manager(om=self.ontology_manager, llm_tool=self.llm),
    )
    logger.info(
        "Ontology materialize + enrich finished for %d ontolog(ies) in %.2fs "
        "(reindex_concurrency=%d); initialize total %.2fs",
        len(synchronized_ontologies),
        time.perf_counter() - materialize_started,
        concurrency,
        time.perf_counter() - init_started,
    )

    # After materialization, not before: the wipe above is unconditional and
    # the refill is not, so an empty index is only evidence of a fault once
    # the reindex has had its turn.
    await self._check_catalog_ready(
        synchronized_ontologies, required=require_populated_catalog
    )
is_vector_store_ready()
Source code in ontocast/toolbox.py
def is_vector_store_ready(self) -> bool:
    return self.vector_store is not None and self.vector_store_ready
require_patch_retriever()

Return the ontology patch retriever or raise a directive error.

Source code in ontocast/toolbox.py
def require_patch_retriever(self) -> OntologyPatchRetriever:
    """Return the ontology patch retriever or raise a directive error."""
    if self.patch_retriever is None:
        raise RuntimeError(
            "Ontology patch retrieval needs a vector store. Set QDRANT_URI "
            "or LANCEDB_ENABLED=true."
        )
    return self.patch_retriever
require_triple_store_manager()

Return the configured triple store manager or raise a clear error.

Source code in ontocast/toolbox.py
def require_triple_store_manager(self) -> TripleStoreManager:
    """Return the configured triple store manager or raise a clear error."""
    manager = self.triple_store_manager
    if manager is None:
        raise RuntimeError("Triple store backend is not configured")
    return manager
require_vector_store()

Return the configured vector store or raise a directive error.

Source code in ontocast/toolbox.py
def require_vector_store(self) -> VectorStoreManager:
    """Return the configured vector store or raise a directive error."""
    if self.vector_store is None:
        raise RuntimeError(
            "No vector store is configured. Set QDRANT_URI (Qdrant server) "
            "or LANCEDB_ENABLED=true (embedded LanceDB); each needs its "
            "matching extra, ontocast[qdrant] or ontocast[lancedb]."
        )
    return self.vector_store
serialize(state)

Persist the document's ontologies and facts.

Drives :meth:aserialize under a single :func:asyncio.run, so a document with N ontologies costs one event loop and one backend connection rather than N of each -- the per-call sync entry points open (and tear down) a fresh HTTP client every time.

Raises:

Type Description
RuntimeError

If called from inside a running event loop; await :meth:aserialize there instead.

Source code in ontocast/toolbox.py
def serialize(self, state: AgentState) -> None:
    """Persist the document's ontologies and facts.

    Drives :meth:`aserialize` under a single :func:`asyncio.run`, so a
    document with N ontologies costs one event loop and one backend
    connection rather than N of each -- the per-call sync entry points open
    (and tear down) a fresh HTTP client every time.

    Raises:
        RuntimeError: If called from inside a running event loop; await
            :meth:`aserialize` there instead.
    """
    require_no_running_loop("ToolBox.serialize", "ToolBox.aserialize")
    asyncio.run(self.aserialize(state))
shapes_chapter_for_context(context_terms)

Per-unit conformance chapter, joined on the unit's context IRIs.

Source code in ontocast/toolbox.py
def shapes_chapter_for_context(self, context_terms: set[str]) -> str:
    """Per-unit conformance chapter, joined on the unit's context IRIs."""
    cfg = self.config.tool_config.facts_validation
    return self.shapes_catalog.selected_chapter(
        context_terms, max_lines=cfg.shapes_prompt_max_lines
    )
shapes_prompt_contract()

Conformance chapter, term exemptions, and the selection flag.

("", (), False) when the contract is off or the shapes partition is empty -- the prompt is then byte-identical to a shape-less deployment's. Full-catalog modes return the memoized whole chapter with the flag False. When per-unit selection is in force (mode context, or auto with a catalog that outgrew the line cap) the chapter comes back empty with the flag True, and the unit loop fills it via :meth:shapes_chapter_for_context once the unit's ontology snapshot exists. Exemption terms are the full catalog's in every mode -- the gate validates against every shape, and the terms are catalog IRIs either way.

Source code in ontocast/toolbox.py
def shapes_prompt_contract(self) -> tuple[str, tuple[str, ...], bool]:
    """Conformance chapter, term exemptions, and the selection flag.

    ``("", (), False)`` when the contract is off or the shapes partition
    is empty -- the prompt is then byte-identical to a shape-less
    deployment's. Full-catalog modes return the memoized whole chapter
    with the flag False. When per-unit selection is in force (mode
    ``context``, or ``auto`` with a catalog that outgrew the line cap)
    the chapter comes back empty with the flag True, and the unit loop
    fills it via :meth:`shapes_chapter_for_context` once the unit's
    ontology snapshot exists. Exemption terms are the full catalog's in
    every mode -- the gate validates against every shape, and the terms
    are catalog IRIs either way.
    """
    cfg = self.config.tool_config.facts_validation
    mode = cfg.shapes_prompt_contract
    if mode == "off":
        return "", (), False
    max_lines = cfg.shapes_prompt_max_lines
    terms = self.shapes_catalog.prompt_contract_terms(max_lines=max_lines)
    select = mode == "context" or (
        mode == "auto" and self.shapes_catalog.needs_selection(max_lines=max_lines)
    )
    if select:
        return "", terms, True
    return (
        self.shapes_catalog.conformance_chapter(max_lines=max_lines),
        terms,
        False,
    )
should_initialize_vector_store(ontology_context_mode)
Source code in ontocast/toolbox.py
def should_initialize_vector_store(
    self, ontology_context_mode: OntologyContextMode | None
) -> bool:
    return (
        self.vector_store is not None
        and ontology_context_mode
        == OntologyContextMode.SELECTED_VECTOR_SEARCH_ONTOLOGY
    )
update_tenancy(tenant, project) async

Retarget Fuseki datasets and Qdrant collections for tenant / project.

Source code in ontocast/toolbox.py
async def update_tenancy(self, tenant: str, project: str) -> None:
    """Retarget Fuseki datasets and Qdrant collections for ``tenant`` / ``project``."""
    await self.update_tenancy_with_vector_mode(
        tenant,
        project,
        initialize_vector_store=True,
        fail_on_vector_store_error=True,
    )
update_tenancy_with_vector_mode(tenant, project, *, initialize_vector_store, fail_on_vector_store_error) async

Retarget tenancy and optionally initialize vector store collections.

Serialized: the body mutates ToolBox-wide state across awaits, and the HTTP layer calls this per request from a ?tenant= query parameter with no concurrency cap. Interleaving two switches leaves the catalog and the store handles describing different tenants.

Source code in ontocast/toolbox.py
async def update_tenancy_with_vector_mode(
    self,
    tenant: str,
    project: str,
    *,
    initialize_vector_store: bool,
    fail_on_vector_store_error: bool,
) -> None:
    """Retarget tenancy and optionally initialize vector store collections.

    Serialized: the body mutates ToolBox-wide state across awaits, and the
    HTTP layer calls this per request from a ``?tenant=`` query parameter
    with no concurrency cap. Interleaving two switches leaves the catalog
    and the store handles describing different tenants.
    """
    async with self._get_tenancy_lock():
        await self._update_tenancy_with_vector_mode_locked(
            tenant,
            project,
            initialize_vector_store=initialize_vector_store,
            fail_on_vector_store_error=fail_on_vector_store_error,
        )

Functions:

render_ontology_summary(ontology, llm_tool) async

Generate a summary of ontology properties using LLM analysis.

This function uses the LLM tool to analyze an RDF graph and generate a structured summary of its properties. Only unset fields are requested.

Parameters:

Name Type Description Default
ontology Ontology

The ontology to analyze (for checking which fields are set).

required
llm_tool LLMTool

The LLM tool instance for analysis.

required

Returns:

Name Type Description
OntologyProperties OntologyProperties

A structured summary containing only the missing properties.

Source code in ontocast/toolbox.py
async def render_ontology_summary(
    ontology: Ontology, llm_tool: LLMTool
) -> OntologyProperties:
    """Generate a summary of ontology properties using LLM analysis.

    This function uses the LLM tool to analyze an RDF graph and generate
    a structured summary of its properties. Only unset fields are requested.

    Args:
        ontology: The ontology to analyze (for checking which fields are set).
        llm_tool: The LLM tool instance for analysis.

    Returns:
        OntologyProperties: A structured summary containing only the missing properties.
    """
    from typing import Any, cast

    from pydantic import create_model

    # Sample the graph intelligently (first 100 sections)
    # This provides context without overwhelming the LLM
    sampled_graph = sample_ontology_graph(ontology.graph, max_triples=100)
    # Serialize with consistent ordering to ensure determinism
    ontology_str = sampled_graph.serialize()

    # Determine which fields are unset and need LLM inference
    unset_fields = {}
    fields_to_fetch = []

    # Fields we want to potentially fetch from LLM (excluding internal fields like created_at)
    fields_to_check = ["title", "description", "ontology_id", "version", "iri"]

    # For Ontology objects, only fetch fields that are unset
    for field in fields_to_check:
        value = getattr(ontology, field, None)
        if value is None or (field == "iri" and value == ONTOLOGY_NULL_IRI):
            fields_to_fetch.append(field)
            # Get the field definition from the base model
            base_field = OntologyProperties.model_fields[field]
            unset_fields[field] = (base_field.annotation, base_field)

    if not unset_fields:
        # All fields are already set, return empty props
        return OntologyProperties()

    # Create a dynamic model with only unset fields
    DynamicProps = create_model("DynamicOntologyProps", **cast(Any, unset_fields))

    # Define the output parser
    parser = PydanticOutputParser(pydantic_object=DynamicProps)

    # Create the prompt template with format instructions
    field_list_str = "\n- ".join(fields_to_fetch)
    format_instructions = parser.get_format_instructions()

    # Build the template - use format_instructions as a separate variable to avoid brace conflicts
    template = (
        "Below is a sample of an ontology in Turtle format:\n\n"
        "```ttl\n{ontology_str}\n```\n\n"
        "Extract ONLY the following properties that are missing:\n"
        f"- {field_list_str}\n\n"
        "{format_instructions}"
    )

    prompt = PromptTemplate(
        template=template,
        input_variables=["ontology_str"],
        partial_variables={"format_instructions": format_instructions},
    )

    response = await llm_tool(prompt.format_prompt(ontology_str=ontology_str))
    dynamic_props = parser.parse(response.content)

    # Convert dynamic props to OntologyProperties
    result = OntologyProperties()
    for field in unset_fields.keys():
        value = getattr(dynamic_props, field, None)
        if value is not None:
            setattr(result, field, value)

    return result

sample_ontology_graph(graph, max_triples=100)

Sample an ontology graph to provide a representative subset.

This function serializes the graph to Turtle format and takes the first N blank-line separated sections. This is deterministic and simpler than complex triple selection logic.

Parameters:

Name Type Description Default
graph RDFGraph

The full ontology graph

required
max_triples int

Maximum number of sections to include in the sample

100

Returns:

Name Type Description
RDFGraph RDFGraph

A sampled version of the ontology with representative triples

Source code in ontocast/toolbox.py
def sample_ontology_graph(graph: RDFGraph, max_triples: int = 100) -> RDFGraph:
    """Sample an ontology graph to provide a representative subset.

    This function serializes the graph to Turtle format and takes the first
    N blank-line separated sections. This is deterministic and simpler than
    complex triple selection logic.

    Args:
        graph: The full ontology graph
        max_triples: Maximum number of sections to include in the sample

    Returns:
        RDFGraph: A sampled version of the ontology with representative triples
    """
    # Serialize to turtle
    turtle_str = graph.serialize_canonical_turtle()

    # Split on blank lines (typical turtle format uses \n\n to separate blocks)
    sections = turtle_str.split("\n\n")

    # Take first max_triples sections (or fewer if graph is smaller)
    num_sections = min(len(sections), max_triples)
    sampled_turtle = "\n\n".join(sections[:num_sections])

    # Parse back into a graph
    sampled = RDFGraph()
    sampled.parse(data=sampled_turtle, format="turtle")

    # Copy namespace bindings from original graph
    for prefix, namespace in graph.namespaces():
        if prefix:
            sampled.bind(prefix, namespace)

    return sampled

update_ontology_manager(om, llm_tool, *, max_concurrency=None) async

Update properties for all ontologies in the manager.

Ontologies that already have title, ontology_id, and description are skipped. Remaining LLM calls run concurrently up to max_concurrency (defaults to the LLM tool's llm_max_inflight).

Parameters:

Name Type Description Default
om OntologyManager

The ontology manager containing ontologies to update.

required
llm_tool LLMTool

The LLM tool instance for analysis.

required
max_concurrency int | None

Optional override for parallel LLM enrich calls.

None
Source code in ontocast/toolbox.py
async def update_ontology_manager(
    om: OntologyManager,
    llm_tool: LLMTool,
    *,
    max_concurrency: int | None = None,
):
    """Update properties for all ontologies in the manager.

    Ontologies that already have title, ontology_id, and description are skipped.
    Remaining LLM calls run concurrently up to ``max_concurrency`` (defaults to
    the LLM tool's ``llm_max_inflight``).

    Args:
        om: The ontology manager containing ontologies to update.
        llm_tool: The LLM tool instance for analysis.
        max_concurrency: Optional override for parallel LLM enrich calls.
    """
    import asyncio
    import time

    pending = [
        o
        for o in om.ontologies
        if (o.title is None) or (o.ontology_id is None) or (o.description is None)
    ]
    if not pending:
        return

    limit = max_concurrency
    if limit is None:
        limit = max(1, getattr(llm_tool.config, "llm_max_inflight", 1))
    semaphore = asyncio.Semaphore(max(1, limit))

    async def _one(ontology: Ontology) -> None:
        async with semaphore:
            await update_ontology_properties(ontology, llm_tool)

    started = time.perf_counter()
    await asyncio.gather(*[_one(o) for o in pending])
    logger.info(
        "Ontology property enrich finished for %d ontolog(ies) in %.2fs",
        len(pending),
        time.perf_counter() - started,
    )

update_ontology_properties(o, llm_tool) async

Update ontology properties using LLM analysis, only if missing.

This function uses the LLM tool to analyze and update the properties of a given ontology based on its graph content, but only if any key property is missing or empty.

Source code in ontocast/toolbox.py
async def update_ontology_properties(o: Ontology, llm_tool: LLMTool):
    """Update ontology properties using LLM analysis, only if missing.

    This function uses the LLM tool to analyze and update the properties
    of a given ontology based on its graph content, but only if any key
    property is missing or empty.
    """
    # Only update if any key property is missing or empty
    if (o.title is None) or (o.ontology_id is None) or (o.description is None):
        props = await render_ontology_summary(o, llm_tool)
        o.set_properties(**props.model_dump())

warn_if_workers_exceed_inflight(config)

Warn when PARALLEL_WORKERS cannot all reach the provider at once.

A unit worker never issues two LLM calls concurrently, so the provider concurrency one document can reach is min(PARALLEL_WORKERS, LLM_MAX_INFLIGHT). Workers past the in-flight cap do not run: they queue behind it -- the wait lands in llm/inflight_wait -- while each still holds a unit slot and its memory. Said once, where the LLM client is built, rather than at every call that queues.

Parameters:

Name Type Description Default
config Config

Fully resolved configuration.

required

Returns:

Type Description
bool

Whether a warning was emitted.

Source code in ontocast/toolbox.py
def warn_if_workers_exceed_inflight(config: Config) -> bool:
    """Warn when PARALLEL_WORKERS cannot all reach the provider at once.

    A unit worker never issues two LLM calls concurrently, so the provider
    concurrency one document can reach is ``min(PARALLEL_WORKERS,
    LLM_MAX_INFLIGHT)``. Workers past the in-flight cap do not run: they queue
    behind it -- the wait lands in ``llm/inflight_wait`` -- while each still
    holds a unit slot and its memory. Said once, where the LLM client is built,
    rather than at every call that queues.

    Args:
        config: Fully resolved configuration.

    Returns:
        Whether a warning was emitted.
    """
    workers = config.server.parallel_workers
    inflight = config.get_tool_config().llm_config.llm_max_inflight
    if workers <= inflight:
        return False
    logger.warning(
        "PARALLEL_WORKERS=%d exceeds LLM_MAX_INFLIGHT=%d: at most %d units of a "
        "document reach the provider at once and the rest queue "
        "(budget.node_durations['llm/inflight_wait']). Lower PARALLEL_WORKERS, "
        "or raise LLM_MAX_INFLIGHT to what the provider tier allows.",
        workers,
        inflight,
        inflight,
    )
    return True