From fa672c298ace8a6eb077619fe65393a14f78a41b Mon Sep 17 00:00:00 2001 From: Yiorgis Gozadinos Date: Fri, 22 May 2026 14:23:22 +0300 Subject: [PATCH] logfire spans for the ingester, nested poller/job/document traces --- CHANGELOG.md | 2 +- haiku_rag_slim/haiku/rag/client/documents.py | 35 ++++++--- .../haiku/rag/ingester/api/server.py | 4 + haiku_rag_slim/haiku/rag/ingester/cli.py | 1 + .../haiku/rag/ingester/pollers/base.py | 75 +++++++++++-------- .../haiku/rag/ingester/workers/pipeline.py | 22 ++++-- haiku_rag_slim/pyproject.toml | 1 + uv.lock | 46 ++++++++++++ 8 files changed, 138 insertions(+), 48 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 7dd1fc92..3c6db2cc 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,7 +3,7 @@ ### Added -- New `haiku-ingester` service for continuous document ingestion: persistent SQLite job queue, async worker pool with retries and a dead-letter queue, FS/HTTP/S3 source adapters with per-source circuit breakers, and a FastAPI control plane (`/health`, `/jobs`, `/sources`, `/dlq`). Configured under `ingester:` in `haiku.rag.yaml`. Shipped behind the `[ingester]` extra. See [docs/ingester.md](docs/ingester.md). +- New `haiku-ingester` service for continuous document ingestion: persistent SQLite job queue, async worker pool with retries and a dead-letter queue, FS/HTTP/S3 source adapters with per-source circuit breakers, and a FastAPI control plane (`/health`, `/jobs`, `/sources`, `/dlq`). Configured under `ingester:` in `haiku.rag.yaml`. Shipped behind the `[ingester]` extra, with Logfire spans (`ingester.poller.sweep` → `ingester.job` → `document.{fetch,convert,chunk,embed,store}`) for traceable ingestion. See [docs/ingester.md](docs/ingester.md). ### Removed diff --git a/haiku_rag_slim/haiku/rag/client/documents.py b/haiku_rag_slim/haiku/rag/client/documents.py index 973e123f..f66325bc 100644 --- a/haiku_rag_slim/haiku/rag/client/documents.py +++ b/haiku_rag_slim/haiku/rag/client/documents.py @@ -3,6 +3,8 @@ from pathlib import Path from typing import TYPE_CHECKING from urllib.parse import unquote, urlparse +import logfire + from haiku.rag.client.processing import ( ensure_chunks_embedded, get_extension_from_content_type_or_url, @@ -247,9 +249,13 @@ async def _ingest_fetch_result( cleanup_path = target_path try: - docling_document = await client.convert(target_path, source_uri=result.uri) - chunks = await client.chunk(docling_document) - embedded_chunks = await embed_chunks(chunks, client._config) + with logfire.span("document.convert", uri=result.uri): + docling_document = await client.convert(target_path, source_uri=result.uri) + with logfire.span("document.chunk", uri=result.uri) as chunk_span: + chunks = await client.chunk(docling_document) + chunk_span.set_attribute("chunks_created", len(chunks)) + with logfire.span("document.embed", uri=result.uri): + embedded_chunks = await embed_chunks(chunks, client._config) finally: if cleanup_path is not None: cleanup_path.unlink(missing_ok=True) @@ -267,9 +273,12 @@ async def _ingest_fetch_result( existing_doc.title = await resolve_title( client._config, docling_document, stored_content ) - return await _update_document_with_chunks( - client, existing_doc, embedded_chunks, docling_document - ) + with logfire.span("document.store", uri=result.uri, op="update") as store_span: + updated = await _update_document_with_chunks( + client, existing_doc, embedded_chunks, docling_document + ) + store_span.set_attribute("document_id", updated.id) + return updated if title is None: title = await resolve_title(client._config, docling_document, stored_content) @@ -280,9 +289,12 @@ async def _ingest_fetch_result( metadata=final_metadata, ) document.set_docling(docling_document) - return await _store_document_with_chunks( - client, document, embedded_chunks, docling_document - ) + with logfire.span("document.store", uri=result.uri, op="create") as store_span: + created = await _store_document_with_chunks( + client, document, embedded_chunks, docling_document + ) + store_span.set_attribute("document_id", created.id) + return created async def create_document_from_source( @@ -388,7 +400,10 @@ async def create_document_from_source( source_metadata=None, ) - result = await fetcher.fetch(source_str) + with logfire.span("document.fetch", uri=source_str) as fetch_span: + result = await fetcher.fetch(source_str) + fetch_span.set_attribute("bytes", len(result.body)) + fetch_span.set_attribute("content_hash", result.content_hash) # MD5 short-circuit: the bytes are unchanged even if the revision wasn't. # Refresh the source-derived metadata (etag may have rolled) but skip diff --git a/haiku_rag_slim/haiku/rag/ingester/api/server.py b/haiku_rag_slim/haiku/rag/ingester/api/server.py index 6fd301a1..c956eabe 100644 --- a/haiku_rag_slim/haiku/rag/ingester/api/server.py +++ b/haiku_rag_slim/haiku/rag/ingester/api/server.py @@ -1,6 +1,7 @@ from dataclasses import dataclass from typing import TYPE_CHECKING +import logfire from fastapi import Depends, FastAPI, Request from haiku.rag.ingester.api.auth import require_auth @@ -49,4 +50,7 @@ def build_app( app.include_router(jobs.router, dependencies=auth_dep) app.include_router(sources.router, dependencies=auth_dep) app.include_router(dlq.router, dependencies=auth_dep) + + # Every request becomes a span when logfire is configured; no-op otherwise. + logfire.instrument_fastapi(app) return app diff --git a/haiku_rag_slim/haiku/rag/ingester/cli.py b/haiku_rag_slim/haiku/rag/ingester/cli.py index 96b961bf..395ecc17 100644 --- a/haiku_rag_slim/haiku/rag/ingester/cli.py +++ b/haiku_rag_slim/haiku/rag/ingester/cli.py @@ -45,6 +45,7 @@ def _configure_logfire() -> None: import logfire logfire.configure( + service_name="haiku-ingester", send_to_logfire="if-token-present", console=False, ) diff --git a/haiku_rag_slim/haiku/rag/ingester/pollers/base.py b/haiku_rag_slim/haiku/rag/ingester/pollers/base.py index 5c39dd52..41ad66e0 100644 --- a/haiku_rag_slim/haiku/rag/ingester/pollers/base.py +++ b/haiku_rag_slim/haiku/rag/ingester/pollers/base.py @@ -2,6 +2,8 @@ import asyncio import logging from datetime import UTC, datetime +import logfire + from haiku.rag.config import SourceConfig from haiku.rag.ingester.pollers.circuit_breaker import CircuitBreaker from haiku.rag.ingester.queue.models import JobOp @@ -17,7 +19,9 @@ logger = logging.getLogger(__name__) def _enqueue_extra(cfg: SourceConfig) -> dict | None: """Per-source state worth carrying into the job (so the worker can rebuild - the same fetch context when it processes).""" + the same fetch context when it processes), plus the current logfire trace + context so the worker's `ingester.job` span nests under the + `ingester.poller.sweep` that enqueued it.""" extra: dict = {} storage_options = getattr(cfg, "storage_options", None) if storage_options: @@ -25,6 +29,9 @@ def _enqueue_extra(cfg: SourceConfig) -> dict | None: headers = getattr(cfg, "headers", None) if headers: extra["headers"] = dict(headers) + carrier = logfire.get_context() + if carrier: + extra["_otel"] = dict(carrier) return extra or None @@ -82,36 +89,44 @@ class BasePoller: "Skipping discover() — circuit breaker open for %s", self.source_id ) return False - try: - snapshot = await self._sync.get_snapshot(self.source_id) - counts = { - SourceEventKind.UPSERT: 0, - SourceEventKind.DELETE: 0, - SourceEventKind.UNCHANGED: 0, - } - async for event in self.source.discover(since=snapshot): - counts[event.kind] += 1 - await self._handle_event(event) - self._breaker.record_success() - self._last_polled_at = datetime.now(UTC) - if counts[SourceEventKind.UPSERT] or counts[SourceEventKind.DELETE]: - logger.info( - "Swept %s: %d upsert, %d delete, %d unchanged", - self.source_id, - counts[SourceEventKind.UPSERT], - counts[SourceEventKind.DELETE], - counts[SourceEventKind.UNCHANGED], + with logfire.span("ingester.poller.sweep", source_id=self.source_id) as span: + try: + snapshot = await self._sync.get_snapshot(self.source_id) + counts = { + SourceEventKind.UPSERT: 0, + SourceEventKind.DELETE: 0, + SourceEventKind.UNCHANGED: 0, + } + async for event in self.source.discover(since=snapshot): + counts[event.kind] += 1 + await self._handle_event(event) + self._breaker.record_success() + self._last_polled_at = datetime.now(UTC) + span.set_attribute("upsert", counts[SourceEventKind.UPSERT]) + span.set_attribute("delete", counts[SourceEventKind.DELETE]) + span.set_attribute("unchanged", counts[SourceEventKind.UNCHANGED]) + if counts[SourceEventKind.UPSERT] or counts[SourceEventKind.DELETE]: + logger.info( + "Swept %s: %d upsert, %d delete, %d unchanged", + self.source_id, + counts[SourceEventKind.UPSERT], + counts[SourceEventKind.DELETE], + counts[SourceEventKind.UNCHANGED], + ) + return True + except Exception as exc: + self._breaker.record_failure() + span.set_attribute( + "consecutive_failures", self._breaker.consecutive_failures ) - return True - except Exception as exc: - self._breaker.record_failure() - logger.exception( - "discover() failed for %s (consecutive=%d): %s", - self.source_id, - self._breaker.consecutive_failures, - exc, - ) - return False + span.record_exception(exc) + logger.exception( + "discover() failed for %s (consecutive=%d): %s", + self.source_id, + self._breaker.consecutive_failures, + exc, + ) + return False async def _handle_event(self, event: SourceEvent) -> None: if event.kind is SourceEventKind.UPSERT: diff --git a/haiku_rag_slim/haiku/rag/ingester/workers/pipeline.py b/haiku_rag_slim/haiku/rag/ingester/workers/pipeline.py index ebf5d728..5f679dd8 100644 --- a/haiku_rag_slim/haiku/rag/ingester/workers/pipeline.py +++ b/haiku_rag_slim/haiku/rag/ingester/workers/pipeline.py @@ -1,4 +1,5 @@ import asyncio +from contextlib import nullcontext from typing import TYPE_CHECKING import httpx @@ -74,14 +75,21 @@ async def run_job(client: "HaikuRAG", job: Job) -> JobResult: extra = job.extra or {} storage_options = extra.get("storage_options") user_metadata = extra.get("metadata", {}) + # Restore the poller's trace context (if any) so the job span nests + # under the `ingester.poller.sweep` that enqueued it. + parent_ctx = extra.get("_otel") + attach = logfire.attach_context(parent_ctx) if parent_ctx else nullcontext() - with logfire.span( - "ingester.job", - job_id=job.id, - source_id=job.source_id, - uri=job.uri, - op=job.op.value, - attempt=job.attempts, + with ( + attach, + logfire.span( + "ingester.job", + job_id=job.id, + source_id=job.source_id, + uri=job.uri, + op=job.op.value, + attempt=job.attempts, + ), ): try: if job.op is JobOp.DELETE: diff --git a/haiku_rag_slim/pyproject.toml b/haiku_rag_slim/pyproject.toml index 56af1a97..88f923c8 100644 --- a/haiku_rag_slim/pyproject.toml +++ b/haiku_rag_slim/pyproject.toml @@ -58,6 +58,7 @@ ingester = [ "fastapi>=0.125", "uvicorn[standard]>=0.32", "aiosqlite>=0.20", + "logfire[fastapi]>=4.30", "haiku.rag-slim[s3]", ] # TUI (chat and inspect commands) diff --git a/uv.lock b/uv.lock index 315266e6..c0dfdf44 100644 --- a/uv.lock +++ b/uv.lock @@ -243,6 +243,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/38/0e/27be9fdef66e72d64c0cdc3cc2823101b80585f8119b5c112c2e8f5f7dab/anyio-4.12.1-py3-none-any.whl", hash = "sha256:d405828884fc140aa80a3c667b8beed277f1dfedec42ba031bd6ac3db606ab6c", size = 113592, upload-time = "2026-01-06T11:45:19.497Z" }, ] +[[package]] +name = "asgiref" +version = "3.11.1" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/63/40/f03da1264ae8f7cfdbf9146542e5e7e8100a4c66ab48e791df9a03d3f6c0/asgiref-3.11.1.tar.gz", hash = "sha256:5f184dc43b7e763efe848065441eac62229c9f7b0475f41f80e207a114eda4ce", size = 38550, upload-time = "2026-02-03T13:30:14.33Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/5c/0a/a72d10ed65068e115044937873362e6e32fab1b7dce0046aeb224682c989/asgiref-3.11.1-py3-none-any.whl", hash = "sha256:e8667a091e69529631969fd45dc268fa79b99c92c5fcdda727757e52146ec133", size = 24345, upload-time = "2026-02-03T13:30:13.039Z" }, +] + [[package]] name = "attrs" version = "26.1.0" @@ -1596,6 +1605,7 @@ groq = [ ingester = [ { name = "aiosqlite" }, { name = "fastapi" }, + { name = "logfire", extra = ["fastapi"] }, { name = "obstore" }, { name = "uvicorn", extra = ["standard"] }, ] @@ -1642,6 +1652,7 @@ requires-dist = [ { name = "jinja2", specifier = ">=3.1.0" }, { name = "jsonpatch", specifier = ">=1.33" }, { name = "lancedb", specifier = "==0.30.2" }, + { name = "logfire", extras = ["fastapi"], marker = "extra == 'ingester'", specifier = ">=4.30" }, { name = "mxbai-rerank", marker = "extra == 'mxbai'", specifier = ">=0.1.6" }, { name = "obstore", marker = "extra == 's3'", specifier = ">=0.9,<0.10" }, { name = "opencv-python-headless", marker = "extra == 'docling'", specifier = ">=4.13.0.92" }, @@ -2249,6 +2260,9 @@ wheels = [ ] [package.optional-dependencies] +fastapi = [ + { name = "opentelemetry-instrumentation-fastapi" }, +] httpx = [ { name = "opentelemetry-instrumentation-httpx" }, ] @@ -3103,6 +3117,38 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/77/d2/6788e83c5c86a2690101681aeef27eeb2a6bf22df52d3f263a22cee20915/opentelemetry_instrumentation-0.60b1-py3-none-any.whl", hash = "sha256:04480db952b48fb1ed0073f822f0ee26012b7be7c3eac1a3793122737c78632d", size = 33096, upload-time = "2025-12-11T13:35:33.067Z" }, ] +[[package]] +name = "opentelemetry-instrumentation-asgi" +version = "0.60b1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "asgiref" }, + { name = "opentelemetry-api" }, + { name = "opentelemetry-instrumentation" }, + { name = "opentelemetry-semantic-conventions" }, + { name = "opentelemetry-util-http" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/77/db/851fa88db7441da82d50bd80f2de5ee55213782e25dc858e04d0c9961d60/opentelemetry_instrumentation_asgi-0.60b1.tar.gz", hash = "sha256:16bfbe595cd24cda309a957456d0fc2523f41bc7b076d1f2d7e98a1ad9876d6f", size = 26107, upload-time = "2025-12-11T13:36:47.015Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/76/76/1fb94367cef64420d2171157a6b9509582873bd09a6afe08a78a8d1f59d9/opentelemetry_instrumentation_asgi-0.60b1-py3-none-any.whl", hash = "sha256:d48def2dbed10294c99cfcf41ebbd0c414d390a11773a41f472d20000fcddc25", size = 16933, upload-time = "2025-12-11T13:35:40.462Z" }, +] + +[[package]] +name = "opentelemetry-instrumentation-fastapi" +version = "0.60b1" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "opentelemetry-api" }, + { name = "opentelemetry-instrumentation" }, + { name = "opentelemetry-instrumentation-asgi" }, + { name = "opentelemetry-semantic-conventions" }, + { name = "opentelemetry-util-http" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/9c/e7/e7e5e50218cf488377209d85666b182fa2d4928bf52389411ceeee1b2b60/opentelemetry_instrumentation_fastapi-0.60b1.tar.gz", hash = "sha256:de608955f7ff8eecf35d056578346a5365015fd7d8623df9b1f08d1c74769c01", size = 24958, upload-time = "2025-12-11T13:36:59.35Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/7d/cc/6e808328ba54662e50babdcab21138eae4250bc0fddf67d55526a615a2ca/opentelemetry_instrumentation_fastapi-0.60b1-py3-none-any.whl", hash = "sha256:af94b7a239ad1085fc3a820ecf069f67f579d7faf4c085aaa7bd9b64eafc8eaf", size = 13478, upload-time = "2025-12-11T13:36:00.811Z" }, +] + [[package]] name = "opentelemetry-instrumentation-httpx" version = "0.60b1"