diff --git a/.gitignore b/.gitignore index 60ca2b9..9134017 100644 --- a/.gitignore +++ b/.gitignore @@ -222,4 +222,8 @@ __marimo__/ data/ -opencode.json \ No newline at end of file +opencode.json + +# Hermes agent working files (per-developer) +.worktrees/ +.worktreeinclude \ No newline at end of file diff --git a/src/api/app.py b/src/api/app.py index 1fc84e7..3c9b8af 100644 --- a/src/api/app.py +++ b/src/api/app.py @@ -19,6 +19,7 @@ health, ingest, ingest_run, + ingest_status, insights, knowledge_gaps, onboarding, @@ -84,6 +85,7 @@ async def unhandled_exception_handler(request: Request, exc: Exception) -> JSONR api_router.include_router(buddy.router) api_router.include_router(ingest.router) api_router.include_router(ingest_run.router) +api_router.include_router(ingest_status.router) api_router.include_router(artifact_projects.router) api_router.include_router(title.router) api_router.include_router(vector_db.router) diff --git a/src/api/routes/ingest_status.py b/src/api/routes/ingest_status.py new file mode 100644 index 0000000..2349d9c --- /dev/null +++ b/src/api/routes/ingest_status.py @@ -0,0 +1,78 @@ +"""GET /api/v1/ingest/status -- read-only AI index state for many artifacts.""" + +from typing import Annotated + +from fastapi import APIRouter, Depends, Query + +from api.dependencies import get_ingestion_metadata_store +from api.schemas import ( + ArtifactIndexStatus, + ArtifactIngestStatusResponse, + IngestStatusResponse, + StatusArtifactId, +) +from ingestion.metadata_store import ( + ArtifactRecord, + IngestionMetadataStore, + IngestionStatus, +) + +router = APIRouter() + +# The knowledge base asks once per visible page, and a page never shows more. +MAX_STATUS_IDS = 100 + +_STATUS_BY_RECORDED: dict[IngestionStatus, ArtifactIndexStatus] = { + "completed": "indexed", + "processing": "processing", + "failed": "failed", + "deindexed": "deindexed", +} + + +def _to_item( + artifact_id: str, + record: ArtifactRecord | None, +) -> ArtifactIngestStatusResponse: + if record is None: + return ArtifactIngestStatusResponse(artifact_id=artifact_id, status="unknown") + + return ArtifactIngestStatusResponse( + artifact_id=artifact_id, + # A status this service does not know how to name is reported as + # unknown rather than guessed at: the chip must never claim more + # than the record says. + status=_STATUS_BY_RECORDED.get(record.status, "unknown"), + updated_at=record.updated_at, + chunk_count=record.chunk_count, + ) + + +@router.get( + "/ingest/status", + response_model=IngestStatusResponse, + summary="Get the AI index status of artifacts", + description=( + "Read-only. Returns one entry per distinct requested artifact id, in " + "request order: indexed, processing, failed, deindexed, or unknown " + "for an id this service holds no record of. Pass ids as a repeated " + f"query parameter, at most {MAX_STATUS_IDS} per request." + ), +) +def get_ingest_status( + metadata_store: Annotated[ + IngestionMetadataStore, + Depends(get_ingestion_metadata_store), + ], + artifact_ids: Annotated[ + list[StatusArtifactId], + Query(default_factory=list, max_length=MAX_STATUS_IDS), + ], +) -> IngestStatusResponse: + requested = list(dict.fromkeys(artifact_ids)) + records = metadata_store.get_artifacts(requested) + return IngestStatusResponse( + items=[ + _to_item(artifact_id, records.get(artifact_id)) for artifact_id in requested + ] + ) diff --git a/src/api/schemas.py b/src/api/schemas.py index 0b9b28c..ec2c4f6 100644 --- a/src/api/schemas.py +++ b/src/api/schemas.py @@ -908,6 +908,41 @@ class RunArtifactsSyncResponse(BaseModel): ) +# What the knowledge base shows as an artifact's AI index state. ``indexed`` +# is the metadata store's ``completed``, renamed for the reader: it means the +# chatbot can answer from the artifact. ``unknown`` means this service holds +# no record of the id at all. +ArtifactIndexStatus = Literal["indexed", "processing", "failed", "deindexed", "unknown"] + +# One requested id. Non-empty so that a stray ``artifact_ids=`` is a client +# error rather than a lookup that can only ever answer ``unknown``. +StatusArtifactId = Annotated[str, StringConstraints(min_length=1)] + + +class ArtifactIngestStatusResponse(BaseModel): + """One artifact's index state, as the AI service recorded it.""" + + artifact_id: str = Field(description="The requested artifact id, verbatim.") + status: ArtifactIndexStatus + updated_at: str | None = Field( + default=None, + description="ISO timestamp of the last recorded change; null if unknown.", + ) + chunk_count: int | None = Field( + default=None, + description="Chunks recorded for the artifact; null if unknown.", + ) + + +class IngestStatusResponse(BaseModel): + items: list[ArtifactIngestStatusResponse] = Field( + description=( + "One entry per distinct requested id, in the order the ids were " + "first requested. An id this service never saw is 'unknown'." + ), + ) + + # ── Connector / source enable-disable ─────────────────────────────────────── diff --git a/src/ingestion/metadata_store.py b/src/ingestion/metadata_store.py index a9dd0c8..0a26247 100644 --- a/src/ingestion/metadata_store.py +++ b/src/ingestion/metadata_store.py @@ -38,6 +38,11 @@ "labels", ) +# SQLite builds before 3.32 reject statements binding more than 999 +# parameters. Staying below that keeps batch lookups portable to any +# interpreter's bundled SQLite. +_MAX_IN_PARAMS = 900 + @dataclass(frozen=True) class ArtifactRecord: @@ -544,6 +549,35 @@ def get_artifact(self, artifact_id: str) -> ArtifactRecord | None: return self._row_to_record(row) + def get_artifacts(self, artifact_ids: Sequence[str]) -> dict[str, ArtifactRecord]: + """Look up many artifacts at once, keyed by id. + + Ids with no record are simply absent from the result, so callers can + tell "never ingested" from every recorded status. The result carries + no order; a caller that must answer in request order walks its own + list. Duplicates collapse, and the lookup is batched into as few + ``IN`` queries as SQLite's bound-parameter limit allows -- one for any + request of up to ``_MAX_IN_PARAMS`` distinct ids. + """ + normalized = tuple(dict.fromkeys(artifact_ids)) + if not normalized: + return {} + + rows: list[sqlite3.Row] = [] + with self._lock: + for start in range(0, len(normalized), _MAX_IN_PARAMS): + batch = normalized[start : start + _MAX_IN_PARAMS] + placeholders = ", ".join("?" for _ in batch) + cursor = self._connection.execute( + f"SELECT {', '.join(_COLUMNS)} FROM artifacts " + f"WHERE id IN ({placeholders})", + batch, + ) + rows.extend(cast(list[sqlite3.Row], cursor.fetchall())) + + records = (self._row_to_record(row) for row in rows) + return {record.id: record for record in records} + def list_artifacts( self, status: IngestionStatus | None = "completed", diff --git a/tests/api/test_ingest_status.py b/tests/api/test_ingest_status.py new file mode 100644 index 0000000..94acf0f --- /dev/null +++ b/tests/api/test_ingest_status.py @@ -0,0 +1,184 @@ +from collections.abc import Iterable +from pathlib import Path + +import pytest +from fastapi.testclient import TestClient + +from api.app import app +from api.dependencies import get_ingestion_metadata_store +from ingestion.metadata_store import ArtifactRecord, IngestionMetadataStore + +_URL = "/api/v1/ingest/status" +_CREATED = "2026-01-01T00:00:00+00:00" +_UPDATED = "2026-01-02T00:00:00+00:00" + + +@pytest.fixture +def metadata_store(tmp_path: Path) -> Iterable[IngestionMetadataStore]: + store = IngestionMetadataStore(path=str(tmp_path / "metadata.db")) + yield store + store.close() + + +@pytest.fixture +def client(metadata_store: IngestionMetadataStore) -> Iterable[TestClient]: + app.dependency_overrides[get_ingestion_metadata_store] = lambda: metadata_store + + yield TestClient(app) + + app.dependency_overrides.clear() + + +def _save( + store: IngestionMetadataStore, + artifact_id: str, + chunk_count: int = 2, +) -> None: + store.save_completed_artifact( + ArtifactRecord( + id=artifact_id, + filename=f"{artifact_id}.md", + content_type="text/markdown", + source_type="github", + size_bytes=10, + chunk_count=chunk_count, + status="completed", + created_at=_CREATED, + updated_at=_UPDATED, + ) + ) + + +def test_completed_artifact_is_reported_as_indexed( + client: TestClient, + metadata_store: IngestionMetadataStore, +) -> None: + _save(metadata_store, "a1", chunk_count=3) + + response = client.get(_URL, params={"artifact_ids": ["a1"]}) + + assert response.status_code == 200 + assert response.json() == { + "items": [ + { + "artifact_id": "a1", + "status": "indexed", + "updated_at": _UPDATED, + "chunk_count": 3, + } + ] + } + + +def test_unknown_id_is_reported_with_null_details(client: TestClient) -> None: + response = client.get(_URL, params={"artifact_ids": ["never-ingested"]}) + + assert response.status_code == 200 + assert response.json() == { + "items": [ + { + "artifact_id": "never-ingested", + "status": "unknown", + "updated_at": None, + "chunk_count": None, + } + ] + } + + +def test_every_recorded_status_keeps_its_name_except_completed( + client: TestClient, + metadata_store: IngestionMetadataStore, +) -> None: + for artifact_id in ("done", "broken", "gone"): + _save(metadata_store, artifact_id) + metadata_store.mark_failed("broken", "parse error", _UPDATED) + metadata_store.mark_deindexed("gone", _UPDATED) + metadata_store.save_artifact( + ArtifactRecord( + id="busy", + filename="busy.md", + content_type="text/markdown", + source_type="github", + size_bytes=10, + chunk_count=0, + status="processing", + created_at=_CREATED, + updated_at=_CREATED, + ) + ) + + response = client.get( + _URL, params={"artifact_ids": ["done", "broken", "gone", "busy"]} + ) + + statuses = {i["artifact_id"]: i["status"] for i in response.json()["items"]} + assert statuses == { + "done": "indexed", + "broken": "failed", + "gone": "deindexed", + "busy": "processing", + } + + +def test_items_follow_request_order_without_duplicates( + client: TestClient, + metadata_store: IngestionMetadataStore, +) -> None: + _save(metadata_store, "a1") + _save(metadata_store, "a2") + + response = client.get( + _URL, params={"artifact_ids": ["a2", "x", "a1", "a2", "x", "a1"]} + ) + + assert response.status_code == 200 + ids = [item["artifact_id"] for item in response.json()["items"]] + assert ids == ["a2", "x", "a1"] + + +def test_batch_of_exactly_the_limit_is_accepted(client: TestClient) -> None: + ids = [f"id-{i}" for i in range(100)] + + response = client.get(_URL, params={"artifact_ids": ids}) + + assert response.status_code == 200 + assert [i["artifact_id"] for i in response.json()["items"]] == ids + + +def test_batch_above_the_limit_is_rejected(client: TestClient) -> None: + ids = [f"id-{i}" for i in range(101)] + + response = client.get(_URL, params={"artifact_ids": ids}) + + assert response.status_code == 422 + assert response.json()["detail"][0]["loc"] == ["query", "artifact_ids"] + + +def test_no_ids_returns_no_items(client: TestClient) -> None: + response = client.get(_URL) + + assert response.status_code == 200 + assert response.json() == {"items": []} + + +def test_blank_id_is_rejected(client: TestClient) -> None: + response = client.get(f"{_URL}?artifact_ids=a1&artifact_ids=") + + assert response.status_code == 422 + + +def test_status_lookup_writes_nothing( + client: TestClient, + metadata_store: IngestionMetadataStore, +) -> None: + _save(metadata_store, "a1") + revision = metadata_store.corpus_revision() + changes = metadata_store._connection.total_changes + + response = client.get(_URL, params={"artifact_ids": ["a1", "missing"]}) + + assert response.status_code == 200 + assert metadata_store._connection.total_changes == changes + assert metadata_store.corpus_revision() == revision + assert metadata_store.get_artifact("missing") is None diff --git a/tests/ingestion/test_metadata_store.py b/tests/ingestion/test_metadata_store.py index f6e3e15..9f4f433 100644 --- a/tests/ingestion/test_metadata_store.py +++ b/tests/ingestion/test_metadata_store.py @@ -1,6 +1,9 @@ import sqlite3 from pathlib import Path +import pytest + +from ingestion import metadata_store from ingestion.metadata_store import ArtifactRecord, IngestionMetadataStore _NOW = "2026-01-01T00:00:00+00:00" @@ -432,3 +435,50 @@ def test_only_failed_revocation_operations_are_recoverable(tmp_path: Path) -> No ) finally: store.close() + + +def test_get_artifacts_returns_known_ids_and_omits_unknown() -> None: + store = IngestionMetadataStore(":memory:") + store.save_completed_artifact(_record("a1", chunk_count=3)) + store.save_completed_artifact(_record("a2", status="failed")) + + records = store.get_artifacts(["a2", "missing", "a1", "a2"]) + + assert set(records) == {"a1", "a2"} + assert records["a1"].chunk_count == 3 + assert records["a2"].status == "failed" + + +def test_get_artifacts_without_ids_returns_empty() -> None: + store = IngestionMetadataStore(":memory:") + store.save_completed_artifact(_record("a1")) + + assert store.get_artifacts([]) == {} + + +def test_get_artifacts_reads_a_batch_in_one_query() -> None: + store = IngestionMetadataStore(":memory:") + ids = [f"a{i}" for i in range(100)] + for artifact_id in ids: + store.save_completed_artifact(_record(artifact_id)) + statements: list[str] = [] + store._connection.set_trace_callback(statements.append) + + records = store.get_artifacts(ids) + + assert set(records) == set(ids) + selects = [s for s in statements if s.lstrip().upper().startswith("SELECT")] + assert len(selects) == 1 + + +def test_get_artifacts_splits_batches_above_the_parameter_limit( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setattr(metadata_store, "_MAX_IN_PARAMS", 2) + store = IngestionMetadataStore(":memory:") + for artifact_id in ("a1", "a2", "a3", "a4", "a5"): + store.save_completed_artifact(_record(artifact_id)) + + records = store.get_artifacts(["a5", "a1", "missing", "a3", "a2", "a4"]) + + assert set(records) == {"a1", "a2", "a3", "a4", "a5"}