fix(11-07): roll back the savepoint when a swallowed read failed

- beachcomber and lighthouse released their savepoint whenever the inner
  function reported no error; a Gate that counts a failed read as off,
  or a channel function or delete snapshot that swallows one, left the
  caller's Postgres transaction aborted (25P02) and failed the write
- a failed RELEASE now rolls back to the savepoint, as the READMEs promise
- beachcomber gets its testcontainers harness and sync tests
  (TestSyncGates, TestSyncAfterCommit, TestSyncDeleteAndSoftDelete,
  TestSyncFailuresNonFatal, TestServiceSetup); lighthouse gets
  TestBroadcastSwallowedReadFailure
This commit is contained in:
Jakub Zych
2026-09-30 14:26:51 +02:00
parent 6dadbf6957
commit 6df43d45b8
7 changed files with 890 additions and 4 deletions

View File

@@ -31,7 +31,7 @@ The package itself knows no search server. An engine package registers itself fr
## Sync semantics
- **After commit.** The callbacks register the sync with `lagoon.AfterCommit`. Inside `lagoon.Transaction` it runs after that transaction commits, and not at all when it rolls back. A single-statement write, for which GORM opens its own transaction, syncs after that commit and not when the write fails. Inside a plain `gorm` transaction there is no commit hook, so the sync runs immediately through the transaction's handle. Its reads run in a savepoint, so a failed read never aborts the caller's transaction.
- **After commit.** The callbacks register the sync with `lagoon.AfterCommit`. Inside `lagoon.Transaction` it runs after that transaction commits, and not at all when it rolls back. A single-statement write, for which GORM opens its own transaction, syncs after that commit and not when the write fails. Inside a plain `gorm` transaction there is no commit hook, so the sync runs immediately through the transaction's handle. Its reads run in a savepoint, so a failed read never aborts the caller's transaction, including a read that the application Gate swallows and counts as off.
- **Inline and non-fatal.** The sync runs in the writing goroutine, after the commit, so a create followed by a search sees the document. Every engine request is bounded by the engine's timeout (`search.typesense.connection_timeout_seconds`), and the caller's context cancellation does not abandon it. A failure, a timeout or a panic is logged at Warn as `search: sync failed` with the index, key and operation. The write is already committed and stays so. The log never carries the document or the API key.
- **Three gates, before any request.** Nothing is sent when:
1. the engine is not configured (the `null` engine, or Typesense with an empty `search.typesense.api_key`);

View File

@@ -0,0 +1,148 @@
package beachcomber
import (
"context"
"database/sql"
"fmt"
"net/url"
"os"
"strings"
"sync/atomic"
"testing"
"time"
"git.golem15.com/golem15/summercms/modules/lagoon"
_ "github.com/jackc/pgx/v5/stdlib"
"github.com/testcontainers/testcontainers-go"
"github.com/testcontainers/testcontainers-go/modules/postgres"
)
var (
bcPG *postgres.PostgresContainer
bcSQL *sql.DB
bcDSN string
bcPGErr error
dbSeq atomic.Int64
)
func TestMain(m *testing.M) {
code := 1
if !testShort() {
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
bcPGErr = startBeachcomberPostgres(ctx)
cancel()
if bcPGErr != nil {
fmt.Fprintf(os.Stderr, "beachcomber: testcontainers postgres: %v\n", bcPGErr)
stopBeachcomberPostgres()
os.Exit(1)
}
}
code = m.Run()
stopBeachcomberPostgres()
os.Exit(code)
}
func testShort() bool {
for _, a := range os.Args {
if a == "-test.short" {
return true
}
}
return false
}
func startBeachcomberPostgres(ctx context.Context) error {
ctr, err := postgres.Run(ctx,
"postgres:16-alpine",
postgres.WithDatabase("beachcomber"),
postgres.WithUsername("beachcomber"),
postgres.WithPassword("beachcomber"),
postgres.BasicWaitStrategies(),
testcontainers.WithEnv(map[string]string{
"POSTGRES_INITDB_ARGS": "--locale-provider=icu --icu-locale=pl-PL --encoding=UTF8",
}),
)
if err != nil {
return err
}
bcPG = ctr
dsn, err := ctr.ConnectionString(ctx, "sslmode=disable")
if err != nil {
return err
}
db, err := sql.Open("pgx", dsn)
if err != nil {
return err
}
if err := db.PingContext(ctx); err != nil {
_ = db.Close()
return err
}
bcSQL = db
bcDSN = dsn
return nil
}
func stopBeachcomberPostgres() {
if bcSQL != nil {
_ = bcSQL.Close()
}
if bcPG != nil {
_ = testcontainers.TerminateContainer(bcPG)
}
}
func adminDB(t *testing.T) *sql.DB {
t.Helper()
if testing.Short() {
t.Skip("requires testcontainers postgres")
}
if bcPGErr != nil {
t.Fatalf("postgres unavailable: %v", bcPGErr)
}
if bcSQL == nil {
t.Fatal("postgres unavailable: container was not started")
}
return bcSQL
}
// migratedDB returns a dedicated ICU pl-PL database migrated with
// lagoon.Migrate (River v7 and summer_jobs) plus its DSN.
func migratedDB(t *testing.T) (*sql.DB, string) {
t.Helper()
admin := adminDB(t)
ctx := t.Context()
name := fmt.Sprintf("beachcomber_%d", dbSeq.Add(1))
if _, err := admin.ExecContext(ctx, `CREATE DATABASE `+name+` TEMPLATE template0 ENCODING 'UTF8' LOCALE_PROVIDER icu ICU_LOCALE 'pl-PL'`); err != nil && !strings.Contains(err.Error(), "already exists") {
t.Fatalf("create %s: %v", name, err)
}
dsn, err := dsnWithDB(bcDSN, name)
if err != nil {
t.Fatal(err)
}
db, err := sql.Open("pgx", dsn)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() {
_ = db.Close()
_, _ = admin.ExecContext(context.Background(), `DROP DATABASE IF EXISTS `+name+` WITH (FORCE)`)
})
gdb, err := lagoon.Use(ctx, db)
if err != nil {
t.Fatal(err)
}
if err := lagoon.Migrate(gdb, nil); err != nil {
t.Fatal(err)
}
return db, dsn
}
func dsnWithDB(dsn, name string) (string, error) {
u, err := url.Parse(dsn)
if err != nil {
return "", err
}
u.Path = "/" + name
return u.String(), nil
}

