From d9b951fe29243cf1855b680e746778336271c6b2 Mon Sep 17 00:00:00 2001 From: Jakub Zych Date: Tue, 29 Sep 2026 13:40:53 +0200 Subject: [PATCH] docs(11): map phase patterns --- .../11-PATTERNS.md | 268 ++++++++++++++++++ 1 file changed, 268 insertions(+) create mode 100644 .planning/phases/11-jobs-realtime-and-search-infrastructure/11-PATTERNS.md diff --git a/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-PATTERNS.md b/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-PATTERNS.md new file mode 100644 index 0000000..674a777 --- /dev/null +++ b/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-PATTERNS.md @@ -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(": ...: %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