Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
123 changes: 120 additions & 3 deletions clio-agentic-search/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ Part of [**CLIO Kit**](https://github.com/iowarp/clio-kit) — the IoWarp platfo

---

Hybrid retrieval engine for scientific computing corpora. Indexes documents into namespace-specific backends and supports lexical (BM25), vector, graph, metadata, and scientific-operator retrieval in one pipeline. DuckDB storage, FastAPI server, async job queue, OpenTelemetry tracing, Prometheus metrics.
Agentic hybrid retrieval engine for scientific computing corpora. Indexes documents into namespace-specific backends and supports lexical (BM25), vector, graph, metadata, and scientific-operator retrieval in one pipeline, with an optional multi-hop agentic loop that rewrites queries and adapts to each corpus. DuckDB storage, FastAPI server, async job queue, OpenTelemetry tracing, Prometheus metrics.

## Quick start

Expand All @@ -33,15 +33,68 @@ uv run clio query --namespace local_fs --q "pressure > 200 kPa"
uv run clio index --namespace local_fs
```

### Optional extras

The core install is lightweight; heavier or backend-specific dependencies ship as extras (`uv sync --extra <name>`):

| Extra | Pulls in | Enables |
|-------|----------|---------|
| `semantic` | sentence-transformers | Transformer embeddings (otherwise a hash embedder is used) |
| `ann` | numpy, hnswlib | Approximate nearest-neighbour vector backend (`CLIO_ANN_BACKEND=hnsw`) |
| `hdf5` | h5py | HDF5 connector (`hdf5_data` namespace) |
| `netcdf` | xarray, netCDF4 | NetCDF connector (`netcdf_data` namespace) |
| `llm` | anthropic, openai | LLM-based query rewriting (`--llm-rewrite`); without it, a rule-based fallback is used |
| `telemetry` | opentelemetry, prometheus-client | Tracing + `/metrics` exposition |
| `eval` | claude-agent-sdk, anthropic | SC26 evaluation harness |

## Features

- **Multi-namespace registry** with runtime/auth config bundles
- **Connectors**: filesystem + DuckDB (`local_fs`), S3 object store, Qdrant vector store, Neo4j graph, Redis KV log
- **Hybrid retrieval** across lexical (BM25), vector, graph and metadata branches in one pipeline
- **Scientific retrieval operators**: numeric range (`unit`, `min`, `max`), unit matching, formula targeting (normalized signatures)
- **Agentic retrieval**: optional multi-hop loop with LLM query rewriting (with a no-LLM fallback) and SI-unit variant inference
- **Corpus-adaptive strategy**: schema/metadata profiling drives per-query branch selection and content-quality filtering
- **Structured ingestion**: CSV/tabular detection and table-aware chunking alongside text
- **Nine connectors** spanning POSIX, object, vector, graph, KV and science formats — see [Connectors](#connectors)
- **Background indexing** job API with cancellation tokens and per-namespace serialized execution
- **Retry/backoff** wrappers for connect/index operations
- **Telemetry**: OpenTelemetry tracing (opt-in), Prometheus metrics at `/metrics`

## Retrieval pipeline

```
Query → Namespace registry → Retrieval coordinator → parallel branches
├── Lexical (BM25)
├── Vector (embeddings; hash or transformer)
├── Graph (BFS)
├── Metadata (schema-aware filters)
└── Scientific (SI unit conversion + formula normalization)
→ Merge + rerank → Citations + trace events
```

With `--agentic`, the coordinator runs inside an observe–decide–act loop: it
inspects results, rewrites the query (LLM or rule-based), and re-runs branches
until it converges or hits `--max-hops`. A corpus profiler inspects what
metadata each namespace actually provides and adapts branch selection and
quality filtering per query.

## Connectors

| Connector | Namespace | Default registry | Extra required |
|-----------|-----------|:---------------:|----------------|
| Filesystem + DuckDB | `local_fs` | ✅ | — |
| S3 object store | `object_s3` | ✅ | — |
| Qdrant vector store | `vector_qdrant` | ✅ | — |
| HDF5 | `hdf5_data` | ✅ | `hdf5` (h5py is also a core dep) |
| NetCDF | `netcdf_data` | ✅ | `netcdf` |
| Neo4j graph | (configurable) | — | — |
| Redis KV log | (configurable) | — | — |
| IOWarp content store | (configurable) | — | `iowarp_core` wheel |
| NDP datasets | (configurable) | — | `mcp` (for MCP-backed discovery) |

`build_default_registry()` provisions the first five namespaces; the remaining
connectors are available to register explicitly.

## API endpoints

| Method | Path | Description |
Expand All @@ -59,12 +112,76 @@ uv run clio index --namespace local_fs

| Command | Description |
|---------|-------------|
| `clio query` | Run retrieval queries against a namespace |
| `clio query` | Run retrieval queries against a namespace (add `--agentic --max-hops N` for the multi-hop loop, `--llm-rewrite` for LLM query rewriting) |
| `clio index` | Index documents into a namespace |
| `clio list` | List indexed documents |
| `clio seed` | Seed sample data for testing |
| `clio serve` | Start the FastAPI server |

Agentic retrieval is opt-in — a plain `clio query` behaves exactly as before:

```bash
# Single-shot (default)
clio query --namespace local_fs --q "pressure 200 kPa"

# Multi-hop agentic loop (max 3 hops), with LLM query rewriting
clio query --namespace local_fs --q "pressure 200 kPa" --agentic --max-hops 3 --llm-rewrite
```

## Examples

Point the filesystem connector at a folder, index it, then run the queries below.

```bash
export CLIO_LOCAL_ROOT=./docs # folder of .txt/.md/.csv files
export CLIO_STORAGE_PATH=./clio.duckdb
clio index --namespace local_fs # build the index
clio list --namespace local_fs # show indexed docs + chunk counts
```

**Scientific numeric-range** — match by real unit math, not keywords. Only
documents whose measurements fall in the range are returned:

```bash
# "pressure between 300 and 400 kPa"
clio query --namespace local_fs --q "pressure" --numeric-range "300:400:kPa"

# Same physical range expressed in Pa — finds the same 320 kPa document,
# because values are canonicalized to SI base units before matching.
clio query --namespace local_fs --q "pressure" --numeric-range "300000:400000:Pa"
```

**Formula targeting** — match normalized equation signatures:

```bash
clio query --namespace local_fs --q "newton law" --formula "F=ma"
```

**Agentic multi-hop** — the loop rewrites/expands the query between hops
(here, `kPa` is auto-expanded to its SI variants):

```bash
clio query --namespace local_fs --q "pressure 320 kPa" --agentic --max-hops 3
```

**Science-format connectors** — index HDF5 / NetCDF datasets:

```bash
CLIO_HDF5_ROOT=./h5_files clio index --namespace hdf5_data
CLIO_NETCDF_ROOT=./nc_files clio index --namespace netcdf_data # needs the `netcdf` extra
clio query --namespace hdf5_data --q "compressor pressure"
```

**HTTP API** — start the server and query over HTTP:

```bash
clio serve & # FastAPI on :8000
curl -s localhost:8000/health
curl -s -X POST localhost:8000/query \
-H "Content-Type: application/json" \
-d '{"namespace":"local_fs","query":"turbine pressure","top_k":3}'
```

## Environment variables

| Variable | Default | Description |
Expand Down
33 changes: 33 additions & 0 deletions clio-agentic-search/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,10 @@ dependencies = [
"fastapi>=0.115.0,<1.0.0",
"uvicorn>=0.30.0,<1.0.0",
"tenacity>=8.2.0,<10.0.0",
"httpx>=0.28.1",
"h5py>=3.16.0",
"aiohttp>=3.13.5",
"matplotlib>=3.10.8",
]

[project.scripts]
Expand All @@ -31,6 +35,22 @@ ann = [
"numpy>=1.26.0",
"hnswlib>=0.8.0",
]
hdf5 = ["h5py>=3.10.0"]
netcdf = [
"xarray>=2024.1.0",
"netCDF4>=1.6.0",
]
llm = ["anthropic>=0.40.0", "openai>=1.0.0"]
eval = [
# Evaluation harness for the SC26 submission: Claude Agent SDK,
# Anthropic API client. The ndp-mcp server is installed separately
# from the clio-kit repository (not on PyPI) via
# uv add /path/to/clio-kit/clio-kit-mcp-servers/ndp
"claude-agent-sdk>=0.1.56",
"anthropic>=0.87.0",
"aiohttp>=3.10.0",
"matplotlib>=3.8.0",
]
telemetry = [
"opentelemetry-api>=1.20.0",
"opentelemetry-sdk>=1.20.0",
Expand Down Expand Up @@ -75,12 +95,25 @@ python_version = "3.11"
strict = true
mypy_path = "src"
packages = ["clio_agentic_search"]
# Optional-backend imports carry `# type: ignore` that is needed when the
# typed package is installed (dev/all-extras) but unused when it is absent
# (CI/--ignore-missing-imports). Tolerate both so `mypy src/` is clean in
# either environment.
warn_unused_ignores = false

[[tool.mypy.overrides]]
module = [
"duckdb",
"sentence_transformers",
"opentelemetry.*",
"prometheus_client",
"h5py",
"xarray",
"netCDF4",
"anthropic",
"openai",
"clio_cte_core_ext",
"iowarp_core",
]
ignore_missing_imports = true

71 changes: 70 additions & 1 deletion clio-agentic-search/src/clio_agentic_search/cli/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,22 @@ def build_parser() -> argparse.ArgumentParser:
default="",
help="Formula-targeted retrieval expression.",
)
query_parser.add_argument(
"--agentic",
action="store_true",
help="Enable multi-hop agentic retrieval with query rewriting.",
)
query_parser.add_argument(
"--max-hops",
type=int,
default=3,
help="Maximum retrieval hops in agentic mode (default: 3).",
)
query_parser.add_argument(
"--llm-rewrite",
action="store_true",
help="Use LLM-based query rewriting (requires anthropic). Falls back to SI expansion.",
)

seed_parser = subparsers.add_parser("seed", help="Seed explicit demo/test records.")
seed_parser.add_argument(
Expand Down Expand Up @@ -138,6 +154,9 @@ def _run_query(
numeric_range: str,
unit_match: str,
formula: str,
agentic: bool = False,
max_hops: int = 3,
llm_rewrite: bool = False,
) -> int:
registry = build_default_registry()
filters = _parse_filters(filter_pairs)
Expand Down Expand Up @@ -175,7 +194,54 @@ def _run_query(
)

coordinator = RetrievalCoordinator()
if len(connectors) == 1:

if agentic:
from clio_agentic_search.retrieval.agentic import AgenticRetriever
from clio_agentic_search.retrieval.query_rewriter import (
FallbackQueryRewriter,
QueryRewriter,
)

if llm_rewrite:
try:
rewriter: QueryRewriter | FallbackQueryRewriter = QueryRewriter()
except RuntimeError:
print(
"Warning: anthropic not installed, falling back to SI expansion",
file=sys.stderr,
)
rewriter = FallbackQueryRewriter()
else:
rewriter = FallbackQueryRewriter()

agentic_retriever = AgenticRetriever(
coordinator=coordinator,
rewriter=rewriter,
max_hops=max_hops,
)
if len(connectors) == 1:
agentic_result = agentic_retriever.query(
connector=connectors[0],
query=text,
top_k=top_k,
metadata_filters=filters,
scientific_operators=scientific_operators,
)
else:
agentic_result = agentic_retriever.query_namespaces(
connectors=connectors,
query=text,
top_k=top_k,
metadata_filters=filters,
scientific_operators=scientific_operators,
)
result_namespaces = agentic_result.namespace.split(",")
citations = agentic_result.citations
trace = agentic_result.trace
print(
f"agentic_hops={agentic_result.total_hops},final_query={agentic_result.final_query}"
)
elif len(connectors) == 1:
result = coordinator.query(
connector=connectors[0],
query=text,
Expand Down Expand Up @@ -407,6 +473,9 @@ def main(argv: Sequence[str] | None = None) -> int:
numeric_range=args.numeric_range,
unit_match=args.unit_match,
formula=args.formula,
agentic=args.agentic,
max_hops=args.max_hops,
llm_rewrite=args.llm_rewrite,
)
except ValueError as error:
print(str(error), file=sys.stderr)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
)
from clio_agentic_search.retrieval.ann import ANNAdapter, AnnResult, build_ann_adapter
from clio_agentic_search.retrieval.capabilities import ScoredChunk
from clio_agentic_search.retrieval.corpus_profile import CorpusProfile, build_corpus_profile
from clio_agentic_search.retrieval.scientific import (
ScientificQueryOperators,
score_scientific_metadata,
Expand Down Expand Up @@ -365,7 +366,7 @@ def search_lexical(self, query: str, top_k: int) -> list[ScoredChunk]:
chunk_id=match.chunk.chunk_id,
document_id=match.chunk.document_id,
text=match.chunk.text,
lexical_score=match.overlap_count / len(query_tokens),
lexical_score=match.bm25_score,
)
for match in matches
]
Expand Down Expand Up @@ -453,6 +454,12 @@ def search_scientific(

candidate_ids: set[str] | None = None

# When a quality filter is active, push the acceptable-flag list down
# to the SQL layer so we skip bad/missing rows without loading them.
acceptable_quality: tuple[str, ...] | None = None
if operators.quality_filter is not None:
acceptable_quality = operators.quality_filter.acceptable_strings()

if operators.numeric_range is not None:
try:
canonical_min = None
Expand All @@ -469,7 +476,11 @@ def search_scientific(
except ValueError:
return []
range_chunks = self.storage.query_chunks_by_measurement_range(
self.namespace, canonical_unit, canonical_min, canonical_max
self.namespace,
canonical_unit,
canonical_min,
canonical_max,
acceptable_quality=acceptable_quality,
)
range_ids = {c.chunk_id for c in range_chunks}
candidate_ids = range_ids if candidate_ids is None else candidate_ids & range_ids
Expand Down Expand Up @@ -520,6 +531,11 @@ def build_citation(self, chunk: ScoredChunk) -> CitationRecord:
score=round(chunk.combined_score, 6),
)

def corpus_profile(self) -> CorpusProfile:
"""Return a statistical profile of the indexed namespace."""
self._ensure_connected()
return build_corpus_profile(self.storage, self.namespace)

def _build_chunks(self, *, document_id: str, text: str) -> ScientificChunkPlan:
return build_structure_aware_chunk_plan(
namespace=self.namespace,
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
"""HDF5 connector package."""

from clio_agentic_search.connectors.hdf5.connector import HDF5Connector

__all__ = ["HDF5Connector"]
Loading
Loading