View File

@@ -282,7 +282,13 @@ func inSavepoint(db *gorm.DB, fn func(tx *gorm.DB) error) error {
db.RollbackTo(syncSavepoint)
return err
}
if db.Exec("RELEASE SAVEPOINT "+syncSavepoint).Error != nil {
// A statement inside fn failed although fn did not report it (a
// Gate counts a failed read as off): the transaction is aborted
// and only a rollback to the savepoint makes it usable again.
db.RollbackTo(syncSavepoint)
db.Exec("RELEASE SAVEPOINT " + syncSavepoint)
}
return nil
}

View File

@@ -0,0 +1,696 @@
package beachcomber
import (
"context"
"database/sql"
"errors"
"fmt"
"log/slog"
"reflect"
"strings"
"sync"
"testing"
"git.golem15.com/golem15/summercms/modules/backpack"
"git.golem15.com/golem15/summercms/modules/compass"
"git.golem15.com/golem15/summercms/modules/lagoon"
"gorm.io/gorm"
)
// Doc is a soft-deletable Searchable model whose document reads a related
// row through the handle it is given.
type Doc struct {
ID uint `gorm:"column:id;primaryKey"`
Title string `gorm:"column:title"`
Hidden bool `gorm:"column:hidden"`
OwnerID uint `gorm:"column:owner_id"`
DeletedAt gorm.DeletedAt `gorm:"column:deleted_at"`
}
func (Doc) TableName() string { return "acme_docs" }
func (Doc) SearchableAs() string { return "acme_docs" }
func (d *Doc) ShouldBeSearchable() bool { return !d.Hidden }
func (d *Doc) ToSearchableArray(ctx context.Context, db *gorm.DB) (map[string]any, error) {
switch d.Title {
case "boom":
return nil, errors.New("document failed")
case "panic":
panic("document panicked")
case "empty":
return map[string]any{}, nil
}
var owners []string
if err := db.WithContext(ctx).Table("acme_owners").Where("id = ?", d.OwnerID).Pluck("name", &owners).Error; err != nil {
return nil, err
}
owner := ""
if len(owners) == 1 {
owner = owners[0]
}
return map[string]any{"title": d.Title, "owner": owner, "collection_id": d.OwnerID}, nil
}
func (Doc) SearchIndexSchema() map[string]any {
return map[string]any{"fields": []map[string]any{{"name": "title", "type": "string"}}}
}
// Note is a hard-deleted Searchable model with its own document key.
type Note struct {
ID uint `gorm:"column:id;primaryKey"`
Body string `gorm:"column:body"`
}
func (Note) TableName() string { return "acme_notes" }
func (Note) SearchableAs() string { return "acme_notes" }
func (n *Note) ShouldBeSearchable() bool { return true }
func (n Note) SearchKey() string { return fmt.Sprintf("note-%d", n.ID) }
func (n *Note) ToSearchableArray(context.Context, *gorm.DB) (map[string]any, error) {
return map[string]any{"body": n.Body}, nil
}
// Plain is not Searchable.
type Plain struct {
ID uint `gorm:"column:id;primaryKey"`
Name string `gorm:"column:name"`
}
func (Plain) TableName() string { return "acme_plain" }
const syncTables = `
CREATE TABLE acme_owners (id SERIAL PRIMARY KEY, name TEXT NOT NULL);
CREATE TABLE acme_docs (id SERIAL PRIMARY KEY, title TEXT NOT NULL UNIQUE, hidden BOOLEAN NOT NULL DEFAULT FALSE, owner_id INT NOT NULL DEFAULT 0, deleted_at TIMESTAMPTZ NULL);
CREATE TABLE acme_notes (id SERIAL PRIMARY KEY, body TEXT NOT NULL);
CREATE TABLE acme_plain (id SERIAL PRIMARY KEY, name TEXT NOT NULL);
CREATE TABLE acme_settings (id SERIAL PRIMARY KEY, search_enabled BOOLEAN NOT NULL);
INSERT INTO acme_owners (name) VALUES ('Owner One');
INSERT INTO acme_settings (search_enabled) VALUES (true);
`
// engineCall is one recorded engine call.
type engineCall struct {
op string
index string
schema map[string]any
docs []map[string]any
ids []string
ctxErr error
}
// recEngine records calls and fails on demand. It is configured when
// search.acme_key is set.
type recEngine struct {
configured bool
mu sync.Mutex
calls []engineCall
fail string // "upsert", "delete", "panic"
}
func (e *recEngine) Name() string { return "acme-recording" }
func (e *recEngine) Configured() bool { return e.configured }
func (e *recEngine) record(c engineCall) error {
e.mu.Lock()
defer e.mu.Unlock()
e.calls = append(e.calls, c)
switch e.fail {
case "panic":
panic("engine panicked")
case c.op:
return errors.New(c.op + " failed with status 500")
}
return nil
}
func (e *recEngine) Upsert(ctx context.Context, index string, schema map[string]any, docs []map[string]any) error {
return e.record(engineCall{op: "upsert", index: index, schema: schema, docs: docs, ctxErr: ctx.Err()})
}
func (e *recEngine) Delete(ctx context.Context, index string, ids []string) error {
return e.record(engineCall{op: "delete", index: index, ids: ids, ctxErr: ctx.Err()})
}
func (e *recEngine) Flush(context.Context, string) error { return nil }
func (e *recEngine) SearchIDs(context.Context, string, Query) ([]string, error) { return nil, nil }
func (e *recEngine) take() []engineCall {
e.mu.Lock()
defer e.mu.Unlock()
out := e.calls
e.calls = nil
return out
}
func (e *recEngine) setFail(mode string) {
e.mu.Lock()
e.fail = mode
e.mu.Unlock()
}
func init() {
RegisterEngine("acme-recording", func(app *backpack.App) (Engine, error) {
return &recEngine{configured: app != nil && app.Config != nil && app.Config.String("search.acme_key") != ""}, nil
})
RegisterEngine("acme-broken", func(*backpack.App) (Engine, error) { return nil, errors.New("engine config invalid") })
RegisterEngine("acme-nil", func(*backpack.App) (Engine, error) { return nil, nil })
}
// logRecorder records log messages with their attributes.
type logRecorder struct {
mu sync.Mutex
records []map[string]string
}
func (h *logRecorder) Enabled(context.Context, slog.Level) bool { return true }
func (h *logRecorder) WithAttrs([]slog.Attr) slog.Handler { return h }
func (h *logRecorder) WithGroup(string) slog.Handler { return h }
func (h *logRecorder) Handle(_ context.Context, r slog.Record) error {
rec := map[string]string{"msg": r.Message, "level": r.Level.String()}
r.Attrs(func(a slog.Attr) bool {
rec[a.Key] = a.Value.String()
return true
})
h.mu.Lock()
h.records = append(h.records, rec)
h.mu.Unlock()
return nil
}
func (h *logRecorder) take(msg string) []map[string]string {
h.mu.Lock()
defer h.mu.Unlock()
var out, rest []map[string]string
for _, r := range h.records {
if r["msg"] == msg {
out = append(out, r)
} else {
rest = append(rest, r)
}
}
h.records = rest
return out
}
type syncEnv struct {
app *backpack.App
svc *Service
eng *recEngine
db *sql.DB
gdb *gorm.DB
logs *logRecorder
}
// newSyncEnv builds an app in the production boot order: the service
// first, then the database published (which installs the callbacks).
func newSyncEnv(t *testing.T, kv map[string]any, gate bool) syncEnv {
t.Helper()
db, _ := migratedDB(t)
cfg, err := compass.Open(compass.Options{Dir: t.TempDir(), Env: "testing", Environ: []string{}})
if err != nil {
t.Fatal(err)
}
settings := map[string]any{"search.driver": "acme-recording", "search.acme_key": "test-only-key", "search.prefix": "dev_"}
for k, v := range kv {
settings[k] = v
}
for k, v := range settings {
if err := cfg.Set(k, v); err != nil {
t.Fatal(err)
}
}
app := backpack.New(cfg)
logs := &logRecorder{}
if err := app.Publish(slog.New(logs)); err != nil {
t.Fatal(err)
}
svc, err := From(app)
if err != nil {
t.Fatal(err)
}
if gate {
svc.SetGate(GateFunc(func(ctx context.Context, db *gorm.DB) bool {
var on []bool
if err := db.WithContext(ctx).Table("acme_settings").Order("id").Limit(1).Pluck("search_enabled", &on).Error; err != nil {
return false
}
return len(on) == 1 && on[0]
}))
}
gdb, err := lagoon.Use(t.Context(), db)
if err != nil {
t.Fatal(err)
}
if err := gdb.Exec(syncTables).Error; err != nil {
t.Fatal(err)
}
if err := lagoon.Publish(app, db, gdb); err != nil {
t.Fatal(err)
}
eng, ok := svc.Engine().(*recEngine)
if !ok {
eng = nil
}
return syncEnv{app: app, svc: svc, eng: eng, db: db, gdb: gdb, logs: logs}
}
func (e syncEnv) setGate(t *testing.T, on bool) {
t.Helper()
if err := e.gdb.Exec(`UPDATE acme_settings SET search_enabled = ?`, on).Error; err != nil {
t.Fatal(err)
}
}
// TestSyncGates covers D-20 and T-11-25: with the null engine, an engine
// without a key, no published database or the application gate off,
// nothing is sent.
func TestSyncGates(t *testing.T) {
ctx := context.Background()
t.Run("null_engine_by_default", func(t *testing.T) {
env := newSyncEnv(t, map[string]any{"search.driver": ""}, false)
if env.svc.Engine().Name() != NullEngine || env.svc.Engine().Configured() {
t.Fatalf("engine = %s", env.svc.Engine().Name())
}
if err := env.gdb.WithContext(ctx).Create(&Doc{Title: "null"}).Error; err != nil {
t.Fatal(err)
}
if err := env.svc.Sync(ctx, nil, &Doc{ID: 1}); err != nil {
t.Fatal(err)
}
ne := env.svc.Engine()
if ne.Upsert(ctx, "i", nil, nil) != nil || ne.Delete(ctx, "i", nil) != nil || ne.Flush(ctx, "i") != nil {
t.Fatal("null engine returned an error")
}
if ids, err := ne.SearchIDs(ctx, "i", Query{}); err != nil || ids == nil || len(ids) != 0 {
t.Fatalf("null SearchIDs = %v, %v", ids, err)
}
})
t.Run("empty_key_sends_nothing", func(t *testing.T) {
env := newSyncEnv(t, map[string]any{"search.acme_key": ""}, false)
if err := env.gdb.WithContext(ctx).Create(&Doc{Title: "no-key"}).Error; err != nil {
t.Fatal(err)
}
if err := env.svc.Sync(ctx, env.gdb, &Doc{ID: 1}); err != nil {
t.Fatal(err)
}
if n := len(env.eng.take()); n != 0 {
t.Fatalf("%d engine calls without a key", n)
}
})
t.Run("gate_off_then_on", func(t *testing.T) {
env := newSyncEnv(t, nil, true)
env.setGate(t, false)
if err := env.gdb.WithContext(ctx).Create(&Doc{Title: "gated"}).Error; err != nil {
t.Fatal(err)
}
if n := len(env.eng.take()); n != 0 {
t.Fatalf("%d engine calls with the gate off", n)
}
env.setGate(t, true)
if err := env.gdb.WithContext(ctx).Create(&Doc{Title: "open"}).Error; err != nil {
t.Fatal(err)
}
if calls := env.eng.take(); len(calls) != 1 || calls[0].op != "upsert" {
t.Fatalf("calls with the gate on = %+v", calls)
}
})
t.Run("no_published_database", func(t *testing.T) {
env := newSyncEnv(t, nil, false)
bare := backpack.New(env.app.Config)
svc := &Service{app: bare, engine: env.eng}
if err := svc.Sync(ctx, nil, &Doc{ID: 1}); err != nil {
t.Fatal(err)
}
if err := svc.syncOne(ctx, env.gdb, pending{op: opUpsert}); err != nil {
t.Fatal(err)
}
if n := len(env.eng.take()); n != 0 {
t.Fatalf("%d engine calls without a published database", n)
}
})
}
// TestSyncAfterCommit covers D-20: a write syncs only after it commits
// (lagoon.Transaction, a single statement's implicit transaction, a plain
// gorm transaction through its handle), a rollback syncs nothing, and the
// document is built from a reload on a clean statement.
func TestSyncAfterCommit(t *testing.T) {
env := newSyncEnv(t, nil, true)
ctx := context.Background()
t.Run("lagoon_transaction_commit", func(t *testing.T) {
var d Doc
err := lagoon.Transaction(ctx, env.gdb, func(ctx context.Context, tx *gorm.DB) error {
d = Doc{Title: "committed", OwnerID: 1}
if err := tx.WithContext(ctx).Create(&d).Error; err != nil {
return err
}
if n := len(env.eng.take()); n != 0 {
t.Errorf("%d engine calls before commit", n)
}
return nil
})
if err != nil {
t.Fatal(err)
}
calls := env.eng.take()
if len(calls) != 1 {
t.Fatalf("calls = %+v", calls)
}
c := calls[0]
want := map[string]any{"id": fmt.Sprint(d.ID), "title": "committed", "owner": "Owner One", "collection_id": uint(1)}
if c.op != "upsert" || c.index != "dev_acme_docs" || !reflect.DeepEqual(c.docs[0], want) || c.schema["fields"] == nil || c.ctxErr != nil {
t.Fatalf("call = %+v, want doc %v", c, want)
}
})
t.Run("rollback_sends_nothing", func(t *testing.T) {
err := lagoon.Transaction(ctx, env.gdb, func(ctx context.Context, tx *gorm.DB) error {
if err := tx.WithContext(ctx).Create(&Doc{Title: "rolled-back"}).Error; err != nil {
return err
}
return errors.New("rollback")
})
if err == nil {
t.Fatal("transaction did not fail")
}
if n := len(env.eng.take()); n != 0 {
t.Fatalf("%d engine calls after a rollback", n)
}
})
t.Run("single_statement_implicit_commit", func(t *testing.T) {
d := Doc{Title: "implicit", OwnerID: 1}
if err := env.gdb.WithContext(ctx).Create(&d).Error; err != nil {
t.Fatal(err)
}
calls := env.eng.take()
if len(calls) != 1 || calls[0].docs[0]["owner"] != "Owner One" {
t.Fatalf("calls = %+v", calls)
}
d.Title = "implicit-updated"
if err := env.gdb.WithContext(ctx).Save(&d).Error; err != nil {
t.Fatal(err)
}
if calls := env.eng.take(); len(calls) != 1 || calls[0].docs[0]["title"] != "implicit-updated" {
t.Fatalf("update calls = %+v", calls)
}
})
t.Run("plain_gorm_transaction_syncs_through_the_tx", func(t *testing.T) {
err := env.gdb.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := tx.Create(&Doc{Title: "plain-tx", OwnerID: 1}).Error; err != nil {
return err
}
calls := env.eng.take()
if len(calls) != 1 || calls[0].docs[0]["title"] != "plain-tx" {
t.Errorf("inside a plain transaction: calls = %+v, want an immediate sync from the uncommitted row", calls)
}
// The transaction is still usable.
return tx.Create(&Plain{Name: "after"}).Error
})
if err != nil {
t.Fatal(err)
}
})
t.Run("zero_key_batch_update_and_plain_model_are_skipped", func(t *testing.T) {
if err := env.gdb.WithContext(ctx).Model(&Doc{}).Where("owner_id = ?", 1).Update("hidden", false).Error; err != nil {
t.Fatal(err)
}
if err := env.gdb.WithContext(ctx).Create(&Plain{Name: "plain"}).Error; err != nil {
t.Fatal(err)
}
if n := len(env.eng.take()); n != 0 {
t.Fatalf("%d engine calls for a batch update or a plain model", n)
}
})
t.Run("batch_create_syncs_each_row", func(t *testing.T) {
docs := []Doc{{Title: "batch-1"}, {Title: "batch-2"}}
if err := env.gdb.WithContext(ctx).Create(&docs).Error; err != nil {
t.Fatal(err)
}
if n := len(env.eng.take()); n != 2 {
t.Fatalf("batch create calls = %d, want 2", n)
}
})
}
// TestSyncDeleteAndSoftDelete covers deletes: a hard delete removes the
// document under its SearchKey, a soft delete or a row that should not be
// searchable removes it, a restore re-indexes it, and Sync/Remove follow
// the same gated path.
func TestSyncDeleteAndSoftDelete(t *testing.T) {
env := newSyncEnv(t, nil, true)
ctx := context.Background()
n := Note{Body: "note"}
if err := env.gdb.WithContext(ctx).Create(&n).Error; err != nil {
t.Fatal(err)
}
if calls := env.eng.take(); len(calls) != 1 || calls[0].docs[0]["id"] != fmt.Sprintf("note-%d", n.ID) || calls[0].schema != nil {
t.Fatalf("note upsert = %+v", calls)
}
if err := env.gdb.WithContext(ctx).Delete(&Note{ID: n.ID}).Error; err != nil {
t.Fatal(err)
}
if calls := env.eng.take(); len(calls) != 1 || calls[0].op != "delete" || calls[0].ids[0] != fmt.Sprintf("note-%d", n.ID) || calls[0].index != "dev_acme_notes" {
t.Fatalf("note delete = %+v", calls)
}
d := Doc{Title: "soft", OwnerID: 1}
if err := env.gdb.WithContext(ctx).Create(&d).Error; err != nil {
t.Fatal(err)
}
env.eng.take()
if err := env.gdb.WithContext(ctx).Delete(&d).Error; err != nil {
t.Fatal(err)
}
if calls := env.eng.take(); len(calls) != 1 || calls[0].op != "delete" || calls[0].ids[0] != fmt.Sprint(d.ID) {
t.Fatalf("soft delete = %+v", calls)
}
if err := env.gdb.WithContext(ctx).Unscoped().Model(&d).Update("deleted_at", nil).Error; err != nil {
t.Fatal(err)
}
if calls := env.eng.take(); len(calls) != 1 || calls[0].op != "upsert" {
t.Fatalf("restore = %+v", calls)
}
if err := env.gdb.WithContext(ctx).Model(&d).Update("hidden", true).Error; err != nil {
t.Fatal(err)
}
if calls := env.eng.take(); len(calls) != 1 || calls[0].op != "delete" {
t.Fatalf("hidden row = %+v, want a delete", calls)
}
t.Run("explicit_sync_and_remove", func(t *testing.T) {
if err := env.gdb.Exec(`UPDATE acme_docs SET hidden = false WHERE id = ?`, d.ID).Error; err != nil {
t.Fatal(err)
}
if err := env.svc.Sync(ctx, nil, &Doc{ID: d.ID}); err != nil {
t.Fatal(err)
}
if err := env.svc.Remove(ctx, env.gdb, &Doc{ID: d.ID}); err != nil {
t.Fatal(err)
}
if err := env.svc.Sync(nil, env.gdb, &Doc{ID: 999999}); err != nil {
t.Fatal(err)
}
calls := env.eng.take()
if len(calls) != 3 || calls[0].op != "upsert" || calls[1].op != "delete" || calls[2].op != "delete" || calls[2].ids[0] != "999999" {
t.Fatalf("explicit calls = %+v", calls)
}
for _, bad := range []any{Doc{ID: 1}, &Doc{}, &Plain{ID: 1}, (*Doc)(nil), 42} {
if err := env.svc.Sync(ctx, env.gdb, bad); err == nil {
t.Errorf("Sync(%T %v) accepted", bad, bad)
}
}
var nilSvc *Service
if err := nilSvc.Sync(ctx, env.gdb, &Doc{ID: 1}); err == nil {
t.Fatal("nil service accepted")
}
})
}
// TestSyncFailuresNonFatal covers T-11-26: engine errors and panics, a
// failing or panicking document builder, and a failed read inside a
// caller's transaction are each one Warn log with index, key and
// operation, and never touch the write.
func TestSyncFailuresNonFatal(t *testing.T) {
env := newSyncEnv(t, nil, true)
ctx := context.Background()
warnings := func() []map[string]string { return env.logs.take("search: sync failed") }
for _, mode := range []string{"upsert", "panic"} {
t.Run("engine_"+mode, func(t *testing.T) {
env.eng.setFail(mode)
defer env.eng.setFail("")
d := Doc{Title: "engine-" + mode, OwnerID: 1}
if err := env.gdb.WithContext(ctx).Create(&d).Error; err != nil {
t.Fatalf("an engine failure failed the write: %v", err)
}
ws := warnings()
if len(ws) != 1 || ws[0]["index"] != "dev_acme_docs" || ws[0]["key"] != fmt.Sprint(d.ID) || ws[0]["operation"] != "upsert" || ws[0]["level"] != "WARN" {
t.Fatalf("warnings = %v", ws)
}
if strings.Contains(ws[0]["error"], "test-only-key") {
t.Fatal("the key reached the log")
}
})
}
t.Run("delete_error", func(t *testing.T) {
d := Doc{Title: "delete-error"}
if err := env.gdb.WithContext(ctx).Create(&d).Error; err != nil {
t.Fatal(err)
}
env.eng.setFail("delete")
defer env.eng.setFail("")
if err := env.gdb.WithContext(ctx).Delete(&d).Error; err != nil {
t.Fatal(err)
}
if ws := warnings(); len(ws) != 1 || ws[0]["operation"] != "delete" {
t.Fatalf("warnings = %v", ws)
}
})
for _, title := range []string{"boom", "panic"} {
t.Run("document_"+title, func(t *testing.T) {
env.eng.take()
if err := env.gdb.WithContext(ctx).Create(&Doc{Title: title}).Error; err != nil {
t.Fatal(err)
}
if ws := warnings(); len(ws) != 1 {
t.Fatalf("warnings = %v", ws)
}
if calls := env.eng.take(); len(calls) != 0 {
t.Fatalf("a failed document was sent: %+v", calls)
}
})
}
t.Run("empty_document_is_skipped", func(t *testing.T) {
if err := env.gdb.WithContext(ctx).Create(&Doc{Title: "empty"}).Error; err != nil {
t.Fatal(err)
}
if ws, calls := warnings(), env.eng.take(); len(ws) != 0 || len(calls) != 0 {
t.Fatalf("empty document: warnings %v calls %v", ws, calls)
}
})
t.Run("failed_gate_read_keeps_the_callers_transaction", func(t *testing.T) {
env.svc.SetGate(GateFunc(func(ctx context.Context, db *gorm.DB) bool {
var n int
return db.WithContext(ctx).Raw(`SELECT count(*) FROM acme_missing_settings`).Scan(&n).Error == nil
}))
defer env.svc.SetGate(nil)
err := env.gdb.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := tx.Create(&Doc{Title: "gate-read-fails"}).Error; err != nil {
return err
}
return tx.Create(&Plain{Name: "still-usable"}).Error
})
if err != nil {
t.Fatalf("a failed gate read aborted the caller's transaction: %v", err)
}
if calls := env.eng.take(); len(calls) != 0 {
t.Fatalf("gate off still sent %+v", calls)
}
})
t.Run("failed_document_read_inside_a_transaction", func(t *testing.T) {
if err := env.gdb.Exec(`ALTER TABLE acme_owners RENAME TO acme_owners_gone`).Error; err != nil {
t.Fatal(err)
}
defer env.gdb.Exec(`ALTER TABLE acme_owners_gone RENAME TO acme_owners`)
err := env.gdb.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := tx.Create(&Doc{Title: "owner-read-fails"}).Error; err != nil {
return err
}
return tx.Create(&Plain{Name: "still-usable-2"}).Error
})
if err != nil {
t.Fatalf("a failed document read aborted the caller's transaction: %v", err)
}
if ws := warnings(); len(ws) != 1 || !strings.Contains(ws[0]["error"], "build document") {
t.Fatalf("warnings = %v", ws)
}
})
}
// TestServiceSetup covers From, the engine registry and the accessors.
func TestServiceSetup(t *testing.T) {
cfg := func(t *testing.T, kv map[string]any) *backpack.App {
t.Helper()
c, err := compass.Open(compass.Options{Dir: t.TempDir(), Env: "testing", Environ: []string{}})
if err != nil {
t.Fatal(err)
}
for k, v := range kv {
if err := c.Set(k, v); err != nil {
t.Fatal(err)
}
}
return backpack.New(c)
}
if _, err := From(nil); err == nil {
t.Fatal("From(nil) succeeded")
}
_, err := From(cfg(t, map[string]any{"search.driver": "nope"}))
if err == nil || !strings.Contains(err.Error(), `unknown search.driver "nope"`) || !strings.Contains(err.Error(), "acme-broken, acme-nil, acme-recording, null") {
t.Fatalf("unknown driver: %v", err)
}
if _, err := From(cfg(t, map[string]any{"search.driver": "acme-broken"})); err == nil || !strings.Contains(err.Error(), "engine config invalid") {
t.Fatalf("broken engine: %v", err)
}
if _, err := From(cfg(t, map[string]any{"search.driver": "acme-nil"})); err == nil || !strings.Contains(err.Error(), "returned nil") {
t.Fatalf("nil engine: %v", err)
}
app := cfg(t, map[string]any{"search.driver": " ACME-Recording ", "search.prefix": " p_ "})
svc, err := From(app)
if err != nil {
t.Fatal(err)
}
if again, err := From(app); err != nil || again != svc {
t.Fatal("From is not idempotent")
}
if svc.Prefix() != "p_" || svc.IndexName(&Doc{}) != "p_acme_docs" || svc.IndexName(nil) != "p_" || svc.Logger() == nil {
t.Fatalf("prefix %q index %q", svc.Prefix(), svc.IndexName(&Doc{}))
}
var nilSvc *Service
nilSvc.SetGate(nil)
if nilSvc.Engine() != nil || nilSvc.Prefix() != "" || nilSvc.Logger() == nil {
t.Fatal("nil service accessors")
}
if GateFunc(nil).Enabled(context.Background(), nil) {
t.Fatal("a nil GateFunc is enabled")
}
for name, fn := range map[string]func(){
"empty": func() { RegisterEngine("", func(*backpack.App) (Engine, error) { return nullEngine{}, nil }) },
"nil": func() { RegisterEngine("acme-x", nil) },
"duplicate": func() { RegisterEngine(NullEngine, func(*backpack.App) (Engine, error) { return nullEngine{}, nil }) },
} {
func() {
defer func() {
if recover() == nil {
t.Errorf("RegisterEngine %s did not panic", name)
}
}()
fn()
}()
}
for pk, want := range map[any]string{uint(1): "1", uint32(2): "2", uint64(3): "3", 4: "4", int32(5): "5", int64(6): "6", "k": "k"} {
if got := searchKey(reflect.ValueOf(Plain{}), pk); got != want {
t.Errorf("searchKey(%T) = %q", pk, got)
}
}
if got := searchKey(reflect.ValueOf(Note{ID: 7}), uint(7)); got != "note-7" {
t.Errorf("searchKey of an unaddressable SearchKeyer = %q", got)
}
}

