18 KiB
Phase 11: Jobs, realtime and search infrastructure - Pattern Map
Mapped: 2026-09-29 Files analyzed: 22 (new or modified file groups) Analogs found: 19 / 22
Package names follow RESEARCH.md's suggested layout: conga (jobs), lighthouse (realtime) + lighthouse/centrifugo, flare (web push), beachcomber (search) + beachcomber/typesense. All are discretionary (CONTEXT Claude's Discretion). Paths below are relative to summercms.go/ unless prefixed fonoteka.go/. All analogs are git-tracked source.
Plan-time user decisions applied: River uses riverdatabasesql.NewWithPgxListener (one client); summer_jobs gets an internal nullable river_job_id BIGINT; Web Push ports the seams (Pusher interface, stdlib VAPID driver, generate-vapid-keys, health, test-push via an app-provided SubscriptionSource).
File Classification
| New/Modified File | Role | Data Flow | Closest Analog | Match Quality |
|---|---|---|---|---|
modules/conga/client.go (River client build, roles) |
service | event-driven | modules/lagoon/connection.go (shared *sql.DB) |
partial |
modules/conga/manager.go (Dispatch/Start/Update/Complete/Fail/Cancel/CheckIfCanceled/GetMetadata) |
service | CRUD | modules/lagoon/connection.go + fonoteka classes/*_write_service.go |
role-match |
modules/conga/job.go (conga.Job[T] adapter over pact.Job) |
adapter | event-driven | modules/pact/capabilities.go:84-99 |
contract source |
modules/conga/migrations.go (River migrate + summer_jobs) |
migration | batch | modules/lagoon/backend_admin_migrations.go |
exact |
modules/conga/commands.go (queue:work, queue:clear, schedule:run) |
command | batch | modules/lagoon/commands.go + modules/surf/serve.go |
exact |
modules/conga/schedule.go (Daily/Every PeriodicSchedule) |
utility | event-driven | none | no analog |
modules/conga/activate.go (Activate before Boot) |
provider | request-response | modules/postcard/mailer.go:85-134 |
exact |
modules/lighthouse/{publisher,registry,channel,route}.go |
service/interface | pub-sub | modules/postcard/mailer.go + drivers.go |
role-match |
modules/lighthouse/{memory,log,null}.go drivers |
driver | pub-sub | modules/postcard/drivers.go:42-110 |
exact |
modules/lighthouse/broadcast.go (Broadcastable, callbacks, WithoutBroadcasting[T], Emit, BroadcastWorker) |
middleware/hook | event-driven | fonoteka.go/plugins/golem15/fonoteka/classes/artist_resolver.go:45-55 |
role-match |
modules/lighthouse/mount.go (Mount(driver, Surfaces)) |
route | request-response | fonoteka.go/plugins/golem15/fonoteka/routes.go (Group/GroupRaw) |
role-match |
modules/lighthouse/centrifugo/client.go (HTTP API) |
service | request-response | none in-repo (hand-rolled net/http) | no analog |
modules/lighthouse/centrifugo/token.go (5 generators) |
service | transform | modules/wristband/token.go (JWT) |
role-match |
modules/lighthouse/centrifugo/proxy.go + token handler |
controller | request-response | modules/wristband/server.go:180-199 (writeExactJSON) |
exact |
modules/lighthouse/centrifugo/commands.go (health) |
command | request-response | modules/lagoon/commands.go |
exact |
modules/flare/* (Pusher, VAPID driver, generate-vapid-keys, test-push) |
service/command | request-response | modules/postcard/mailer.go (driver select) |
role-match (VAPID crypto: no analog) |
modules/beachcomber/{searchable,engine,sync,gate}.go |
service | event-driven | modules/postcard/mailer.go + artist_resolver callbacks |
role-match |
modules/beachcomber/typesense/engine.go |
service | request-response | none | no analog (RESEARCH Pattern 10) |
modules/lagoon/{ondatabase,aftercommit,transaction}.go |
utility | event-driven | modules/lagoon/connection.go:108-121 (Publish) |
exact |
modules/pact/capabilities.go (+HasSchedule, ScheduledCommand, Cadence) |
config/contract | — | same file 84-99, 368-372 | exact |
modules/bonfire/call.go (Call) |
utility | request-response | modules/bonfire/root.go (NewRootIO) |
exact |
modules/party/registry.go (Activate slot) + internal/build/build.go + cmd/summer/main.go |
wiring | — | party/registry.go:121-126, build.go:113-129, cmd/summer/main.go:41-44 |
exact |
modules/surf/serve.go (start worker in-process) |
command | event-driven | same file 23-83 | exact |
fonoteka.go/.../fonoteka/classes/ws/{collection,wishlist}_authorizer.go |
service | request-response | fonoteka classes/*_write_service.go DB reads via app.Lookup[*gorm.DB] |
role-match |
fonoteka.go/.../fonoteka/models/album.go (+Broadcastable, Searchable) |
model | transform | same file | exact |
fonoteka.go/.../fonoteka/plugin.go (ws-api bucket, Schedule, registry.Register) |
config | — | same file 251-298, 66-160 | exact |
fonoteka.go/.../fonoteka/routes.go (lighthouse.Mount) |
route | request-response | same file (GroupRaw block) |
exact |
fonoteka.go/parity/* realtime tests + manifest entries, tide Centrifugo capture |
test | request-response | fonoteka.go/parity/oauth_flow_test.go, manifest.yaml |
role-match |
modules/*/..._test.go (last plan) |
test | — | modules/lagoon/postgres_test.go |
exact |
Pattern Assignments
modules/conga/activate.go, modules/lighthouse driver select, modules/beachcomber driver select, modules/flare (provider, driver selection)
Analog: modules/postcard/mailer.go
Activate before Boot, publishing the service on the app (lines 85-104):
func Activate[P interface{ ID() string }](app *backpack.App, _ []P) error {
if app == nil {
return fmt.Errorf("postcard: app is nil")
}
driver, opts, err := driverFromApp(app)
if err != nil {
return err
}
m := &mailer{catalog: NewCatalog(), driver: driver, opts: opts}
if err := app.Publish[Mailer](m); err != nil {
return fmt.Errorf("postcard: %w", err)
}
return nil
}
Driver selection by config key, unknown value is a boot error (lines 106-134). Copy for realtime.driver (centrifugo|memory|log|null), search.driver (typesense|null), push.*:
name := "memory"
if app != nil && app.Config != nil {
c := app.Config
if v := strings.TrimSpace(c.String("mail.driver")); v != "" {
name = strings.ToLower(v)
}
...
}
switch name {
case "memory":
return NewMemoryDriver(), opts, nil
case "log":
return NewLogDriver(loggerFromApp(app)), opts, nil
...
default:
return nil, opts, fmt.Errorf("postcard: unknown mail.driver %q (use memory, log, or smtp)", name)
}
Durations from config accept "5s" or an int of seconds (lines 145-151) — reuse for broadcast_timeout, token_ttl. Logger lookup (155-162): app.Lookup[*slog.Logger]() falling back to slog.Default(). Plugin capability type-assert at BootPlugin (164-180): any(p).(pact.HasMailTemplates) — same shape for pact.HasSchedule / pact.HasJobs.
Wire the Activate call in modules/party/registry.go:121-126 next to phrasebook.Activate / postcard.Activate.
modules/lighthouse/{memory,log,null}.go (drivers)
Analog: modules/postcard/drivers.go:42-110. Memory driver is a mutex-guarded slice with a copy-returning accessor:
type MemoryDriver struct {
mu sync.Mutex
messages []RenderedMessage
}
func (d *MemoryDriver) Send(_ context.Context, msg RenderedMessage) error {
if d == nil { return fmt.Errorf("postcard: memory driver is nil") }
d.mu.Lock(); defer d.mu.Unlock()
d.messages = append(d.messages, cloneRendered(msg))
return nil
}
LogDriver at 84-110 (NewLogDriver(log *slog.Logger)); FailDriver at ~170-195 is the model for a failing publisher in the "publish failure never affects write" test (D-07).
modules/conga/client.go, modules/lagoon/ondatabase.go (River on the shared pool)
Analog: modules/lagoon/connection.go
Reserved seam, lines 21-26: "Phase 11 owns a separate pgxpool.Pool for River LISTEN/NOTIFY. Do not create that listener pool here". Build the listener pool from lagoon.DSN(app.Config) (lines 100-106). Publish pattern to extend with OnDatabase callbacks (lines 110-121):
func Publish(app *backpack.App, sqlDB *sql.DB, gdb *gorm.DB) error {
if app == nil { return fmt.Errorf("lagoon: app is nil") }
if sqlDB == nil || gdb == nil { return fmt.Errorf("lagoon: database handles are nil") }
if err := app.Publish(sqlDB); err != nil { return err }
return app.Publish(gdb)
}
OnDatabase(app, fn) runs fn immediately if app.Lookup[*gorm.DB]() succeeds, else queues it and Publish drains the queue after publishing. Error prefix convention: "lagoon: ...", "conga: ...". River client code: RESEARCH Pattern 1 (riverdatabasesql.NewWithPgxListener(sqlDB, listener), listener MaxConns=1, MinConns=0, MaxAttempts: 3, JobTimeout: 5*time.Minute).
modules/conga/migrations.go (summer_jobs + River schema)
Analog: modules/lagoon/backend_admin_migrations.go:1-50
var BackendAdminMigrations = []*gormigrate.Migration{
{
ID: "202609240001_backend_admin_identity",
Migrate: func(tx *gorm.DB) error {
stmts := []string{
`CREATE TABLE IF NOT EXISTS backend_user_roles (
id SERIAL PRIMARY KEY,
...
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW()
)`,
Framework-owned history is isolated under a framework plugin id (comment line 13: "History is isolated under the summercms.cabana plugin id"). Use e.g. summercms.conga. Columns per RESEARCH Pattern 3 (id SERIAL, label, status INT DEFAULT 0, progress, progress_max, user_id INT NULL, is_admin, is_canceled, metadata TEXT NOT NULL, timestamps) plus internal river_job_id BIGINT NULL. River schema via rivermigrate inside Migrate/Rollback on tx's *sql.DB.
modules/conga/commands.go, modules/lighthouse/centrifugo/commands.go, modules/flare/commands.go
Analog: modules/lagoon/commands.go:16-78
func RuntimeCommands(app *backpack.App, plugins []party.Plugin) []bonfire.Command {
return []bonfire.Command{
{
Name: "migrate:rollback",
Description: "Roll back the last migration of a plugin",
Flags: []bonfire.Flag{{
Name: "plugin",
Description: "Plugin ID whose last migration to roll back",
}},
Run: func(ctx context.Context, in bonfire.Input, out bonfire.Output) error {
plugin, _ := in.Flag("plugin")
return withDB(ctx, app, func(gdb *gorm.DB) error { ... out.Success(...) })
},
},
withDB (line 78) opens + publishes the DB for CLI commands. --queue filters use bonfire.Flag{Repeatable: true} read via in.Flags (modules/bonfire/command.go:24-38); --once uses Bare: true. Long-running queue:work / schedule:run copy the signal loop from modules/surf/serve.go:54-83:
ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM)
defer stop()
... select { case <-ctx.Done(): shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) ... }
Register in the generated main: internal/build/build.go:113-116 (append conga.RuntimeCommands(app, plugins)... beside lagoon.RuntimeCommands / cabana.RuntimeCommands), and in cmd/summer/main.go:41-44 add delegateCommand("queue:work", ...), etc.
modules/surf/serve.go (in-process worker)
Analog: itself, lines 37-47: after lagoon.Publish(app, sqlDB, gdb) start the conga worker client unless queue.work_in_serve is false; stop it in the shutdown branch before srv.Shutdown returns. Keep surf free of River imports by calling a conga hook published on the app (app.Lookup[conga.Runner]()) — surf must not import conga if conga imports surf (check cycle).
modules/lighthouse/broadcast.go, modules/beachcomber/sync.go (GORM callbacks)
Analog: fonoteka.go/plugins/golem15/fonoteka/classes/artist_resolver.go:45-55 + classes/registry.go:1-21
func init() {
RegisterHook(func(gdb *gorm.DB) error {
if err := gdb.Callback().Create().Before("gorm:create").Register("fonoteka:album_artist_resolver", albumBeforeSaveCallback); err != nil {
return err
}
...
})
}
Framework version: register through lagoon.OnDatabase (fixes the boot gap documented at fonoteka/plugin.go:143-154), names prefixed lighthouse: / beachcomber:, positions per RESEARCH Patterns 5 and 10 (After("gorm:after_create"), Delete().Before("gorm:before_delete") snapshot, search flush After("gorm:commit_or_rollback_transaction")). Suppression read from db.Statement.Context; tx from db.Statement.ConnPool.(*sql.Tx). No package-level mutable state beyond callback registration.
modules/lighthouse/centrifugo/proxy.go + token handler (controller)
Analog: modules/wristband/server.go:180-199 — raw, non-house JSON writer (no trailing newline, no HTML escaping):
func writeExactJSON(w http.ResponseWriter, status int, v any, extraHeaders map[string]string) {
var buf bytes.Buffer
enc := json.NewEncoder(&buf)
enc.SetEscapeHTML(false)
if err := enc.Encode(v); err != nil { w.WriteHeader(http.StatusInternalServerError); return }
h := w.Header()
h.Set("Content-Type", "application/json")
...
w.WriteHeader(status)
_, _ = w.Write(bytes.TrimSuffix(buf.Bytes(), []byte("\n")))
}
Deny body is HTTP 200 {"error":{"code":403,"message":"Access denied"}}; allow {"result":{"info":[]}} (empty info must encode as []); 503 {"error":"WebSocket not configured"}; 401 {"error":"Unauthorized"}. Secret check with crypto/subtle.ConstantTimeCompare. Framework handlers are app-mounted (Phase 8 D-09 precedent: p.oauthServer.Token mounted in fonoteka routes.go GroupRaw).
modules/lighthouse/centrifugo/token.go
Analog: modules/wristband/token.go (golang-jwt/v5 usage in the framework). Claims verbatim per RESEARCH Pattern 7; info.name may be JSON null (fonoteka.go/plugins/golem15/user/models/user.go:15, Name *string).
modules/pact/capabilities.go (+HasSchedule)
Analog: same file lines 84-99 (JobArgs/Job/HasJobs, River-free) — add ScheduledCommand{Command string; Args []string; Cadence Cadence} and HasSchedule{ Schedule() []ScheduledCommand } in the same style, and update the "Future capability families" comment at 368-379 (remove HasSchedule from the future list, add to the kernel type-assert list). Update modules/pact/README.md in the same commit.
fonoteka.go/.../fonoteka/plugin.go
Analog: itself. Add ws-api to Buckets() (lines 251-298), copying the id-else-IP key form:
"fonoteka-api-token": {
Max: 60, Decay: time.Minute,
Key: func(r *http.Request) string {
if cred, ok := bouncer.Credential(r.Context()); ok { ... return "tok:" + ... }
return surf.ClientIP(r, trusted)
},
},
(ws-api: Max 120, key by bouncer principal ID, else surf.ClientIP). Authorizer registration in Boot follows the bouncer.Registry lookup-or-publish block (lines ~124-134): reg, ok := app.Lookup[*bouncer.Registry](); if !ok { reg = NewRegistry(); app.Publish(reg) }; reg.Register(...). Lazy DB access via app.Lookup[*gorm.DB]() per call (see lazyInvTokenGuard, line 173), since Boot runs before serve publishes the DB.
fonoteka.go/.../fonoteka/routes.go
Analog: itself — r.GroupRaw("/", surf.Use(), ...) for server-to-server and r.Group(..., surf.Use("jwt.auth", ...)) for UserAuth. lighthouse.Mount(driver, lighthouse.Surfaces{...}) is called once in Routes(r pact.Router), with throttle:ws-api as group middleware.
Tests (last plan)
Analog: modules/lagoon/postgres_test.go:1-60 — testcontainers Postgres in TestMain, skipped under -short, ICU pl-PL initdb args:
ctr, err := postgres.Run(ctx, "postgres:16-alpine",
postgres.WithDatabase("lagoon"), ...,
postgres.BasicWaitStrategies(),
testcontainers.WithEnv(map[string]string{
"POSTGRES_INITDB_ARGS": "--locale-provider=icu --icu-locale=pl-PL --encoding=UTF8",
}),
Timed LISTEN pickup test (SC-1): two River clients on one container, set FetchPollInterval high so poll cannot mask a broken notify (RESEARCH Pitfall 5). Fake Centrifugo/Typesense via httptest.NewServer. Parity: extend fonoteka.go/parity/manifest.yaml with /api/realtime/token and /api/realtime/subscribe entries in the existing format.
Shared Patterns
- Service publication:
app.Publish[T](v)/app.Lookup[T]()on*backpack.App; never package globals (postcard/mailer.go:100,fonoteka/plugin.go:124-134). - Error prefixes:
fmt.Errorf("<pkg>: ...: %w", err). - Config reads:
app.Config.String/Int/Bool("dotted.key"), trimmed, with defaults set before reading (postcard/mailer.go:106-153). - Logging:
*slog.Loggerfrom the app (postcard/mailer.go:155-162); broadcast/search failures areWarn, never returned to the write. - Raw JSON responses:
writeExactJSONtechnique (wristband/server.go:180-199) for all PHP-parity bodies. - Boot order: anything needing
*gorm.DBat Boot goes throughlagoon.OnDatabaseor lazyapp.Lookupper call. - Docs: each new module gets a README per CLAUDE.md structure plus a root README modules row; modified modules (lagoon, pact, bonfire, party, surf) update their README in the same commit.
No Analog Found
| File | Role | Data Flow | Reason |
|---|---|---|---|
modules/conga/schedule.go |
utility | event-driven | No wall-clock scheduling code exists; use RESEARCH Pattern 6 |
modules/lighthouse/centrifugo/client.go |
service | request-response | No outbound HTTP API client in the framework; bytes per RESEARCH Pattern 8 |
modules/beachcomber/typesense/engine.go |
service | request-response | Same; wire contract per RESEARCH Pattern 10 |
modules/flare VAPID/RFC 8291 encryption |
service | transform | No crypto of this kind in repo; stdlib crypto/ecdh, crypto/hkdf, AES-GCM |
Metadata
Analog search scope: summercms.go/modules/{postcard,lagoon,pact,party,surf,bonfire,wristband}, summercms.go/{cmd/summer,internal/build}, fonoteka.go/plugins/golem15/fonoteka, fonoteka.go/parity
Files scanned: ~20
Pattern extraction date: 2026-09-29