Skip to content

windlass.interfaces.retriever

retriever

The retriever interface.

A retriever answers "which chunks are relevant to this query?". That is a different job from the vector store, which only answers "which vectors are closest to this one" — the distinction is what lets BM25, hybrid fusion, contextual retrieval and reranking all be retrievers while only some of them touch a vector database.

Implementers override one coroutine, :meth:Retriever.aretrieve_chunks.

Example

from windlass.providers.retrievers.bm25 import BM25Retriever from windlass.core.types import Chunk r = BM25Retriever() r.index([Chunk(content="the cat sat"), Chunk(content="dogs bark")]) 2 r.retrieve("cat").hits[0].chunk.content 'the cat sat'

Retriever

Retriever(
    *,
    top_k: int = 5,
    score_threshold: float | None = None,
    rerank: Any = None,
    fetch_k: int | None = None,
    name: str | None = None,
    **config: Any
)

Bases: Component

Abstract retrieval strategy.

Parameters:

Name Type Description Default
top_k int

Default number of chunks to return.

5
score_threshold float | None

Drop hits scoring below this. None keeps everything. Useful for "answer only if we actually found something" flows.

None
rerank Any

Optional reranker applied to the candidate set before truncation. Any object with an arerank(query, hits, k) coroutine.

None
fetch_k int | None

How many candidates to pull before reranking/filtering. Defaults to top_k when no reranker is configured, 4 * top_k when one is.

None
name str | None

Component name for traces.

None
**config Any

Strategy-specific options.

{}

Attributes:

Name Type Description
top_k

The configured result count.

requires_index bool

Whether :meth:aindex must be called before searching.

Example

Implementing a retriever takes one method::

class RandomRetriever(Retriever):
    provider_name = "random"

    async def aretrieve_chunks(self, query, k, *, filters=None, **kw):
        picks = random.sample(self.corpus, k)
        return [ScoredChunk(chunk=c, score=1.0) for c in picks]
Source code in src\windlass\interfaces\retriever.py
def __init__(
    self,
    *,
    top_k: int = 5,
    score_threshold: float | None = None,
    rerank: Any = None,
    fetch_k: int | None = None,
    name: str | None = None,
    **config: Any,
) -> None:
    if top_k <= 0:
        raise ValueError("top_k must be positive")
    super().__init__(
        name=name or self.provider_name,
        top_k=top_k,
        score_threshold=score_threshold,
        fetch_k=fetch_k,
        **config,
    )
    self.top_k = top_k
    self.score_threshold = score_threshold
    self.reranker = rerank
    self.fetch_k = fetch_k or (top_k * 4 if rerank is not None else top_k)

aretrieve_chunks abstractmethod async

aretrieve_chunks(
    query: str, k: int, *, filters: MetadataFilter | None = None, **kwargs: Any
) -> list[ScoredChunk]

Return candidate chunks for query.

The only method a strategy must implement. Thresholding, reranking, truncation and timing are applied by :meth:aretrieve.

Parameters:

Name Type Description Default
query str

The search query.

required
k int

How many candidates to produce. This is fetch_k, not top_k, when a reranker is configured.

required
filters MetadataFilter | None

Metadata constraints.

None
**kwargs Any

Strategy-specific options.

{}

Returns:

Type Description
list[ScoredChunk]

Scored chunks, ideally already sorted by descending score.

Source code in src\windlass\interfaces\retriever.py
@abc.abstractmethod
async def aretrieve_chunks(
    self,
    query: str,
    k: int,
    *,
    filters: MetadataFilter | None = None,
    **kwargs: Any,
) -> list[ScoredChunk]:
    """Return candidate chunks for ``query``.

    The only method a strategy must implement. Thresholding, reranking,
    truncation and timing are applied by :meth:`aretrieve`.

    Args:
        query: The search query.
        k: How many candidates to produce. This is ``fetch_k``, not
            ``top_k``, when a reranker is configured.
        filters: Metadata constraints.
        **kwargs: Strategy-specific options.

    Returns:
        Scored chunks, ideally already sorted by descending score.
    """

aindex async

aindex(chunks: Sequence[Chunk]) -> int

Add chunks to whatever index this retriever maintains.

The default is a no-op, which is right for retrievers that read from a shared vector store. Lexical retrievers override it.

Parameters:

Name Type Description Default
chunks Sequence[Chunk]

Chunks to index.

required

Returns:

Type Description
int

How many chunks were indexed.

Source code in src\windlass\interfaces\retriever.py
async def aindex(self, chunks: Sequence[Chunk]) -> int:
    """Add chunks to whatever index this retriever maintains.

    The default is a no-op, which is right for retrievers that read from a
    shared vector store. Lexical retrievers override it.

    Args:
        chunks: Chunks to index.

    Returns:
        How many chunks were indexed.
    """
    return 0

index

index(chunks: Sequence[Chunk]) -> int

Blocking :meth:aindex.

Source code in src\windlass\interfaces\retriever.py
def index(self, chunks: Sequence[Chunk]) -> int:
    """Blocking :meth:`aindex`."""
    return run_sync(self.aindex(chunks))

aretrieve async

aretrieve(
    query: str,
    k: int | None = None,
    *,
    filters: MetadataFilter | None = None,
    **kwargs: Any
) -> SearchResult

