test(11-07): pin the Typesense engine wire contract

- TestEngineWire: API key and Accept on every request, create on 404 with
  the schema or auto fields, 409 as success, JSONL import with
  success:false and unreadable lines as errors (message capped, no
  document), 404-tolerant delete/flush with id escaping, SearchIDs
  parameters and id parsing, typed errors without bodies, transport
  errors without the URL, timeout
- TestEngineConfig: search.typesense.* parsing and the registered engine
  (coverage 95.6%)
This commit is contained in:
Jakub Zych
2026-09-30 14:28:02 +02:00
parent 6df43d45b8
commit 18d3097be5

View File

@@ -0,0 +1,340 @@
package typesense
import (
"context"
"encoding/json"
"errors"
"io"
"net/http"
"net/http/httptest"
"net/url"
"strconv"
"strings"
"sync"
"testing"
"time"
"git.golem15.com/golem15/summercms/modules/backpack"
"git.golem15.com/golem15/summercms/modules/beachcomber"
"git.golem15.com/golem15/summercms/modules/compass"
)
const testKey = "test-only-typesense-key-5e1a"
type tsCall struct {
method, path, rawQuery, key, contentType, accept, body string
}
// fakeTypesense answers by "METHOD path" from a table; unlisted requests
// get 404.
type fakeTypesense struct {
mu sync.Mutex
calls []tsCall
answers map[string]fakeAnswer
}
type fakeAnswer struct {
status int
body string
}
func newFake(t *testing.T) (*fakeTypesense, *Engine) {
t.Helper()
f := &fakeTypesense{answers: map[string]fakeAnswer{}}
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
f.mu.Lock()
f.calls = append(f.calls, tsCall{r.Method, r.URL.EscapedPath(), r.URL.RawQuery, r.Header.Get(apiKeyHeader), r.Header.Get("Content-Type"), r.Header.Get("Accept"), string(body)})
a, ok := f.answers[r.Method+" "+r.URL.EscapedPath()]
f.mu.Unlock()
if !ok {
w.WriteHeader(http.StatusNotFound)
_, _ = io.WriteString(w, `{"message":"Not Found"}`)
return
}
w.WriteHeader(a.status)
_, _ = io.WriteString(w, a.body)
}))
t.Cleanup(srv.Close)
u, err := url.Parse(srv.URL)
if err != nil {
t.Fatal(err)
}
port, _ := strconv.Atoi(u.Port())
return f, New(Config{APIKey: testKey, Host: u.Hostname(), Port: port, Protocol: "http", Path: "/ts"})
}
func (f *fakeTypesense) answer(key string, status int, body string) {
f.mu.Lock()
f.answers[key] = fakeAnswer{status, body}
f.mu.Unlock()
}
func (f *fakeTypesense) take() []tsCall {
f.mu.Lock()
defer f.mu.Unlock()
out := f.calls
f.calls = nil
return out
}
func statusOf(err error) int {
var se *StatusError
if errors.As(err, &se) {
return se.StatusCode()
}
return 0
}
// TestEngineWire covers the Scout Typesense wire contract (SRCH-01, D-19):
// the API key on every request, collection create on 404 (with the schema
// or an auto-typed field set, 409 as success), JSON-lines import with
// success:false lines as errors, 404-tolerant delete and flush, search
// parameters and id parsing, escaping, and typed errors without bodies.
func TestEngineWire(t *testing.T) {
ctx := context.Background()
schema := map[string]any{"default_sorting_field": "created_at", "fields": []map[string]any{{"name": "title", "type": "string"}}}
docs := []map[string]any{{"id": "1", "title": "a/b"}, {"id": "2", "title": "c"}}
t.Run("upsert_creates_missing_collection", func(t *testing.T) {
f, e := newFake(t)
f.answer("POST /ts/collections", 201, `{}`)
f.answer("POST /ts/collections/dev_albums/documents/import", 200, "{\"success\":true}\n{\"success\":true}\n")
if err := e.Upsert(ctx, "dev_albums", schema, docs); err != nil {
t.Fatal(err)
}
calls := f.take()
if len(calls) != 3 {
t.Fatalf("calls = %+v", calls)
}
if calls[0].method != "GET" || calls[0].path != "/ts/collections/dev_albums" {
t.Fatalf("first call = %+v", calls[0])
}
var created map[string]any
if err := json.Unmarshal([]byte(calls[1].body), &created); err != nil {
t.Fatal(err)
}
if calls[1].path != "/ts/collections" || calls[1].contentType != "application/json" || created["name"] != "dev_albums" || created["default_sorting_field"] != "created_at" {
t.Fatalf("create = %+v", calls[1])
}
if schema["name"] != nil {
t.Fatal("Upsert mutated the caller's schema")
}
imp := calls[2]
if imp.method != "POST" || imp.rawQuery != "action=upsert" || imp.contentType != "text/plain" || imp.body != "{\"id\":\"1\",\"title\":\"a/b\"}\n{\"id\":\"2\",\"title\":\"c\"}\n" {
t.Fatalf("import = %+v", imp)
}
for _, c := range calls {
if c.key != testKey || c.accept != "application/json" {
t.Fatalf("call without the key or Accept header: %+v", c)
}
}
})
t.Run("existing_collection_and_auto_schema", func(t *testing.T) {
f, e := newFake(t)
f.answer("GET /ts/collections/idx", 200, `{"name":"idx"}`)
f.answer("POST /ts/collections/idx/documents/import", 200, `{"success":true}`)
if err := e.Upsert(ctx, "idx", nil, docs[:1]); err != nil {
t.Fatal(err)
}
if calls := f.take(); len(calls) != 2 {
t.Fatalf("calls = %+v, want no create", calls)
}
f2, e2 := newFake(t)
f2.answer("POST /ts/collections", 409, `{"message":"exists"}`)
f2.answer("POST /ts/collections/auto/documents/import", 200, `{"success":true}`)
if err := e2.Upsert(ctx, "auto", nil, docs[:1]); err != nil {
t.Fatalf("409 on create: %v", err)
}
if body := f2.take()[1].body; !strings.Contains(body, `"fields":[{"name":".*","type":"auto"}]`) {
t.Fatalf("auto schema = %s", body)
}
if err := e2.Upsert(ctx, "auto", nil, nil); err != nil || len(f2.take()) != 0 {
t.Fatal("empty upsert sent requests")
}
})
t.Run("upsert_failures", func(t *testing.T) {
f, e := newFake(t)
f.answer("GET /ts/collections/g", 500, `{"message":"`+testKey+`"}`)
if err := e.Upsert(ctx, "g", nil, docs); statusOf(err) != 500 || strings.Contains(err.Error(), testKey) {
t.Fatalf("GET 500: %v", err)
}
f.answer("GET /ts/collections/c", 404, `{}`)
f.answer("POST /ts/collections", 400, `{"message":"bad schema"}`)
if err := e.Upsert(ctx, "c", nil, docs); statusOf(err) != 400 {
t.Fatalf("create 400: %v", err)
}
f.answer("GET /ts/collections/i", 200, `{}`)
f.answer("POST /ts/collections/i/documents/import", 503, `busy`)
if err := e.Upsert(ctx, "i", nil, docs); statusOf(err) != 503 || !strings.Contains(err.Error(), "POST /collections/i/documents/import: status 503") {
t.Fatalf("import 503: %v", err)
}
long := strings.Repeat("x", 300)
f.answer("POST /ts/collections/i/documents/import", 200, "{\"success\":true}\n\n{\"success\":false,\"error\":\""+long+"\",\"document\":\"{\\\"title\\\":\\\"private\\\"}\"}\n")
err := e.Upsert(ctx, "i", nil, docs)
if err == nil || !strings.Contains(err.Error(), "1 of 2 documents failed") || strings.Contains(err.Error(), long) || strings.Contains(err.Error(), "private") {
t.Fatalf("success:false: %v", err)
}
f.answer("POST /ts/collections/i/documents/import", 200, "not json\n")
if err := e.Upsert(ctx, "i", nil, docs); err == nil || !strings.Contains(err.Error(), "unreadable answer line") {
t.Fatalf("unreadable line: %v", err)
}
if err := e.Upsert(ctx, "i", nil, []map[string]any{{"bad": make(chan int)}}); err == nil {
t.Fatal("unencodable document accepted")
}
if err := e.Upsert(ctx, "c", map[string]any{"bad": make(chan int)}, docs); err == nil {
t.Fatal("unencodable schema accepted")
}
})
t.Run("delete_and_flush", func(t *testing.T) {
f, e := newFake(t)
f.answer("DELETE /ts/collections/idx/documents/1", 200, `{}`)
if err := e.Delete(ctx, "idx", []string{"1", "a/b"}); err != nil {
t.Fatalf("delete with a 404 for the second id: %v", err)
}
calls := f.take()
if len(calls) != 2 || calls[1].path != "/ts/collections/idx/documents/a%2Fb" {
t.Fatalf("delete calls = %+v", calls)
}
f.answer("DELETE /ts/collections/idx/documents/9", 500, `{}`)
if err := e.Delete(ctx, "idx", []string{"9"}); statusOf(err) != 500 {
t.Fatalf("delete 500: %v", err)
}
if err := e.Flush(ctx, "gone"); err != nil {
t.Fatalf("flush 404: %v", err)
}
f.answer("DELETE /ts/collections/idx", 200, `{}`)
if err := e.Flush(ctx, "idx"); err != nil {
t.Fatal(err)
}
f.answer("DELETE /ts/collections/bad", 500, `{}`)
if err := e.Flush(ctx, "bad"); statusOf(err) != 500 {
t.Fatalf("flush 500: %v", err)
}
})
t.Run("search_ids", func(t *testing.T) {
f, e := newFake(t)
f.answer("GET /ts/collections/idx/documents/search", 200, `{"hits":[{"document":{"id":"12"}},{"document":{"id":7}},{"document":{"id":null}},{"document":{}},{"document":{"id":"3"}}]}`)
ids, err := e.SearchIDs(ctx, "idx", beachcomber.Query{Q: "blue note", QueryBy: []string{"name", "artist"}, FilterBy: "collection_id:=5", SortBy: "created_at:desc", Page: 2, PerPage: 50})
if err != nil {
t.Fatal(err)
}
if strings.Join(ids, ",") != "12,7,3" {
t.Fatalf("ids = %v", ids)
}
q, err := url.ParseQuery(f.take()[0].rawQuery)
if err != nil {
t.Fatal(err)
}
want := url.Values{"q": {"blue note"}, "query_by": {"name,artist"}, "filter_by": {"collection_id:=5"}, "sort_by": {"created_at:desc"}, "page": {"2"}, "per_page": {"50"}}
if q.Encode() != want.Encode() {
t.Fatalf("query = %v, want %v", q, want)
}
f.answer("GET /ts/collections/empty/documents/search", 200, `{"hits":[]}`)
ids, err = e.SearchIDs(ctx, "empty", beachcomber.Query{})
if err != nil || ids == nil || len(ids) != 0 {
t.Fatalf("no hits = %v, %v", ids, err)
}
if q, _ := url.ParseQuery(f.take()[0].rawQuery); q.Encode() != "q=%2A" {
t.Fatalf("default query = %v", q)
}
f.answer("GET /ts/collections/bad/documents/search", 200, `nope`)
if _, err := e.SearchIDs(ctx, "bad", beachcomber.Query{}); err == nil || !strings.Contains(err.Error(), "unreadable answer") {
t.Fatalf("unreadable: %v", err)
}
if _, err := e.SearchIDs(ctx, "missing", beachcomber.Query{}); statusOf(err) != 404 {
t.Fatalf("404: %v", err)
}
})
t.Run("transport_errors_hide_the_url", func(t *testing.T) {
e := New(Config{APIKey: testKey, Host: "127.0.0.1", Port: 1, Protocol: "http", ConnectionTimeout: time.Second})
err := e.Flush(ctx, "idx")
if err == nil || strings.Contains(err.Error(), "http://127.0.0.1:1") || strings.Contains(err.Error(), testKey) {
t.Fatalf("transport error = %v", err)
}
bad := New(Config{APIKey: testKey, Host: "bad host", Port: 1, Protocol: "http"})
if err := bad.Flush(ctx, "idx"); err == nil {
t.Fatal("malformed base URL accepted")
}
hang := make(chan struct{})
slow := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { <-hang }))
defer slow.Close()
defer close(hang)
u, _ := url.Parse(slow.URL)
port, _ := strconv.Atoi(u.Port())
timed := New(Config{APIKey: testKey, Host: u.Hostname(), Port: port, Protocol: "http", ConnectionTimeout: 200 * time.Millisecond})
start := time.Now()
if err := timed.Flush(ctx, "idx"); err == nil || time.Since(start) > 3*time.Second {
t.Fatalf("timeout: %v after %s", err, time.Since(start))
}
})
}
// TestEngineConfig covers search.typesense.* parsing, the defaults and the
// registered engine.
func TestEngineConfig(t *testing.T) {
def := LoadConfig(nil)
if def.BaseURL() != "http://localhost:8181" || def.ConnectionTimeout != DefaultConnectionTimeout || def.ImportAction != DefaultImportAction || def.APIKey != "" {
t.Fatalf("defaults = %+v", def)
}
cfg, err := compass.Open(compass.Options{Dir: t.TempDir(), Env: "testing", Environ: []string{}})
if err != nil {
t.Fatal(err)
}
for k, v := range map[string]any{
"search.driver": "typesense",
"search.typesense.api_key": " k ",
"search.typesense.host": "ts.internal",
"search.typesense.port": "8108",
"search.typesense.protocol": "HTTPS",
"search.typesense.path": "typesense/",
"search.typesense.connection_timeout_seconds": "1.5",
"search.typesense.import_action": "emplace",
} {
if err := cfg.Set(k, v); err != nil {
t.Fatal(err)
}
}
got := LoadConfig(cfg)
if got.BaseURL() != "https://ts.internal:8108/typesense" || got.APIKey != "k" || got.ConnectionTimeout != 1500*time.Millisecond || got.ImportAction != "emplace" {
t.Fatalf("config = %+v (%s)", got, got.BaseURL())
}
for raw, want := range map[string]time.Duration{"3s": 3 * time.Second, "nope": DefaultConnectionTimeout, "-1": DefaultConnectionTimeout} {
if err := cfg.Set("search.typesense.connection_timeout_seconds", raw); err != nil {
t.Fatal(err)
}
if d := LoadConfig(cfg).ConnectionTimeout; d != want {
t.Errorf("timeout %q = %s, want %s", raw, d, want)
}
}
if err := cfg.Set("search.typesense.port", "not-a-port"); err != nil {
t.Fatal(err)
}
if LoadConfig(cfg).Port != DefaultPort {
t.Fatal("invalid port did not keep the default")
}
svc, err := beachcomber.From(backpack.New(cfg))
if err != nil {
t.Fatal(err)
}
e, ok := svc.Engine().(*Engine)
if !ok || e.Name() != DriverName || !e.Configured() || e.Config().Host != "ts.internal" {
t.Fatalf("engine = %#v", svc.Engine())
}
bare := New(Config{})
if bare.Configured() || bare.Config().ConnectionTimeout != DefaultConnectionTimeout || bare.Config().ImportAction != DefaultImportAction {
t.Fatal("New does not fill the timeout and import action")
}
var nilEngine *Engine
if nilEngine.Configured() {
t.Fatal("nil engine is configured")
}
if unwrapURLError(errors.New("plain")).Error() != "plain" {
t.Fatal("unwrapURLError changed a plain error")
}
}