API reference¶
Generated from the apogee-ai-rag source with mkdocstrings. Every symbol below is exported from apogee_ai_rag, so it is part of the supported public surface.
Other¶
AdaptiveRagPipeline
¶
AdaptiveRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, classifier: Callable[[RagQuery], Complexity] | None = None, max_iterations: int = 3)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/adaptive/adaptive_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
classifier: Callable[[RagQuery], Complexity] | None = None,
max_iterations: int = 3,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"AdaptiveRagPipeline requires chunker, embedder, store and generator"
)
if max_iterations < 1:
raise PipelineConfigError("max_iterations must be >= 1")
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._classifier = classifier or heuristic_complexity
self._max_iterations = max_iterations
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/adaptive/adaptive_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
await self._emit(
"query.received", query_id=query.id, text=query.text, pipeline_type=self.name,
)
complexity, retrieved = await self._route(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
response.traces.append({"pipeline": self.name, "complexity": complexity})
await self._emit(
"query.generated", query_id=query.id, complexity=complexity,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/adaptive/adaptive_rag_pipeline.py
AdvancedRagPipeline
¶
AdvancedRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, query_rewriter: IQueryRewriter | None = None, reranker: IReranker | None = None, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, dedupe_by_parent: bool = True)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/advanced/advanced_rag_pipeline.py
def __init__( # noqa: PLR0913 — DI surface
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
query_rewriter: IQueryRewriter | None = None,
reranker: IReranker | None = None,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
dedupe_by_parent: bool = True,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"AdvancedRagPipeline requires chunker, embedder, vector_store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._rewriter = query_rewriter or PassThroughRewriter()
self._reranker = reranker or NoOpReranker()
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._dedupe = dedupe_by_parent
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/advanced/advanced_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.query_id = query.id
await self._emit(
"query.generated",
query_id=query.id,
tokens=response.usage.get("total_tokens", 0),
finish_reason=response.finish_reason,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/advanced/advanced_rag_pipeline.py
AgenticRagPipeline
¶
AgenticRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, max_steps: int = 5, extra_tools: dict[str, Callable[[str], str]] | None = None)
Bases: ReActRagPipeline
Source code in apogee_ai_rag/infrastructure/pipelines/agentic/agentic_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
max_steps: int = 5,
extra_tools: dict[str, Callable[[str], str]] | None = None,
) -> None:
super().__init__(
chunker=chunker,
embedder=embedder,
vector_store=vector_store,
generator=generator,
observability=observability,
generation_config=generation_config,
max_steps=max_steps,
)
self._extra_tools = extra_tools or {}
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/agentic/agentic_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
response = await super().run(query)
if self._extra_tools:
response.traces.append({
"pipeline": self.name,
"extra_tools": sorted(self._extra_tools.keys()),
})
else:
response.traces.append({"pipeline": self.name, "extra_tools": []})
return response
ApogeeProvidersGenerator
¶
ApogeeProvidersGenerator(chat_provider: IChatCompletionProvider, default_model: str = 'gpt-4o-mini', system_prompt: str | None = None)
Bridges :class:IGenerator to apogee-ai-providers chat providers.
The chat provider implements IChatCompletionProvider from
apogee_ai_providers and may be backed by any of the 14 LLMs supported
by that package (OpenAI, Anthropic, Gemini, Bedrock, OpenRouter, …).
Source code in apogee_ai_rag/infrastructure/generators/apogee_providers_generator.py
generate
async
¶
generate(prompt: str, context: list[RetrievedChunk], config: GenerationConfig | None = None) -> RagResponse
Source code in apogee_ai_rag/infrastructure/generators/apogee_providers_generator.py
async def generate(
self,
prompt: str,
context: list[RetrievedChunk],
config: GenerationConfig | None = None,
) -> RagResponse:
try:
from apogee_ai_providers import ChatRequest
except ImportError as exc: # pragma: no cover - covered indirectly
raise GenerationError(
"apogee-ai-providers is required for ApogeeProvidersGenerator"
) from exc
messages, model = self._build_messages(prompt, context, config)
request = ChatRequest(
model=model,
messages=messages,
temperature=(config.temperature if config else 0.2),
max_tokens=(config.max_tokens if config else 1024),
)
try:
chat_response = await self._chat.complete(request)
except Exception as exc: # pragma: no cover - real provider errors
raise GenerationError("chat provider failure", stage="generate", cause=exc) from exc
choice = chat_response.choices[0] if chat_response.choices else None
answer = (choice.message.content if choice and choice.message else "") or ""
finish = choice.finish_reason.value if choice and choice.finish_reason else None
usage_obj = getattr(chat_response, "usage", None)
usage: dict = {}
if usage_obj is not None:
usage = {
"prompt_tokens": getattr(usage_obj, "prompt_tokens", 0),
"completion_tokens": getattr(usage_obj, "completion_tokens", 0),
"total_tokens": getattr(usage_obj, "total_tokens", 0),
}
return RagResponse(
answer=answer,
sources=context,
finish_reason=finish,
usage=usage,
)
stream
async
¶
stream(prompt: str, context: list[RetrievedChunk], config: GenerationConfig | None = None) -> AsyncIterator[RagChunk]
Source code in apogee_ai_rag/infrastructure/generators/apogee_providers_generator.py
async def stream(
self,
prompt: str,
context: list[RetrievedChunk],
config: GenerationConfig | None = None,
) -> AsyncIterator[RagChunk]:
try:
from apogee_ai_providers import ChatRequest
except ImportError as exc: # pragma: no cover
raise GenerationError(
"apogee-ai-providers is required for ApogeeProvidersGenerator"
) from exc
messages, model = self._build_messages(prompt, context, config)
request = ChatRequest(
model=model,
messages=messages,
temperature=(config.temperature if config else 0.2),
max_tokens=(config.max_tokens if config else 1024),
stream=True,
)
try:
async for chunk in self._chat.stream(request):
yield RagChunk(
delta=chunk.delta or "",
sources=[],
finish_reason=(chunk.finish_reason.value if chunk.finish_reason else None),
)
except Exception as exc: # pragma: no cover
raise GenerationError("chat provider stream failure", stage="stream", cause=exc) from exc
yield RagChunk(delta="", sources=context, finish_reason="stop")
BM25Okapi
dataclass
¶
BM25Okapi(k1: float = 1.5, b: float = 0.75, _docs: list[list[str]] = list(), _doc_lengths: list[int] = list(), _avg_dl: float = 0.0, _df: Counter[str] = Counter(), _idf: dict[str, float] = dict())
In-memory BM25 index over a collection of token lists.
fit
¶
Source code in apogee_ai_rag/infrastructure/search/bm25.py
def fit(self, corpus: list[list[str]]) -> None:
self._docs = list(corpus)
self._doc_lengths = [len(doc) for doc in self._docs]
self._avg_dl = (sum(self._doc_lengths) / len(self._docs)) if self._docs else 0.0
self._df = Counter()
for doc in self._docs:
for term in set(doc):
self._df[term] += 1
n = len(self._docs)
self._idf = {
term: math.log(1.0 + (n - df + 0.5) / (df + 0.5))
for term, df in self._df.items()
}
score
¶
Source code in apogee_ai_rag/infrastructure/search/bm25.py
def score(self, query_tokens: list[str], doc_index: int) -> float:
if not self._docs or doc_index >= len(self._docs):
return 0.0
doc = self._docs[doc_index]
if not doc:
return 0.0
tf = Counter(doc)
dl = self._doc_lengths[doc_index]
norm = self.k1 * (1.0 - self.b + self.b * dl / self._avg_dl) if self._avg_dl else 0.0
score = 0.0
for term in query_tokens:
if term not in self._idf:
continue
f = tf[term]
if f == 0:
continue
score += self._idf[term] * (f * (self.k1 + 1.0)) / (f + norm)
return score
search
¶
BM25Search
¶
Convenience wrapper that owns a chunk index and returns RetrievedChunk.
Source code in apogee_ai_rag/infrastructure/search/bm25.py
search
¶
search(query: str, top_k: int = 10) -> list[RetrievedChunk]
Source code in apogee_ai_rag/infrastructure/search/bm25.py
CAGPipeline
¶
CAGPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, cache_store: ICacheStore | None = None, ttl_seconds: int = 3600)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/cag/cag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
cache_store: ICacheStore | None = None,
ttl_seconds: int = 3600,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"CAGPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._cache: ICacheStore = cache_store or InMemoryCacheStore()
self._ttl = ttl_seconds
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/cag/cag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
key = _cache_key(query.text, retrieved)
cached = await self._cache.get(key)
if cached is not None:
await self._emit("cag.hit", query_id=query.id, key=key)
return RagResponse(
answer=cached.value.decode("utf-8"),
sources=retrieved,
query_id=query.id,
traces=[{"pipeline": self.name, "cache_hit": True}],
)
await self._emit("cag.miss", query_id=query.id, key=key)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
response.traces.append({"pipeline": self.name, "cache_hit": False})
await self._cache.set(
CacheEntry(
key=key,
value=response.answer.encode("utf-8"),
ttl_seconds=self._ttl,
),
)
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
CRagPipeline
¶
CRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, upper_threshold: float = 0.6, lower_threshold: float = 0.2, web_search_fn: WebSearchFn | None = None, max_web_results: int = 4)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/crag/crag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
upper_threshold: float = 0.6,
lower_threshold: float = 0.2,
web_search_fn: WebSearchFn | None = None,
max_web_results: int = 4,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"CRagPipeline requires chunker, embedder, store and generator"
)
if not 0.0 <= lower_threshold <= upper_threshold <= 1.0:
raise PipelineConfigError(
"thresholds must satisfy 0 <= lower <= upper <= 1"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._upper = upper_threshold
self._lower = lower_threshold
self._web = web_search_fn
self._max_web = max_web_results
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/crag/crag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved, verdict = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
response.traces.append({"pipeline": self.name, "verdict": verdict})
await self._emit(
"query.generated",
query_id=query.id,
verdict=verdict,
finish_reason=response.finish_reason,
)
return response
stream
async
¶
CacheConfig
dataclass
¶
CacheConfig(enabled: bool = False, backend: str = 'in_memory', ttl_seconds: int = 3600, max_size: int = 1024)
CacheEntry
dataclass
¶
Chunk
dataclass
¶
Chunk(text: str, parent_id: str, position: int = 0, id: str = _new_chunk_id(), metadata: dict = dict(), modality: Modality = TEXT, embedding: list[float] = list())
embedding
class-attribute
instance-attribute
¶
ChunkStrategy
dataclass
¶
ChunkStrategy(strategy: ChunkingStrategy = RECURSIVE, chunk_size: int = 400, chunk_overlap: int = 40, separators: tuple[str, ...] = ('\n\n', '\n', '. ', ' ', ''), metadata: dict = dict())
separators
class-attribute
instance-attribute
¶
ChunkingStrategy
¶
Bases: StrEnum
ContextualRagPipeline
¶
ContextualRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, max_context_chars: int = 200)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/contextual/contextual_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
max_context_chars: int = 200,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"ContextualRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._max_context = max_context_chars
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
Source code in apogee_ai_rag/infrastructure/pipelines/contextual/contextual_rag_pipeline.py
async def ingest(self, job: IngestionJob) -> IngestionResult:
await self._emit(
"ingestion.started",
job_id=job.id,
n_documents=len(job.documents),
chunker=self._chunker.name,
embedder=self._embedder.model,
store=self._store.name,
)
chunks = await self._chunker.chunk(job.documents)
await self._emit("ingestion.chunked", job_id=job.id, n_chunks=len(chunks))
if not chunks:
return IngestionResult(
job_id=job.id,
documents_ingested=len(job.documents),
chunks_produced=0,
vectors_upserted=0,
)
contextualised: list[Chunk] = []
for chunk in chunks:
contextualised.append(await self._contextualise(chunk))
await self._emit(
"contextual.contextualised",
job_id=job.id,
n_chunks=sum(1 for c in contextualised if c.metadata.get("contextualised")),
)
embeddings = await self._embedder.embed([c.text for c in contextualised])
for chunk, vec in zip(contextualised, embeddings, strict=False):
chunk.embedding = vec
upserted = await self._store.upsert(contextualised)
await self._emit("ingestion.upserted", job_id=job.id, n_upserted=upserted)
return IngestionResult(
job_id=job.id,
documents_ingested=len(job.documents),
chunks_produced=len(contextualised),
vectors_upserted=upserted,
)
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/contextual/contextual_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/contextual/contextual_rag_pipeline.py
ConversationalRagPipeline
¶
ConversationalRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, history_window: int = 6)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/conversational/conversational_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
history_window: int = 6,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"ConversationalRagPipeline requires chunker, embedder, store and generator"
)
if history_window < 1:
raise PipelineConfigError("history_window must be >= 1")
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._history_window = history_window
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/conversational/conversational_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved, expanded = await self._retrieve(query)
response = await self._generator.generate(
expanded, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
response.traces.append({"pipeline": self.name, "history_turns": len(query.history)})
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/conversational/conversational_rag_pipeline.py
DiskCacheStore
¶
Source code in apogee_ai_rag/infrastructure/cache_stores/disk_cache_store.py
get
async
¶
get(key: str) -> CacheEntry | None
Source code in apogee_ai_rag/infrastructure/cache_stores/disk_cache_store.py
async def get(self, key: str) -> CacheEntry | None:
path = self._path(key)
if not path.is_file():
return None
meta = self._meta(key)
ttl = None
if meta.is_file():
try:
ttl = float(meta.read_text(encoding="utf-8").strip())
except ValueError:
ttl = None
if ttl is not None and ttl < time.time():
for p in (path, meta, self._key_file(key)):
await asyncio.to_thread(p.unlink, missing_ok=True)
return None
value = await asyncio.to_thread(path.read_bytes)
return CacheEntry(key=key, value=value, ttl_seconds=None)
set
async
¶
set(entry: CacheEntry) -> None
Source code in apogee_ai_rag/infrastructure/cache_stores/disk_cache_store.py
async def set(self, entry: CacheEntry) -> None:
await asyncio.to_thread(self._path(entry.key).write_bytes, entry.value)
await asyncio.to_thread(
self._key_file(entry.key).write_text, entry.key, encoding="utf-8",
)
if entry.ttl_seconds is not None:
await asyncio.to_thread(
self._meta(entry.key).write_text,
str(time.time() + entry.ttl_seconds),
encoding="utf-8",
)
invalidate
async
¶
Source code in apogee_ai_rag/infrastructure/cache_stores/disk_cache_store.py
async def invalidate(self, prefix: str) -> int:
removed = 0
for kfile in list(self._dir.glob("*.key")):
try:
stored_key = kfile.read_text(encoding="utf-8")
except OSError:
continue
if not stored_key.startswith(prefix):
continue
stem = kfile.stem
for suffix in (".bin", ".meta", ".key"):
p = self._dir / f"{stem}{suffix}"
await asyncio.to_thread(p.unlink, missing_ok=True)
removed += 1
return removed
Document
dataclass
¶
Document(text: str = '', metadata: dict = dict(), id: str = _new_doc_id(), modality: Modality = TEXT, image_b64: str | None = None, audio_url: str | None = None, video_url: str | None = None, embedding: list[float] = list())
embedding
class-attribute
instance-attribute
¶
EchoGenerator
¶
Deterministic generator that mirrors the prompt and cites sources.
Used by the smoke E2E test (no API key needed) and as a sane default for
pipelines created without an explicit IGenerator.
generate
async
¶
generate(prompt: str, context: list[RetrievedChunk], config: GenerationConfig | None = None) -> RagResponse
Source code in apogee_ai_rag/infrastructure/generators/echo_generator.py
async def generate(
self,
prompt: str,
context: list[RetrievedChunk],
config: GenerationConfig | None = None,
) -> RagResponse:
bullets = "\n".join(
f"- [{i + 1}] {c.chunk.text}" for i, c in enumerate(context)
)
body = (
f"Question: {prompt}\n\n"
f"Context:\n{bullets}" if bullets else f"Question: {prompt}\n\n(no context)"
)
return RagResponse(
answer=body,
sources=context,
confidence=1.0 if context else 0.0,
finish_reason="stop",
usage={"prompt_chars": len(prompt), "context_chunks": len(context)},
)
stream
async
¶
stream(prompt: str, context: list[RetrievedChunk], config: GenerationConfig | None = None) -> AsyncIterator[RagChunk]
Source code in apogee_ai_rag/infrastructure/generators/echo_generator.py
async def stream(
self,
prompt: str,
context: list[RetrievedChunk],
config: GenerationConfig | None = None,
) -> AsyncIterator[RagChunk]:
response = await self.generate(prompt, context, config)
for word in response.answer.split(" "):
yield RagChunk(delta=word + " ", sources=[])
yield RagChunk(delta="", sources=response.sources, finish_reason="stop")
EmbeddingConfig
dataclass
¶
EmbeddingConfig(provider: EmbedderProvider = HASHING, model: str = 'hashing-128', dims: int = 128, batch_size: int = 32)
EvalCase
dataclass
¶
EvalCase(question: str, ground_truth: str | None = None, expected_substrings: list[str] = list(), contexts: list[str] = list(), metadata: dict = dict())
EvalConfig
dataclass
¶
EvalConfig(evaluator: Evaluator = RAGAS, metrics: tuple[str, ...] = ('faithfulness', 'answer_relevancy', 'context_precision', 'context_recall'), threshold: float = 0.7, options: dict = dict())
metrics
class-attribute
instance-attribute
¶
metrics: tuple[str, ...] = ('faithfulness', 'answer_relevancy', 'context_precision', 'context_recall')
EvalReport
dataclass
¶
EvalReport(metrics: dict[str, float] = dict(), per_case: list[EvalResult] = list(), passed: int = 0, failed: int = 0)
metrics
class-attribute
instance-attribute
¶
per_case
class-attribute
instance-attribute
¶
per_case: list[EvalResult] = field(default_factory=list)
EvalResult
dataclass
¶
EvalResult(case: EvalCase, answer: str, metrics: dict = dict(), passed: bool = False, error: str | None = None)
Evaluator
¶
Bases: StrEnum
FLAREPipeline
¶
FLAREPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, max_re_retrievals: int = 2, min_segment_chars: int = 32)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/flare/flare_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
max_re_retrievals: int = 2,
min_segment_chars: int = 32,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"FLAREPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._max_re = max_re_retrievals
self._min_segment = min_segment_chars
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/flare/flare_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
await self._emit(
"query.received", query_id=query.id, text=query.text, pipeline_type=self.name,
)
retrieved = await self._retrieve(query.text, query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
triggers = 0
while (
triggers < self._max_re
and _looks_uncertain(response.answer, self._min_segment)
):
triggers += 1
await self._emit(
"flare.uncertain", query_id=query.id, attempt=triggers,
)
extra = await self._retrieve(
f"{query.text} {response.answer}", query,
)
seen = {r.chunk.id for r in retrieved}
retrieved.extend(r for r in extra if r.chunk.id not in seen)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
response.traces.append({"pipeline": self.name, "re_retrievals": triggers})
await self._emit(
"query.generated", query_id=query.id, re_retrievals=triggers,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/flare/flare_rag_pipeline.py
FeedbackInput
dataclass
¶
Framework
¶
Bases: StrEnum
GenerationConfig
dataclass
¶
GenerationConfig(model: str = 'gpt-4o-mini', temperature: float = 0.2, max_tokens: int | None = 1024, system_prompt: str | None = None, stream: bool = False)
GraphConfig
dataclass
¶
GraphConfig(enable_communities: bool = True, community_algorithm: str = 'leiden', max_depth: int = 2, entity_extractor: str = 'llm')
GraphRagPipeline
¶
GraphRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, graph_store: IGraphStore | None = None, traversal_depth: int = 1)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/graph/graph_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
graph_store: IGraphStore | None = None,
traversal_depth: int = 1,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"GraphRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._graph: IGraphStore = graph_store or NetworkXGraphStore()
self._depth = traversal_depth
self._community_summaries: list[tuple[set[str], str]] = []
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
Source code in apogee_ai_rag/infrastructure/pipelines/graph/graph_rag_pipeline.py
async def ingest(self, job: IngestionJob) -> IngestionResult:
result = await self._ingest_default(job)
chunks = list(getattr(self._store, "_chunks", {}).values()) # type: ignore[attr-defined]
# Build the entity graph.
for chunk in chunks:
entities = _extract_entities(chunk.text)
for entity in entities:
await self._graph.add_node(
KnowledgeNode(label=entity, properties={"chunk_id": chunk.id}),
)
for i, src in enumerate(entities):
for dst in entities[i + 1 :]:
await self._graph.add_edge(
KnowledgeEdge(source_id=src, target_id=dst, label="co_occurs"),
)
# Cache one short summary per community.
communities: list[list[str]] = getattr(self._graph, "communities", lambda: [])()
self._community_summaries = []
for component in communities:
tokens = Counter()
for entity in component:
for chunk in chunks:
if entity in chunk.text:
tokens.update(chunk.text.lower().split())
top = " ".join(t for t, _ in tokens.most_common(10))
self._community_summaries.append((set(component), top))
all_nodes_fn = getattr(self._graph, "all_nodes", None)
n_entities = 0
if callable(all_nodes_fn):
try:
nodes_value = all_nodes_fn()
n_entities = sum(1 for _ in nodes_value) # type: ignore[unused-ignore]
except Exception:
n_entities = 0
await self._emit(
"graph.indexed",
job_id=job.id,
communities=len(self._community_summaries),
entities=n_entities,
)
return result
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/graph/graph_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
response.traces.append({"pipeline": self.name, "n_sources": len(retrieved)})
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/graph/graph_rag_pipeline.py
HashingEmbedder
¶
Deterministic, dependency-free token-hashing embedder.
Useful for tests, offline development and smoke pipelines without API keys. Output vectors are L2-normalised so cosine similarity matches dot product.
Source code in apogee_ai_rag/infrastructure/embeddings/hashing_embedder.py
embed
async
¶
Source code in apogee_ai_rag/infrastructure/embeddings/hashing_embedder.py
HeuristicQueryRewriter
¶
Generates a small set of paraphrases without calling any LLM.
Strategies:
- the original query;
- keyword-only form (drop interrogative prefixes and stopwords);
- a contextual phrasing prepended with about.
Good enough as a baseline before plugging :class:LLMQueryRewriter.
Source code in apogee_ai_rag/infrastructure/query_rewriters/heuristic_query_rewriter.py
rewrite
async
¶
rewrite(query: RagQuery) -> list[str]
Source code in apogee_ai_rag/infrastructure/query_rewriters/heuristic_query_rewriter.py
async def rewrite(self, query: RagQuery) -> list[str]:
text = query.text.strip()
if not text:
return [""]
variants: list[str] = [text]
lowered = text.lower()
for prefix in _QUESTION_PREFIXES:
if lowered.startswith(prefix):
variants.append(text[len(prefix) :].strip(" ?."))
break
keywords = [
tok for tok in tokenize(text)
if tok not in self._STOPWORDS and len(tok) > 2
]
if keywords:
variants.append(" ".join(keywords))
variants.append("about " + " ".join(keywords[:5]))
seen: set[str] = set()
deduped: list[str] = []
for variant in variants:
if variant and variant not in seen:
seen.add(variant)
deduped.append(variant)
return deduped[: self._max]
HierarchicalRagPipeline
¶
HierarchicalRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/hierarchical/hierarchical_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"HierarchicalRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/hierarchical/hierarchical_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/hierarchical/hierarchical_rag_pipeline.py
HyDERagPipeline
¶
HyDERagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, hypothesis_max_chars: int = 600)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/hyde/hyde_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
hypothesis_max_chars: int = 600,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"HyDERagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._hypothesis_cap = hypothesis_max_chars
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/hyde/hyde_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.query_id = query.id
await self._emit(
"query.generated", query_id=query.id, finish_reason=response.finish_reason,
)
return response
stream
async
¶
HybridRagPipeline
¶
HybridRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, rrf_k: int = 60)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/hybrid/hybrid_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
rrf_k: int = 60,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"HybridRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._bm25 = BM25Search()
self._chunks: list[Chunk] = []
self._rrf_k = rrf_k
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
Source code in apogee_ai_rag/infrastructure/pipelines/hybrid/hybrid_rag_pipeline.py
async def ingest(self, job: IngestionJob) -> IngestionResult:
result = await self._ingest_default(job)
# Rebuild BM25 from the in-memory store when available; otherwise
# accumulate the freshly chunked input.
chunks = list(getattr(self._store, "_chunks", {}).values()) # type: ignore[attr-defined]
if not chunks:
chunks = self._chunks + await self._chunker.chunk(job.documents)
self._chunks = chunks
self._bm25.fit(self._chunks)
return result
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/hybrid/hybrid_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.query_id = query.id
await self._emit(
"query.generated", query_id=query.id, finish_reason=response.finish_reason,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/hybrid/hybrid_rag_pipeline.py
InMemoryCacheStore
¶
Source code in apogee_ai_rag/infrastructure/cache_stores/in_memory_cache_store.py
get
async
¶
get(key: str) -> CacheEntry | None
set
async
¶
set(entry: CacheEntry) -> None
Source code in apogee_ai_rag/infrastructure/cache_stores/in_memory_cache_store.py
async def set(self, entry: CacheEntry) -> None:
if len(self._store) >= self._max:
# Drop the oldest entry to keep the cache bounded.
oldest_key = next(iter(self._store))
self._store.pop(oldest_key, None)
expiry = (
time.time() + entry.ttl_seconds if entry.ttl_seconds is not None else None
)
self._store[entry.key] = (entry, expiry)
invalidate
async
¶
size
async
¶
InMemoryObservabilityEmitter
¶
InMemoryVectorStore
¶
Cosine-similarity store for the naive RAG MVP.
Holds chunks in a dict keyed by chunk.id. Supports filters via exact
match on chunk.metadata and a lightweight hybrid mode that mixes
cosine similarity with a normalised token-overlap score (sparse stand-in).
Source code in apogee_ai_rag/infrastructure/vector_stores/in_memory_vector_store.py
search
async
¶
search(query_embedding: list[float], top_k: int = 5, threshold: float = 0.0, filters: dict | None = None) -> list[RetrievedChunk]
Source code in apogee_ai_rag/infrastructure/vector_stores/in_memory_vector_store.py
async def search(
self,
query_embedding: list[float],
top_k: int = 5,
threshold: float = 0.0,
filters: dict | None = None,
) -> list[RetrievedChunk]:
scored: list[RetrievedChunk] = []
for chunk in self._chunks.values():
if not chunk.embedding:
continue
if not _matches_filters(chunk.metadata, filters):
continue
score = _cosine(query_embedding, chunk.embedding)
scored.append(RetrievedChunk(chunk=chunk, score=score, retriever="vector"))
scored.sort(key=lambda r: r.score, reverse=True)
return [r for r in scored[:top_k] if r.score >= threshold]
hybrid_search
async
¶
hybrid_search(query_text: str, query_embedding: list[float], top_k: int = 5, threshold: float = 0.0, filters: dict | None = None, alpha: float = 0.5) -> list[RetrievedChunk]
Source code in apogee_ai_rag/infrastructure/vector_stores/in_memory_vector_store.py
async def hybrid_search(
self,
query_text: str,
query_embedding: list[float],
top_k: int = 5,
threshold: float = 0.0,
filters: dict | None = None,
alpha: float = 0.5,
) -> list[RetrievedChunk]:
query_tokens = set(_TOKEN_RE.findall(query_text.lower()))
scored: list[RetrievedChunk] = []
for chunk in self._chunks.values():
if not _matches_filters(chunk.metadata, filters):
continue
dense = _cosine(query_embedding, chunk.embedding) if chunk.embedding else 0.0
sparse = _bm25ish_score(query_tokens, set(_TOKEN_RE.findall(chunk.text.lower())))
blended = alpha * dense + (1.0 - alpha) * sparse
scored.append(
RetrievedChunk(chunk=chunk, score=blended, retriever="hybrid"),
)
scored.sort(key=lambda r: r.score, reverse=True)
return [r for r in scored[:top_k] if r.score >= threshold]
delete
async
¶
count
async
¶
IngestionJob
dataclass
¶
IngestionJob(documents: list[Document], batch_size: int = 64, idempotency_key: str | None = None, id: str = (lambda: f'job_{hex[:10]}')(), metadata: dict = dict())
id
class-attribute
instance-attribute
¶
IngestionResult
dataclass
¶
IngestionResult(job_id: str, documents_ingested: int, chunks_produced: int, vectors_upserted: int, duration_ms: float = 0.0, errors: list[str] = list())
IterativeRagPipeline
¶
IterativeRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, n_rounds: int = 2)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/iterative/iterative_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
n_rounds: int = 2,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"IterativeRagPipeline requires chunker, embedder, store and generator"
)
if n_rounds < 1:
raise PipelineConfigError("n_rounds must be >= 1")
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._n_rounds = n_rounds
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/iterative/iterative_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
await self._emit(
"query.received", query_id=query.id, text=query.text, pipeline_type=self.name,
)
seen: set[str] = set()
retrieved: list[RetrievedChunk] = []
last_response: RagResponse | None = None
current_text = query.text
for round_idx in range(1, self._n_rounds + 1):
[embedding] = await self._embedder.embed([current_text])
results = await self._store.search(
embedding,
top_k=query.top_k,
threshold=query.threshold,
filters=query.filters or None,
)
new_items = [r for r in results if r.chunk.id not in seen]
for r in new_items:
seen.add(r.chunk.id)
retrieved.extend(new_items)
last_response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
await self._emit(
"iterative.round",
query_id=query.id,
round=round_idx,
added=len(new_items),
)
if not new_items:
break
current_text = f"{query.text} {last_response.answer}"
retrieved.append(
RetrievedChunk(
chunk=Chunk(
text=last_response.answer[:400],
parent_id=f"iterative_thought_{round_idx}",
),
score=0.0,
retriever="iterative_thought",
),
)
assert last_response is not None
last_response.sources = [
r for r in retrieved if r.retriever != "iterative_thought"
]
last_response.query_id = query.id
await self._emit("query.generated", query_id=query.id)
return last_response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/iterative/iterative_rag_pipeline.py
KAGPipeline
¶
KAGPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, graph_store: IGraphStore | None = None)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/kag/kag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
graph_store: IGraphStore | None = None,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"KAGPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._graph: IGraphStore = graph_store or NetworkXGraphStore()
self._entity_index: dict[str, list[str]] = {}
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
Source code in apogee_ai_rag/infrastructure/pipelines/kag/kag_pipeline.py
async def ingest(self, job: IngestionJob) -> IngestionResult:
result = await self._ingest_default(job)
chunks = list(getattr(self._store, "_chunks", {}).values()) # type: ignore[attr-defined]
self._entity_index = {}
for chunk in chunks:
for entity in _extract_entities(chunk.text):
await self._graph.add_node(
KnowledgeNode(
label=entity, properties={"chunk_id": chunk.id, "kind": "entity"},
),
)
self._entity_index.setdefault(entity, []).append(chunk.id)
await self._emit(
"kag.indexed", job_id=job.id, entities=len(self._entity_index),
)
return result
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/kag/kag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
response.traces.append(
{"pipeline": self.name, "entities_in_query": _extract_entities(query.text)},
)
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
KnowledgeEdge
dataclass
¶
KnowledgeEdge(source_id: str, target_id: str, label: str, weight: float = 1.0, properties: dict = dict(), id: str = (lambda: f'e_{hex[:10]}')())
KnowledgeNode
dataclass
¶
LLMQueryRewriter
¶
LLMQueryRewriter(generator: IGenerator, n: int = 3, config: GenerationConfig | None = None)
Generates rewrites via an :class:IGenerator (typically an LLM).
The default prompt asks the model for n numbered alternatives. This
works with any provider supported by apogee-ai-providers once you
inject :class:ApogeeProvidersGenerator. With EchoGenerator it
falls back to returning the original query (the echo just mirrors back).
Source code in apogee_ai_rag/infrastructure/query_rewriters/llm_query_rewriter.py
rewrite
async
¶
rewrite(query: RagQuery) -> list[str]
Source code in apogee_ai_rag/infrastructure/query_rewriters/llm_query_rewriter.py
async def rewrite(self, query: RagQuery) -> list[str]:
prompt = (
f"Generate {self._n} alternative phrasings of the following user "
f"question that would help retrieve relevant context. Output one "
f"alternative per line, no numbering or extra commentary.\n\n"
f"Question: {query.text}"
)
response = await self._generator.generate(prompt, [], self._config)
candidates = [
line.strip().lstrip("-•0123456789.) ").strip()
for line in (response.answer or "").splitlines()
]
candidates = [c for c in candidates if c]
rewrites = [query.text] + candidates
seen: set[str] = set()
deduped: list[str] = []
for r in rewrites:
if r not in seen:
seen.add(r)
deduped.append(r)
return deduped[: self._n + 1]
LongRagPipeline
¶
LongRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, long_chunk_size: int = 4000, long_chunk_overlap: int = 200, long_top_k: int = 2)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/long/long_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
long_chunk_size: int = 4000,
long_chunk_overlap: int = 200,
long_top_k: int = 2,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"LongRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._strategy = ChunkStrategy(
chunk_size=long_chunk_size, chunk_overlap=long_chunk_overlap,
)
self._long_top_k = long_top_k
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
Source code in apogee_ai_rag/infrastructure/pipelines/long/long_rag_pipeline.py
async def ingest(self, job: IngestionJob) -> IngestionResult:
await self._emit(
"ingestion.started",
job_id=job.id,
n_documents=len(job.documents),
chunker=self._chunker.name,
embedder=self._embedder.model,
store=self._store.name,
chunk_size=self._strategy.chunk_size,
)
chunks = await self._chunker.chunk(job.documents, self._strategy)
await self._emit("ingestion.chunked", job_id=job.id, n_chunks=len(chunks))
if not chunks:
return IngestionResult(
job_id=job.id,
documents_ingested=len(job.documents),
chunks_produced=0,
vectors_upserted=0,
)
embeddings = await self._embedder.embed([c.text for c in chunks])
for chunk, vec in zip(chunks, embeddings, strict=False):
chunk.embedding = vec
upserted = await self._store.upsert(chunks)
await self._emit("ingestion.upserted", job_id=job.id, n_upserted=upserted)
return IngestionResult(
job_id=job.id,
documents_ingested=len(job.documents),
chunks_produced=len(chunks),
vectors_upserted=upserted,
)
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/long/long_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
await self._emit("query.generated", query_id=query.id, n_sources=len(retrieved))
return response
stream
async
¶
Modality
¶
ModularRagPipeline
¶
ModularRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, stages: list[Stage] | None = None, reranker: IReranker | None = None)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/modular/modular_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
stages: list[Stage] | None = None,
reranker: IReranker | None = None,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"ModularRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._reranker = reranker or NoOpReranker()
self._stages = stages or self.default_stages()
default_stages
¶
Source code in apogee_ai_rag/infrastructure/pipelines/modular/modular_rag_pipeline.py
def default_stages(self) -> list[Stage]:
async def embed_query(state: ModularState) -> ModularState:
[vec] = await self._embedder.embed([state.query.text])
state.embedding = vec
return state
async def retrieve(state: ModularState) -> ModularState:
state.retrieved = await self._store.search(
state.embedding,
top_k=state.query.top_k,
threshold=state.query.threshold,
filters=state.query.filters or None,
)
return state
async def rerank(state: ModularState) -> ModularState:
state.retrieved = await self._reranker.rerank(
state.query.text, state.retrieved, top_n=state.query.top_k,
)
return state
async def generate(state: ModularState) -> ModularState:
state.response = await self._generator.generate(
state.query.text, state.retrieved, self._generation_config,
)
return state
return [embed_query, retrieve, rerank, generate]
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/modular/modular_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
await self._emit(
"query.received", query_id=query.id, text=query.text, pipeline_type=self.name,
)
state = ModularState(query=query)
for stage in self._stages:
state = await stage(state)
await self._emit(
"modular.stage",
query_id=query.id,
stage=getattr(stage, "__name__", repr(stage)),
)
if state.response is None:
raise PipelineConfigError("Modular pipeline finished without producing a response")
state.response.sources = state.retrieved
state.response.query_id = query.id
await self._emit("query.generated", query_id=query.id)
return state.response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/modular/modular_rag_pipeline.py
ModularState
dataclass
¶
ModularState(query: RagQuery, embedding: list[float] = list(), retrieved: list[RetrievedChunk] = list(), response: RagResponse | None = None, bag: dict[str, Any] = dict())
embedding
class-attribute
instance-attribute
¶
retrieved
class-attribute
instance-attribute
¶
retrieved: list[RetrievedChunk] = field(default_factory=list)
MultiHopRagPipeline
¶
MultiHopRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, max_hops: int = 3)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/multi_hop/multi_hop_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
max_hops: int = 3,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"MultiHopRagPipeline requires chunker, embedder, store and generator"
)
if max_hops < 1:
raise PipelineConfigError("max_hops must be >= 1")
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._max_hops = max_hops
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/multi_hop/multi_hop_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/multi_hop/multi_hop_rag_pipeline.py
MultiModalRagPipeline
¶
MultiModalRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, image_caption_max_chars: int = 256)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/multi_modal/multi_modal_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
image_caption_max_chars: int = 256,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"MultiModalRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._image_caption_max = image_caption_max_chars
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/multi_modal/multi_modal_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/multi_modal/multi_modal_rag_pipeline.py
NaiveRagPipeline
¶
NaiveRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, reranker: IReranker | None = None, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None)
The simplest RAG pipeline — embed, search, optionally rerank, generate.
Source code in apogee_ai_rag/infrastructure/pipelines/naive/naive_rag_pipeline.py
def __init__(
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
reranker: IReranker | None = None,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
) -> None:
if chunker is None or embedder is None or vector_store is None or generator is None:
raise PipelineConfigError(
"NaiveRagPipeline requires chunker, embedder, vector_store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._reranker = reranker
self._emitter = observability or NoOpObservabilityEmitter()
self._generation_config = generation_config
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
Source code in apogee_ai_rag/infrastructure/pipelines/naive/naive_rag_pipeline.py
async def ingest(self, job: IngestionJob) -> IngestionResult:
started = time.perf_counter()
await self._emit(
"ingestion.started",
job_id=job.id,
n_documents=len(job.documents),
chunker=self._chunker.name,
embedder=self._embedder.model,
store=self._store.name,
)
chunks = await self._chunker.chunk(job.documents)
await self._emit("ingestion.chunked", job_id=job.id, n_chunks=len(chunks))
if not chunks:
return IngestionResult(
job_id=job.id,
documents_ingested=len(job.documents),
chunks_produced=0,
vectors_upserted=0,
duration_ms=(time.perf_counter() - started) * 1000.0,
)
embeddings = await self._embedder.embed([c.text for c in chunks])
embedded: list[Chunk] = []
for c, e in zip(chunks, embeddings, strict=False):
c.embedding = e
embedded.append(c)
await self._emit(
"ingestion.embedded",
job_id=job.id,
n_vectors=len(embedded),
dims=self._embedder.dims,
)
upserted = 0
batch = max(job.batch_size, 1)
for start in range(0, len(embedded), batch):
upserted += await self._store.upsert(embedded[start : start + batch])
await self._emit("ingestion.upserted", job_id=job.id, n_upserted=upserted)
return IngestionResult(
job_id=job.id,
documents_ingested=len(job.documents),
chunks_produced=len(chunks),
vectors_upserted=upserted,
duration_ms=(time.perf_counter() - started) * 1000.0,
)
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/naive/naive_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(query.text, retrieved, self._generation_config)
response.query_id = query.id
await self._emit(
"query.generated",
query_id=query.id,
tokens=response.usage.get("total_tokens", 0),
finish_reason=response.finish_reason,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/naive/naive_rag_pipeline.py
Neo4jGraphStore
¶
Source code in apogee_ai_rag/infrastructure/graph_stores/neo4j_graph_store.py
def __init__(self, *, uri: str, auth: tuple[str, str] | None = None) -> None:
try:
import neo4j # type: ignore[import-untyped] # noqa: F401
except ImportError as exc:
raise ProviderNotInstalledError("neo4j", "neo4j") from exc
# pragma: no cover — full driver wiring is intentionally deferred.
self._uri = uri
self._auth = auth
add_node
async
¶
add_node(node: KnowledgeNode) -> str
add_edge
async
¶
add_edge(edge: KnowledgeEdge) -> str
traverse
async
¶
traverse(start_id: str, depth: int = 2) -> tuple[list[KnowledgeNode], list[KnowledgeEdge]]
NetworkXGraphStore
¶
Source code in apogee_ai_rag/infrastructure/graph_stores/networkx_graph_store.py
add_node
async
¶
add_node(node: KnowledgeNode) -> str
add_edge
async
¶
add_edge(edge: KnowledgeEdge) -> str
traverse
async
¶
traverse(start_id: str, depth: int = 2) -> tuple[list[KnowledgeNode], list[KnowledgeEdge]]
Source code in apogee_ai_rag/infrastructure/graph_stores/networkx_graph_store.py
async def traverse(
self, start_id: str, depth: int = 2,
) -> tuple[list[KnowledgeNode], list[KnowledgeEdge]]:
if start_id not in self._nodes:
return [], []
visited_nodes: set[str] = {start_id}
visited_edges: list[KnowledgeEdge] = []
queue: deque[tuple[str, int]] = deque([(start_id, 0)])
while queue:
current, level = queue.popleft()
if level >= depth:
continue
for edge in self._adj.get(current, []):
if edge.id in {e.id for e in visited_edges}:
continue
visited_edges.append(edge)
neighbour = (
edge.target_id if edge.source_id == current else edge.source_id
)
if neighbour not in visited_nodes:
visited_nodes.add(neighbour)
queue.append((neighbour, level + 1))
nodes = [self._nodes[nid] for nid in visited_nodes if nid in self._nodes]
return nodes, visited_edges
all_nodes
¶
all_nodes() -> list[KnowledgeNode]
all_edges
¶
all_edges() -> list[KnowledgeEdge]
Source code in apogee_ai_rag/infrastructure/graph_stores/networkx_graph_store.py
communities
¶
Returns weakly-connected components as a community partition.
Source code in apogee_ai_rag/infrastructure/graph_stores/networkx_graph_store.py
def communities(self) -> list[list[str]]:
"""Returns weakly-connected components as a community partition."""
visited: set[str] = set()
communities: list[list[str]] = []
for node_id in self._nodes:
if node_id in visited:
continue
stack = [node_id]
component: list[str] = []
while stack:
current = stack.pop()
if current in visited:
continue
visited.add(current)
component.append(current)
for edge in self._adj.get(current, []):
neighbour = (
edge.target_id if edge.source_id == current else edge.source_id
)
if neighbour not in visited:
stack.append(neighbour)
communities.append(component)
return communities
NoOpObservabilityEmitter
¶
NoOpReranker
¶
Default reranker that preserves the upstream order, only applying top_n.
rerank
async
¶
rerank(query: str, chunks: list[RetrievedChunk], top_n: int = 5) -> list[RetrievedChunk]
ObservabilityProvider
¶
Bases: StrEnum
PassThroughRewriter
¶
Returns the original query verbatim — used as the default.
PipelineRuntimeConfig
dataclass
¶
PipelineRuntimeConfig(max_iterations: int = 5, timeout_seconds: float = 60.0, enable_streaming: bool = False, enable_caching: bool = False, enable_observability: bool = True)
PipelineSpec
dataclass
¶
PipelineSpec(rag_type: RagType = NAIVE, chunker: Any | None = None, embedder: Any | None = None, vector_store: Any | None = None, retriever: Any | None = None, generator: Any | None = None, reranker: Any | None = None, query_rewriter: Any | None = None, evaluator: Any | None = None, observability: Any | None = None, graph_store: Any | None = None, cache_store: Any | None = None, options: dict = dict())
Aggregate of components that compose a RAG pipeline.
Each field is a Protocol implementation (chunker, embedder, vector store,
generator, etc.). Components left as None are either optional or filled
in by the factory based on rag_type.
QueryRewritingRagPipeline
¶
QueryRewritingRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, query_rewriter: IQueryRewriter | None = None, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/query_rewriting/query_rewriting_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
query_rewriter: IQueryRewriter | None = None,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"QueryRewritingRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._rewriter = query_rewriter or HeuristicQueryRewriter()
self._emitter = default_emitter(observability)
self._generation_config = generation_config
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/query_rewriting/query_rewriting_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.query_id = query.id
await self._emit(
"query.generated", query_id=query.id, finish_reason=response.finish_reason,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/query_rewriting/query_rewriting_rag_pipeline.py
RAPTORPipeline
¶
RAPTORPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, max_levels: int = 2, cluster_similarity: float = 0.4, min_cluster_size: int = 2)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/raptor/raptor_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
max_levels: int = 2,
cluster_similarity: float = 0.4,
min_cluster_size: int = 2,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"RAPTORPipeline requires chunker, embedder, store and generator"
)
if max_levels < 1:
raise PipelineConfigError("max_levels must be >= 1")
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._max_levels = max_levels
self._cluster_threshold = cluster_similarity
self._min_cluster_size = min_cluster_size
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
Source code in apogee_ai_rag/infrastructure/pipelines/raptor/raptor_rag_pipeline.py
async def ingest(self, job: IngestionJob) -> IngestionResult:
result = await self._ingest_default(job)
# Build summary levels on top of the freshly upserted chunks.
current_chunks = list(getattr(self._store, "_chunks", {}).values()) # type: ignore[attr-defined]
for level in range(1, self._max_levels + 1):
clusters = _greedy_cluster(current_chunks, self._cluster_threshold)
qualifying = [c for c in clusters if len(c.members) >= self._min_cluster_size]
await self._emit(
"raptor.level",
job_id=job.id,
level=level,
clusters=len(qualifying),
)
if not qualifying:
break
summaries: list[Chunk] = []
for idx, cluster in enumerate(qualifying):
merged_text = "\n".join(m.text for m in cluster.members)
response = await self._generator.generate(
f"Summarise the following passages in 2-3 sentences:\n\n{merged_text}",
[], self._generation_config,
)
summary_text = (response.answer or merged_text[:400]).strip()
summaries.append(
Chunk(
text=summary_text,
parent_id=f"raptor_l{level}_c{idx}",
metadata={"raptor_level": level},
),
)
embeddings = await self._embedder.embed([s.text for s in summaries])
for summary, vec in zip(summaries, embeddings, strict=False):
summary.embedding = vec
await self._store.upsert(summaries)
result.chunks_produced += len(summaries)
result.vectors_upserted += len(summaries)
current_chunks = summaries
return result
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/raptor/raptor_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/raptor/raptor_rag_pipeline.py
RagChunk
dataclass
¶
RagChunk(delta: str, sources: list[RetrievedChunk] = list(), finish_reason: str | None = None)
Streaming delta returned by IRagPipeline.stream.
sources
class-attribute
instance-attribute
¶
sources: list[RetrievedChunk] = field(default_factory=list)
RagEvent
dataclass
¶
RagFactory
¶
Selects and instantiates a concrete pipeline implementing :class:IRagPipeline.
build
staticmethod
¶
build(rag_type: RagType, spec: PipelineSpec) -> IRagPipeline
RagFusionPipeline
¶
RagFusionPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, query_rewriter: IQueryRewriter | None = None, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, rrf_k: int = 60)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/rag_fusion/rag_fusion_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
query_rewriter: IQueryRewriter | None = None,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
rrf_k: int = 60,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"RagFusionPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._rewriter = query_rewriter or HeuristicQueryRewriter(max_variants=4)
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._rrf_k = rrf_k
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/rag_fusion/rag_fusion_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.query_id = query.id
await self._emit(
"query.generated", query_id=query.id, finish_reason=response.finish_reason,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/rag_fusion/rag_fusion_pipeline.py
RagProviderCredentials
dataclass
¶
RagProviderCredentials(api_key: str | None = None, endpoint: str | None = None, region: str | None = None, project: str | None = None, namespace: str | None = None, extra: dict = dict())
RagQuery
dataclass
¶
RagQuery(text: str, filters: dict = dict(), top_k: int = 5, threshold: float = 0.0, modality: Modality = TEXT, history: list[dict] = list(), conversation_id: str | None = None, locale: str | None = None, image_b64: str | None = None, id: str = (lambda: f'q_{hex[:10]}')())
RagResponse
dataclass
¶
RagResponse(answer: str, sources: list[RetrievedChunk] = list(), confidence: float | None = None, finish_reason: str | None = None, usage: dict = dict(), traces: list[dict] = list(), query_id: str | None = None)
sources
class-attribute
instance-attribute
¶
sources: list[RetrievedChunk] = field(default_factory=list)
ReActRagPipeline
¶
ReActRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, max_steps: int = 4)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/react/react_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
max_steps: int = 4,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"ReActRagPipeline requires chunker, embedder, store and generator"
)
if max_steps < 1:
raise PipelineConfigError("max_steps must be >= 1")
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._max_steps = max_steps
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/react/react_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
await self._emit(
"query.received", query_id=query.id, text=query.text, pipeline_type=self.name,
)
seen_ids: set[str] = set()
evidence: list[RetrievedChunk] = []
scratchpad = f"Question: {query.text}\n"
last_answer = ""
for step in range(1, self._max_steps + 1):
thought = await self._generator.generate(
f"{scratchpad}\nThought:", [], self._generation_config,
)
action_query = (thought.answer or query.text).split("\n")[0][:200]
results = await self._search(action_query, query)
new_items = [r for r in results if r.chunk.id not in seen_ids]
if not new_items:
break
for r in new_items:
seen_ids.add(r.chunk.id)
evidence.extend(new_items)
observation = "\n".join(f"- {r.chunk.text}" for r in new_items[:2])
scratchpad += (
f"\nThought: {thought.answer.strip()[:200]}\n"
f"Action: search[{action_query}]\n"
f"Observation:\n{observation}\n"
)
await self._emit(
"react.step",
query_id=query.id,
step=step,
action=action_query,
added=len(new_items),
)
response = await self._generator.generate(
f"{scratchpad}\nAnswer:", evidence, self._generation_config,
)
last_answer = response.answer or ""
looks_complete = (
"Answer:" in last_answer or last_answer.strip().endswith(".")
)
if looks_complete and (step >= 2 or len(last_answer) > 60):
# Heuristic stop: the model produced a complete sentence.
break
if not last_answer:
response = await self._generator.generate(
query.text, evidence, self._generation_config,
)
last_answer = response.answer or ""
traces = [
{"pipeline": self.name, "scratchpad": scratchpad, "evidence": len(evidence)},
]
# Surface scratchpad as a synthetic source so callers can inspect it.
evidence.append(
RetrievedChunk(
chunk=Chunk(text=scratchpad[-1000:], parent_id="react_scratchpad"),
score=0.0,
retriever="react_scratchpad",
),
)
await self._emit("query.generated", query_id=query.id, evidence=len(evidence))
return RagResponse(
answer=last_answer,
sources=[r for r in evidence if r.retriever != "react_scratchpad"],
traces=traces,
query_id=query.id,
)
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/react/react_rag_pipeline.py
RecursiveRagPipeline
¶
RecursiveRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, max_depth: int = 3)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/recursive/recursive_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
max_depth: int = 3,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"RecursiveRagPipeline requires chunker, embedder, store and generator"
)
if max_depth < 1:
raise PipelineConfigError("max_depth must be >= 1")
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._max_depth = max_depth
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/recursive/recursive_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.sources = retrieved
response.query_id = query.id
await self._emit("query.generated", query_id=query.id)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/recursive/recursive_rag_pipeline.py
RecursiveTextChunker
¶
RecursiveTextChunker(chunk_size: int = 400, chunk_overlap: int = 40, separators: tuple[str, ...] | None = None)
Splits documents recursively along configurable separators.
Mirrors the behaviour of LangChain's RecursiveCharacterTextSplitter
while staying dependency-free. Adds a per-chunk overlap by prepending the
tail of the previous chunk.
Source code in apogee_ai_rag/infrastructure/chunkers/recursive_text_chunker.py
def __init__(
self,
chunk_size: int = 400,
chunk_overlap: int = 40,
separators: tuple[str, ...] | None = None,
) -> None:
if chunk_size <= 0:
raise ValueError("chunk_size must be > 0")
if chunk_overlap < 0 or chunk_overlap >= chunk_size:
raise ValueError("chunk_overlap must satisfy 0 <= overlap < chunk_size")
self._chunk_size = chunk_size
self._chunk_overlap = chunk_overlap
self._separators = separators or _DEFAULT_STRATEGY.separators
chunk
async
¶
chunk(documents: list[Document], strategy: ChunkStrategy | None = None) -> list[Chunk]
Source code in apogee_ai_rag/infrastructure/chunkers/recursive_text_chunker.py
async def chunk(
self, documents: list[Document], strategy: ChunkStrategy | None = None
) -> list[Chunk]:
size = strategy.chunk_size if strategy else self._chunk_size
overlap = strategy.chunk_overlap if strategy else self._chunk_overlap
seps = strategy.separators if strategy else self._separators
chunks: list[Chunk] = []
for doc in documents:
if not doc.text:
continue
parts = _recursive_split(doc.text, seps, size)
parts = _apply_overlap(parts, overlap)
for position, text in enumerate(parts):
chunks.append(
Chunk(
text=text,
parent_id=doc.id,
position=position,
metadata=dict(doc.metadata),
modality=doc.modality,
)
)
return chunks
RedisCacheStore
¶
Source code in apogee_ai_rag/infrastructure/cache_stores/redis_cache_store.py
get
async
¶
get(key: str) -> CacheEntry | None
set
async
¶
set(entry: CacheEntry) -> None
invalidate
async
¶
Source code in apogee_ai_rag/infrastructure/cache_stores/redis_cache_store.py
async def invalidate(self, prefix: str) -> int: # pragma: no cover - I/O
cursor = 0
deleted = 0
pattern = self._k(prefix) + "*"
while True:
cursor, keys = await self._client.scan(cursor, match=pattern, count=128)
if keys:
deleted += await self._client.delete(*keys)
if cursor == 0:
break
return deleted
RerankConfig
dataclass
¶
RerankConfig(provider: RerankerProvider = NONE, model: str | None = None, top_n: int = 5)
RerankerProvider
¶
Bases: StrEnum
RerankingRagPipeline
¶
RerankingRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, reranker: IReranker, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, over_fetch: int = 4)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/reranking/reranking_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
reranker: IReranker,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
over_fetch: int = 4,
) -> None:
if reranker is None:
raise PipelineConfigError("RerankingRagPipeline requires a reranker")
if over_fetch < 1:
raise PipelineConfigError("over_fetch must be >= 1")
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._reranker = reranker
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._over_fetch = over_fetch
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/reranking/reranking_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
response = await self._generator.generate(
query.text, retrieved, self._generation_config,
)
response.query_id = query.id
await self._emit(
"query.generated", query_id=query.id, finish_reason=response.finish_reason,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/reranking/reranking_rag_pipeline.py
RetrievalConfig
dataclass
¶
RetrievalConfig(top_k: int = 5, threshold: float = 0.0, algorithm: SearchAlgorithm = COSINE, filters: dict = dict(), hybrid_alpha: float = 0.5)
RetrievedChunk
dataclass
¶
RetrievedChunk(chunk: Chunk, score: float, retriever: str = 'vector')
ScoreThresholdReranker
¶
Drops every chunk whose score falls below min_score and truncates to top_n.
Source code in apogee_ai_rag/infrastructure/rerankers/score_threshold_reranker.py
rerank
async
¶
rerank(query: str, chunks: list[RetrievedChunk], top_n: int = 5) -> list[RetrievedChunk]
Source code in apogee_ai_rag/infrastructure/rerankers/score_threshold_reranker.py
SearchAlgorithm
¶
Bases: StrEnum
SelfRagPipeline
¶
SelfRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, retrieve_min_words: int = 4, relevance_threshold: float = 0.15)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/self_rag/self_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
retrieve_min_words: int = 4,
relevance_threshold: float = 0.15,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"SelfRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._generator = generator
self._emitter = default_emitter(observability)
self._generation_config = generation_config
self._retrieve_min_words = retrieve_min_words
self._relevance_threshold = relevance_threshold
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/self_rag/self_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
await self._emit(
"query.received", query_id=query.id, text=query.text, pipeline_type=self.name,
)
sources, verdicts = await self._retrieve_if_needed(query)
response = await self._generator.generate(
query.text, sources, self._generation_config,
)
supported = self._is_supported(response.answer, sources)
response.sources = sources
response.query_id = query.id
response.confidence = (
0.95 if supported and verdicts["relevant"]
else 0.5 if verdicts["retrieve"]
else 0.3
)
response.traces.append({
"pipeline": self.name,
"verdicts": verdicts,
"supported": supported,
})
await self._emit(
"query.generated",
query_id=query.id,
confidence=response.confidence,
supported=supported,
)
return response
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/self_rag/self_rag_pipeline.py
SpeculativeRagPipeline
¶
SpeculativeRagPipeline(chunker: IChunker, embedder: IEmbedder, vector_store: IVectorStore, generator: IGenerator, observability: IRagObservabilityEmitter | None = None, generation_config: GenerationConfig | None = None, drafter: IGenerator | None = None, drafter_config: GenerationConfig | None = None)
Bases: IngestionMixin
Source code in apogee_ai_rag/infrastructure/pipelines/speculative/speculative_rag_pipeline.py
def __init__( # noqa: PLR0913
self,
chunker: IChunker,
embedder: IEmbedder,
vector_store: IVectorStore,
generator: IGenerator,
observability: IRagObservabilityEmitter | None = None,
generation_config: GenerationConfig | None = None,
drafter: IGenerator | None = None,
drafter_config: GenerationConfig | None = None,
) -> None:
if not all([chunker, embedder, vector_store, generator]):
raise PipelineConfigError(
"SpeculativeRagPipeline requires chunker, embedder, store and generator"
)
self._chunker = chunker
self._embedder = embedder
self._store = vector_store
self._verifier = generator
self._drafter = drafter or generator
self._emitter = default_emitter(observability)
self._verifier_config = generation_config
self._drafter_config = drafter_config or GenerationConfig(
temperature=0.0, max_tokens=192,
)
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
Source code in apogee_ai_rag/infrastructure/pipelines/speculative/speculative_rag_pipeline.py
async def run(self, query: RagQuery) -> RagResponse:
retrieved = await self._retrieve(query)
draft_context = retrieved[:1]
draft = await self._drafter.generate(
query.text, draft_context, self._drafter_config,
)
await self._emit(
"speculative.drafted",
query_id=query.id,
draft_chars=len(draft.answer or ""),
)
verifier_prompt = (
f"Question: {query.text}\n\n"
f"Draft answer (from a small model): {draft.answer}\n\n"
"Using the supplied context, ratify or rewrite the draft. "
"Quote the brackets when referencing sources."
)
verified = await self._verifier.generate(
verifier_prompt, retrieved, self._verifier_config,
)
verified.sources = retrieved
verified.query_id = query.id
verified.traces.append({
"pipeline": self.name,
"draft_answer": draft.answer,
"draft_sources": len(draft_context),
})
# Track the draft as an auxiliary chunk for downstream UIs.
verified.sources = list(verified.sources) + [
RetrievedChunk(
chunk=Chunk(text=draft.answer or "", parent_id="speculative_draft"),
score=0.0,
retriever="speculative_draft",
),
]
await self._emit("query.generated", query_id=query.id)
return verified
stream
async
¶
Source code in apogee_ai_rag/infrastructure/pipelines/speculative/speculative_rag_pipeline.py
TokenOverlapReranker
¶
Cheap dense+sparse re-ranker — blends dense score with token overlap.
Useful as a baseline when no Cohere/BGE/ColBERT key is available. The
final score is alpha * dense + (1 - alpha) * overlap, both
normalised to [0, 1]. Mirrors how :class:InMemoryVectorStore.hybrid_search
blends signals, but is applied on top of an already-retrieved list.
Source code in apogee_ai_rag/infrastructure/rerankers/token_overlap_reranker.py
rerank
async
¶
rerank(query: str, chunks: list[RetrievedChunk], top_n: int = 5) -> list[RetrievedChunk]
Source code in apogee_ai_rag/infrastructure/rerankers/token_overlap_reranker.py
async def rerank(
self, query: str, chunks: list[RetrievedChunk], top_n: int = 5
) -> list[RetrievedChunk]:
if not chunks:
return []
query_tokens = set(tokenize(query))
max_dense = max((c.score for c in chunks), default=1.0) or 1.0
rescored: list[RetrievedChunk] = []
for retrieved in chunks:
doc_tokens = set(tokenize(retrieved.chunk.text))
overlap = (
len(query_tokens & doc_tokens) / max(len(query_tokens), 1)
if query_tokens
else 0.0
)
dense_norm = retrieved.score / max_dense if max_dense else 0.0
blended = self._alpha * dense_norm + (1.0 - self._alpha) * overlap
rescored.append(
RetrievedChunk(
chunk=retrieved.chunk,
score=blended,
retriever=f"{retrieved.retriever}+overlap",
),
)
rescored.sort(key=lambda r: r.score, reverse=True)
return rescored[:top_n]
VectorStoreProvider
¶
Bases: StrEnum
reciprocal_rank_fusion
¶
reciprocal_rank_fusion(rankings: Iterable[list[RetrievedChunk]], *, k: int = 60, top_k: int | None = None, retriever_label: str = 'rrf') -> list[RetrievedChunk]
Fuses several ranked lists using reciprocal rank fusion.
Score for a chunk = Σ 1/(k + rank_i) across every list it appears in.
Returns the merged list sorted by descending RRF score, optionally
truncated to top_k items.
Source code in apogee_ai_rag/infrastructure/search/rrf.py
def reciprocal_rank_fusion(
rankings: Iterable[list[RetrievedChunk]],
*,
k: int = 60,
top_k: int | None = None,
retriever_label: str = "rrf",
) -> list[RetrievedChunk]:
"""Fuses several ranked lists using reciprocal rank fusion.
Score for a chunk = Σ 1/(k + rank_i) across every list it appears in.
Returns the merged list sorted by descending RRF score, optionally
truncated to ``top_k`` items.
"""
scores: dict[str, float] = {}
repr_chunk: dict[str, RetrievedChunk] = {}
for ranking in rankings:
for rank, retrieved in enumerate(ranking, start=1):
cid = retrieved.chunk.id
scores[cid] = scores.get(cid, 0.0) + 1.0 / (k + rank)
repr_chunk.setdefault(cid, retrieved)
fused = [
RetrievedChunk(
chunk=repr_chunk[cid].chunk,
score=score,
retriever=retriever_label,
)
for cid, score in scores.items()
]
fused.sort(key=lambda r: r.score, reverse=True)
if top_k is not None:
fused = fused[:top_k]
return fused
tokenize
¶
Splits text into word tokens.
Locale-agnostic, dependency-free. lowercase=True matches the
behaviour expected by the BM25 implementation and by most embedders.
Source code in apogee_ai_rag/infrastructure/search/tokenizer.py
def tokenize(text: str, lowercase: bool = True) -> list[str]:
"""Splits text into word tokens.
Locale-agnostic, dependency-free. ``lowercase=True`` matches the
behaviour expected by the BM25 implementation and by most embedders.
"""
text = text or ""
tokens = _TOKEN_RE.findall(text)
if lowercase:
tokens = [t.lower() for t in tokens]
return tokens
Other · DTOs¶
ChunkDTO
¶
Bases: BaseModel
DocumentDTO
¶
Bases: BaseModel
EvalCaseDTO
¶
EvalReportDTO
¶
Bases: BaseModel
metrics
class-attribute
instance-attribute
¶
per_case
class-attribute
instance-attribute
¶
per_case: list[EvalResultDTO] = Field(default_factory=list)
IngestionJobDTO
¶
Bases: BaseModel
IngestionResultDTO
¶
Bases: BaseModel
RagChunkDTO
¶
Bases: BaseModel
sources
class-attribute
instance-attribute
¶
sources: list[RetrievedChunkDTO] = Field(default_factory=list)
RagResponseDTO
¶
Bases: BaseModel
sources
class-attribute
instance-attribute
¶
sources: list[RetrievedChunkDTO] = Field(default_factory=list)
RetrievedChunkDTO
¶
Other · Enums¶
ChunkerKind
¶
Bases: StrEnum
RagType
¶
Bases: StrEnum
Other · Exceptions¶
PipelineConfigError
¶
PipelineConfigError(message: str, *, stage: str | None = None, cause: Exception | None = None)
ProviderNotInstalledError
¶
Bases: RagError
Raised when an optional provider is selected but its extra is missing.
Carries the suggested pip install apogee-ai-rag[<extra>] command so the
caller can recover automatically (e.g. CLI surfaces it as the only message).
Source code in apogee_ai_rag/domain/exceptions/provider_not_installed_error.py
RagError
¶
Bases: Exception
Source code in apogee_ai_rag/domain/exceptions/rag_error.py
UnsupportedModalityError
¶
UnsupportedModalityError(message: str, *, stage: str | None = None, cause: Exception | None = None)
Other · Protocols (ports)¶
ICacheStore
¶
Bases: Protocol
get
async
¶
get(key: str) -> CacheEntry | None
set
async
¶
set(entry: CacheEntry) -> None
invalidate
async
¶
IChunker
¶
Bases: Protocol
chunk
async
¶
chunk(documents: list[Document], strategy: ChunkStrategy | None = None) -> list[Chunk]
IEmbedder
¶
Bases: Protocol
IGenerator
¶
Bases: Protocol
generate
async
¶
generate(prompt: str, context: list[RetrievedChunk], config: GenerationConfig | None = None) -> RagResponse
stream
¶
stream(prompt: str, context: list[RetrievedChunk], config: GenerationConfig | None = None) -> AsyncIterator[RagChunk]
IGraphStore
¶
Bases: Protocol
add_node
async
¶
add_node(node: KnowledgeNode) -> str
add_edge
async
¶
add_edge(edge: KnowledgeEdge) -> str
traverse
async
¶
traverse(start_id: str, depth: int = 2) -> tuple[list[KnowledgeNode], list[KnowledgeEdge]]
IIngestor
¶
Bases: Protocol
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
IQueryRewriter
¶
IRagEvaluator
¶
Bases: Protocol
evaluate
async
¶
evaluate(suite: list[EvalCase], pipeline: IRagPipeline) -> EvalReport
IRagObservabilityEmitter
¶
IRagPipeline
¶
Bases: Protocol
ingest
async
¶
ingest(job: IngestionJob) -> IngestionResult
run
async
¶
run(query: RagQuery) -> RagResponse
IReranker
¶
Bases: Protocol
rerank
async
¶
rerank(query: str, chunks: list[RetrievedChunk], top_n: int = 5) -> list[RetrievedChunk]
IRetriever
¶
Bases: Protocol
retrieve
async
¶
retrieve(query: RagQuery) -> list[RetrievedChunk]
IVectorStore
¶
Bases: Protocol
search
async
¶
search(query_embedding: list[float], top_k: int = 5, threshold: float = 0.0, filters: dict | None = None) -> list[RetrievedChunk]
hybrid_search
async
¶
hybrid_search(query_text: str, query_embedding: list[float], top_k: int = 5, threshold: float = 0.0, filters: dict | None = None, alpha: float = 0.5) -> list[RetrievedChunk]
delete
async
¶
count
async
¶
Other · Use cases¶
DeleteDocumentsUseCase
¶
DeleteDocumentsUseCase(store: IVectorStore)
Source code in apogee_ai_rag/application/use_cases/delete_documents_use_case.py
execute
async
¶
EvaluateRagUseCase
¶
EvaluateRagUseCase(evaluator: IRagEvaluator, pipeline: IRagPipeline)
Source code in apogee_ai_rag/application/use_cases/evaluate_rag_use_case.py
execute
async
¶
execute(suite: list[EvalCaseDTO]) -> EvalReportDTO
FeedbackUseCase
¶
FeedbackUseCase(emitter: IRagObservabilityEmitter)
Source code in apogee_ai_rag/application/use_cases/feedback_use_case.py
execute
async
¶
execute(feedback: FeedbackInput) -> None
Source code in apogee_ai_rag/application/use_cases/feedback_use_case.py
async def execute(self, feedback: FeedbackInput) -> None:
attributes = {
"query_id": feedback.query_id,
"score": feedback.score,
}
if feedback.comment is not None:
attributes["comment"] = feedback.comment
if feedback.metadata:
attributes.update(feedback.metadata)
await self._emitter.emit(RagEvent(name="feedback.submitted", attributes=attributes))
IngestDocumentsUseCase
¶
IngestDocumentsUseCase(pipeline: IRagPipeline)
Source code in apogee_ai_rag/application/use_cases/ingest_documents_use_case.py
execute
async
¶
execute(dto: IngestionJobDTO) -> IngestionResultDTO
QueryRagUseCase
¶
QueryRagUseCase(pipeline: IRagPipeline)
Source code in apogee_ai_rag/application/use_cases/query_rag_use_case.py
execute
async
¶
execute(dto: RagQueryDTO) -> RagResponseDTO
RefreshIndexUseCase
¶
RefreshIndexUseCase(pipeline: IRagPipeline, store: IVectorStore)
Re-ingests a batch of documents after deleting their previous chunks.
The caller supplies an IVectorStore used to find and remove the chunks
that share a parent_id with the incoming documents. Implementation is
pluggable via the store's delete and count operations.
Source code in apogee_ai_rag/application/use_cases/refresh_index_use_case.py
execute
async
¶
execute(dto: IngestionJobDTO) -> IngestionResultDTO
Source code in apogee_ai_rag/application/use_cases/refresh_index_use_case.py
StreamRagUseCase
¶
StreamRagUseCase(pipeline: IRagPipeline)
Source code in apogee_ai_rag/application/use_cases/stream_rag_use_case.py
execute
async
¶
execute(dto: RagQueryDTO) -> AsyncIterator[RagChunkDTO]