From ab19f78507f942105579509d1f10b6602c56af9d Mon Sep 17 00:00:00 2001 From: Yiorgis Gozadinos Date: Thu, 20 Aug 2026 11:46:55 +0300 Subject: [PATCH] Move source adapters out of the ingester package MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit haiku.rag.ingester.sources was never ingester-only: one-shot client ingestion resolves adapters through it (create_document_from_source), and convert() now fetches through HTTPSource, so the core client imported into the ingester package to reach them. Move the package to haiku.rag.sources and update every import. No shims: haiku.rag.ingester.sources is gone. The haiku.rag.sources plugin entry-point group is unchanged, so third-party source packages need no edit — the group name now matches the module path it always implied. Source unit tests move to tests/sources/. test_source_plugins.py stays in tests/ingester/: it drives a PeriodicPoller against the job repo, so it is plugin wiring through ingester machinery rather than a source test. --- CHANGELOG.md | 1 + docs/ingester.md | 4 ++-- haiku_rag_slim/haiku/rag/client/__init__.py | 2 +- haiku_rag_slim/haiku/rag/client/documents.py | 10 +++++----- haiku_rag_slim/haiku/rag/client/processing.py | 2 +- haiku_rag_slim/haiku/rag/ingester/metadata.py | 2 +- haiku_rag_slim/haiku/rag/ingester/pollers/base.py | 2 +- .../haiku/rag/ingester/pollers/factory.py | 4 ++-- haiku_rag_slim/haiku/rag/ingester/pollers/fs.py | 4 ++-- .../haiku/rag/ingester/pollers/manager.py | 4 ++-- .../haiku/rag/ingester/workers/pipeline.py | 6 +++--- haiku_rag_slim/haiku/rag/ingester/workers/pool.py | 2 +- .../haiku/rag/{ingester => }/sources/__init__.py | 14 +++++++------- .../haiku/rag/{ingester => }/sources/base.py | 0 .../haiku/rag/{ingester => }/sources/filter.py | 0 .../haiku/rag/{ingester => }/sources/fs.py | 4 ++-- .../haiku/rag/{ingester => }/sources/http.py | 2 +- .../haiku/rag/{ingester => }/sources/plugins.py | 2 +- .../haiku/rag/{ingester => }/sources/registry.py | 8 ++++---- .../haiku/rag/{ingester => }/sources/s3.py | 4 ++-- .../haiku/rag/{ingester => }/sources/webdav.py | 4 ++-- tests/ingester/test_api.py | 2 +- tests/ingester/test_filter.py | 2 +- tests/ingester/test_metadata.py | 2 +- tests/ingester/test_pipeline.py | 4 ++-- tests/ingester/test_pollers.py | 8 ++++---- tests/ingester/test_revision_round_trip.py | 4 ++-- tests/ingester/test_run_batch.py | 4 +--- tests/ingester/test_serve_integration.py | 2 +- tests/ingester/test_source_plugins.py | 8 ++++---- tests/ingester/test_workers.py | 2 +- tests/sources/__init__.py | 0 tests/{ingester => sources}/test_fs_source.py | 6 +++--- tests/{ingester => sources}/test_http_source.py | 4 ++-- .../{ingester => sources}/test_resolve_fetcher.py | 8 ++++---- tests/{ingester => sources}/test_s3_source.py | 4 ++-- tests/{ingester => sources}/test_sources_base.py | 2 +- tests/{ingester => sources}/test_webdav_source.py | 6 +++--- tests/test_client.py | 4 ++-- tests/test_pdf_attachments.py | 2 +- 40 files changed, 77 insertions(+), 78 deletions(-) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/__init__.py (52%) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/base.py (100%) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/filter.py (100%) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/fs.py (98%) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/http.py (99%) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/plugins.py (96%) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/registry.py (91%) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/s3.py (98%) rename haiku_rag_slim/haiku/rag/{ingester => }/sources/webdav.py (99%) create mode 100644 tests/sources/__init__.py rename tests/{ingester => sources}/test_fs_source.py (98%) rename tests/{ingester => sources}/test_http_source.py (99%) rename tests/{ingester => sources}/test_resolve_fetcher.py (95%) rename tests/{ingester => sources}/test_s3_source.py (98%) rename tests/{ingester => sources}/test_sources_base.py (97%) rename tests/{ingester => sources}/test_webdav_source.py (99%) diff --git a/CHANGELOG.md b/CHANGELOG.md index a70e556d..ce46cde0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ ### Changed +- Source adapters moved from `haiku.rag.ingester.sources` to `haiku.rag.sources`: `FetchResult`, `SourceEvent`, `SourceEventKind`, `RevisionSnapshot`, `Source` and the `FSSource`/`HTTPSource`/`S3Source`/`WebDAVSource` adapters. One-shot client ingestion uses them too, so they were never ingester-only. Update imports; the `haiku.rag.sources` plugin entry-point group is unchanged. - `HaikuRAG.convert(url)` fetches through `HTTPSource`, the same adapter the ingester uses, instead of its own httpx client. `_write_fetch_body` moved from `client.documents` to `client.processing`. - Chunk embedding is owned by the persistence funnels: `create_document`, `update_document` and source ingestion no longer embed eagerly before handing chunks to a check that would embed them anyway. The `document.embed` span moved onto `ensure_chunks_embedded`, so every path is instrumented rather than only ingest. - One-shot directory ingestion and `FSSource.discover` share `walk_files`, so the symlink-escape guard lives in one place. diff --git a/docs/ingester.md b/docs/ingester.md index 81bb8cf5..37f816e9 100644 --- a/docs/ingester.md +++ b/docs/ingester.md @@ -206,7 +206,7 @@ so a class is its own factory: # example_pkg/__init__.py from urllib.parse import urlparse -from haiku.rag.ingester.sources import FetchResult +from haiku.rag.sources import FetchResult class Provider: @@ -306,7 +306,7 @@ class Source(Protocol): ``` `FetchResult`, `SourceEvent`, `SourceEventKind`, and `RevisionSnapshot` -live in `haiku.rag.ingester.sources`. +live in `haiku.rag.sources`. ```toml # in the source package's pyproject.toml diff --git a/haiku_rag_slim/haiku/rag/client/__init__.py b/haiku_rag_slim/haiku/rag/client/__init__.py index 308987b5..4cf6ba48 100644 --- a/haiku_rag_slim/haiku/rag/client/__init__.py +++ b/haiku_rag_slim/haiku/rag/client/__init__.py @@ -34,9 +34,9 @@ if TYPE_CHECKING: from haiku.rag.embeddings import EmbedderWrapper from haiku.rag.ingester.metadata import MetadataProvider - from haiku.rag.ingester.sources.base import Source from haiku.rag.reranking.base import RerankerBase from haiku.rag.sandbox import AnalysisResult + from haiku.rag.sources.base import Source from haiku.rag.store.models.citation import Citation logger = logging.getLogger(__name__) diff --git a/haiku_rag_slim/haiku/rag/client/documents.py b/haiku_rag_slim/haiku/rag/client/documents.py index 39b93781..13fc65b0 100644 --- a/haiku_rag_slim/haiku/rag/client/documents.py +++ b/haiku_rag_slim/haiku/rag/client/documents.py @@ -26,7 +26,7 @@ if TYPE_CHECKING: from haiku.rag.client import HaikuRAG from haiku.rag.ingester.metadata import MetadataProvider - from haiku.rag.ingester.sources.base import FetchResult, Source + from haiku.rag.sources.base import FetchResult, Source logger = logging.getLogger(__name__) @@ -582,7 +582,7 @@ async def _reconcile_pdf_attachments( ): continue - from haiku.rag.ingester.sources.base import FetchResult + from haiku.rag.sources.base import FetchResult child_fr = FetchResult( uri=child_uri, @@ -668,8 +668,8 @@ async def create_document_from_source( "uri override is not supported for directory sources; each file " "produces its own document with its own auto-derived URI." ) - from haiku.rag.ingester.sources.filter import FileFilter - from haiku.rag.ingester.sources.fs import walk_files + from haiku.rag.sources.filter import FileFilter + from haiku.rag.sources.fs import walk_files # One-shot CLI directory ingest uses the converter's supported # extensions but no include/ignore patterns. For pattern-based @@ -707,7 +707,7 @@ async def create_document_from_source( # renamed/removed source surfaces as a DLQ instead of silently dropping # credentials. Ad-hoc CLI calls (no source_id) fall back to scheme-based # adapters when no configured source matches. - from haiku.rag.ingester.sources import ( + from haiku.rag.sources import ( resolve_adhoc_fetcher, resolve_configured_source, ) diff --git a/haiku_rag_slim/haiku/rag/client/processing.py b/haiku_rag_slim/haiku/rag/client/processing.py index 407b4020..ae2dfd85 100644 --- a/haiku_rag_slim/haiku/rag/client/processing.py +++ b/haiku_rag_slim/haiku/rag/client/processing.py @@ -131,7 +131,7 @@ async def convert( if parsed.scheme in ("http", "https"): # One HTTP acquisition path: the same adapter the ingester fetches with. - from haiku.rag.ingester.sources.http import HTTPSource + from haiku.rag.sources.http import HTTPSource fetcher = HTTPSource(source_id="convert") try: diff --git a/haiku_rag_slim/haiku/rag/ingester/metadata.py b/haiku_rag_slim/haiku/rag/ingester/metadata.py index 39a98d3b..2eabf291 100644 --- a/haiku_rag_slim/haiku/rag/ingester/metadata.py +++ b/haiku_rag_slim/haiku/rag/ingester/metadata.py @@ -3,7 +3,7 @@ from importlib.metadata import entry_points from typing import TYPE_CHECKING, Protocol, runtime_checkable if TYPE_CHECKING: - from haiku.rag.ingester.sources.base import FetchResult + from haiku.rag.sources.base import FetchResult ENTRY_POINT_GROUP = "haiku.rag.metadata_providers" diff --git a/haiku_rag_slim/haiku/rag/ingester/pollers/base.py b/haiku_rag_slim/haiku/rag/ingester/pollers/base.py index 26e992fb..b71b053b 100644 --- a/haiku_rag_slim/haiku/rag/ingester/pollers/base.py +++ b/haiku_rag_slim/haiku/rag/ingester/pollers/base.py @@ -8,7 +8,7 @@ from haiku.rag.config import SourceConfig from haiku.rag.ingester.batch import BatchChange, BatchSourceSummary from haiku.rag.ingester.queue.models import JobOp, SyncRow from haiku.rag.ingester.queue.repository import JobRepo, SyncStateRepo -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( Source, SourceEvent, SourceEventKind, diff --git a/haiku_rag_slim/haiku/rag/ingester/pollers/factory.py b/haiku_rag_slim/haiku/rag/ingester/pollers/factory.py index 693905fc..362c6ae9 100644 --- a/haiku_rag_slim/haiku/rag/ingester/pollers/factory.py +++ b/haiku_rag_slim/haiku/rag/ingester/pollers/factory.py @@ -6,14 +6,14 @@ from haiku.rag.config import ( SourceConfig, WebDAVSourceConfig, ) -from haiku.rag.ingester.sources import ( +from haiku.rag.sources import ( FSSource, HTTPSource, S3Source, Source, WebDAVSource, ) -from haiku.rag.ingester.sources.plugins import ( +from haiku.rag.sources.plugins import ( ENTRY_POINT_GROUP, load_source_factories, ) diff --git a/haiku_rag_slim/haiku/rag/ingester/pollers/fs.py b/haiku_rag_slim/haiku/rag/ingester/pollers/fs.py index f30f5111..265488c4 100644 --- a/haiku_rag_slim/haiku/rag/ingester/pollers/fs.py +++ b/haiku_rag_slim/haiku/rag/ingester/pollers/fs.py @@ -7,13 +7,13 @@ from watchfiles import Change, awatch from haiku.rag.ingester.pollers.base import BasePoller, _enqueue_extra from haiku.rag.ingester.queue.models import JobOp -from haiku.rag.ingester.sources.filter import FileFilter +from haiku.rag.sources.filter import FileFilter from haiku.rag.telemetry import logfire if TYPE_CHECKING: from haiku.rag.circuit_breaker import CircuitBreaker from haiku.rag.config import FSSourceConfig - from haiku.rag.ingester.sources.fs import FSSource + from haiku.rag.sources.fs import FSSource logger = logging.getLogger(__name__) diff --git a/haiku_rag_slim/haiku/rag/ingester/pollers/manager.py b/haiku_rag_slim/haiku/rag/ingester/pollers/manager.py index 7067e215..babfb532 100644 --- a/haiku_rag_slim/haiku/rag/ingester/pollers/manager.py +++ b/haiku_rag_slim/haiku/rag/ingester/pollers/manager.py @@ -14,7 +14,7 @@ from haiku.rag.ingester.pollers.periodic import PeriodicPoller if TYPE_CHECKING: from haiku.rag.ingester.queue.repository import JobRepo, SyncStateRepo - from haiku.rag.ingester.sources.base import Source + from haiku.rag.sources.base import Source logger = logging.getLogger(__name__) @@ -48,7 +48,7 @@ class PollerManager: source = build_source(cfg, supported_extensions=self._supported_extensions) breaker = CircuitBreaker(cfg.circuit_breaker) if isinstance(cfg, FSSourceConfig): - from haiku.rag.ingester.sources.fs import FSSource + from haiku.rag.sources.fs import FSSource assert isinstance(source, FSSource) return FSPoller( diff --git a/haiku_rag_slim/haiku/rag/ingester/workers/pipeline.py b/haiku_rag_slim/haiku/rag/ingester/workers/pipeline.py index 4b3fe3b3..068b1d2d 100644 --- a/haiku_rag_slim/haiku/rag/ingester/workers/pipeline.py +++ b/haiku_rag_slim/haiku/rag/ingester/workers/pipeline.py @@ -14,8 +14,8 @@ from pydantic import BaseModel from haiku.rag.client.exceptions import UnsupportedSourceError from haiku.rag.ingester.exceptions import PermanentError, TransientError from haiku.rag.ingester.queue.models import Job, JobOp -from haiku.rag.ingester.sources.base import FileTooLargeError -from haiku.rag.ingester.sources.registry import resolve_configured_source +from haiku.rag.sources.base import FileTooLargeError +from haiku.rag.sources.registry import resolve_configured_source from haiku.rag.telemetry import attach_context, logfire if TYPE_CHECKING: @@ -23,7 +23,7 @@ if TYPE_CHECKING: from haiku.rag.client import HaikuRAG from haiku.rag.ingester.metadata import MetadataProvider - from haiku.rag.ingester.sources.base import Source + from haiku.rag.sources.base import Source class JobResult(BaseModel): diff --git a/haiku_rag_slim/haiku/rag/ingester/workers/pool.py b/haiku_rag_slim/haiku/rag/ingester/workers/pool.py index a7bbcf63..8aa3ca14 100644 --- a/haiku_rag_slim/haiku/rag/ingester/workers/pool.py +++ b/haiku_rag_slim/haiku/rag/ingester/workers/pool.py @@ -19,7 +19,7 @@ if TYPE_CHECKING: from haiku.rag.client import HaikuRAG from haiku.rag.ingester.metadata import MetadataProvider - from haiku.rag.ingester.sources.base import Source + from haiku.rag.sources.base import Source logger = logging.getLogger(__name__) diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/__init__.py b/haiku_rag_slim/haiku/rag/sources/__init__.py similarity index 52% rename from haiku_rag_slim/haiku/rag/ingester/sources/__init__.py rename to haiku_rag_slim/haiku/rag/sources/__init__.py index 15e1072f..0b6a1543 100644 --- a/haiku_rag_slim/haiku/rag/ingester/sources/__init__.py +++ b/haiku_rag_slim/haiku/rag/sources/__init__.py @@ -1,19 +1,19 @@ -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( FetchResult, RevisionSnapshot, Source, SourceEvent, SourceEventKind, ) -from haiku.rag.ingester.sources.filter import FileFilter -from haiku.rag.ingester.sources.fs import FSSource -from haiku.rag.ingester.sources.http import HTTPSource -from haiku.rag.ingester.sources.registry import ( +from haiku.rag.sources.filter import FileFilter +from haiku.rag.sources.fs import FSSource +from haiku.rag.sources.http import HTTPSource +from haiku.rag.sources.registry import ( resolve_adhoc_fetcher, resolve_configured_source, ) -from haiku.rag.ingester.sources.s3 import S3Source -from haiku.rag.ingester.sources.webdav import WebDAVSource +from haiku.rag.sources.s3 import S3Source +from haiku.rag.sources.webdav import WebDAVSource __all__ = [ "FetchResult", diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/base.py b/haiku_rag_slim/haiku/rag/sources/base.py similarity index 100% rename from haiku_rag_slim/haiku/rag/ingester/sources/base.py rename to haiku_rag_slim/haiku/rag/sources/base.py diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/filter.py b/haiku_rag_slim/haiku/rag/sources/filter.py similarity index 100% rename from haiku_rag_slim/haiku/rag/ingester/sources/filter.py rename to haiku_rag_slim/haiku/rag/sources/filter.py diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/fs.py b/haiku_rag_slim/haiku/rag/sources/fs.py similarity index 98% rename from haiku_rag_slim/haiku/rag/ingester/sources/fs.py rename to haiku_rag_slim/haiku/rag/sources/fs.py index 7c77b3d5..ac0cab90 100644 --- a/haiku_rag_slim/haiku/rag/ingester/sources/fs.py +++ b/haiku_rag_slim/haiku/rag/sources/fs.py @@ -8,14 +8,14 @@ from pathlib import Path from urllib.parse import unquote, urlparse from haiku.rag.client.exceptions import UnsupportedSourceError -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( FetchResult, RevisionSnapshot, SourceEvent, SourceEventKind, check_file_size, ) -from haiku.rag.ingester.sources.filter import ( +from haiku.rag.sources.filter import ( FileFilter, _default_supported_extensions, ) diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/http.py b/haiku_rag_slim/haiku/rag/sources/http.py similarity index 99% rename from haiku_rag_slim/haiku/rag/ingester/sources/http.py rename to haiku_rag_slim/haiku/rag/sources/http.py index a54ce811..9bb97ae5 100644 --- a/haiku_rag_slim/haiku/rag/ingester/sources/http.py +++ b/haiku_rag_slim/haiku/rag/sources/http.py @@ -6,7 +6,7 @@ from urllib.parse import urlparse import httpx -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( FetchResult, RevisionSnapshot, SourceEvent, diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/plugins.py b/haiku_rag_slim/haiku/rag/sources/plugins.py similarity index 96% rename from haiku_rag_slim/haiku/rag/ingester/sources/plugins.py rename to haiku_rag_slim/haiku/rag/sources/plugins.py index 65bb2f86..f70573aa 100644 --- a/haiku_rag_slim/haiku/rag/ingester/sources/plugins.py +++ b/haiku_rag_slim/haiku/rag/sources/plugins.py @@ -1,7 +1,7 @@ from importlib.metadata import entry_points from typing import Any, Protocol, runtime_checkable -from haiku.rag.ingester.sources.base import Source +from haiku.rag.sources.base import Source @runtime_checkable diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/registry.py b/haiku_rag_slim/haiku/rag/sources/registry.py similarity index 91% rename from haiku_rag_slim/haiku/rag/ingester/sources/registry.py rename to haiku_rag_slim/haiku/rag/sources/registry.py index 793c74c0..76b9620b 100644 --- a/haiku_rag_slim/haiku/rag/ingester/sources/registry.py +++ b/haiku_rag_slim/haiku/rag/sources/registry.py @@ -3,10 +3,10 @@ from pathlib import Path from urllib.parse import urlparse from haiku.rag.client.exceptions import UnsupportedSourceError -from haiku.rag.ingester.sources.base import Source -from haiku.rag.ingester.sources.fs import FSSource -from haiku.rag.ingester.sources.http import HTTPSource -from haiku.rag.ingester.sources.s3 import S3Source +from haiku.rag.sources.base import Source +from haiku.rag.sources.fs import FSSource +from haiku.rag.sources.http import HTTPSource +from haiku.rag.sources.s3 import S3Source def resolve_configured_source( diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/s3.py b/haiku_rag_slim/haiku/rag/sources/s3.py similarity index 98% rename from haiku_rag_slim/haiku/rag/ingester/sources/s3.py rename to haiku_rag_slim/haiku/rag/sources/s3.py index fc44474f..3c069582 100644 --- a/haiku_rag_slim/haiku/rag/ingester/sources/s3.py +++ b/haiku_rag_slim/haiku/rag/sources/s3.py @@ -5,14 +5,14 @@ from datetime import UTC, datetime from urllib.parse import urlparse from haiku.rag.client.exceptions import UnsupportedSourceError -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( FetchResult, RevisionSnapshot, SourceEvent, SourceEventKind, check_file_size, ) -from haiku.rag.ingester.sources.filter import ( +from haiku.rag.sources.filter import ( FileFilter, _default_supported_extensions, ) diff --git a/haiku_rag_slim/haiku/rag/ingester/sources/webdav.py b/haiku_rag_slim/haiku/rag/sources/webdav.py similarity index 99% rename from haiku_rag_slim/haiku/rag/ingester/sources/webdav.py rename to haiku_rag_slim/haiku/rag/sources/webdav.py index 4437de31..b0dab48c 100644 --- a/haiku_rag_slim/haiku/rag/ingester/sources/webdav.py +++ b/haiku_rag_slim/haiku/rag/sources/webdav.py @@ -7,14 +7,14 @@ from xml.etree.ElementTree import Element, fromstring import httpx -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( FetchResult, RevisionSnapshot, SourceEvent, SourceEventKind, check_file_size, ) -from haiku.rag.ingester.sources.filter import ( +from haiku.rag.sources.filter import ( FileFilter, _default_supported_extensions, ) diff --git a/tests/ingester/test_api.py b/tests/ingester/test_api.py index fa5b365a..301435fc 100644 --- a/tests/ingester/test_api.py +++ b/tests/ingester/test_api.py @@ -7,7 +7,7 @@ from httpx import ASGITransport from haiku.rag.config import AppConfig from haiku.rag.ingester.api.server import APIState, build_app from haiku.rag.ingester.queue.models import JobOp, JobStatus -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( FetchResult, SourceEvent, SourceEventKind, diff --git a/tests/ingester/test_filter.py b/tests/ingester/test_filter.py index c880fd14..85fc373a 100644 --- a/tests/ingester/test_filter.py +++ b/tests/ingester/test_filter.py @@ -1,6 +1,6 @@ from watchfiles import Change -from haiku.rag.ingester.sources.filter import FileFilter, _default_supported_extensions +from haiku.rag.sources.filter import FileFilter, _default_supported_extensions def test_default_supported_extensions_returns_nonempty_list(): diff --git a/tests/ingester/test_metadata.py b/tests/ingester/test_metadata.py index bd0adb20..d9bf2c30 100644 --- a/tests/ingester/test_metadata.py +++ b/tests/ingester/test_metadata.py @@ -7,7 +7,7 @@ from haiku.rag.ingester.metadata import ( build_providers, load_metadata_providers, ) -from haiku.rag.ingester.sources.base import FetchResult +from haiku.rag.sources.base import FetchResult class _Provider: diff --git a/tests/ingester/test_pipeline.py b/tests/ingester/test_pipeline.py index 2d52873f..d756e903 100644 --- a/tests/ingester/test_pipeline.py +++ b/tests/ingester/test_pipeline.py @@ -15,8 +15,8 @@ from obstore.exceptions import ( from haiku.rag.client import HaikuRAG from haiku.rag.ingester.exceptions import PermanentError, TransientError from haiku.rag.ingester.queue.models import Job, JobOp, JobStatus -from haiku.rag.ingester.sources.base import FetchResult, FileTooLargeError, Source from haiku.rag.ingester.workers.pipeline import run_job +from haiku.rag.sources.base import FetchResult, FileTooLargeError, Source from haiku.rag.store.models.document import Document @@ -80,7 +80,7 @@ async def test_upsert_calls_create_document_from_source_and_returns_metadata(): async def test_upsert_threads_configured_sources_to_client(): """The list of configured Source adapters reaches the client so resolve_fetcher can pick the right one by source_id.""" - from haiku.rag.ingester.sources.http import HTTPSource + from haiku.rag.sources.http import HTTPSource client = _mock_client() client.create_document_from_source.return_value = Document( diff --git a/tests/ingester/test_pollers.py b/tests/ingester/test_pollers.py index b33bc2d3..f1868b97 100644 --- a/tests/ingester/test_pollers.py +++ b/tests/ingester/test_pollers.py @@ -15,7 +15,7 @@ from haiku.rag.config import ( from haiku.rag.ingester.pollers.manager import PollerManager from haiku.rag.ingester.pollers.periodic import PeriodicPoller from haiku.rag.ingester.queue.models import JobOp, JobStatus -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( FetchResult, SourceEvent, SourceEventKind, @@ -416,7 +416,7 @@ async def test_manager_sources_available_at_construction(tmp_path, jobs, sync): """PollerManager builds Sources eagerly so callers (WorkerPool) can receive them by plain construction order.""" from haiku.rag.config import SourceConfig - from haiku.rag.ingester.sources.http import HTTPSource + from haiku.rag.sources.http import HTTPSource configs: list[SourceConfig] = [ FSSourceConfig(type="fs", id="docs", root=tmp_path), @@ -491,7 +491,7 @@ def _fs_poller(tmp_path, jobs, sync): """Construct an FSPoller for unit-testing the watch-change handler. Doesn't start the watch loop — tests call `_handle_watch_change` directly.""" from haiku.rag.ingester.pollers.fs import FSPoller - from haiku.rag.ingester.sources.fs import FSSource + from haiku.rag.sources.fs import FSSource cfg = FSSourceConfig( type="fs", @@ -713,7 +713,7 @@ async def test_watch_deleted_skipped_when_delete_orphans_false(tmp_path, jobs, s from watchfiles import Change from haiku.rag.ingester.pollers.fs import FSPoller - from haiku.rag.ingester.sources.fs import FSSource + from haiku.rag.sources.fs import FSSource cfg = FSSourceConfig( type="fs", diff --git a/tests/ingester/test_revision_round_trip.py b/tests/ingester/test_revision_round_trip.py index e71e5750..6e01c116 100644 --- a/tests/ingester/test_revision_round_trip.py +++ b/tests/ingester/test_revision_round_trip.py @@ -10,8 +10,8 @@ from pathlib import Path import pytest from haiku.rag.client import HaikuRAG -from haiku.rag.ingester.sources.base import FetchResult, SourceEventKind -from haiku.rag.ingester.sources.fs import FSSource +from haiku.rag.sources.base import FetchResult, SourceEventKind +from haiku.rag.sources.fs import FSSource @pytest.fixture(scope="module") diff --git a/tests/ingester/test_run_batch.py b/tests/ingester/test_run_batch.py index bf87f851..eb4ee956 100644 --- a/tests/ingester/test_run_batch.py +++ b/tests/ingester/test_run_batch.py @@ -240,9 +240,7 @@ async def test_run_batch_reports_failed_sweep( raise RuntimeError("discover blew up") yield # unreachable; makes this an async generator - monkeypatch.setattr( - "haiku.rag.ingester.sources.fs.FSSource.discover", _failing_discover - ) + monkeypatch.setattr("haiku.rag.sources.fs.FSSource.discover", _failing_discover) with caplog.at_level("ERROR", logger="haiku.rag.ingester.pollers.base"): report = await IngesterApp( diff --git a/tests/ingester/test_serve_integration.py b/tests/ingester/test_serve_integration.py index 34962d98..7fa811ae 100644 --- a/tests/ingester/test_serve_integration.py +++ b/tests/ingester/test_serve_integration.py @@ -9,8 +9,8 @@ from haiku.rag.client import HaikuRAG from haiku.rag.config import FSSourceConfig, HTTPSourceConfig from haiku.rag.ingester.pollers.manager import PollerManager from haiku.rag.ingester.queue.models import JobOp -from haiku.rag.ingester.sources.http import HTTPSource from haiku.rag.ingester.workers.pool import WorkerPool +from haiku.rag.sources.http import HTTPSource from haiku.rag.store.models.document import Document diff --git a/tests/ingester/test_source_plugins.py b/tests/ingester/test_source_plugins.py index 1a8ccf50..9ee55fc2 100644 --- a/tests/ingester/test_source_plugins.py +++ b/tests/ingester/test_source_plugins.py @@ -7,15 +7,15 @@ from haiku.rag.config import PluginSourceConfig from haiku.rag.ingester.pollers.factory import build_source from haiku.rag.ingester.pollers.periodic import PeriodicPoller from haiku.rag.ingester.queue.models import JobOp -from haiku.rag.ingester.sources import plugins as plugins_module -from haiku.rag.ingester.sources import resolve_configured_source -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources import plugins as plugins_module +from haiku.rag.sources import resolve_configured_source +from haiku.rag.sources.base import ( FetchResult, Source, SourceEvent, SourceEventKind, ) -from haiku.rag.ingester.sources.plugins import ( +from haiku.rag.sources.plugins import ( ENTRY_POINT_GROUP, load_source_factories, ) diff --git a/tests/ingester/test_workers.py b/tests/ingester/test_workers.py index bf6f6c28..82036372 100644 --- a/tests/ingester/test_workers.py +++ b/tests/ingester/test_workers.py @@ -325,7 +325,7 @@ async def test_drain_passes_configured_sources_to_client(client, jobs, sync): """The pool's `sources` list flows through run_job to client.create_document_from_source so resolve_fetcher can pick the configured authenticated source over an adhoc adapter.""" - from haiku.rag.ingester.sources.http import HTTPSource + from haiku.rag.sources.http import HTTPSource client.create_document_from_source.return_value = Document( id="d", content="x", uri="u", metadata={"md5": "m", "source_revision": "r"} diff --git a/tests/sources/__init__.py b/tests/sources/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/tests/ingester/test_fs_source.py b/tests/sources/test_fs_source.py similarity index 98% rename from tests/ingester/test_fs_source.py rename to tests/sources/test_fs_source.py index 1ce4a40e..d64a0625 100644 --- a/tests/ingester/test_fs_source.py +++ b/tests/sources/test_fs_source.py @@ -4,8 +4,8 @@ from pathlib import Path import pytest from haiku.rag.client.exceptions import UnsupportedSourceError -from haiku.rag.ingester.sources.base import FileTooLargeError, SourceEventKind -from haiku.rag.ingester.sources.fs import FSSource +from haiku.rag.sources.base import FileTooLargeError, SourceEventKind +from haiku.rag.sources.fs import FSSource @pytest.fixture @@ -378,7 +378,7 @@ async def test_discover_skips_symlink_to_missing_in_root_target(tmp_path): def test_walk_files_drops_links_escaping_the_root(tmp_path): import os - from haiku.rag.ingester.sources.fs import walk_files + from haiku.rag.sources.fs import walk_files tree = tmp_path / "tree" outside = tmp_path / "outside" diff --git a/tests/ingester/test_http_source.py b/tests/sources/test_http_source.py similarity index 99% rename from tests/ingester/test_http_source.py rename to tests/sources/test_http_source.py index 856ad50c..40e874a3 100644 --- a/tests/ingester/test_http_source.py +++ b/tests/sources/test_http_source.py @@ -3,8 +3,8 @@ import hashlib import httpx import pytest -from haiku.rag.ingester.sources.base import FileTooLargeError, SourceEventKind -from haiku.rag.ingester.sources.http import HTTPSource +from haiku.rag.sources.base import FileTooLargeError, SourceEventKind +from haiku.rag.sources.http import HTTPSource def _transport(routes: dict[tuple[str, str], httpx.Response]) -> httpx.MockTransport: diff --git a/tests/ingester/test_resolve_fetcher.py b/tests/sources/test_resolve_fetcher.py similarity index 95% rename from tests/ingester/test_resolve_fetcher.py rename to tests/sources/test_resolve_fetcher.py index ddf198c3..efb64558 100644 --- a/tests/ingester/test_resolve_fetcher.py +++ b/tests/sources/test_resolve_fetcher.py @@ -2,13 +2,13 @@ from pathlib import Path import pytest -from haiku.rag.ingester.sources import ( +from haiku.rag.sources import ( resolve_adhoc_fetcher, resolve_configured_source, ) -from haiku.rag.ingester.sources.fs import FSSource -from haiku.rag.ingester.sources.http import HTTPSource -from haiku.rag.ingester.sources.s3 import S3Source +from haiku.rag.sources.fs import FSSource +from haiku.rag.sources.http import HTTPSource +from haiku.rag.sources.s3 import S3Source # --- resolve_adhoc_fetcher --- diff --git a/tests/ingester/test_s3_source.py b/tests/sources/test_s3_source.py similarity index 98% rename from tests/ingester/test_s3_source.py rename to tests/sources/test_s3_source.py index 25d420f7..463ad010 100644 --- a/tests/ingester/test_s3_source.py +++ b/tests/sources/test_s3_source.py @@ -3,8 +3,8 @@ from unittest.mock import AsyncMock, MagicMock import pytest -from haiku.rag.ingester.sources.base import FileTooLargeError, SourceEventKind -from haiku.rag.ingester.sources.s3 import S3Source +from haiku.rag.sources.base import FileTooLargeError, SourceEventKind +from haiku.rag.sources.s3 import S3Source @pytest.fixture diff --git a/tests/ingester/test_sources_base.py b/tests/sources/test_sources_base.py similarity index 97% rename from tests/ingester/test_sources_base.py rename to tests/sources/test_sources_base.py index ebaaebae..57291e90 100644 --- a/tests/ingester/test_sources_base.py +++ b/tests/sources/test_sources_base.py @@ -1,6 +1,6 @@ from datetime import UTC, datetime -from haiku.rag.ingester.sources.base import ( +from haiku.rag.sources.base import ( FetchResult, Source, SourceEvent, diff --git a/tests/ingester/test_webdav_source.py b/tests/sources/test_webdav_source.py similarity index 99% rename from tests/ingester/test_webdav_source.py rename to tests/sources/test_webdav_source.py index 8659c017..8139c0cb 100644 --- a/tests/ingester/test_webdav_source.py +++ b/tests/sources/test_webdav_source.py @@ -3,8 +3,8 @@ import hashlib import httpx import pytest -from haiku.rag.ingester.sources.base import FileTooLargeError, SourceEventKind -from haiku.rag.ingester.sources.webdav import WebDAVSource, _strip_etag +from haiku.rag.sources.base import FileTooLargeError, SourceEventKind +from haiku.rag.sources.webdav import WebDAVSource, _strip_etag @pytest.mark.parametrize( @@ -716,7 +716,7 @@ async def test_head_returns_none_for_empty_multistatus(): def test_entry_with_status_but_no_prop_has_no_revision(): """A 200 propstat carrying no still yields an entry, without a revision — distinct from the malformed bodies that yield no entry at all.""" - from haiku.rag.ingester.sources.webdav import _parse_multistatus + from haiku.rag.sources.webdav import _parse_multistatus entries = _parse_multistatus(_raw_multistatus(_STATUS_WITHOUT_PROP)) diff --git a/tests/test_client.py b/tests/test_client.py index e543f9c2..a5661750 100644 --- a/tests/test_client.py +++ b/tests/test_client.py @@ -19,7 +19,7 @@ from haiku.rag.client.documents import ( from haiku.rag.client.processing import _write_fetch_body from haiku.rag.config import get_config from haiku.rag.embeddings import EmbedderWrapper -from haiku.rag.ingester.sources.base import FetchResult +from haiku.rag.sources.base import FetchResult from haiku.rag.store.compression import decompress_json from haiku.rag.store.models.chunk import Chunk from haiku.rag.store.models.document import Document @@ -2532,7 +2532,7 @@ class _CountingSource: async def test_adhoc_source_is_closed_after_ingest(temp_db_path, monkeypatch): """An ad-hoc fetcher is built for this one call, so this call has to close it — HTTP and WebDAV adapters hold an httpx connection pool.""" - from haiku.rag.ingester import sources as sources_module + from haiku.rag import sources as sources_module uri = "https://example.com/counting.md" fetcher = _CountingSource(uri, b"# Counting\n\nAd-hoc fetched body.") diff --git a/tests/test_pdf_attachments.py b/tests/test_pdf_attachments.py index 69851570..f283aba3 100644 --- a/tests/test_pdf_attachments.py +++ b/tests/test_pdf_attachments.py @@ -10,7 +10,7 @@ from haiku.rag.client.documents import ( _reconcile_pdf_attachments, parent_uri_filter, ) -from haiku.rag.ingester.sources import FetchResult +from haiku.rag.sources import FetchResult from haiku.rag.store.models.document import Document