Retrieve chunks for a query.

Parameters:

Name Type Description Default
query str

The search query.

required
k int | None

Override for :attr:top_k.

None
filters MetadataFilter | None

Metadata constraints.

None
**kwargs Any

Strategy-specific options.

{}

Returns:

Name Type Description
A SearchResult

class:~windlass.core.types.SearchResult with ranked hits,

SearchResult

candidate count and latency.

Raises:

Type Description
RetrievalError

When the underlying strategy fails.

Performance

With a reranker configured, fetch_k candidates are retrieved and then narrowed to k. Raising fetch_k improves recall at the cost of one larger rerank call.

Example

import asyncio from windlass.providers.retrievers.bm25 import BM25Retriever from windlass.core.types import Chunk r = BM25Retriever() _ = r.index([Chunk(content="vector search rocks")]) asyncio.run(r.aretrieve("vector")).hits[0].score > 0 True

Source code in src\windlass\interfaces\retriever.py
async def aretrieve(
    self,
    query: str,
    k: int | None = None,
    *,
    filters: MetadataFilter | None = None,
    **kwargs: Any,
) -> SearchResult:
    """Retrieve chunks for a query.

    Args:
        query: The search query.
        k: Override for :attr:`top_k`.
        filters: Metadata constraints.
        **kwargs: Strategy-specific options.

    Returns:
        A :class:`~windlass.core.types.SearchResult` with ranked hits,
        candidate count and latency.

    Raises:
        RetrievalError: When the underlying strategy fails.

    Performance:
        With a reranker configured, ``fetch_k`` candidates are retrieved and
        then narrowed to ``k``. Raising ``fetch_k`` improves recall at the
        cost of one larger rerank call.

    Example:
        >>> import asyncio
        >>> from windlass.providers.retrievers.bm25 import BM25Retriever
        >>> from windlass.core.types import Chunk
        >>> r = BM25Retriever()
        >>> _ = r.index([Chunk(content="vector search rocks")])
        >>> asyncio.run(r.aretrieve("vector")).hits[0].score > 0
        True
    """
    limit = k or self.top_k
    candidates = max(limit, self.fetch_k) if self.reranker is not None else limit
    started = time.perf_counter()

    hits = await self.aretrieve_chunks(query, candidates, filters=filters, **kwargs)

    if self.reranker is not None and hits:
        hits = await self.reranker.arerank(query, hits, limit)

    if self.score_threshold is not None:
        hits = [h for h in hits if h.score >= self.score_threshold]

    hits = sorted(hits, key=lambda h: h.score, reverse=True)[:limit]
    for position, hit in enumerate(hits, start=1):
        hit.rank = position
        if not hit.retriever:
            hit.retriever = self.name

    return SearchResult(
        query=query,
        hits=hits,
        total=len(hits),
        latency_ms=(time.perf_counter() - started) * 1000,
    )

retrieve

retrieve(
    query: str,
    k: int | None = None,
    *,
    filters: MetadataFilter | None = None,
    **kwargs: Any
) -> SearchResult

Blocking :meth:aretrieve.

Source code in src\windlass\interfaces\retriever.py
def retrieve(
    self,
    query: str,
    k: int | None = None,
    *,
    filters: MetadataFilter | None = None,
    **kwargs: Any,
) -> SearchResult:
    """Blocking :meth:`aretrieve`."""
    return run_sync(self.aretrieve(query, k, filters=filters, **kwargs))

abatch_retrieve async

abatch_retrieve(
    queries: Sequence[str],
    k: int | None = None,
    *,
    filters: MetadataFilter | None = None,
    concurrency: int | None = None
) -> list[SearchResult]

Retrieve for many queries concurrently.

Parameters:

Name Type Description Default
queries Sequence[str]

The queries to run.

required
k int | None

Override for :attr:top_k.

None
filters MetadataFilter | None

Metadata constraints applied to every query.

None
concurrency int | None

Maximum simultaneous retrievals.

None

Returns:

Type Description
list[SearchResult]

One result per query, in input order.

Source code in src\windlass\interfaces\retriever.py
async def abatch_retrieve(
    self,
    queries: Sequence[str],
    k: int | None = None,
    *,
    filters: MetadataFilter | None = None,
    concurrency: int | None = None,
) -> list[SearchResult]:
    """Retrieve for many queries concurrently.

    Args:
        queries: The queries to run.
        k: Override for :attr:`top_k`.
        filters: Metadata constraints applied to every query.
        concurrency: Maximum simultaneous retrievals.

    Returns:
        One result per query, in input order.
    """
    from windlass.core.config import settings

    limit = concurrency or settings().max_concurrency
    return await gather_bounded(
        [self.aretrieve(q, k, filters=filters) for q in queries], limit=limit
    )

batch_retrieve

batch_retrieve(queries: Sequence[str], k: int | None = None) -> list[SearchResult]

Blocking :meth:abatch_retrieve.

Source code in src\windlass\interfaces\retriever.py
def batch_retrieve(self, queries: Sequence[str], k: int | None = None) -> list[SearchResult]:
    """Blocking :meth:`abatch_retrieve`."""
    return run_sync(self.abatch_retrieve(queries, k))