Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
6 changes: 5 additions & 1 deletion .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -222,4 +222,8 @@ __marimo__/

data/

opencode.json
opencode.json

# Hermes agent working files (per-developer)
.worktrees/
.worktreeinclude
2 changes: 2 additions & 0 deletions src/api/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
health,
ingest,
ingest_run,
ingest_status,
insights,
knowledge_gaps,
onboarding,
Expand Down Expand Up @@ -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)
Expand Down
78 changes: 78 additions & 0 deletions src/api/routes/ingest_status.py
Original file line number Diff line number Diff line change
@@ -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
]
)
35 changes: 35 additions & 0 deletions src/api/schemas.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 ───────────────────────────────────────


Expand Down
34 changes: 34 additions & 0 deletions src/ingestion/metadata_store.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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",
Expand Down
184 changes: 184 additions & 0 deletions tests/api/test_ingest_status.py
Original file line number Diff line number Diff line change
@@ -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
Loading
Loading