feat(11-05): add the beachcomber search sync package and its Typesense engine
- Searchable, Engine, Gate and an init-time engine registry with the null engine - GORM callbacks installed through lagoon.OnDatabase register an after-commit sync that reloads the row and upserts or deletes its document - Gates run before any request: engine configured, database published, app Gate - hand-rolled net/http Typesense engine following the Scout wire contract - module README and root modules row
This commit is contained in:
89
modules/beachcomber/typesense/config.go
Normal file
89
modules/beachcomber/typesense/config.go
Normal file
@@ -0,0 +1,89 @@
|
||||
// Package typesense is the Typesense engine of beachcomber: a hand-rolled
|
||||
// net/http client for the collection, import, delete and search endpoints.
|
||||
// Import it for its side effect to register the "typesense" search.driver.
|
||||
package typesense
|
||||
|
||||
import (
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"git.golem15.com/golem15/summercms/modules/compass"
|
||||
)
|
||||
|
||||
// Default values of the search.typesense.* keys.
|
||||
const (
|
||||
DefaultHost = "localhost"
|
||||
DefaultPort = 8181
|
||||
DefaultProtocol = "http"
|
||||
DefaultConnectionTimeout = 2 * time.Second
|
||||
DefaultImportAction = "upsert"
|
||||
)
|
||||
|
||||
// Config is the search.typesense.* configuration.
|
||||
type Config struct {
|
||||
// APIKey is sent as X-TYPESENSE-API-KEY; empty means not configured,
|
||||
// and nothing is ever sent.
|
||||
APIKey string
|
||||
// Host, Port, Protocol and Path form the base URL
|
||||
// {protocol}://{host}:{port}{path}.
|
||||
Host string
|
||||
Port int
|
||||
Protocol string
|
||||
Path string
|
||||
// ConnectionTimeout bounds each request, including reading the answer.
|
||||
ConnectionTimeout time.Duration
|
||||
// ImportAction is the action query parameter of document imports.
|
||||
ImportAction string
|
||||
}
|
||||
|
||||
// LoadConfig reads search.typesense.* from c, filling the defaults.
|
||||
// connection_timeout_seconds is an integer number of seconds or a duration
|
||||
// string.
|
||||
func LoadConfig(c *compass.Config) Config {
|
||||
cfg := Config{
|
||||
Host: DefaultHost,
|
||||
Port: DefaultPort,
|
||||
Protocol: DefaultProtocol,
|
||||
ConnectionTimeout: DefaultConnectionTimeout,
|
||||
ImportAction: DefaultImportAction,
|
||||
}
|
||||
if c == nil {
|
||||
return cfg
|
||||
}
|
||||
str := func(key string) string { return strings.TrimSpace(c.String("search.typesense." + key)) }
|
||||
cfg.APIKey = str("api_key")
|
||||
if v := str("host"); v != "" {
|
||||
cfg.Host = v
|
||||
}
|
||||
if v := str("port"); v != "" {
|
||||
if n, err := strconv.Atoi(v); err == nil && n > 0 {
|
||||
cfg.Port = n
|
||||
}
|
||||
}
|
||||
if v := str("protocol"); v != "" {
|
||||
cfg.Protocol = strings.ToLower(v)
|
||||
}
|
||||
if v := strings.TrimSuffix(str("path"), "/"); v != "" {
|
||||
if !strings.HasPrefix(v, "/") {
|
||||
v = "/" + v
|
||||
}
|
||||
cfg.Path = v
|
||||
}
|
||||
if v := str("connection_timeout_seconds"); v != "" {
|
||||
if n, err := strconv.ParseFloat(v, 64); err == nil && n > 0 {
|
||||
cfg.ConnectionTimeout = time.Duration(n * float64(time.Second))
|
||||
} else if d, err := time.ParseDuration(v); err == nil && d > 0 {
|
||||
cfg.ConnectionTimeout = d
|
||||
}
|
||||
}
|
||||
if v := str("import_action"); v != "" {
|
||||
cfg.ImportAction = v
|
||||
}
|
||||
return cfg
|
||||
}
|
||||
|
||||
// BaseURL is {protocol}://{host}:{port}{path}.
|
||||
func (c Config) BaseURL() string {
|
||||
return c.Protocol + "://" + c.Host + ":" + strconv.Itoa(c.Port) + c.Path
|
||||
}
|
||||
322
modules/beachcomber/typesense/engine.go
Normal file
322
modules/beachcomber/typesense/engine.go
Normal file
@@ -0,0 +1,322 @@
|
||||
package typesense
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"maps"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"git.golem15.com/golem15/summercms/modules/backpack"
|
||||
"git.golem15.com/golem15/summercms/modules/beachcomber"
|
||||
)
|
||||
|
||||
// DriverName is the search.driver value of this engine.
|
||||
const DriverName = "typesense"
|
||||
|
||||
// apiKeyHeader authenticates every request.
|
||||
const apiKeyHeader = "X-TYPESENSE-API-KEY"
|
||||
|
||||
// maxResponseBytes caps how much of an answer is read.
|
||||
const maxResponseBytes = 32 << 20
|
||||
|
||||
func init() {
|
||||
beachcomber.RegisterEngine(DriverName, func(app *backpack.App) (beachcomber.Engine, error) {
|
||||
var cfg Config
|
||||
if app != nil {
|
||||
cfg = LoadConfig(app.Config)
|
||||
} else {
|
||||
cfg = LoadConfig(nil)
|
||||
}
|
||||
return New(cfg), nil
|
||||
})
|
||||
}
|
||||
|
||||
// Engine is the Typesense beachcomber.Engine. It follows the Laravel Scout
|
||||
// TypesenseEngine wire contract.
|
||||
type Engine struct {
|
||||
cfg Config
|
||||
base string
|
||||
client *http.Client
|
||||
}
|
||||
|
||||
var _ beachcomber.Engine = (*Engine)(nil)
|
||||
|
||||
// New returns an engine for cfg. Every request is bounded by
|
||||
// cfg.ConnectionTimeout.
|
||||
func New(cfg Config) *Engine {
|
||||
if cfg.ConnectionTimeout <= 0 {
|
||||
cfg.ConnectionTimeout = DefaultConnectionTimeout
|
||||
}
|
||||
if cfg.ImportAction == "" {
|
||||
cfg.ImportAction = DefaultImportAction
|
||||
}
|
||||
return &Engine{
|
||||
cfg: cfg,
|
||||
base: strings.TrimSuffix(cfg.BaseURL(), "/"),
|
||||
client: &http.Client{Timeout: cfg.ConnectionTimeout},
|
||||
}
|
||||
}
|
||||
|
||||
// Name returns "typesense".
|
||||
func (e *Engine) Name() string { return DriverName }
|
||||
|
||||
// Configured reports whether an API key is set. Without one the engine is
|
||||
// never called.
|
||||
func (e *Engine) Configured() bool { return e != nil && e.cfg.APIKey != "" }
|
||||
|
||||
// Config returns the engine's configuration.
|
||||
func (e *Engine) Config() Config { return e.cfg }
|
||||
|
||||
// Upsert makes sure the collection exists (GET /collections/{index}, then
|
||||
// POST /collections with schema and the name when it answers 404), then
|
||||
// imports docs as JSON lines with POST
|
||||
// /collections/{index}/documents/import?action={import_action}. Typesense
|
||||
// answers 200 even when documents fail, so every answer line is checked
|
||||
// and any "success":false line is an error.
|
||||
func (e *Engine) Upsert(ctx context.Context, index string, schema map[string]any, docs []map[string]any) error {
|
||||
if len(docs) == 0 {
|
||||
return nil
|
||||
}
|
||||
if err := e.ensureCollection(ctx, index, schema); err != nil {
|
||||
return err
|
||||
}
|
||||
var body bytes.Buffer
|
||||
for _, doc := range docs {
|
||||
line, err := json.Marshal(doc)
|
||||
if err != nil {
|
||||
return fmt.Errorf("typesense: encode document: %w", err)
|
||||
}
|
||||
body.Write(line)
|
||||
body.WriteByte('\n')
|
||||
}
|
||||
q := url.Values{"action": {e.cfg.ImportAction}}
|
||||
path := "/collections/" + url.PathEscape(index) + "/documents/import"
|
||||
code, answer, err := e.do(ctx, http.MethodPost, path, q, &body, "text/plain")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !ok2xx(code) {
|
||||
return statusError(http.MethodPost, path, code)
|
||||
}
|
||||
return checkImport(index, answer, len(docs))
|
||||
}
|
||||
|
||||
// ensureCollection creates index from schema when Typesense does not have
|
||||
// it. A missing schema creates an auto-typed collection.
|
||||
func (e *Engine) ensureCollection(ctx context.Context, index string, schema map[string]any) error {
|
||||
path := "/collections/" + url.PathEscape(index)
|
||||
code, _, err := e.do(ctx, http.MethodGet, path, nil, nil, "")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if ok2xx(code) {
|
||||
return nil
|
||||
}
|
||||
if code != http.StatusNotFound {
|
||||
return statusError(http.MethodGet, path, code)
|
||||
}
|
||||
create := map[string]any{}
|
||||
maps.Copy(create, schema)
|
||||
if _, ok := create["fields"]; !ok {
|
||||
create["fields"] = []map[string]any{{"name": ".*", "type": "auto"}}
|
||||
}
|
||||
create["name"] = index
|
||||
raw, err := json.Marshal(create)
|
||||
if err != nil {
|
||||
return fmt.Errorf("typesense: encode collection schema: %w", err)
|
||||
}
|
||||
code, _, err = e.do(ctx, http.MethodPost, "/collections", nil, bytes.NewReader(raw), "application/json")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// 409: created concurrently by another writer.
|
||||
if ok2xx(code) || code == http.StatusConflict {
|
||||
return nil
|
||||
}
|
||||
return statusError(http.MethodPost, "/collections", code)
|
||||
}
|
||||
|
||||
// importLine is one line of an import answer.
|
||||
type importLine struct {
|
||||
Success bool `json:"success"`
|
||||
Error string `json:"error"`
|
||||
}
|
||||
|
||||
// checkImport returns an error when any answer line reports a failure. The
|
||||
// error carries the Typesense message, never the document.
|
||||
func checkImport(index string, answer []byte, total int) error {
|
||||
failed := 0
|
||||
first := ""
|
||||
sc := bufio.NewScanner(bytes.NewReader(answer))
|
||||
sc.Buffer(make([]byte, 0, 64*1024), maxResponseBytes)
|
||||
for sc.Scan() {
|
||||
line := bytes.TrimSpace(sc.Bytes())
|
||||
if len(line) == 0 {
|
||||
continue
|
||||
}
|
||||
var l importLine
|
||||
if err := json.Unmarshal(line, &l); err != nil {
|
||||
return fmt.Errorf("typesense: import into %s: unreadable answer line", index)
|
||||
}
|
||||
if !l.Success {
|
||||
failed++
|
||||
if first == "" {
|
||||
first = l.Error
|
||||
}
|
||||
}
|
||||
}
|
||||
if err := sc.Err(); err != nil {
|
||||
return fmt.Errorf("typesense: import into %s: read answer: %w", index, err)
|
||||
}
|
||||
if failed == 0 {
|
||||
return nil
|
||||
}
|
||||
if len(first) > 200 {
|
||||
first = first[:200]
|
||||
}
|
||||
return fmt.Errorf("typesense: import into %s: %d of %d documents failed: %s", index, failed, total, first)
|
||||
}
|
||||
|
||||
// Delete sends DELETE /collections/{index}/documents/{id} per id. A 404
|
||||
// (document or collection already gone) counts as success.
|
||||
func (e *Engine) Delete(ctx context.Context, index string, ids []string) error {
|
||||
for _, id := range ids {
|
||||
path := "/collections/" + url.PathEscape(index) + "/documents/" + url.PathEscape(id)
|
||||
code, _, err := e.do(ctx, http.MethodDelete, path, nil, nil, "")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !ok2xx(code) && code != http.StatusNotFound {
|
||||
return statusError(http.MethodDelete, path, code)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Flush drops the collection with DELETE /collections/{index}. A 404
|
||||
// counts as success.
|
||||
func (e *Engine) Flush(ctx context.Context, index string) error {
|
||||
path := "/collections/" + url.PathEscape(index)
|
||||
code, _, err := e.do(ctx, http.MethodDelete, path, nil, nil, "")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !ok2xx(code) && code != http.StatusNotFound {
|
||||
return statusError(http.MethodDelete, path, code)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// searchAnswer is the part of a search answer SearchIDs reads.
|
||||
type searchAnswer struct {
|
||||
Hits []struct {
|
||||
Document struct {
|
||||
ID json.RawMessage `json:"id"`
|
||||
} `json:"document"`
|
||||
} `json:"hits"`
|
||||
}
|
||||
|
||||
// SearchIDs sends GET /collections/{index}/documents/search with q
|
||||
// (default "*"), query_by, filter_by, sort_by, page and per_page, and
|
||||
// returns hits[].document.id in order. The ids are candidates: the caller
|
||||
// must re-check each one in SQL before exposing it. No hits is an empty
|
||||
// list.
|
||||
func (e *Engine) SearchIDs(ctx context.Context, index string, q beachcomber.Query) ([]string, error) {
|
||||
params := url.Values{}
|
||||
text := q.Q
|
||||
if text == "" {
|
||||
text = "*"
|
||||
}
|
||||
params.Set("q", text)
|
||||
if len(q.QueryBy) > 0 {
|
||||
params.Set("query_by", strings.Join(q.QueryBy, ","))
|
||||
}
|
||||
if q.FilterBy != "" {
|
||||
params.Set("filter_by", q.FilterBy)
|
||||
}
|
||||
if q.SortBy != "" {
|
||||
params.Set("sort_by", q.SortBy)
|
||||
}
|
||||
if q.Page > 0 {
|
||||
params.Set("page", strconv.Itoa(q.Page))
|
||||
}
|
||||
if q.PerPage > 0 {
|
||||
params.Set("per_page", strconv.Itoa(q.PerPage))
|
||||
}
|
||||
path := "/collections/" + url.PathEscape(index) + "/documents/search"
|
||||
code, answer, err := e.do(ctx, http.MethodGet, path, params, nil, "")
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if !ok2xx(code) {
|
||||
return nil, statusError(http.MethodGet, path, code)
|
||||
}
|
||||
var res searchAnswer
|
||||
if err := json.Unmarshal(answer, &res); err != nil {
|
||||
return nil, fmt.Errorf("typesense: search %s: unreadable answer", index)
|
||||
}
|
||||
ids := make([]string, 0, len(res.Hits))
|
||||
for _, h := range res.Hits {
|
||||
raw := bytes.TrimSpace(h.Document.ID)
|
||||
if len(raw) == 0 || bytes.Equal(raw, []byte("null")) {
|
||||
continue
|
||||
}
|
||||
var s string
|
||||
if err := json.Unmarshal(raw, &s); err == nil {
|
||||
ids = append(ids, s)
|
||||
continue
|
||||
}
|
||||
ids = append(ids, string(raw))
|
||||
}
|
||||
return ids, nil
|
||||
}
|
||||
|
||||
// do sends one request with the API key header and returns the status and
|
||||
// the (capped) body.
|
||||
func (e *Engine) do(ctx context.Context, method, path string, query url.Values, body io.Reader, contentType string) (int, []byte, error) {
|
||||
u := e.base + path
|
||||
if len(query) > 0 {
|
||||
u += "?" + query.Encode()
|
||||
}
|
||||
req, err := http.NewRequestWithContext(ctx, method, u, body)
|
||||
if err != nil {
|
||||
return 0, nil, fmt.Errorf("typesense: %s %s: build request: %w", method, path, err)
|
||||
}
|
||||
req.Header.Set(apiKeyHeader, e.cfg.APIKey)
|
||||
req.Header.Set("Accept", "application/json")
|
||||
if contentType != "" {
|
||||
req.Header.Set("Content-Type", contentType)
|
||||
}
|
||||
resp, err := e.client.Do(req)
|
||||
if err != nil {
|
||||
return 0, nil, fmt.Errorf("typesense: %s %s: %w", method, path, unwrapURLError(err))
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
answer, err := io.ReadAll(io.LimitReader(resp.Body, maxResponseBytes))
|
||||
if err != nil {
|
||||
return resp.StatusCode, nil, fmt.Errorf("typesense: %s %s: read answer: %w", method, path, err)
|
||||
}
|
||||
return resp.StatusCode, answer, nil
|
||||
}
|
||||
|
||||
// unwrapURLError drops the *url.Error wrapper, whose message repeats the
|
||||
// full URL, and keeps the cause (a timeout, a refused connection).
|
||||
func unwrapURLError(err error) error {
|
||||
if ue, ok := err.(*url.Error); ok && ue.Err != nil {
|
||||
return ue.Err
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func ok2xx(code int) bool { return code >= 200 && code < 300 }
|
||||
|
||||
func statusError(method, path string, code int) error {
|
||||
return fmt.Errorf("typesense: %s %s: status %d", method, path, code)
|
||||
}
|
||||
Reference in New Issue
Block a user