Documentation
¶
Overview ¶
Package elasticsearch implements the Core vector-store capabilities using a borrowed native Elasticsearch SDK client. The host provisions the concrete index and owns the client, credentials, topology, retries, timeouts and lifetime. NewStore reads the native mapping and settings and validates current records; it never creates or modifies an index. The native index owns dimensions, similarity, vector index tuning and result-window policy. Native KNN also owns its candidate defaults; Scope does not reproduce a candidate multiplier.
The index must declare exactly content (text), embedding (indexed FLOAT32 dense_vector with explicit dimensions), and metadata_json (keyword with index:false and doc_values:false), under dynamic:strict. Complete stored _source is required. Aliases, required or per-document custom routing, runtime fields, dynamic templates and ingest pipelines are incompatible. The host must preserve these namespace and storage guarantees while Store is in use.
Native cosine, l2_norm, dot_product and max_inner_product similarities are supported. The adapter validates the complete FLOAT32 vectors, including native finite-magnitude and unit-vector constraints, before publication. Native scores are projected through Core; group merging uses native scores before Core normalization so saturated scores do not change ranking.
Each source record contains exactly three fields and uses the native _id as its sole identity. Metadata is one Core JSON string rather than native mapped terms. Nil and empty metadata remain distinct; arbitrary Core-supported keys, nested values, large integers and decimal values preserve their exact meaning. Media is unsupported. Native IDs are limited to 512 UTF-8 bytes. Existing deployments must provision the current schema and reindex their source documents; obsolete fields, configuration and record formats are not accepted.
Index prepares all records, embedding batches, native vectors and NDJSON before its first bulk publication. A validation or model failure publishes nothing. Native I/O failures may leave earlier bulk records published; no cross-record transaction is claimed. Every bulk acknowledgment must identify the requested operation, index and document, with a valid status and no hidden failure.
Search reads a complete native scroll snapshot before KNN retrieval. Every record is validated, and Core filter.Match alone decides metadata membership. Filtered KNN queries use bounded native ID selections, then merge native ranks before TopK and the Core score threshold. This costs O(N) source reads and document-ID bookkeeping. Enumeration and later KNN retrieval are separate requests; returned membership is revalidated and no cross-request snapshot is promised. Native TopK limits are checked before I/O.
Timeout, failed or missing shard acknowledgments, malformed records, repeated identities, truncated pages, missing concurrency tokens and unreadable hits fail the operation. Search returns no response on any error. Scroll cleanup uses a bounded context even after caller cancellation; cleanup failures remain visible. StoreConfig.MaxResponseBytes bounds every response and defaults to DefaultMaxResponseBytes; response bodies are always closed.
DeleteIDs deduplicates IDs and treats unknown IDs as successful no-ops. DeleteWhere enumerates the complete metadata snapshot before writing, then uses native sequence-number and primary-term conditional deletes. A changed or vanished snapshot member fails rather than deleting a newer version or silently succeeding. Earlier deletions may remain applied after a failure.
Default tests are offline. Tests selected with -tags=integration require SCOPE_ELASTICSEARCH_ENDPOINT and optionally SCOPE_ELASTICSEARCH_USERNAME and SCOPE_ELASTICSEARCH_PASSWORD. They create and delete unique native indexes; run them only against an isolated instance. Native verification exercises Elasticsearch 8.19.7 and its current KNN and stored-source contracts.
See https://www.elastic.co/guide/en/elasticsearch/reference/8.19/dense-vector.html and https://www.elastic.co/guide/en/elasticsearch/reference/8.19/search-search.html.
Index ¶
- Constants
- Variables
- type APIClient
- type Store
- func (s *Store) DeleteIDs(ctx context.Context, ids []string) error
- func (s *Store) DeleteWhere(ctx context.Context, predicate filter.Predicate) error
- func (s *Store) Index(ctx context.Context, request *vectorstore.IndexRequest) error
- func (s *Store) Search(ctx context.Context, request *vectorstore.SearchRequest) (response *vectorstore.SearchResponse, err error)
- type StoreConfig
Examples ¶
Constants ¶
const ( Provider = "Elasticsearch" DefaultIndexName = "scope-vector-index" DefaultMaxResponseBytes int64 = 16 << 20 )
Variables ¶
var ( ErrIndexMissing = errors.New("elasticsearch: index not found") ErrIncompatibleIndex = errors.New("elasticsearch: native index is incompatible") )
Functions ¶
This section is empty.
Types ¶
type APIClient ¶ added in v0.43.0
APIClient borrows the native SDK's transport. The host owns authentication, topology, retries, timeouts and lifetime; Store closes every response body.
type Store ¶
type Store struct {
// contains filtered or unexported fields
}
func NewStore ¶
func NewStore(ctx context.Context, config StoreConfig) (*Store, error)
Example ¶
package main
import (
"context"
"errors"
"net/http"
"os"
"slices"
"strings"
"time"
elasticsearchsdk "github.com/elastic/go-elasticsearch/v8"
"github.com/Tangerg/scope/core/document"
"github.com/Tangerg/scope/core/embedding"
"github.com/Tangerg/scope/core/vectorstore"
"github.com/Tangerg/scope/vectorstores/elasticsearch"
)
type exampleBatcher struct{}
func (e exampleBatcher) Batch(_ context.Context, docs []*document.Document) ([][]*document.Document, error) {
return slices.Collect(slices.Chunk(docs, 32)), nil
}
func main() {
ctx, cancel := context.WithTimeout(context.Background(), time.Minute)
defer cancel()
transport := http.DefaultTransport.(*http.Transport).Clone()
defer transport.CloseIdleConnections()
client, err := elasticsearchsdk.NewClient(elasticsearchsdk.Config{Addresses: []string{os.Getenv("SCOPE_ELASTICSEARCH_ENDPOINT")}, Username: os.Getenv("SCOPE_ELASTICSEARCH_USERNAME"), Password: os.Getenv("SCOPE_ELASTICSEARCH_PASSWORD"), Transport: transport})
if err != nil {
panic(err)
}
defer func() {
cleanup, stop := context.WithTimeout(context.WithoutCancel(ctx), 5*time.Second)
defer stop()
if closeErr := client.Close(cleanup); closeErr != nil {
panic(closeErr)
}
}()
// Provision once through the native host API before constructing Store.
created, err := client.Indices.Create("documents", client.Indices.Create.WithContext(ctx), client.Indices.Create.WithBody(strings.NewReader(`{"mappings":{"dynamic":"strict","properties":{"content":{"type":"text"},"embedding":{"type":"dense_vector","dims":2,"similarity":"cosine"},"metadata_json":{"type":"keyword","index":false,"doc_values":false}}}}`)))
if err != nil {
panic(err)
}
closeErr := created.Body.Close()
if created.IsError() || closeErr != nil {
panic(errors.Join(errors.New("native index provisioning failed"), closeErr))
}
model := embedding.ModelFunc(func(_ context.Context, request *embedding.Request) (*embedding.Response, error) {
outputs := make([]*embedding.Output, len(request.Texts))
for i := range outputs {
outputs[i] = &embedding.Output{Embedding: []float64{1, 0}}
}
return embedding.NewResponse(outputs, nil)
})
store, err := elasticsearch.NewStore(ctx, elasticsearch.StoreConfig{Client: client, IndexName: "documents", EmbeddingModel: model, DocumentBatcher: exampleBatcher{}})
if err != nil {
panic(err)
}
if err = store.Index(ctx, &vectorstore.IndexRequest{Documents: []*document.Document{{ID: "one", Text: "example"}}}); err != nil {
panic(err)
}
if _, err = store.Search(ctx, &vectorstore.SearchRequest{Query: "example"}); err != nil {
panic(err)
}
if err = store.DeleteIDs(ctx, []string{"one"}); err != nil {
panic(err)
}
}
Output:
func (*Store) DeleteWhere ¶
func (*Store) Index ¶
func (s *Store) Index(ctx context.Context, request *vectorstore.IndexRequest) error
func (*Store) Search ¶
func (s *Store) Search(ctx context.Context, request *vectorstore.SearchRequest) (response *vectorstore.SearchResponse, err error)
type StoreConfig ¶
type StoreConfig struct {
Client APIClient
IndexName string
EmbeddingModel embedding.Model
DocumentBatcher vectorstore.Batcher
MaxResponseBytes int64
}
func (StoreConfig) Validate ¶
func (s StoreConfig) Validate() error