View File

@@ -39,7 +39,7 @@ The `centrifugo` sub-package is the Centrifugo driver. It has a hand-rolled `net
- `lighthouse.BroadcastTTLer` or `Binding.TTL` replaces the ttl.
The event name is `{action}.{alias}` lowercased: `lighthouse.ActionCreated`, `lighthouse.ActionUpdated` or `lighthouse.ActionDeleted`, then an alias that defaults to `<plugin>.<model>` (the Go package name, or the parent directory of a `models` package, and the type name). The payload builder receives a `lighthouse.Event` with the action, the `lighthouse.Actor`, the timestamp and the ttl. A soft delete counts as a delete. A delete's channels and payload are computed from a fresh read of the row before it is deleted, so deleting a model that holds only its id still broadcasts. An empty channel list means no broadcast.
- Transactional delivery. GORM callbacks (`lighthouse.CallbackAfterCreate`, `lighthouse.CallbackAfterUpdate`, `lighthouse.CallbackSnapshot` and `lighthouse.CallbackAfterDelete`) are installed through `lagoon.OnDatabase`. The after-write callbacks run after the model's own after hook and before GORM commits the transaction it opens for a single-statement write, so they enqueue a `lighthouse.BroadcastArgs` job on the write's `*sql.Tx` in every case (an explicit transaction or a single `Create`, `Save` or `Delete`), on the `realtime.broadcast_queue` queue with MaxAttempts 1 and the `realtime.broadcast_timeout` timeout. Channel and payload queries and the enqueue run inside a savepoint, so a failure is rolled back to it, logged at Warn with channels and event (never the payload), and the write goes on. A write with a zero primary key, such as `Model(&T{}).Where(…).Updates(…)`, is not broadcast; bulk paths suppress and emit instead. The null driver, or a driver whose `Enabled` reports false (Centrifugo without an API key), gets no jobs.
- Transactional delivery. GORM callbacks (`lighthouse.CallbackAfterCreate`, `lighthouse.CallbackAfterUpdate`, `lighthouse.CallbackSnapshot` and `lighthouse.CallbackAfterDelete`) are installed through `lagoon.OnDatabase`. The after-write callbacks run after the model's own after hook and before GORM commits the transaction it opens for a single-statement write, so they enqueue a `lighthouse.BroadcastArgs` job on the write's `*sql.Tx` in every case (an explicit transaction or a single `Create`, `Save` or `Delete`), on the `realtime.broadcast_queue` queue with MaxAttempts 1 and the `realtime.broadcast_timeout` timeout. Channel and payload queries and the enqueue run inside a savepoint, so a failure (also one a channel or payload function swallows) is rolled back to it, logged at Warn with channels and event (never the payload), and the write goes on. A write with a zero primary key, such as `Model(&T{}).Where(…).Updates(…)`, is not broadcast; bulk paths suppress and emit instead. The null driver, or a driver whose `Enabled` reports false (Centrifugo without an API key), gets no jobs.
- The broadcast job lowercases the channels and adds the `realtime.broadcast_namespace` prefix unless a channel already has it. It then publishes to one channel or broadcasts to several. A failure is logged as `realtime: broadcast failed` and is not retried. Delivery order across separate jobs is not guaranteed. The payload travels inside the job as a JSON string, so its key order survives Postgres JSONB.
- Suppression: `lighthouse.WithoutBroadcasting` silences one model type for writes made with the context it hands to its function. Other types still broadcast, and a write through an outer context is not suppressed. `lighthouse.Service.Emit` enqueues one `lighthouse.Broadcast` on the caller's transaction and returns its error. Together they turn N row events into one summary event.
- Centrifugo driver (`centrifugo.Driver`, driver name `centrifugo`):

View File

@@ -469,7 +469,12 @@ func (s *Service) inSavepoint(db *gorm.DB, fn func(tx *gorm.DB) error) {
}
return
}
if inTx {
if inTx && tx.Exec("RELEASE SAVEPOINT "+savepoint).Error != nil {
// A statement inside fn failed although fn did not report it (a
// channel or payload function that treats a failed read as "no
// broadcast", or the delete snapshot's reload): the transaction is
// aborted and only a rollback to the savepoint keeps the write alive.
tx.RollbackTo(savepoint)
tx.Exec("RELEASE SAVEPOINT " + savepoint)
}
}

View File

@@ -379,6 +379,11 @@ func (s *Sprocket) BroadcastChannels(_ context.Context, tx *gorm.DB) ([]string,
if s.Channels == "fail" {
return nil, errors.New("channels failed")
}
if s.Channels == "swallow" {
// A channel function that treats a failed read as "no channels".
_ = tx.Exec(`SELECT * FROM acme_missing_table`).Error
return nil, nil
}
if s.Channels == "abort-tx" {
// A failed statement aborts a Postgres transaction unless it runs
// inside a savepoint.
@@ -750,6 +755,32 @@ func TestBroadcastEdges(t *testing.T) {
})
}
// TestBroadcastSwallowedReadFailure covers a channel function that
// swallows a failed read: the savepoint must still be rolled back, or the
// failed statement leaves the caller's transaction aborted.
func TestBroadcastSwallowedReadFailure(t *testing.T) {
env := newLHEnv(t, nil, nil)
ctx := t.Context()
err := lagoon.Transaction(ctx, env.gdb, func(ctx context.Context, tx *gorm.DB) error {
if err := tx.WithContext(ctx).Create(&Sprocket{Name: "swallowed", Channels: "swallow"}).Error; err != nil {
return err
}
return tx.WithContext(ctx).Create(&Widget{Name: "after-swallow", OwnerID: 0}).Error
})
if err != nil {
t.Fatalf("a swallowed read failure in the broadcast savepoint aborted the write: %v", err)
}
// A single-statement write runs its callbacks inside GORM's own
// transaction; the commit must still succeed.
if err := env.gdb.WithContext(ctx).Create(&Sprocket{Name: "swallowed-implicit", Channels: "swallow"}).Error; err != nil {
t.Fatalf("single-statement write with a swallowed read failure: %v", err)
}
var n int64
if err := env.gdb.Model(&Sprocket{}).Where("name LIKE ?", "swallowed%").Count(&n).Error; err != nil || n != 2 {
t.Fatalf("committed sprockets = %d (err %v), want 2", n, err)
}
}
// TestBroadcastPublishFailure covers D-09: a failed publish is logged at
// Warn without the payload and the one-attempt job is not retried.
func TestBroadcastPublishFailure(t *testing.T) {