- pin the sync callbacks before gorm:commit_or_rollback_transaction: an After-only anchor is appended past the commit and lagoon's after-commit flush, so single-statement writes never synced - give the gate and the document builder a clean session: Session with NewDB and a Context clones the write's statement, and a later WithContext queried through the written model's table - typesense.StatusError carries method, path and status, never the body - README: sync semantics, the three gates, delete on soft delete, and the SQL re-gate required of SearchIDs callers
184 lines
5.5 KiB
Go
184 lines
5.5 KiB
Go
// Package beachcomber keeps search indexes in step with GORM models: after
|
|
// a write commits, the written row is reloaded and its document upserted
|
|
// into, or deleted from, the index of a pluggable engine. A failed sync is
|
|
// logged and never touches the write.
|
|
//
|
|
// Models implement Searchable. An engine package (for example
|
|
// beachcomber/typesense) is imported for its side effect of registering
|
|
// itself, and is chosen with search.driver.
|
|
package beachcomber
|
|
|
|
import (
|
|
"database/sql"
|
|
"fmt"
|
|
"log/slog"
|
|
"strings"
|
|
"sync"
|
|
|
|
"git.golem15.com/golem15/summercms/modules/backpack"
|
|
"git.golem15.com/golem15/summercms/modules/lagoon"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
// DefaultDriver is the engine used when search.driver is empty.
|
|
const DefaultDriver = NullEngine
|
|
|
|
// Names of the GORM callbacks that sync Searchable models.
|
|
const (
|
|
CallbackAfterCreate = "beachcomber:after_create"
|
|
CallbackAfterUpdate = "beachcomber:after_update"
|
|
CallbackAfterDelete = "beachcomber:after_delete"
|
|
)
|
|
|
|
// Service is the app-scoped search service. Get it with From.
|
|
type Service struct {
|
|
app *backpack.App
|
|
engine Engine
|
|
prefix string
|
|
log *slog.Logger
|
|
|
|
mu sync.RWMutex
|
|
gate Gate
|
|
}
|
|
|
|
// From returns the app's Service, building and publishing it on first use.
|
|
// The first call reads search.driver (default "null") and search.prefix
|
|
// (default ""), builds the engine through its registered factory, and
|
|
// installs the sync GORM callbacks through lagoon.OnDatabase. An unknown
|
|
// driver name is an error.
|
|
func From(app *backpack.App) (*Service, error) {
|
|
if app == nil {
|
|
return nil, fmt.Errorf("beachcomber: app is nil")
|
|
}
|
|
if s, ok := app.Lookup[*Service](); ok && s != nil {
|
|
return s, nil
|
|
}
|
|
svc := &Service{app: app, log: loggerFromApp(app)}
|
|
name := DefaultDriver
|
|
if app.Config != nil {
|
|
if v := strings.TrimSpace(app.Config.String("search.driver")); v != "" {
|
|
name = strings.ToLower(v)
|
|
}
|
|
svc.prefix = strings.TrimSpace(app.Config.String("search.prefix"))
|
|
}
|
|
factory, ok := engineFactory(name)
|
|
if !ok {
|
|
return nil, fmt.Errorf("beachcomber: unknown search.driver %q (registered: %s; an engine package must be imported to register itself)", name, strings.Join(engineNames(), ", "))
|
|
}
|
|
engine, err := factory(app)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("beachcomber: engine %s: %w", name, err)
|
|
}
|
|
if engine == nil {
|
|
return nil, fmt.Errorf("beachcomber: engine %s returned nil", name)
|
|
}
|
|
svc.engine = engine
|
|
// The callbacks need the GORM handle, which serve publishes after
|
|
// plugins boot.
|
|
if err := lagoon.OnDatabase(app, func(_ *sql.DB, gdb *gorm.DB) error {
|
|
return svc.installCallbacks(gdb)
|
|
}); err != nil {
|
|
return nil, fmt.Errorf("beachcomber: install sync callbacks: %w", err)
|
|
}
|
|
if err := app.Publish(svc); err != nil {
|
|
if existing, ok := app.Lookup[*Service](); ok && existing != nil {
|
|
return existing, nil
|
|
}
|
|
return nil, fmt.Errorf("beachcomber: %w", err)
|
|
}
|
|
return svc, nil
|
|
}
|
|
|
|
// SetGate installs the application kill-switch. With no gate, sync runs
|
|
// whenever the engine is configured and a database is published.
|
|
func (s *Service) SetGate(g Gate) {
|
|
if s == nil {
|
|
return
|
|
}
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
s.gate = g
|
|
}
|
|
|
|
func (s *Service) currentGate() Gate {
|
|
s.mu.RLock()
|
|
defer s.mu.RUnlock()
|
|
return s.gate
|
|
}
|
|
|
|
// Engine returns the engine selected by search.driver.
|
|
func (s *Service) Engine() Engine {
|
|
if s == nil {
|
|
return nil
|
|
}
|
|
return s.engine
|
|
}
|
|
|
|
// Prefix returns search.prefix.
|
|
func (s *Service) Prefix() string {
|
|
if s == nil {
|
|
return ""
|
|
}
|
|
return s.prefix
|
|
}
|
|
|
|
// IndexName is search.prefix followed by m.SearchableAs().
|
|
func (s *Service) IndexName(m Searchable) string {
|
|
if m == nil {
|
|
return s.Prefix()
|
|
}
|
|
return s.Prefix() + m.SearchableAs()
|
|
}
|
|
|
|
// Logger returns the app logger sync failures are logged through.
|
|
func (s *Service) Logger() *slog.Logger {
|
|
if s == nil || s.log == nil {
|
|
return slog.Default()
|
|
}
|
|
return s.log
|
|
}
|
|
|
|
// commitCallback is GORM's commit of a transaction it opened itself.
|
|
const commitCallback = "gorm:commit_or_rollback_transaction"
|
|
|
|
// installCallbacks registers the sync callbacks on gdb, replacing earlier
|
|
// ones so a handle shared by several apps syncs through the most recent
|
|
// service. Each runs after the model's own after hook and before GORM
|
|
// commits a single-statement write: GORM appends a callback that names only
|
|
// an After anchor to the end of the chain, past the commit and past
|
|
// lagoon's after-commit flush, where the registered work would never run.
|
|
func (s *Service) installCallbacks(gdb *gorm.DB) error {
|
|
cb := gdb.Callback()
|
|
if cb.Create().Get(CallbackAfterCreate) == nil {
|
|
if err := cb.Create().After("gorm:after_create").Before(commitCallback).Register(CallbackAfterCreate, s.afterCreate); err != nil {
|
|
return err
|
|
}
|
|
} else if err := cb.Create().Replace(CallbackAfterCreate, s.afterCreate); err != nil {
|
|
return err
|
|
}
|
|
if cb.Update().Get(CallbackAfterUpdate) == nil {
|
|
if err := cb.Update().After("gorm:after_update").Before(commitCallback).Register(CallbackAfterUpdate, s.afterUpdate); err != nil {
|
|
return err
|
|
}
|
|
} else if err := cb.Update().Replace(CallbackAfterUpdate, s.afterUpdate); err != nil {
|
|
return err
|
|
}
|
|
if cb.Delete().Get(CallbackAfterDelete) == nil {
|
|
if err := cb.Delete().After("gorm:after_delete").Before(commitCallback).Register(CallbackAfterDelete, s.afterDelete); err != nil {
|
|
return err
|
|
}
|
|
} else if err := cb.Delete().Replace(CallbackAfterDelete, s.afterDelete); err != nil {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func loggerFromApp(app *backpack.App) *slog.Logger {
|
|
if app != nil {
|
|
if log, ok := app.Lookup[*slog.Logger](); ok && log != nil {
|
|
return log
|
|
}
|
|
}
|
|
return slog.Default()
|
|
}
|