From 6df43d45b86095f68f2528c2ddc7527f90dede2b Mon Sep 17 00:00:00 2001 From: Jakub Zych Date: Wed, 30 Sep 2026 14:26:51 +0200 Subject: [PATCH] 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 --- modules/beachcomber/README.md | 2 +- modules/beachcomber/postgres_test.go | 148 ++++++ modules/beachcomber/sync.go | 8 +- modules/beachcomber/sync_test.go | 696 +++++++++++++++++++++++++++ modules/lighthouse/README.md | 2 +- modules/lighthouse/broadcast.go | 7 +- modules/lighthouse/broadcast_test.go | 31 ++ 7 files changed, 890 insertions(+), 4 deletions(-) create mode 100644 modules/beachcomber/postgres_test.go create mode 100644 modules/beachcomber/sync_test.go diff --git a/modules/beachcomber/README.md b/modules/beachcomber/README.md index be9d2dd..6b41cf2 100644 --- a/modules/beachcomber/README.md +++ b/modules/beachcomber/README.md @@ -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`); diff --git a/modules/beachcomber/postgres_test.go b/modules/beachcomber/postgres_test.go new file mode 100644 index 0000000..46d8eb1 --- /dev/null +++ b/modules/beachcomber/postgres_test.go @@ -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 +} diff --git a/modules/beachcomber/sync.go b/modules/beachcomber/sync.go index 8f57423..a8d21a4 100644 --- a/modules/beachcomber/sync.go +++ b/modules/beachcomber/sync.go @@ -282,7 +282,13 @@ func inSavepoint(db *gorm.DB, fn func(tx *gorm.DB) error) error { db.RollbackTo(syncSavepoint) return err } - db.Exec("RELEASE SAVEPOINT " + syncSavepoint) + 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 } diff --git a/modules/beachcomber/sync_test.go b/modules/beachcomber/sync_test.go new file mode 100644 index 0000000..47cbb15 --- /dev/null +++ b/modules/beachcomber/sync_test.go @@ -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) + } +} diff --git a/modules/lighthouse/README.md b/modules/lighthouse/README.md index dfa8241..b43e0d5 100644 --- a/modules/lighthouse/README.md +++ b/modules/lighthouse/README.md @@ -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 `.` (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`): diff --git a/modules/lighthouse/broadcast.go b/modules/lighthouse/broadcast.go index 8ba51c1..b2b2cd5 100644 --- a/modules/lighthouse/broadcast.go +++ b/modules/lighthouse/broadcast.go @@ -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) } } diff --git a/modules/lighthouse/broadcast_test.go b/modules/lighthouse/broadcast_test.go index 72e1c33..7c5b2fc 100644 --- a/modules/lighthouse/broadcast_test.go +++ b/modules/lighthouse/broadcast_test.go @@ -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) {