Move source adapters out of the ingester package
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.
This commit is contained in:
parent
a896ac9eec
commit
ab19f78507
40 changed files with 77 additions and 78 deletions
|
|
@ -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.
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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__)
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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"
|
||||
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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__)
|
||||
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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):
|
||||
|
|
|
|||
|
|
@ -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__)
|
||||
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
@ -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,
|
||||
)
|
||||
|
|
@ -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,
|
||||
|
|
@ -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
|
||||
|
|
@ -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(
|
||||
|
|
@ -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,
|
||||
)
|
||||
|
|
@ -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,
|
||||
)
|
||||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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():
|
||||
|
|
|
|||
|
|
@ -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:
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
|
|
@ -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")
|
||||
|
|
|
|||
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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"}
|
||||
|
|
|
|||
0
tests/sources/__init__.py
Normal file
0
tests/sources/__init__.py
Normal file
|
|
@ -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"
|
||||
|
|
@ -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:
|
||||
|
|
@ -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 ---
|
||||
|
||||
|
|
@ -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
|
||||
|
|
@ -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,
|
||||
|
|
@ -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 <prop> 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))
|
||||
|
||||
|
|
@ -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.")
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue