elasticsearch

package module
v0.44.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: Apache-2.0 Imports: 23 Imported by: 0

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

Examples

Constants

View Source
const (
	Provider                      = "Elasticsearch"
	DefaultIndexName              = "scope-vector-index"
	DefaultMaxResponseBytes int64 = 16 << 20
)

Variables

View Source
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

type APIClient interface {
	Perform(*http.Request) (*http.Response, error)
}

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)
	}
}

func (*Store) DeleteIDs

func (s *Store) DeleteIDs(ctx context.Context, ids []string) error

func (*Store) DeleteWhere

func (s *Store) DeleteWhere(ctx context.Context, predicate filter.Predicate) error

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

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL