docs(11): map phase patterns

This commit is contained in:
Jakub Zych
2026-09-29 13:40:53 +02:00
parent a1ed5b3539
commit d9b951fe29

View File

@@ -0,0 +1,268 @@
# 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):
```go
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.*`:
```go
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:
```go
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):
```go
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`
```go
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`
```go
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`:
```go
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`
```go
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):
```go
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:
```go
"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:
```go
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.Logger` from the app (`postcard/mailer.go:155-162`); broadcast/search failures are `Warn`, never returned to the write.
- **Raw JSON responses:** `writeExactJSON` technique (`wristband/server.go:180-199`) for all PHP-parity bodies.
- **Boot order:** anything needing `*gorm.DB` at Boot goes through `lagoon.OnDatabase` or lazy `app.Lookup` per 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