Documentation
¶
Index ¶
- Variables
- type Pipeline
- func NewPipeline(store domain.VectorStore, llm domain.LLMProvider, embedder domain.Embedder) (*Pipeline, error)
- func NewPipelineWithChunker(store domain.VectorStore, llm domain.LLMProvider, embedder domain.Embedder, ...) (*Pipeline, error)
- func NewPipelineWithConfig(store domain.VectorStore, llm domain.LLMProvider, embedder domain.Embedder, ...) (*Pipeline, error)
- func (p *Pipeline) Answer(ctx context.Context, question string, topK int) (string, error)
- func (p *Pipeline) AnswerHyDE(ctx context.Context, question string, topK int) (string, error)
- func (p *Pipeline) AnswerHyDEStream(ctx context.Context, question string, topK int) (<-chan string, error)
- func (p *Pipeline) AnswerHyDEStreamWithInlineCitations(ctx context.Context, question string, topK int) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerHyDEStreamWithSources(ctx context.Context, question string, topK int) (<-chan string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerHyDEWithCitations(ctx context.Context, question string, topK int) (string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerHyDEWithInlineCitations(ctx context.Context, question string, topK int) (string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerHybrid(ctx context.Context, question string, topK int, config domain.HybridConfig) (string, error)
- func (p *Pipeline) AnswerHybridStream(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (<-chan string, error)
- func (p *Pipeline) AnswerHybridStreamWithInlineCitations(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerHybridStreamWithSources(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (<-chan string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerHybridWithCitations(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerHybridWithInlineCitations(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerMulti(ctx context.Context, question string, n, topK int) (string, error)
- func (p *Pipeline) AnswerMultiStream(ctx context.Context, question string, n, topK int) (<-chan string, error)
- func (p *Pipeline) AnswerMultiStreamWithInlineCitations(ctx context.Context, question string, n, topK int) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerMultiStreamWithSources(ctx context.Context, question string, n, topK int) (<-chan string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerMultiWithCitations(ctx context.Context, question string, n, topK int) (string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerMultiWithInlineCitations(ctx context.Context, question string, n, topK int) (string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerStream(ctx context.Context, question string, topK int) (<-chan string, error)
- func (p *Pipeline) AnswerStreamWithInlineCitations(ctx context.Context, question string, topK int) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerStreamWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (<-chan string, error)
- func (p *Pipeline) AnswerStreamWithMetadataFilterWithInlineCitations(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerStreamWithMetadataFilterWithSources(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (<-chan string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerStreamWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (<-chan string, error)
- func (p *Pipeline) AnswerStreamWithParentIDsWithInlineCitations(ctx context.Context, question string, topK int, parentIDs []string) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerStreamWithParentIDsWithSources(ctx context.Context, question string, topK int, parentIDs []string) (<-chan string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerStreamWithSources(ctx context.Context, question string, topK int) (<-chan string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerSubDecompose(ctx context.Context, question string, topK int, ...) (string, error)
- func (p *Pipeline) AnswerSubDecomposeStream(ctx context.Context, question string, topK int, ...) (<-chan string, error)
- func (p *Pipeline) AnswerSubDecomposeStreamWithInlineCitations(ctx context.Context, question string, topK int, ...) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerSubDecomposeStreamWithSources(ctx context.Context, question string, topK int, ...) (<-chan string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerSubDecomposeWithCitations(ctx context.Context, question string, topK int, ...) (string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerSubDecomposeWithInlineCitations(ctx context.Context, question string, topK int, ...) (string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerWithCitations(ctx context.Context, question string, topK int) (string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerWithCitationsWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerWithCitationsWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerWithInlineCitations(ctx context.Context, question string, topK int) (string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerWithInlineCitationsWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerWithInlineCitationsWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (string, error)
- func (p *Pipeline) AnswerWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (string, error)
- func (p *Pipeline) AnswerWithQueries(ctx context.Context, originalQuery string, queries []string, topK int) (string, error)
- func (p *Pipeline) AnswerWithQueriesAndCitations(ctx context.Context, originalQuery string, queries []string, topK int) (string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerWithQueriesStream(ctx context.Context, originalQuery string, queries []string, topK int) (<-chan string, error)
- func (p *Pipeline) AnswerWithQueriesStreamWithInlineCitations(ctx context.Context, originalQuery string, queries []string, topK int) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) AnswerWithQueriesStreamWithSources(ctx context.Context, originalQuery string, queries []string, topK int) (<-chan string, domain.RetrievalResult, error)
- func (p *Pipeline) AnswerWithQueriesWithInlineCitations(ctx context.Context, originalQuery string, queries []string, topK int) (string, domain.RetrievalResult, []domain.InlineCitation, error)
- func (p *Pipeline) Close() error
- func (p *Pipeline) DeleteDocument(ctx context.Context, docID string) error
- func (p *Pipeline) ExecuteWithTools(ctx context.Context, systemPrompt, userMessage string, ...) (string, error)
- func (p *Pipeline) Health(ctx context.Context) error
- func (p *Pipeline) Index(ctx context.Context, docs []domain.Document) error
- func (p *Pipeline) IndexBatch(ctx context.Context, docs []domain.Document, batchSize int) (*domain.IndexBatchResult, error)
- func (p *Pipeline) MaxContextChars() int
- func (p *Pipeline) MaxContextChunks() int
- func (p *Pipeline) Query(ctx context.Context, question string, topK int) (domain.RetrievalResult, error)
- func (p *Pipeline) QueryHyDE(ctx context.Context, question string, topK int) (domain.RetrievalResult, error)
- func (p *Pipeline) QueryHybrid(ctx context.Context, question string, topK int, config domain.HybridConfig) (domain.RetrievalResult, error)
- func (p *Pipeline) QueryMulti(ctx context.Context, question string, n, topK int) (domain.RetrievalResult, error)
- func (p *Pipeline) QuerySubDecompose(ctx context.Context, question string, topK int, ...) (domain.RetrievalResult, error)
- func (p *Pipeline) QueryWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (domain.RetrievalResult, error)
- func (p *Pipeline) QueryWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (domain.RetrievalResult, error)
- func (p *Pipeline) QueryWithQueries(ctx context.Context, queries []string, topK int) (domain.RetrievalResult, error)
- func (p *Pipeline) RedactRetrievalResult(r domain.RetrievalResult) domain.RetrievalResult
- func (p *Pipeline) SystemPrompt() string
- func (p *Pipeline) UpdateDocument(ctx context.Context, doc domain.Document) error
- type PipelineOptions
Constants ¶
This section is empty.
Variables ¶
var ( ErrFiltersNotSupported = errors.New("vector store does not support filters") ErrStreamingNotSupported = errors.New("LLM provider does not support streaming") ErrDeleteNotSupported = errors.New("vector store does not support DeleteByParentID") // @sk-task sub-query-decomposition#T3.3: sentinel for nil decomposer guard (AC-005, AC-006) ErrSubDecomposeNotSupported = errors.New("sub-query decomposition not supported: no QueryDecomposer configured") // @sk-task arch-issues#T1.2: sentinel для Pipeline.Close (AC-008) ErrPipelineClosed = errors.New("pipeline is closed") )
Sentinel errors returned by Pipeline operations.
var ErrHybridNotSupported = errors.New("vector store does not support hybrid search")
ErrHybridNotSupported is returned when a hybrid search method is called but the underlying VectorStore does not implement HybridSearcher.
Functions ¶
This section is empty.
Types ¶
type Pipeline ¶
type Pipeline struct {
// contains filtered or unexported fields
}
Pipeline is the core RAG pipeline coordinating store, LLM, and embedder.
@sk-task arch-issues#T3.1: Health + Close (AC-007, AC-008) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task pii-guardrails#T2.1: Pipeline.PIIDetector (RQ-001, RQ-002)
func NewPipeline ¶
func NewPipeline(store domain.VectorStore, llm domain.LLMProvider, embedder domain.Embedder) (*Pipeline, error)
NewPipeline creates a Pipeline with required dependencies.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task arch-quality-pass#T2.1: error return вместо panic для конфигурации (AC-002)
func NewPipelineWithChunker ¶
func NewPipelineWithChunker( store domain.VectorStore, llm domain.LLMProvider, embedder domain.Embedder, chunker domain.Chunker, ) (*Pipeline, error)
NewPipelineWithChunker creates a Pipeline with an optional chunker.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task arch-quality-pass#T2.1: error return вместо panic для конфигурации (AC-002)
func NewPipelineWithConfig ¶
func NewPipelineWithConfig( store domain.VectorStore, llm domain.LLMProvider, embedder domain.Embedder, cfg PipelineOptions, ) (*Pipeline, error)
NewPipelineWithConfig creates a Pipeline with the given configuration.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task arch-quality-pass#T2.1: error return вместо panic для конфигурации (AC-002) @sk-task arch-quality-pass#T3.2: принимает PipelineOptions вместо PipelineConfig (AC-004)
func (*Pipeline) Answer ¶
Answer выполняет полный RAG-цикл: retrieval (Embed+Search) → prompt → LLM.Generate.
@sk-task arch-issues#T2.1: PII redaction в Answer (AC-001, AC-002) @sk-task hierarchical-indices#T3.3: parent context attach in Answer (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @sk-task middleware-chain#T3.2: userMessage from d.Query, return stageResult.Answer (AC-004) @sk-task arch-issues#T3.1: closed guard в Answer (AC-008)
func (*Pipeline) AnswerHyDE ¶
AnswerHyDE генерирует ответ, используя HyDE для retrieval.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHyDEStream ¶
func (p *Pipeline) AnswerHyDEStream(ctx context.Context, question string, topK int) (<-chan string, error)
AnswerHyDEStream выполняет HyDE retrieval и streaming генерацию.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHyDEStreamWithInlineCitations ¶
func (p *Pipeline) AnswerHyDEStreamWithInlineCitations(ctx context.Context, question string, topK int) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerHyDEStreamWithInlineCitations выполняет HyDE retrieval с streaming и inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHyDEStreamWithSources ¶
func (p *Pipeline) AnswerHyDEStreamWithSources(ctx context.Context, question string, topK int) (<-chan string, domain.RetrievalResult, error)
AnswerHyDEStreamWithSources выполняет HyDE retrieval с streaming и возвращает источники.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHyDEWithCitations ¶
func (p *Pipeline) AnswerHyDEWithCitations(ctx context.Context, question string, topK int) (string, domain.RetrievalResult, error)
AnswerHyDEWithCitations выполняет HyDE retrieval и генерирует ответ с цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHyDEWithInlineCitations ¶
func (p *Pipeline) AnswerHyDEWithInlineCitations(ctx context.Context, question string, topK int) (string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerHyDEWithInlineCitations выполняет HyDE retrieval и генерирует ответ с inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHybrid ¶
func (p *Pipeline) AnswerHybrid(ctx context.Context, question string, topK int, config domain.HybridConfig) (string, error)
AnswerHybrid выполняет гибридный поиск (BM25 + semantic) и генерирует ответ.
Если store не реализует HybridSearcher — возвращает ErrHybridNotSupported.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003)
func (*Pipeline) AnswerHybridStream ¶
func (p *Pipeline) AnswerHybridStream(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (<-chan string, error)
AnswerHybridStream выполняет Hybrid retrieval и streaming генерацию.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHybridStreamWithInlineCitations ¶
func (p *Pipeline) AnswerHybridStreamWithInlineCitations(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerHybridStreamWithInlineCitations выполняет Hybrid retrieval с streaming и inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHybridStreamWithSources ¶
func (p *Pipeline) AnswerHybridStreamWithSources(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (<-chan string, domain.RetrievalResult, error)
AnswerHybridStreamWithSources выполняет Hybrid retrieval с streaming и возвращает источники.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHybridWithCitations ¶
func (p *Pipeline) AnswerHybridWithCitations(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (string, domain.RetrievalResult, error)
AnswerHybridWithCitations выполняет Hybrid retrieval и генерирует ответ с цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerHybridWithInlineCitations ¶
func (p *Pipeline) AnswerHybridWithInlineCitations(ctx context.Context, question string, topK int, cfg domain.HybridConfig) (string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerHybridWithInlineCitations выполняет Hybrid retrieval и генерирует ответ с inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerMulti ¶
AnswerMulti генерирует ответ используя multi-query retrieval.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerMultiStream ¶
func (p *Pipeline) AnswerMultiStream(ctx context.Context, question string, n, topK int) (<-chan string, error)
AnswerMultiStream выполняет MultiQuery retrieval и streaming генерацию.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerMultiStreamWithInlineCitations ¶
func (p *Pipeline) AnswerMultiStreamWithInlineCitations(ctx context.Context, question string, n, topK int) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerMultiStreamWithInlineCitations выполняет MultiQuery retrieval с streaming и inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerMultiStreamWithSources ¶
func (p *Pipeline) AnswerMultiStreamWithSources(ctx context.Context, question string, n, topK int) (<-chan string, domain.RetrievalResult, error)
AnswerMultiStreamWithSources выполняет MultiQuery retrieval с streaming и возвращает источники.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerMultiWithCitations ¶
func (p *Pipeline) AnswerMultiWithCitations(ctx context.Context, question string, n, topK int) (string, domain.RetrievalResult, error)
AnswerMultiWithCitations выполняет MultiQuery retrieval и генерирует ответ с цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerMultiWithInlineCitations ¶
func (p *Pipeline) AnswerMultiWithInlineCitations(ctx context.Context, question string, n, topK int) (string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerMultiWithInlineCitations выполняет MultiQuery retrieval и генерирует ответ с inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerStream ¶
func (p *Pipeline) AnswerStream( ctx context.Context, question string, topK int, ) (<-chan string, error)
AnswerStream выполняет RAG-цикл с streaming генерацией ответа. Возвращает канал для чтения текстовых чанков; канал закрывается при завершении или ошибке. Retrieval выполняется синхронно перед началом streaming'а.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @ds-task T2.3: Реализовать AnswerStream в application Pipeline (AC-001, DEC-003) @sk-task middleware-chain#T3.2: userMessage from d.Query (AC-004)
func (*Pipeline) AnswerStreamWithInlineCitations ¶
func (p *Pipeline) AnswerStreamWithInlineCitations( ctx context.Context, question string, topK int, ) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerStreamWithInlineCitations выполняет RAG-цикл с streaming генерацией и inline-цитатами. Возвращает канал для чтения текстовых чанков и слайс цитат (заполняется синхронно перед streaming'ом).
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @ds-task T2.4: Реализовать AnswerStreamWithInlineCitations в application Pipeline (AC-002) @sk-task middleware-chain#T-concern: add middleware on embed/search/generate (AC-003)
func (*Pipeline) AnswerStreamWithMetadataFilter ¶
func (p *Pipeline) AnswerStreamWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (<-chan string, error)
AnswerStreamWithMetadataFilter выполняет retrieval с фильтром по метаданным и streaming генерацию.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerStreamWithMetadataFilterWithInlineCitations ¶
func (p *Pipeline) AnswerStreamWithMetadataFilterWithInlineCitations(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerStreamWithMetadataFilterWithInlineCitations выполняет retrieval с фильтром по метаданным с streaming и inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerStreamWithMetadataFilterWithSources ¶
func (p *Pipeline) AnswerStreamWithMetadataFilterWithSources(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (<-chan string, domain.RetrievalResult, error)
AnswerStreamWithMetadataFilterWithSources выполняет retrieval с фильтром по метаданным с streaming и возвращает источники.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerStreamWithParentIDs ¶
func (p *Pipeline) AnswerStreamWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (<-chan string, error)
AnswerStreamWithParentIDs выполняет retrieval с фильтром по ParentIDs и streaming генерацию.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerStreamWithParentIDsWithInlineCitations ¶
func (p *Pipeline) AnswerStreamWithParentIDsWithInlineCitations(ctx context.Context, question string, topK int, parentIDs []string) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerStreamWithParentIDsWithInlineCitations выполняет retrieval с фильтром по ParentIDs с streaming и inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerStreamWithParentIDsWithSources ¶
func (p *Pipeline) AnswerStreamWithParentIDsWithSources(ctx context.Context, question string, topK int, parentIDs []string) (<-chan string, domain.RetrievalResult, error)
AnswerStreamWithParentIDsWithSources выполняет retrieval с фильтром по ParentIDs с streaming и возвращает источники.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerStreamWithSources ¶
func (p *Pipeline) AnswerStreamWithSources(ctx context.Context, question string, topK int) (<-chan string, domain.RetrievalResult, error)
AnswerStreamWithSources выполняет RAG-цикл с streaming и возвращает источники.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerSubDecompose ¶ added in v1.0.0
func (p *Pipeline) AnswerSubDecompose(ctx context.Context, question string, topK int, decomposer domain.QueryDecomposer) (string, error)
@sk-task sub-query-decomposition#T2.2: AnswerSubDecompose (AC-008) AnswerSubDecompose выполняет sub-query decomposition + retrieval + answer.
func (*Pipeline) AnswerSubDecomposeStream ¶ added in v1.0.0
func (p *Pipeline) AnswerSubDecomposeStream(ctx context.Context, question string, topK int, decomposer domain.QueryDecomposer) (<-chan string, error)
@sk-task sub-query-decomposition#T2.2: AnswerSubDecomposeStream (AC-009) AnswerSubDecomposeStream выполняет sub-query decomposition + retrieval + streaming answer.
func (*Pipeline) AnswerSubDecomposeStreamWithInlineCitations ¶ added in v1.0.0
func (p *Pipeline) AnswerSubDecomposeStreamWithInlineCitations(ctx context.Context, question string, topK int, decomposer domain.QueryDecomposer) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
@sk-task sub-query-decomposition#T2.2: AnswerSubDecomposeStreamWithInlineCitations (AC-009) AnswerSubDecomposeStreamWithInlineCitations выполняет sub-query decomposition + retrieval + streaming answer с inline-цитатами.
func (*Pipeline) AnswerSubDecomposeStreamWithSources ¶ added in v1.0.0
func (p *Pipeline) AnswerSubDecomposeStreamWithSources(ctx context.Context, question string, topK int, decomposer domain.QueryDecomposer) (<-chan string, domain.RetrievalResult, error)
@sk-task sub-query-decomposition#T2.2: AnswerSubDecomposeStreamWithSources (AC-009) AnswerSubDecomposeStreamWithSources выполняет sub-query decomposition + retrieval + streaming answer с источниками.
func (*Pipeline) AnswerSubDecomposeWithCitations ¶ added in v1.0.0
func (p *Pipeline) AnswerSubDecomposeWithCitations(ctx context.Context, question string, topK int, decomposer domain.QueryDecomposer) (string, domain.RetrievalResult, error)
@sk-task sub-query-decomposition#T2.2: AnswerSubDecomposeWithCitations (AC-008, AC-009) AnswerSubDecomposeWithCitations выполняет sub-query decomposition + retrieval + answer с источниками.
func (*Pipeline) AnswerSubDecomposeWithInlineCitations ¶ added in v1.0.0
func (p *Pipeline) AnswerSubDecomposeWithInlineCitations(ctx context.Context, question string, topK int, decomposer domain.QueryDecomposer) (string, domain.RetrievalResult, []domain.InlineCitation, error)
@sk-task sub-query-decomposition#T2.2: AnswerSubDecomposeWithInlineCitations (AC-008, AC-009) AnswerSubDecomposeWithInlineCitations выполняет sub-query decomposition + retrieval + answer с inline-цитатами.
func (*Pipeline) AnswerWithCitations ¶
func (p *Pipeline) AnswerWithCitations(ctx context.Context, question string, topK int) (string, domain.RetrievalResult, error)
AnswerWithCitations выполняет полный RAG-цикл и возвращает retrieval evidence вместе с ответом.
Если retrieval уже выполнен, а Generate вернул ошибку, метод возвращает retrieval результат (partial) и ошибку, чтобы упростить диагностику и отображение источников.
@sk-task hierarchical-indices#T3.3: parent context attach in AnswerWithCitations (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @sk-task arch-issues#T2.1: PII redaction в AnswerWithCitations (AC-001)
func (*Pipeline) AnswerWithCitationsWithMetadataFilter ¶
func (p *Pipeline) AnswerWithCitationsWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (string, domain.RetrievalResult, error)
AnswerWithCitationsWithMetadataFilter выполняет retrieval с фильтром по метаданным и генерирует ответ с цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerWithCitationsWithParentIDs ¶
func (p *Pipeline) AnswerWithCitationsWithParentIDs( ctx context.Context, question string, topK int, parentIDs []string, ) (string, domain.RetrievalResult, error)
AnswerWithCitationsWithParentIDs выполняет RAG-цикл с фильтром по ParentIDs и возвращает retrieval evidence.
Если parentIDs пустой — эквивалентно AnswerWithCitations.
@sk-task hierarchical-indices#T3.3: parent context attach in AnswerWithCitationsWithParentIDs (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @sk-task arch-issues#T2.1: PII redaction в AnswerWithCitationsWithParentIDs (AC-001)
func (*Pipeline) AnswerWithInlineCitations ¶
func (p *Pipeline) AnswerWithInlineCitations( ctx context.Context, question string, topK int, ) (string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerWithInlineCitations выполняет полный RAG-цикл и возвращает ответ с inline-цитатами `[n]`, а также retrieval evidence и детерминированный маппинг `n -> chunk`.
Если retrieval уже выполнен, а Generate вернул ошибку, метод возвращает retrieval результат (partial), массив citations и ошибку.
@sk-task hierarchical-indices#T3.3: parent context attach in AnswerWithInlineCitations (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @sk-task arch-issues#T2.1: PII redaction в AnswerWithInlineCitations (AC-001)
func (*Pipeline) AnswerWithInlineCitationsWithMetadataFilter ¶
func (p *Pipeline) AnswerWithInlineCitationsWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerWithInlineCitationsWithMetadataFilter выполняет retrieval с фильтром по метаданным и генерирует ответ с inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerWithInlineCitationsWithParentIDs ¶
func (p *Pipeline) AnswerWithInlineCitationsWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (string, domain.RetrievalResult, []domain.InlineCitation, error)
AnswerWithInlineCitationsWithParentIDs выполняет retrieval с фильтром по ParentIDs и генерирует ответ с inline-цитатами.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerWithMetadataFilter ¶
func (p *Pipeline) AnswerWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (string, error)
AnswerWithMetadataFilter выполняет retrieval с фильтром по метаданным и генерирует ответ.
Если filter.Fields пустой — эквивалентно Answer. Если store не реализует VectorStoreWithFilters — возвращает ErrFiltersNotSupported.
@ds-task T3.1: Добавить AnswerWithMetadataFilter в application.Pipeline (RQ-006, AC-003, DEC-003) @sk-task hierarchical-indices#T3.3: parent context attach in AnswerWithMetadataFilter (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @sk-task arch-issues#T2.1: PII redaction в AnswerWithMetadataFilter (AC-001)
func (*Pipeline) AnswerWithParentIDs ¶
func (p *Pipeline) AnswerWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (string, error)
AnswerWithParentIDs выполняет retrieval с фильтром по ParentIDs и генерирует ответ.
Если parentIDs пустой — эквивалентно Answer.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) AnswerWithQueries ¶ added in v1.0.0
func (p *Pipeline) AnswerWithQueries(ctx context.Context, originalQuery string, queries []string, topK int) (string, error)
@sk-task query-rewriting#T2.2: AnswerWithQueries для pre-generated переформулировок (AC-003) AnswerWithQueries генерирует ответ используя multi-query retrieval из готового списка запросов.
func (*Pipeline) AnswerWithQueriesAndCitations ¶ added in v1.0.0
func (p *Pipeline) AnswerWithQueriesAndCitations(ctx context.Context, originalQuery string, queries []string, topK int) (string, domain.RetrievalResult, error)
@sk-task query-rewriting#T2.2: AnswerWithQueriesAndCitations (AC-003) AnswerWithQueriesAndCitations генерирует ответ с источниками из pre-generated переформулировок.
func (*Pipeline) AnswerWithQueriesStream ¶ added in v1.0.0
func (p *Pipeline) AnswerWithQueriesStream(ctx context.Context, originalQuery string, queries []string, topK int) (<-chan string, error)
@sk-task query-rewriting#T2.2: AnswerWithQueriesStream (AC-003) AnswerWithQueriesStream выполняет multi-query retrieval из готового списка запросов и streaming генерацию.
func (*Pipeline) AnswerWithQueriesStreamWithInlineCitations ¶ added in v1.0.0
func (p *Pipeline) AnswerWithQueriesStreamWithInlineCitations(ctx context.Context, originalQuery string, queries []string, topK int) (<-chan string, domain.RetrievalResult, []domain.InlineCitation, error)
@sk-task query-rewriting#T2.2: AnswerWithQueriesStreamWithInlineCitations (AC-003) AnswerWithQueriesStreamWithInlineCitations выполняет multi-query retrieval из готового списка запросов, streaming генерацию с inline-цитатами.
func (*Pipeline) AnswerWithQueriesStreamWithSources ¶ added in v1.0.0
func (p *Pipeline) AnswerWithQueriesStreamWithSources(ctx context.Context, originalQuery string, queries []string, topK int) (<-chan string, domain.RetrievalResult, error)
@sk-task query-rewriting#T2.2: AnswerWithQueriesStreamWithSources (AC-003) AnswerWithQueriesStreamWithSources выполняет multi-query retrieval из готового списка запросов, streaming генерацию и возвращает источники.
func (*Pipeline) AnswerWithQueriesWithInlineCitations ¶ added in v1.0.0
func (p *Pipeline) AnswerWithQueriesWithInlineCitations(ctx context.Context, originalQuery string, queries []string, topK int) (string, domain.RetrievalResult, []domain.InlineCitation, error)
@sk-task query-rewriting#T2.2: AnswerWithQueriesWithInlineCitations (AC-003) AnswerWithQueriesWithInlineCitations генерирует ответ с inline-цитатами из pre-generated переформулировок.
func (*Pipeline) DeleteDocument ¶
DeleteDocument deletes all chunks belonging to a document by its ID.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003)
func (*Pipeline) ExecuteWithTools ¶ added in v1.0.0
func (p *Pipeline) ExecuteWithTools(ctx context.Context, systemPrompt, userMessage string, tools []domain.ToolDefinition, execTool func(domain.ToolCall) domain.ToolResult) (string, error)
@sk-task arch-issues#T4.2: tool execution loop (AC-003)
func (*Pipeline) Health ¶ added in v1.0.0
@sk-task arch-issues#T3.1: Health() fan-out с таймаутом 1s (AC-007)
func (*Pipeline) Index ¶
Index индексирует набор документов параллельно с ограничением
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T3.1: параллельная обработка Index через processDocsConcurrently (DEC-004, RQ-004, AC-006)
Index индексирует набор документов параллельно с ограничением p.indexConcurrency и rateLimit p.indexBatchRateLimit.
Семантика fail-fast: при первой ошибке обработки cancel-ит in-flight siblings и возвращает оригинальную ошибку (не context.Canceled). Документы, не прошедшие Validate, также прерывают выполнение и возвращают первую такую ошибку.
Параллелизм: реализовано через processDocsConcurrently (T1.2). На каждую документную goroutine — отдельный семафор слот и общий rate-limiter. @sk-task arch-issues#T3.1: closed guard в Index (AC-008)
func (*Pipeline) IndexBatch ¶
func (p *Pipeline) IndexBatch(ctx context.Context, docs []domain.Document, batchSize int) (*domain.IndexBatchResult, error)
IndexBatch индексирует набор документов параллельно и возвращает aggregate результат (успешные + ошибки по документам). batchSize интерпретируется как желаемая concurrency: при batchSize <= 0 используется p.indexConcurrency.
Семантика: best-effort — не отменяет siblings при ошибке отдельного документа. Это отличие от Index, который fail-fast.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T1.2: рефактор — IndexBatch как тонкая обёртка над processDocsConcurrently (DEC-004, RQ-004) @sk-task api-consistency-pass#T3.1: shared processDocumentOp между Index и IndexBatch (DEC-004, RQ-004, AC-006)
func (*Pipeline) MaxContextChars ¶ added in v1.0.0
@sk-task arch-issues#T4.4: MaxContextChars accessor (AC-004)
func (*Pipeline) MaxContextChunks ¶ added in v1.0.0
@sk-task arch-issues#T4.4: MaxContextChunks accessor (AC-004)
func (*Pipeline) Query ¶
func (p *Pipeline) Query(ctx context.Context, question string, topK int) (domain.RetrievalResult, error)
Query выполняет поиск по вопросу и возвращает RetrievalResult.
@sk-task arch-issues#T2.1: PII redaction в Query (AC-001, AC-002) @sk-task hierarchical-indices#T3.3: parent context attach in Query (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @sk-task arch-issues#T3.1: closed guard в Query (AC-008)
func (*Pipeline) QueryHyDE ¶
func (p *Pipeline) QueryHyDE(ctx context.Context, question string, topK int) (domain.RetrievalResult, error)
QueryHyDE выполняет поиск с использованием Hypothetical Document Embeddings.
@sk-task hierarchical-indices#T3.3: parent context attach in QueryHyDE (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) Сначала LLM генерирует гипотетический ответ на вопрос, затем ищем по его embedding. @sk-task arch-issues#T2.1: PII redaction в QueryHyDE (AC-001)
func (*Pipeline) QueryHybrid ¶
func (p *Pipeline) QueryHybrid(ctx context.Context, question string, topK int, config domain.HybridConfig) (domain.RetrievalResult, error)
QueryHybrid выполняет гибридный поиск (BM25 + semantic) по вопросу.
Если store не реализует HybridSearcher — возвращает ErrHybridNotSupported.
@sk-task hierarchical-indices#T3.3: parent context attach in QueryHybrid (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003) @sk-task arch-issues#T2.1: PII redaction в QueryHybrid (AC-001)
func (*Pipeline) QueryMulti ¶
func (p *Pipeline) QueryMulti(ctx context.Context, question string, n, topK int) (domain.RetrievalResult, error)
QueryMulti выполняет multi-query retrieval: генерирует n перефразировок вопроса,
@sk-task hierarchical-indices#T3.3: parent context attach in QueryMulti (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) выполняет поиск по каждой, объединяет результаты через Reciprocal Rank Fusion. @sk-task arch-issues#T2.1: PII redaction в QueryMulti (AC-001)
func (*Pipeline) QuerySubDecompose ¶ added in v1.0.0
func (p *Pipeline) QuerySubDecompose(ctx context.Context, question string, topK int, decomposer domain.QueryDecomposer) (domain.RetrievalResult, error)
@sk-task sub-query-decomposition#T1.3: QuerySubDecompose (AC-002, AC-004, AC-007) QuerySubDecompose разбивает запрос на под-вопросы через decomposer, выполняет параллельный retrieval по каждому и объединяет результаты.
Семантика: - Декомпозиция: decomposer.Decompose → список под-вопросов. - Если под-вопросов <= 1, выполняется обычный Query с исходным запросом. - Параллельный embed+search для каждого под-вопроса (errgroup-style). - Merge: дедупликация по Chunk.ID, max score per chunk, сортировка по score desc. - Graceful degradation: ошибка декомпозиции → single-query по исходному запросу.
func (*Pipeline) QueryWithMetadataFilter ¶
func (p *Pipeline) QueryWithMetadataFilter(ctx context.Context, question string, topK int, filter domain.MetadataFilter) (domain.RetrievalResult, error)
@sk-task arch-issues#T2.1: PII redaction в QueryWithMetadataFilter (AC-001)
QueryWithMetadataFilter выполняет поиск по вопросу с фильтром по метаданным документа.
Если filter.Fields пустой — эквивалентно Query. Если store не реализует VectorStoreWithFilters — возвращает ErrFiltersNotSupported.
@ds-task T3.1: Добавить QueryWithMetadataFilter в application.Pipeline (RQ-005, AC-003, DEC-003) @sk-task hierarchical-indices#T3.3: parent context attach in QueryWithMetadataFilter (AC-002) @sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003)
func (*Pipeline) QueryWithParentIDs ¶
func (p *Pipeline) QueryWithParentIDs(ctx context.Context, question string, topK int, parentIDs []string) (domain.RetrievalResult, error)
@sk-task arch-issues#T2.1: PII redaction в QueryWithParentIDs (AC-001)
QueryWithParentIDs выполняет поиск по вопросу с фильтром по ParentIDs.
Если parentIDs пустой — эквивалентно Query.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T2.1: wrapped domain.ErrEmptyQueryText/ErrInvalidQueryTopK в validation (RQ-003, AC-003)
func (*Pipeline) QueryWithQueries ¶ added in v1.0.0
func (p *Pipeline) QueryWithQueries(ctx context.Context, queries []string, topK int) (domain.RetrievalResult, error)
@sk-task arch-issues#T2.1: PII redaction в QueryWithQueries (AC-001)
@sk-task query-rewriting#T2.2: QueryWithQueries для pre-generated переформулировок (AC-003) @sk-task hierarchical-indices#T3.3: parent context attach in QueryWithQueries (AC-002) QueryWithQueries выполняет multi-query retrieval из уже готового списка запросов. Каждый запрос эмбеддится и ищется, результаты объединяются через RRF.
func (*Pipeline) RedactRetrievalResult ¶ added in v1.0.0
func (p *Pipeline) RedactRetrievalResult(r domain.RetrievalResult) domain.RetrievalResult
@sk-task arch-issues#T2.2: экспортирован для pkg/draftrag делегирования (AC-001, AC-002)
func (*Pipeline) SystemPrompt ¶ added in v1.0.0
@sk-task hierarchical-indices#T3.1: parentEmbeddingOrEmbed helper (AC-001, DEC-003)
parentEmbeddingOrEmbed возвращает embedding для parent-документа. При отсутствии chunker'а возвращает nil — вызывающий код должен использовать embedding единственного чанка. При наличии chunker'а вызывает embedder для полного текста документа. @sk-task arch-issues#T4.4: SystemPrompt accessor для tool route handlers (AC-004)
func (*Pipeline) UpdateDocument ¶
UpdateDocument выполняет атомарное обновление документа через
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task api-consistency-pass#T3.2: делегирует в updateDocumentAtomic (DEC-005, RQ-005, AC-008, AC-009)
@sk-task arch-issues#T3.1: closed guard в UpdateDocument (AC-008)
UpdateDocument выполняет атомарное обновление документа через updateDocumentAtomic, который выбирает transactional или best-effort путь в зависимости от capability underlying store.
type PipelineOptions ¶ added in v0.2.0
type PipelineOptions struct {
SystemPrompt string
Chunker domain.Chunker
MaxContextChars int
MaxContextChunks int
DedupByParentID bool
MMREnabled bool
MMRLambda float64
MMRCandidatePool int
Hooks domain.Hooks
IndexConcurrency int
IndexBatchRateLimit int
IndexBatchRateLimitPerWorker bool
StreamBufferSize int
Reranker domain.Reranker
PIIDetector domain.PIIDetector
ParentContextEnabled *bool
Middleware []domain.Middleware
}
PipelineOptions configures a Pipeline behaviour.
@sk-task hardening-2026q2#T1.1: Разделить pipeline.go на модули (AC-001, AC-003) @sk-task arch-quality-pass#T3.2: единый struct конфигурации (AC-004) @sk-task pii-guardrails#T2.1: PipelineOptions.PIIDetector (RQ-001, RQ-002)