docs(11): research phase domain

This commit is contained in:
Jakub Zych
2026-09-29 12:32:02 +02:00
parent 0bbbce6d27
commit c82950c2a6

View File

@@ -0,0 +1,779 @@
# Phase 11: Jobs, realtime and search infrastructure - Research
**Researched:** 2026-09-29
**Domain:** River job queue on Postgres (dual driver), Centrifugo realtime contract port, Typesense sync, Go scheduler
**Confidence:** HIGH for River API and the PHP contract (both read from source this session); MEDIUM for the Centrifugo/Typesense server behaviour (official docs, no live server probe); LOW only where tagged `[ASSUMED]`
## Summary
The biggest finding is in River itself. River **v0.46.0 (2026-08-29)** added `riverdatabasesql.NewWithPgxListener(dbPool *sql.DB, listenerPool *pgxpool.Pool)`. It is one driver where every query and transaction goes through the shared `*sql.DB`, and only Postgres `LISTEN` goes through a dedicated pgx pool, which the driver wraps with `riverpgxv5.New(listenerPool)` internally. River's official GORM page now recommends exactly this. It is the literal "one shared `*sql.DB` plus a small listener pool" design from STACK.md and the reservation comment in `lagoon/connection.go`. The older "two River clients" pattern in PITFALLS.md Pitfall 14 is no longer required. Recommendation: one `*river.Client[*sql.Tx]` type everywhere. The HTTP process runs it insert-only (`riverdatabasesql.New`). A worker (in `serve` or `queue:work`) runs it with `NewWithPgxListener` and a pgx pool of MaxConns 1. Transactional enqueue via `InsertTx(ctx, *sql.Tx, ...)` emits `pg_notify` inside the transaction (`SupportsListenNotify()` is `true` for the database/sql driver), so a separate worker process wakes on commit. The planner must confirm this still satisfies SC-1's wording, "a separate `riverpgxv5` client" (see Open Question 1).
The PHP contract has several traps that a straight reading of CONTEXT.md would miss:
- `JobManager::dispatch` writes `status = IN_PROGRESS (1)`, never `IN_QUEUE (0)`.
- PHP `cancelJob` only sets `STOPPED`; the CSV controller sets `is_canceled` separately.
- The subscribe-proxy **deny is HTTP 200** with an error body.
- The PHP allow response encodes an empty `info` as `[]`, not `{}`.
- The Centrifugo client authenticates with `Authorization: apikey <key>`. Centrifugo master still accepts that header.
- Album overrides the default broadcast payload entirely: `{id, collection_id, action, actor, timestamp, album?}`, not `{model, actor, timestamp, ttl}`.
- Scout on PHP runs inside the transaction (`after_commit=false`) and deletes the index document on soft delete.
- Web Push is dead code in Płytarium. `minishlink/web-push` is not installed, and `TestPushNotifications` imports a `Golem15.Notifications` plugin that does not exist in the repo.
Two framework seams are missing and should be built in this phase:
1. **A "database ready" hook.** In the real binary, plugin `Boot` runs before `serve` publishes `*gorm.DB` (documented as a known gap in `fonoteka.go/plugins/golem15/fonoteka/plugin.go`). Any GORM callback that realtime or search registers would silently never fire under `summer serve`.
2. **An after-commit hook.** GORM has none. The inline-after-commit search sync (D-20) needs one: a ctx-scoped buffer flushed by a `lagoon` transaction helper, plus a GORM callback after `gorm:commit_or_rollback_transaction` for implicit single-statement transactions.
**Primary recommendation:** Build `conga` on `riverdatabasesql.NewWithPgxListener` with a `conga.Job[T]` typed adapter, not the reflective `pact.Job`. Add `lagoon.OnDatabase` and `lagoon.AfterCommit` seams first. Port every PHP value (status codes, payload shapes, HTTP 200 denies, `[]` vs `{}`) from the file lines quoted below, not from CONTEXT.md paraphrase.
<user_constraints>
## User Constraints (from CONTEXT.md)
### Locked Decisions
#### Job manager and job records
- **D-01:** Apparatus is not ported as a plugin. The job manager is a framework package over River (see `.planning/notes/apparatus-dissolved-into-framework.md`). Its record table is **`summer_jobs`**, with the PHP `golem15_apparatus_jobs` columns: integer auto-increment `id`, `label`, `status`, `progress`, `progress_max`, `user_id`, `is_admin`, `is_canceled`, `metadata`, timestamps. Its migration ships with the framework. — **Reversibility:** one-way — `import_job_id`/`match_job_id` in the CSV API expose these ids, and the Phase 15 cutover copies rows from `golem15_apparatus_jobs` with ids preserved.
- **D-02:** `Dispatch` inserts the `summer_jobs` row and enqueues the River job in the same transaction (`riverdatabasesql` on the shared `*sql.DB`). River only executes. The row is what API responses, progress and cancellation read. The service API mirrors PHP `JobManager`: dispatch, startJob, updateJobState, updateMetadata, completeJob, failJob, cancelJob, checkIfCanceled, getMetadata.
- **D-03:** Outcomes mirror PHP exactly: `IN_QUEUE=0, IN_PROGRESS=1, COMPLETE=2, ERROR=3, STOPPED=4`.
- "Skip" is COMPLETE with metadata `{skipped: true}`, as `WishlistDigestJob` does. No new status value.
- A River retry leaves the row IN_PROGRESS. It becomes ERROR only when River discards the job for good, with the error recorded in metadata.
- **D-04:** Cancellation does both.
- `CancelJob` sets `is_canceled` + STOPPED (PHP semantics the CSV API relies on) and also calls River `JobCancel`. A queued job then never starts, and a running job's ctx is cancelled.
- Long jobs still call `CheckIfCanceled` between items, as PHP does.
- **D-05:** `queue:clear` (the Apparatus `QueueClearCommand` equivalent, a River bulk delete) belongs with this package.
#### Broadcast delivery
- **D-06:** A model broadcast is a River job enqueued in the **same transaction as the write**, on a `broadcasts` queue. The pgx LISTEN client picks it up. It is published only if the write commits, which closes PHP's "rollback still broadcasts" gap. This mirrors PHP's queued `BroadcastEventJob`.
- **D-07:** One attempt, best-effort.
- MaxAttempts 1, timeout from `broadcast_timeout` (default 5 s).
- A publish failure is logged as a warning and never affects the write, as in PHP `tries = 1`.
- **D-08:** Suppression is per model type and ctx-scoped: `WithoutBroadcasting[T](ctx, fn)` silences only T's broadcasts inside fn, matching PHP's per-class `withoutBroadcasting`. There is no package-level state (Phase 1 rule).
- The bulk pattern (suppress, then emit exactly one `collection.bulk_updated`) must be expressible.
- Proving it is success criterion 4.
- **D-09:** The broadcast contract is ported from `BroadcastableModel`:
- Event name `{action}.{plugin}.{model}` (e.g. `created.fonoteka.album`), overridable per model.
- Payload `{model, actor, timestamp, ttl}`, with `ttl` defaulting to 60.
- A delete snapshot is captured before deletion.
- `shouldBroadcast(action)` and a model-supplied channels list.
- Single channel → `publish`, several → `broadcast`.
- `broadcast_namespace` prefixing.
- **D-10:** Payload parity is proven against real PHP output. Phase 2 `tide` tooling is extended to subscribe to the PHP stack's Centrifugo during recorded flows (album create/update/delete, bulk add) and store the publications as goldens. The Go side, running on the `memory` driver or a fake Centrifugo, is diffed against them with `timestamp` and `actor` normalised. Same capture rules as Phase 2: private vars, no live secrets in git.
#### Realtime package split (pluggable transports)
- **D-11:** Realtime is a transport-neutral framework package, following the `postcard` mail-driver pattern. It owns:
- the `Publisher` interface (`Publish`, `Broadcast`);
- the channel-namespace authorizer registry (`Authorize(ctx, userID, channel) → Result{Allowed, Info, Capabilities, Overrides, reason}`);
- channel naming rules (`namespace:entity:id`, an optional single `presence:` prefix, no raw ids exposed);
- the Broadcastable interface, suppression and the River broadcast job.
Drivers are selected by `realtime.driver`: `centrifugo` (v1), `memory` (records publications for tests and parity), `log`, `null`. Models and authorizers never import a driver. — **Reversibility:** costly — every broadcastable model and authorizer is written against this interface.
- **D-12:** The **Centrifugo driver** is a sub-package. It contains:
- A hand-rolled `net/http` HTTP API client: publish, broadcast and presence, with request bytes and the `api_key` header matching the PHP `CentrifugoClient`. gocent is not added.
- A JWT token issuer with **full parity** with PHP `JwtTokenGenerator`: `generateForUser`, `generateSubscriptionToken`, `generateAnonymous`, `generateForIdentifier`, `generateSubscriptionForIdentifier`. Same secret, claims (`sub`, `exp`, `info`, `channel`) and TTL config. The WS-004 rule (no PII in `info`) is kept.
- The subscribe-proxy adapter, which ports `ProxyController`:
- constant-time `X-Centrifugo-Secret` check;
- empty user → deny;
- double-`presence:` and more-than-3-segment rejection;
- namespace lookup in the registry;
- presence defaults `prs`, `force_push_join_leave=false`;
- denials logged in full but returned as a generic `{"error":{"code":403,"message":"Access denied"}}`.
- **D-13:** Routes are declared by the driver and mounted by the app by *surface*.
- The neutral package defines `Route{Name, Method, Path, Surface, Handler}`, with `Surface` ∈ {`UserAuth`, `ServerToServer`, `Public`}.
- `fonoteka.go` calls `realtime.Mount(driver, realtime.Surfaces{UserAuth: <P6 JWT group>, ServerToServer: <raw group>, Middleware: throttle("ws-api")})` once. Switching drivers never edits `routes.go`, and the guard, group and bucket stay app-owned.
- The Centrifugo driver declares `GET /api/realtime/token` (UserAuth) and `POST /api/realtime/subscribe` (ServerToServer), with PHP paths as defaults that can be overridden.
- With an empty secret the routes still exist and return PHP's `503 {"error":"WebSocket not configured"}` (`isConfigured()`). The token route keeps PHP's 401 `{"error":"Unauthorized"}` shape.
- `memory` and `null` declare no routes. Płytarium and its tests always run the Centrifugo driver, pointed at a fake Centrifugo over `httptest` where needed.
- **D-14:** The `ws-api` bucket (120/min keyed by user id, else IP) is registered by the app like the P6 buckets.
- **D-15:** **Web Push/VAPID and all websockets console commands are ported:** `CentrifugoHealthCheck`, `GenerateVapidKeys`, `TestPushNotifications`. Web Push sits behind its own small interface with a VAPID driver. It is a separate channel from realtime, not a realtime driver. Config keys follow the PHP `push.*` block (`enabled`, `public_key`, `private_key`, `subject`).
- **D-16:** fonoteka registers its `collection` and `wishlist` authorizers (ports of `CollectionChannelAuthorizer` and `WishlistChannelAuthorizer`) into the framework registry.
- The `fonoteka.go` `websockets` plugin shrinks to config binding and the `Mount` call, and may not need to exist as a separate plugin.
- This deviates from the PROJECT.md constraint that listed websockets among the app plugins. Record it in PROJECT.md at plan time.
#### Workers and scheduler
- **D-17:** `summer serve` starts the River work client (its own `riverpgxv5` pool for LISTEN/NOTIFY) in-process by default. `queue.work_in_serve: false` disables it when a separate `summer queue:work` process runs, with `--queue` filters (e.g. `broadcasts`). Enqueue always goes through the shared `*sql.DB` driver. The success-criterion-1 timed test asserts that pickup is not poll-interval latency.
- **D-18:** The scheduler uses **River periodic jobs**.
- Plugins declare schedules through `pact.HasSchedule` (already reserved in `pact/capabilities.go`) with daily or interval cadence.
- Each entry runs the named bonfire command in-process. River leader election gives exactly one run across instances.
- The scheduler runs wherever workers run. `summer schedule:run` is a foreground scheduler-only process, and `--once` runs due commands immediately for system-cron setups.
- fonoteka's `fonoteka:prune-notifications` daily schedule is the first consumer. The command itself is ported in Phase 14.
#### Search sync
- **D-19:** Search is pluggable as well.
- The framework `Searchable` model interface covers index name (`searchableAs`), document (`toSearchableArray`) and `shouldBeSearchable`.
- The engine driver interface covers upsert, delete, flush and search-IDs.
- `search.driver: typesense | null`. The Typesense driver is hand-rolled `net/http`; typesense-go is not added.
- The Album document embeds `collection_id` (SRCH-01). The search-IDs method returns candidate ids only; the SQL re-gate (Pitfall 15) is applied by the Phase 12 search endpoint.
- **D-20:** Sync runs **after commit, inline, non-fatal**, matching PHP Scout `queue=false`. Replayed create-then-search flows see the document immediately. A Typesense error is logged and never fails the write, because the SQL re-gate makes a stale index safe. The whole sync is skipped when `search_use_typesense` is off, the Typesense config/API key is absent, or there is no DB (fresh install), matching PHP `Plugin.php` `disableSearchSyncing` and `ReindexAlbums::typesenseIsConfigured`.
### Claude's Discretion
- Package names and layout, following the beach-themed naming. ARCHITECTURE.md suggests `conga` for jobs. Names for the realtime and search packages and the driver sub-packages are open.
- How the ctx-scoped suppression reaches GORM hooks (e.g. via `tx.Statement.Context`), and how lifecycle hooks enqueue into the write transaction.
- How `actor` metadata is captured at dispatch time (PHP captures the authenticated user before queuing).
- River client configuration (queue names and concurrency, the separate listener pool size) and the timed-test harness (existing Postgres test harness vs testcontainers).
- Whether the Phase 8 D-17 expiry sweep moves into a River periodic job now (optional; it must not change recorded replays).
- The `tide` Centrifugo capture mechanics (subscriber client, storage format).
### Deferred Ideas (OUT OF SCOPE)
- **Admin extension point (referred to elsewhere as "Phase 10.1").** It is not on the roadmap. It stays a deferred idea from `10-CONTEXT.md`. Rewording the references that call it a phase (the Phase 10 decisions and the apparatus note) is a follow-up.
- Additional realtime drivers (Soketi/Pusher, Mercure, built-in hub) and search drivers beyond Typesense/null: future projects.
- Jobs admin screen (list of `summer_jobs` with progress): can be plain admin YAML later.
- Backend admin personal API tokens: `.planning/todos/pending/backend-admin-api-tokens.md`.
- Guarded outbound HTTP client and redacting slog handler: pending todos, needed by Phase 14.
</user_constraints>
<phase_requirements>
## Phase Requirements
| ID | Description | Research Support |
|----|-------------|------------------|
| JOBS-01 | River runs on the shared *sql.DB with riverdatabasesql for transactional enqueue and riverpgxv5 for LISTEN/NOTIFY; job outcomes (complete, fail, skip) are recorded and queryable through a job manager service | §Standard Stack (River v0.47.0, `NewWithPgxListener`); Pattern 1 (client roles); Pattern 2 (`conga.Job[T]` adapter); Pattern 3 (job manager semantics from PHP lines); Timed LISTEN test in Validation Architecture |
| CLI-04 | Plugins register recurring commands (daily or interval) that a scheduler runs in-process or via `summer schedule:run` | Pattern 6 (River periodic jobs + custom `Daily`/`Every` schedules, `UniqueOpts.ByPeriod`, `--once` Laravel minute-match semantics, `bonfire.Call`) |
| CLI-06 | Queue worker runs in-process or via `summer queue:work` | Pattern 1 (worker in `serve` vs `queue:work`, `--queue` repeatable flag); generated `main.go` change in `internal/build/build.go`; `summer` CLI delegate commands |
| RT-01 | Publishes to Centrifugo and issues connection and subscription JWTs with the same secret and claims at GET /api/realtime/token | Pattern 7 (claims verbatim from `JwtTokenGenerator.php`); Pattern 8 (HTTP API bytes from `CentrifugoClient.php`); note that the requirement text still says "(gocent)", superseded by D-12 |
| RT-02 | Channel-namespace authorizer registry; server re-validates on every subscribe; channel names never expose raw ids | Pattern 9 (subscribe proxy port, verbatim from `ProxyController.php`, HTTP 200 deny, `[]` info); authorizer ports |
| RT-03 | Broadcastable model interface emits live patches with bulk-write suppression, separate from durable notifications | Pattern 4 (`lagoon.OnDatabase` callback seam), Pattern 5 (ctx-scoped suppression, `realtime.Emit` for summary events), Album payload override |
| SRCH-01 | Album documents sync to Typesense scoped by collection_id, behind a settings kill-switch, degrading gracefully without DB/config | Pattern 10 (after-commit buffer), Typesense HTTP contract from Scout's `TypesenseEngine.php`, kill-switch gates |
</phase_requirements>
## Project Constraints (from CLAUDE.md)
- **Lean planning:** few, large plans. **Checkpoint:** present the plan count with a one-line scope each and wait for confirmation before writing PLAN.md files.
- **Unit tests are always the last plan of the phase.** Earlier plans may include smoke tests.
- **Stdlib first:** a dependency only when the research doc or a phase decision names it. This research names River (already decided) and its sub-modules only. gocent and typesense-go are explicitly excluded (D-12, D-19).
- **Compiled plugins:** no runtime plugin loading.
- **API parity is the acceptance test:** the Nuxt app's requests define the contract. Do not "improve" response shapes.
- **Commits:** no co-author tags. One logical change per commit. Planning docs and code in separate commits.
- **Docs rule:** any change to a `modules/` package's exported API, config keys, CLI commands or dependencies updates that module's `README.md` in the same change. New modules ship a README (H1, one-sentence summary, import line, Overview, Features, Usage, API reference, optional Configuration/CLI commands, Dependencies, Testing) plus a row in the root `README.md` modules table. Framework READMEs never name a consuming application. Every identifier named in a README must pass `go doc ./modules/<name> <Identifier>`.
- **Two repositories:** `summercms.go` holds framework packages only and never mentions Płytarium. `fonoteka.go` holds the wiring (authorizers, Album bindings, schedule, `Mount`, bucket). Planning docs stay in `summercms.go/.planning`.
- **Core plugin contracts:** the PHP user/websockets originals are not changed.
- **Security phases** carry `T-11-xx` threat IDs and tests that fail when the mitigation is broken (Phases 3, 6, 8 precedent).
- `go vet` and `go test ./...` stay green at every commit, **in both repos**. See Pitfall 12: the fonoteka.go schema-diff and migrate tests will break the moment framework migrations add tables.
- User global instruction: core plugins (user, blog, pages, payment) must not get breaking changes.
## Architectural Responsibility Map
| Capability | Primary Tier | Secondary Tier | Rationale |
|------------|-------------|----------------|-----------|
| Transactional job enqueue + `summer_jobs` row | API / Backend (framework `conga`) | Database (River tables, `summer_jobs`) | Must share the business write's `*sql.Tx` |
| Job execution (workers, periodic scheduler) | Backend worker (in `serve` or `queue:work`) | Database (LISTEN via pgx pool) | River producer/leader runs where Queues are configured |
| Token issuing `GET /api/realtime/token` | API / Backend (framework realtime driver) | App (fonoteka `Mount`, JWT group, `ws-api` bucket) | Framework owns claims; app owns guard, group, bucket (D-13) |
| Subscribe authorization `POST /api/realtime/subscribe` | API / Backend (framework proxy adapter + registry) | App (collection/wishlist authorizers) | Server-to-server call from Centrifugo; authorizers need app models |
| Publishing to Centrifugo | Backend worker (`broadcasts` queue job) | External service (Centrifugo HTTP API) | Enqueued in the write tx, published after commit by a worker |
| Live connection / fan-out | External service (Centrifugo, unchanged) | Browser (Nuxt `useCentrifugo.ts`, unchanged) | Out of scope to change |
| Search index sync | API / Backend (after-commit, inline) | External service (Typesense) | D-20: inline after commit, non-fatal |
| Search authorization re-gate | API / Backend (Phase 12 endpoint) | Database | Typesense is only a pre-filter (Pitfall 15) |
| Kill-switch `search_use_typesense` | Database (`golem15_fonoteka_settings`) | App gate implementation | App-owned setting, framework asks through an interface |
## Standard Stack
### Core
| Library | Version | Purpose | Why Standard |
|---------|---------|---------|--------------|
| `github.com/riverqueue/river` | v0.47.0 (2026-08-31T22:21:04Z) | Job queue, periodic jobs, leader election | Already decided (STACK.md). Latest tag `[VERIFIED: go list -m -versions / Go module proxy]` |
| `github.com/riverqueue/river/riverdriver/riverdatabasesql` | v0.47.0 | Driver on the shared `*sql.DB`; `NewWithPgxListener` for LISTEN | `[VERIFIED: source river_database_sql_driver.go @v0.47.0]` + `[CITED: riverqueue.com/docs/gorm]` |
| `github.com/riverqueue/river/riverdriver/riverpgxv5` | v0.47.0 | Used internally by `NewWithPgxListener` for the listener | Pulled transitively; not constructed directly |
| `github.com/riverqueue/river/rivermigrate` | (inside river module) | River schema migrations, run inside gormigrate | `MigrateTx(ctx, tx, DirectionUp, &MigrateOpts{TargetVersion: 7})` `[VERIFIED: rivermigrate/river_migrate.go:347]` |
| `github.com/riverqueue/river/rivertype` | v0.47.0 | `JobRow`, `JobState` | Transitive |
| `github.com/jackc/pgx/v5/pgxpool` | v5.10.0 (already in go.mod) | Listener pool (MaxConns 1, MinConns 0) | River requires pgx v5.10.0, the same version already pinned `[VERIFIED: river go.mod]` |
| stdlib `net/http`, `crypto/subtle`, `crypto/ecdh`, `crypto/hkdf`, `crypto/aes`, `crypto/cipher` | Go 1.27 | Centrifugo client, Typesense client, proxy secret check, Web Push | D-12, D-19 hand-rolled clients |
| `github.com/golang-jwt/jwt/v5` | v5.3.1 (already in go.mod) | HS256 Centrifugo tokens; ES256 VAPID JWT | Already decided |
### Supporting
| Library | Version | Purpose | When to Use |
|---------|---------|---------|-------------|
| `testcontainers-go/modules/postgres` | v0.44.0 (already in go.mod) | Real Postgres for River LISTEN tests | Timed pickup test, migrations, cancel test |
| `modules/wire` `wire.Time` | in repo | Carbon `+00:00` timestamps | Broadcast `timestamp` fields (`[VERIFIED: modules/wire/response.go:32-44]`, quoted: "Time marshals as Carbon's +00:00 form (2006-01-02T15:04:05+00:00), never Go's default "Z"") |
**New transitive modules River brings** (go.sum impact, relevant for any `depguard` config): `github.com/lib/pq v1.12.3` (used by riverdatabasesql for `pq.Array`, not as a driver), `github.com/robfig/cron/v3 v3.0.1`, `github.com/tidwall/gjson`, `github.com/tidwall/sjson`, `github.com/jackc/pgerrcode`. River also requires `testify v1.12.1`, so MVS bumps the repo's indirect `testify v1.11.1` `[VERIFIED: river and riverdatabasesql go.mod @v0.47.0]`.
### Alternatives Considered
| Instead of | Could Use | Tradeoff |
|------------|-----------|----------|
| One client type on `riverdatabasesql.NewWithPgxListener` | Two clients: insert-only `riverdatabasesql.New(sqlDB)` + worker `riverpgxv5.New(pgxPool)` (Pitfall 14 wording) | Two-client version: the worker runs all fetch/complete queries on a second full pgx pool, has two `TTx` types (`*sql.Tx` vs `pgx.Tx`), and `JobCancelTx`/`InsertTx` exist only on the `*sql.Tx` client. Use it only if the user insists on the literal SC-1 wording |
| Custom `Daily`/`Every` `river.PeriodicSchedule` | `robfig/cron/v3` directly (already transitive via River) | cron adds a direct dependency not named by any decision; `Daily`/`Every` is about 20 lines |
| Hand-rolled RFC 8291 Web Push (stdlib) | `github.com/SherClockHolmes/webpush-go` v1.4.0 (2025-01-02) `[ASSUMED]`, not vetted via official docs | Library is about 21 months since its last release, and no decision names it. The hand-rolled version is testable against the RFC 8291 Appendix A vector |
| `tide` fake-Centrifugo HTTP recorder | WebSocket subscriber against a real PHP-side Centrifugo (D-10 wording) | A subscriber needs a WebSocket client (no stdlib client; `golang.org/x/net/websocket` is in go.mod but the Centrifugo protocol is extra work) and only sees `data`. The recorder sees channel, headers and full body. See Open Question 4 |
**Installation:**
```bash
cd summercms.go
go get github.com/riverqueue/river@v0.47.0 \
github.com/riverqueue/river/riverdriver/riverdatabasesql@v0.47.0 \
github.com/riverqueue/river/rivertype@v0.47.0
go mod tidy
```
## Package Legitimacy Audit
The GSD `package-legitimacy` seam supports only `npm|pypi|crates` (error: `Usage: gsd-tools package-legitimacy check --ecosystem <npm|pypi|crates>`). Go modules were checked against the Go module proxy and official docs instead.
| Package | Registry | Age | Downloads | Source Repo | Verdict | Disposition |
|---------|----------|-----|-----------|-------------|---------|-------------|
| github.com/riverqueue/river | proxy.golang.org | v0.47.0 published 2026-08-31; long release history (0.34→0.47 in CHANGELOG since 2026-04) | n/a (Go) | github.com/riverqueue/river | OK (already decided in STACK.md; official docs riverqueue.com) | Approved |
| github.com/riverqueue/river/riverdriver/riverdatabasesql | proxy.golang.org | v0.47.0 | n/a | same repo (sub-module) | OK | Approved |
| github.com/riverqueue/river/rivertype | proxy.golang.org | v0.47.0 | n/a | same repo | OK | Approved |
| github.com/SherClockHolmes/webpush-go | proxy.golang.org | v1.4.0, 2025-01-02 | n/a | github (not inspected) | not run | **Not recommended**, `[ASSUMED]`; do not install without a checkpoint |
**Packages removed due to [SLOP] verdict:** none
**Packages flagged as suspicious [SUS]:** none. `webpush-go` is only an alternative and is not recommended.
## Architecture Patterns
### System Architecture Diagram
```
Nuxt SPA ──GET /api/realtime/token──► [surf: throttle ws-api → jwt.auth] ─► realtime driver TokenHandler ─► {token} (HS256: sub, exp, info{name})
│
└─WS─► Centrifugo (unchanged) ──POST /api/realtime/subscribe (X-Centrifugo-Secret)──► ProxyHandler
▲ │ parse channel → registry[namespace]
│ ▼
│ app authorizer (collection / wishlist) → DB
│ │
│ 200 {"result":{...}} | 200 {"error":{"code":403,...}}
│
HTTP write ─► service ─► lagoon.Transaction(ctx) ─► GORM Create/Update/Delete
│ │
│ ├─ realtime GORM callback (after_create/update/delete; before_delete snapshot)
│ │ └─ suppressed for T in ctx? → skip
│ │ └─ else River InsertTx(*sql.Tx, BroadcastArgs{channels,event,payload}, queue=broadcasts)
│ │ └─ pg_notify(river_insert) inside tx
│ ├─ search GORM callback → record {index,id,op} in ctx after-commit buffer
│ └─ conga.Dispatch → INSERT summer_jobs + River InsertTx (same tx)
▼
COMMIT ─► after-commit flush: search gate (driver configured? settings on? DB?) → reload model → Typesense upsert/delete (errors logged)
│
Postgres NOTIFY
▼
River worker client (serve or queue:work; NewWithPgxListener: queries on shared *sql.DB, LISTEN on pgxpool MaxConns=1)
├─ broadcasts queue → BroadcastWorker → format channels (lowercase, namespace) → 1 channel: POST /publish | N: POST /broadcast → Centrifugo
├─ default queue → conga.Job[T] wrapper → summer_jobs lifecycle (IN_PROGRESS … COMPLETE/ERROR/STOPPED)
└─ leader only: periodic enqueuer → scheduled queue → ScheduledCommandWorker → bonfire.Call(name,args)
```
### Recommended Project Structure
```
summercms.go/modules/
├── conga/ # River client wiring, job manager (summer_jobs), conga.Job[T], queue:work, queue:clear,
│ │ # scheduler (HasSchedule → periodic jobs), schedule:run, migrations (River v7 + summer_jobs)
│ └── README.md
├── lighthouse/ # realtime (name is discretionary): Publisher, Registry, Result, channel rules, Route/Surface/Mount,
│ │ # Broadcastable, WithoutBroadcasting[T], Emit, BroadcastWorker, memory/log/null drivers
│ ├── centrifugo/ # HTTP API client, TokenIssuer (5 generators), ProxyHandler, health command
│ └── README.md
├── flare/ # web push (name discretionary): Pusher interface, VAPID driver (RFC 8291/8292), vapid-keys command, test-push
├── beachcomber/ # search (name discretionary): Searchable, Engine, Gate, after-commit sync, null driver
│ └── typesense/ # hand-rolled net/http engine
├── lagoon/ # + OnDatabase(app, fn), AfterCommit buffer, Transaction(ctx, db, fn)
├── pact/ # + HasSchedule, ScheduledCommand, Cadence
├── bonfire/ # + Call(ctx, commands, name, args, out) (Artisan::call)
└── party/ # + conga/realtime/search Activate before Boot (same slot as postcard.Activate)
internal/build/build.go # generated main.go: + conga.RuntimeCommands, realtime commands, command catalog publish
cmd/summer/main.go # + delegateCommand("queue:work"|"queue:clear"|"schedule:run"|websockets:*)
fonoteka.go/plugins/golem15/fonoteka/
├── classes/ws/ # CollectionChannelAuthorizer, WishlistChannelAuthorizer ports
├── models/album.go # Broadcastable + Searchable bindings
├── routes.go # realtime.Mount(...) once
├── plugin.go # Buckets(): ws-api; Schedule(): prune-notifications daily; Boot: registry.Register(...)
└── config/ or ../../config/realtime.yaml, queue.yaml, search.yaml
```
### Pattern 1: River client roles (one client type, three roles)
**What:** One `*river.Client[*sql.Tx]` built by `conga` from app config.
- **Insert-only** (HTTP process with `queue.work_in_serve: false`, CLI commands): `riverdatabasesql.New(sqlDB)`, `Workers` registered (for kind validation), no `Queues`, never `Start`ed.
- **Worker** (`serve` default, `queue:work`): `riverdatabasesql.NewWithPgxListener(sqlDB, listenerPool)`, `Queues` from config or the `--queue` filter, `PeriodicJobs` from `HasSchedule`, `Start(ctx)`, `Stop` on shutdown.
**Why it works:**
- The database/sql driver reports `SupportsListenNotify() bool { return true }` `[VERIFIED: riverdatabasesql river_database_sql_driver.go:139]`, so `maybeNotifyInsertForQueues` calls `tx.NotifyMany(...)` inside the insert transaction `[VERIFIED: river client.go:2188-2228]`.
- The worker in another process receives it through its listener: "The database/sql pool continues to be used for all database operations other than listening for notifications. The Pgx pool is used only to acquire dedicated connections for Postgres LISTEN commands" `[VERIFIED: river_database_sql_driver.go:57-77]`.
- "A pool dedicated to one River client can generally set MinConns to zero and MaxConns to one." "Applications using PgBouncer must configure the listener pool to use session pooling or connect it directly to Postgres."
**Example:**
```go
// Source: riverqueue.com/docs/gorm + river_database_sql_driver.go @v0.47.0
listenerCfg, err := pgxpool.ParseConfig(lagoon.DSN(app.Config)) // same DB as the shared *sql.DB
if err != nil { return err }
listenerCfg.MaxConns, listenerCfg.MinConns = 1, 0
listener, err := pgxpool.NewWithConfig(ctx, listenerCfg)
if err != nil { return err }
client, err := river.NewClient(riverdatabasesql.NewWithPgxListener(sqlDB, listener), &river.Config{
Queues: map[string]river.QueueConfig{"default": {MaxWorkers: 4}, "broadcasts": {MaxWorkers: 8}, "scheduled": {MaxWorkers: 1}},
Workers: workers,
PeriodicJobs: periodic,
MaxAttempts: 3, // PHP worker: --tries=3 (docs/deploy/systemd/plytarium-queue.service)
JobTimeout: 5 * time.Minute, // PHP worker: --timeout=300; River default is 1m
Logger: log,
})
```
Defaults to override, verbatim from `river@v0.47.0/client.go:46-54`: `FetchCooldownDefault = 100 * time.Millisecond`, `FetchPollIntervalDefault = 1 * time.Second`, `JobTimeoutDefault = 1 * time.Minute`, `MaxAttemptsDefault = rivercommon.MaxAttemptsDefault` (`= 25`, `internal/rivercommon/river_common.go:16`).
### Pattern 2: `conga.Job[T]` typed adapter over `pact.Job`
**What:** River registers workers by static type: `AddWorkerSafely[T](workers, worker)` computes the kind from `var jobArgs T` `[VERIFIED: worker.go:133-136]`. `AddWorkerArgs` (dynamic args) is documented as "its use should be considered internal only" `[VERIFIED: worker.go:114-118]`. The existing reflective contract cannot be registered on River without that internal function:
```go
// modules/pact/capabilities.go:84-99 (verbatim)
type JobArgs interface {
Kind() string
}
type Job interface {
Work(ctx context.Context, args JobArgs) error
}
type HasJobs interface {
Jobs() []Job
}
```
**Recommendation:** keep `pact.Job` River-free, as its own comment requires ("Phase 11 adapts this onto River; the interface itself must not import River", `pact/capabilities.go:90-91`). Add a constructor in `conga`: `conga.Job[T pact.JobArgs](fn func(ctx context.Context, args T) error, opts ...JobOption) pact.Job`. It returns a `*typedJob[T]` that implements `pact.Job` and an unexported `register(*river.Workers) error`, which calls `river.AddWorkerSafely[T]`. Plugins return `conga.Job(...)` values from `Jobs()`. Update the `make:job` stub in `internal/build/stubs/artifacts.tmpl` ("Generate a plugin job without importing River") so it generates a typed func wrapped in `conga.Job`. A `pact.Job` not built by `conga.Job` fails boot with a clear error.
**Job-to-row linkage:** River `InsertOpts.Metadata` carries `{"summer_job_id": N}`; the wrapper reads `job.Metadata` and exposes it as `conga.JobID(ctx)`. This replaces PHP `assignJobId` without polluting typed args. The row stores the River id for `JobCancel` (see Open Question 2 on a `river_job_id` column).
### Pattern 3: Job manager semantics (verbatim PHP)
Status constants, `apparatus/contracts/JobStatus.php:17-21` `[VERIFIED]`:
```
const IN_QUEUE = 0;
const IN_PROGRESS = 1;
const COMPLETE = 2;
const ERROR = 3;
const STOPPED = 4;
```
Table, `apparatus/updates/create_jobs_table.php:13-22` `[VERIFIED]`:
```
$table->increments('id');
$table->string('label');
$table->integer('status')->default(0);
$table->integer('progress')->default(0);
$table->integer('progress_max')->default(0);
$table->integer('user_id')->nullable();
$table->boolean('is_admin')->default(false);
$table->boolean('is_canceled')->default(false);
$table->text('metadata');
$table->timestamps();
```
Dispatch, `apparatus/classes/JobManager.php:76,82` `[VERIFIED]`: `'status' => JobStatus::IN_PROGRESS,` and `'metadata' => json_encode($metadata),`. When no `metadata` parameter is given, `$metadata = ''`, so the stored text is the JSON string `""`.
Method behaviours to copy (JobManager.php, read this session):
| Method | PHP behaviour |
|--------|---------------|
| `dispatch` | Actor = frontend `Auth::getUser()` → `user_id`, `is_admin=false`; then `BackendAuth::getUser()` overrides → `is_admin=true`. `progress_max = parameters['count'] ?? 0`. Delay → `queue->later($delay)`. Returns the id |
| `startJob(id, total)` | `progress=0, progress_max=total, updated_at=now` |
| `updateJobState(id, current, metadata=[])` | Sets only `progress` (no `updated_at`); metadata **replaced** if non-empty |
| `updateMetadata` | Replaces metadata |
| `completeJob(id, metadata=[])` | `status=COMPLETE, progress=progress_max`; metadata replaced only if non-empty. `simpleJob` mode deletes the row |
| `failJob` | `status=ERROR`; metadata replaced if non-empty |
| `cancelJob` | `status=STOPPED` only. `is_canceled` is set by `CsvImportApiController::cancelJob` with a raw `DB::table(...)->update(['is_canceled' => true])` **before** calling `cancelJob` |
| `checkIfCanceled` | Reads `is_canceled` |
| `getMetadata` | `json_decode(...) ?: []` |
In Go, `Principal.ID`/`Principal.Backend` map to `user_id`/`is_admin`. Quoted from `modules/bouncer/context.go:14-22`: `ID uint` and `Backend bool \`json:"-"\``.
Consumers that fix the contract:
- CSV progress reads `progress`, `progress_max`, `status` (`CsvImportApiController.php:548-556`).
- Delayed dispatch is used by `AlbumCsvMatchJob` (rate-limit resume) and `WishlistDigestQueue::enqueue` (`1800`). Map it to `InsertOpts.ScheduledAt`.
### Pattern 4: `lagoon.OnDatabase` seam (fixes the boot-order gap)
The real binary runs `party.Activate` (Boot) **before** `surf.ServeCommand` publishes `*gorm.DB`. `fonoteka.go/plugins/golem15/fonoteka/plugin.go` documents that `RegisterHooks` "only runs when gdb happens to already be published at Boot time (true for every test harness in this repo, false for the real CLI). This is a known gap ... (album/collection slug and credential-encryption hooks would not fire under `summer serve`)". Realtime and search GORM callbacks registered the same way would pass every test and fail in production.
**Recommendation:** `lagoon.OnDatabase(app, func(sqlDB *sql.DB, gdb *gorm.DB) error)` runs immediately if the DB is already published, else when `lagoon.Publish` runs. The realtime and search packages register their `gdb.Callback()` installers through it. GORM callbacks live on the shared `Config`, so all sessions see them. This also closes the fonoteka `RegisterHooks` gap if the app adopts it. Treat that as optional and flag it to the user.
### Pattern 5: Broadcast in the write transaction with ctx-scoped suppression
- Register the callbacks once: `Create().After("gorm:after_create")`, `Update().After("gorm:after_update")`, `Delete().Before("gorm:before_delete")` for the snapshot, and `Delete().After("gorm:after_delete")`. Each checks `db.Statement.Model`/`Dest` for `Broadcastable`.
- Get the transaction from `db.Statement.ConnPool.(*sql.Tx)`. `[VERIFIED: gorm finisher_api.go:699-701]`: `*sql.DB` is a `TxBeginner`, so `Begin` yields `*sql.Tx`. `[VERIFIED: callbacks/transaction.go]`: every single-statement write is wrapped in `gorm:begin_transaction`, since lagoon opens GORM with default config (`modules/lagoon/connection.go:68`, `&gorm.Config{}`).
- Suppression: `ctx = context.WithValue(ctx, suppressKey{}, set ∪ reflect.TypeFor[T]())`. The callback reads `db.Statement.Context`. This only works if the writes inside `fn` use `gdb.WithContext(ctx)` with the ctx handed to `fn`.
- The summary event is an explicit `realtime.Emit(ctx, tx, Broadcast{Channels, Event, Payload})`, enqueued on the same transaction. PHP publishes summary events synchronously after the write (`AlbumApiController.php:153-154`: `'collection.bulk_updated',` / `['reason' => 'bulk_create', 'count' => count($created)]`). Enqueue-in-tx preserves "only if committed".
BroadcastableModel defaults, verbatim `[VERIFIED: websockets/traits/BroadcastableModel.php]`:
- `'model' => $modelData,` / `'actor' => $this->getActorMetadata(),` / `'ttl' => $this->getBroadcastTtl(),` (156-159)
- `return 60; // Default 1 minute` (280)
- actor default `'name' => 'System',` (191)
- event name `return strtolower("{$action}.{$this->getBroadcastAlias()}");` (251)
**Album overrides the whole payload** (`fonoteka/models/Album.php`, read this session): `id`, `collection_id` (int), `action`, `actor`, `timestamp`, plus `album` (serialized fresh with genre/styles/artists/photos) unless `deleted`. There is no `ttl` and no `model`. Channels: `["collection:{id}"]` only when `collection.kind === 'collection'`, else `[]` (no broadcast). Actor comes from `auth()->user()`: `{user_id, name}` or `{user_id: null, name: "System"}`. Go therefore needs an app-provided `ActorFunc(ctx) Actor`, because `bouncer.Principal` has no name.
BroadcastEventJob, verbatim (`websockets/jobs/BroadcastEventJob.php`):
- `public int $tries = 1;` (27)
- `$this->timeout = (int) config('websockets.broadcast_timeout', 5);` (65)
- channels are lowercased with `$name = strtolower((string) $channel);` (148), and the namespace prefix `strtolower($namespace) . ':'` (155) is skipped when already present
- one channel → `publish`, several → `broadcast`
### Pattern 6: Scheduler on River periodic jobs
- `pact.HasSchedule` **does not exist yet**. It is only named in a comment (`pact/capabilities.go:368-372`: "Future capability families are type-asserted when their first consumer packages exist: HasListeners HasSchedule"). This phase defines it, e.g. `type ScheduledCommand struct{ Command string; Args []string; Cadence Cadence }`, `HasSchedule{ Schedule() []ScheduledCommand }`.
- `river.PeriodicSchedule` is `Next(current time.Time) time.Time` `[VERIFIED: periodic_job.go]`. River's scheduler "only retains in-memory state, so anytime a process quits or a new leader is elected, the whole process starts over". `PeriodicInterval(24h)` would therefore drift with every restart. **Use a custom `Daily(hh:mm, loc)` schedule** (next wall-clock occurrence in the app timezone, as Laravel `->daily()` = `0 0 * * *`) and `Every(d)` aligned to wall-clock multiples.
- Leader failover can double-insert. Add `InsertOpts.UniqueOpts{ByPeriod: 24h}` (daily) or `ByPeriod: d` (interval).
- `Start` needs Queues and Workers: "client Queues and Workers must be configured for a client to start working" `[VERIFIED: client.go ~1098]`. So `schedule:run` is a River client with only a `scheduled` queue plus the periodic jobs. Every worker client must carry the same `PeriodicJobs` config, because only the elected leader enqueues.
- `--once` (system cron): skip River and run, in-process, every entry whose cadence matches the current minute (Laravel `schedule:run` semantics, `[ASSUMED]`: training knowledge of Laravel).
- In-process execution needs `bonfire.Call(ctx, commands, name, args, out)`: build a root with `bonfire.NewRootIO(name, []Command{cmd}, …)` and `SetArgs`. That requires the final command list. The generated main builds it last (`internal/build/build.go:113-120`), so the generated main should publish the command catalog on the app before `root.Execute()`.
- fonoteka's first consumer references `fonoteka:prune-notifications`, which does not exist until Phase 14 (Open Question 5).
### Pattern 7: Centrifugo token claims (verbatim, `websockets/classes/JwtTokenGenerator.php`)
| Generator | Claims (file lines) |
|-----------|---------------------|
| `generateForUser` | `'sub' => (string) $user->id,` `'exp' => time() + $this->ttl,` `'info' => [` `'name' => $user->name,` (71-74), HS256 (78) |
| `generateSubscriptionToken` | `'sub' => (string) $user->id,` `'channel' => $channel,` `'exp' => time() + $this->ttl,` (95-97) |
| `generateAnonymous` | `'sub' => '',` `'exp' => time() + 300, // 5 minutes for anonymous` (115-116) |
| `generateForIdentifier` | `'sub' => $identifier,` `'exp' => time() + $this->ttl,` `'info' => $info,` (130-132) |
| `generateSubscriptionForIdentifier` | `'sub' => $identifier,` `'channel' => $channel,` `'exp' => time() + $this->ttl,` (146-148) |
- Config defaults (`websockets/config/config.php`): `'token_ttl' => env('CENTRIFUGO_TOKEN_TTL', 3600)`, `'api_url' => env('CENTRIFUGO_API_URL', 'http://127.0.0.1:8001/api')`.
- Token route (`websockets/routes.php`): prefix `/api/realtime`, group middleware `['throttle:ws-api', 'bindings']`, the `token` route adds `jwt.auth`. `503 {"error":"WebSocket not configured"}` when `token_secret` is empty. Body `{"token": ...}`.
- The Go user model has `Name *string` (`fonoteka.go/plugins/golem15/user/models/user.go:15`), so `info.name` may be `null`.
### Pattern 8: Centrifugo HTTP API bytes (verbatim, `websockets/classes/CentrifugoClient.php`)
- Header: `'Authorization' => 'apikey ' . $this->apiKey,` (99), timeout 5 s, POST `$this->apiUrl . '/publish'` (101) / `'/broadcast'` (159) / `'/presence'` / `'/unsubscribe'`.
- Body: `{"channel": c, "data": {"event": e, "payload": p, "timestamp": now()->toIso8601String()}}` (publish, lines 101-108); `{"channels": [...], "data": {...}}` (broadcast).
- Disabled (`api_key` empty): returns false and never sends.
- Success is `$response->successful()`, i.e. **HTTP 2xx only**. Centrifugo returns API errors as 200 with an `error` body `[CITED: centrifugal.dev/docs/server/server_api]`, so PHP treats those as success. Keep that, and optionally debug-log the body.
- Centrifugo master still accepts `Authorization: apikey <KEY>` besides `X-API-Key` `[CITED: github.com/centrifugal/centrifugo internal/middleware/auth.go]`.
- The locally installed server is `Centrifugo v6.5.1` `[VERIFIED: centrifugo version]`. The production version was not observed.
### Pattern 9: Subscribe proxy (verbatim, `websockets/http/controllers/ProxyController.php`)
- Secret: `if (empty($expectedSecret) || !hash_equals($expectedSecret, $providedSecret ?? ''))` → deny (33-40). An empty configured secret **denies**; it does not return 503.
- `$userId = $request->input('user');`, and `empty($userId)` → deny. In PHP `"0"` is empty, so user `"0"` is denied.
- `parseChannel`: `presence:presence:` → namespace `''`; strip one `presence:`; `count($parts) > 3` → namespace `''`; namespace = first segment.
- Unknown namespace → deny. The authorizer receives the **full original channel** (including any `presence:` prefix) and `(int) $userId`.
- Allow: `$response = ['info' => $auth->info ?? []];`. For presence channels: `allow = capabilities ?? ['prs']`, `override = array_merge(['presence' => ['value' => true], 'join_leave' => ['value' => true], 'force_push_join_leave' => ['value' => false]], overrides ?? [])`. Returns `response()->json(['result' => $response])`.
- Deny: `Log::channel('websockets')->warning('Subscription denied', [...])` and `response()->json(['error' => ['code' => 403, 'message' => 'Access denied']])`. This is **HTTP 200**, the Laravel default status. Centrifugo expects 200 with an `error` body, and treats non-200 as internal error code 100 `[CITED: centrifugal.dev/docs/server/proxy]`.
- `AuthorizationResult::allowed()` has `info = []`. PHP `json_encode([])` is `[]`, so the allow body is `{"result":{"info":[]}}`. A Go `map[string]any{}` would emit `{}`.
- Authorizers: `$registry->register('collection', …CollectionChannelAuthorizer::class)` (fonoteka `Plugin.php:168`) and `register('wishlist', …)` (172). Both parse `explode(':', $channel)[1]` as the id. A `presence:collection:5` channel therefore yields `(int)"collection" = 0` → "malformed channel" deny. Port that behaviour as is.
### Pattern 10: After-commit search sync
- GORM `BeginTransaction` sets `gorm:started_transaction` only when it opened the tx itself. Inside an explicit `db.Transaction`, `Begin()` on a `*sql.Tx` ConnPool returns `ErrInvalidTransaction`, which is ignored `[VERIFIED: gorm callbacks/transaction.go]`. So:
- Search callbacks append `{searchable type, key, op}` to a ctx buffer (`lagoon.WithAfterCommit(ctx)`).
- Implicit single-statement transactions: a callback registered `After("gorm:commit_or_rollback_transaction")` flushes when `InstanceGet("gorm:started_transaction")` is present and `db.Error == nil`.
- Explicit transactions: app code uses `lagoon.Transaction(ctx, gdb, fn)`, which commits and then flushes. Existing `gdb.WithContext(ctx).Transaction(...)` call sites (e.g. `classes/album_write_service.go:55`) must migrate in Phase 12.
- Fallback with no buffer: flush immediately (the PHP `after_commit=false` behaviour).
- Build the document **after commit by reloading** the model. PHP `toSearchableArray` loads genre/styles/artists, and Go `SaveAlbum` syncs the artist pivot *after* `tx.Save(album)`, so a document built in `AfterSave` would carry stale artists.
- Gates, in this order, all non-fatal:
1. `search.driver != typesense` or empty `api_key` → skip.
2. No DB → skip.
3. App `Gate.Enabled(ctx)`: fonoteka reads `golem15_fonoteka_settings.search_use_typesense`. Go model `SearchUseTypesense bool \`gorm:"column:search_use_typesense"\`` (`models/settings.go:9`). An error reading it means "off", like PHP's try/catch.
- When enabled and the model is soft-deleted or deleted: delete the document. Scout's `deleted` → `unsearchable()` because Winter's `SoftDelete` is not Laravel's `SoftDeletes`, and `scout.soft_delete` is `false`.
- Typesense wire contract from Scout `TypesenseEngine.php` (v10.25.0) and typesense-php v4.9.3 `[VERIFIED: vendor source]`:
- upsert: `GET /collections/{name}` (on error `POST /collections` with the model's `collection-schema` + `name`), then `POST /collections/{name}/documents/import?action=upsert` with a JSONL body (`'import_action' => env('TYPESENSE_IMPORT_ACTION', 'upsert')`, `config/scout.php:157`). Any line with `success:false` is an error.
- delete: retrieve then `DELETE /collections/{name}/documents/{id}`, swallowing errors.
- The header is `X-TYPESENSE-API-KEY`. Import returns 200 even when documents fail `[CITED: typesense.org/docs/26.0/api/documents.html]`.
- Index name `return 'golem15_fonoteka_albums';` (`Album.php:498`). The document requires a positive `collection_id` (PHP throws `LogicException`) and field `'collection_id' => $collectionId,` (521). The local server image is `typesense/typesense:26.0` (`scripts/run-typesense.sh:68`).
### Anti-Patterns to Avoid
- **A same-client `Insert` (non-tx) in the timed test.** River's poll-only driver locally triggers a fetch for non-transactional inserts from the same client (`notifyProducerWithoutListenerJobFetch`, `client.go:2064`). The test must use `InsertTx` from a *different* client to prove LISTEN.
- **`PeriodicInterval(24*time.Hour)` for daily jobs.** It restarts from leader start, so a redeploy every <24h means the job never runs.
- **Registering GORM callbacks in `Boot` with `app.Lookup[*gorm.DB]`.** They work in tests and are dead in production (Pattern 4).
- **`map[string]any{}` for PHP empty arrays.** It emits `{}` where PHP emits `[]` (proxy `info`, `generateForIdentifier` with `$info = []`).
- **`time.RFC3339` for broadcast timestamps.** It gives `Z`. PHP gives `+00:00`; use `wire.Time`.
- **Trusting Typesense hits as authorized.** Phase 12 re-gates in SQL (Pitfall 15).
## Don't Hand-Roll
| Problem | Don't Build | Use Instead | Why |
|---------|-------------|-------------|-----|
| Job queue, retries, leader election, periodic enqueue | Custom `summer_jobs` poller | River (row is only the record, D-02) | Leader election, stuck-job rescue, LISTEN wake-ups, unique jobs |
| River schema | Hand-written DDL | `rivermigrate` `MigrateTx` pinned `TargetVersion: 7` inside a gormigrate step | River owns 7 migrations incl. types and functions |
| Constant-time secret compare | `==` | `crypto/subtle.ConstantTimeCompare` | Port of `hash_equals` (PITFALLS "constant-time secret comparison") |
| JWT signing | Manual HMAC + base64 | `golang-jwt/jwt/v5` | Already decided |
| Web Push encryption primitives | Custom ECDH/HKDF/GCM | stdlib `crypto/ecdh`, `crypto/hkdf`, `crypto/cipher` composed per RFC 8291, verified against its Appendix A vector | No undecided dependency; the RFC vector makes it testable |
| Queue bulk delete | Raw `DELETE FROM river_job` | `client.JobDeleteMany(ctx, river.NewJobDeleteManyParams().States(available, scheduled, retryable).First(10000))` in a loop until 0 | River ignores running jobs; default limit is 100 |
**Key insight:** most risk here is contract drift (status codes, JSON shapes, headers), not algorithms. Every hand-rolled client (Centrifugo, Typesense) is small. The shapes must come from the quoted PHP lines.
## Runtime State Inventory
Not a rename phase. Still relevant, because framework migrations add state:
| Category | Items Found | Action Required |
|----------|-------------|------------------|
| Stored data | PHP `golem15_apparatus_jobs` rows (ids exposed as `import_job_id`/`match_job_id`) | Phase 15 cutover copies them into `summer_jobs` with ids preserved (already in the apparatus note). The `summer_jobs` sequence must be bumped past max(id) after the copy (`setval`) |
| Live service config | Centrifugo server config (subscribe proxy endpoint, `X-Centrifugo-Secret` static header, namespaces, `hmac_secret_key`, `http_api.key`) lives on the host, not in git. Local `/etc/centrifugo/config.json` is not readable by this user | None in Phase 11 (Centrifugo unchanged). The cutover must keep the same secret values and proxy URL |
| OS-registered state | `plytarium-queue.service` runs `php artisan queue:work redis --sleep=1 --tries=3 --timeout=300` | Phase 15: replace it with `fonoteka queue:work` or rely on in-serve workers |
| Secrets/env vars | PHP env `CENTRIFUGO_SECRET`, `CENTRIFUGO_API_KEY`, `CENTRIFUGO_API_URL`, `CENTRIFUGO_PROXY_SECRET`, `CENTRIFUGO_TOKEN_TTL`, `BROADCAST_NAMESPACE`, `TYPESENSE_*` | The Go env names become `SUMMER_<FILE>__<KEY>` (compass). A cutover mapping table is needed; the values are unchanged |
| Build artifacts | fonoteka.go `main.go` is generated by `summer build` (`internal/build/build.go:80-130`) | Regenerate after the generator gains conga/realtime commands (also `examples/hello`) |
## Common Pitfalls
### Pitfall 1: Dispatch status is IN_PROGRESS, not IN_QUEUE
**What goes wrong:** A Go `Dispatch` writes `status=0`, and the CSV progress API (`'status' => (int) $row->status`) returns 0 where PHP returns 1.
**How to avoid:** Insert `status = 1`, copying `JobManager.php:76`. `IN_QUEUE` stays defined but unused by dispatch.
**Warning signs:** CSV progress parity diff on `status`.
### Pitfall 2: CancelJob from inside the running job
**What goes wrong:** PHP jobs call `cancelJob($this->jobId)` themselves after `checkIfCanceled` (`AlbumCsvImportJob.php:73-76`). If the Go `CancelJob` also calls River `JobCancel` on the *running* job, the job's own ctx is cancelled mid-cleanup.
**How to avoid:** Split the API: `CancelJob` (external: `is_canceled` + STOPPED + `JobCancel`, D-04) and a worker-side `StopJob` (STOPPED only, the PHP `cancelJob` semantics). The wrapper returns nil once the row is STOPPED, so River neither retries nor marks it discarded.
### Pitfall 3: Retries vs the ERROR status
**What goes wrong:** Calling `FailJob` on every error marks ERROR while River still retries.
**How to avoid:** In the wrapper, `if err != nil && job.Attempt >= job.MaxAttempts` → `FailJob` with the error merged into metadata (`rivertype.JobRow` has `Attempt`, `MaxAttempts`, `Metadata` `[VERIFIED: rivertype @v0.47.0]`). Recover panics in the wrapper too.
### Pitfall 4: River JobTimeout defaults to 1 minute
**What goes wrong:** A long CSV match or import job is ctx-cancelled at 60 s. PHP ran with `--timeout=300`.
**How to avoid:** Set `JobTimeout` from config (default 300 s) and let a per-job `Timeout()` override it.
### Pitfall 5: Timed LISTEN test that proves nothing
**What goes wrong:** It passes on poll latency (default 1 s poll), or the listener is not yet `LISTEN`ing when the insert happens.
**How to avoid:**
- Set `FetchPollInterval: 30*time.Second` on the worker.
- Wait for readiness: `SELECT count(*) FROM pg_stat_activity WHERE query ILIKE 'LISTEN%'` > 0, or a warm-up job completed.
- Insert via `InsertTx` on a *second* insert-only client, commit, and assert pickup < 1 s.
- Negative control with `riverdatabasesql.New` (poll-only) and the same 30 s poll: assert *no* pickup within 2 s.
### Pitfall 6: Boot-order hole for GORM callbacks
See Pattern 4. A test must run callbacks with the DB published *after* `party.Activate`, which is the production order.
### Pitfall 7: Suppression lost through a stale ctx
**What goes wrong:** `WithoutBroadcasting[Album](ctx, fn)`, but `fn` writes through a `gdb` bound to the request ctx, not the ctx passed to `fn`.
**How to avoid:** `fn` receives ctx. Document it, and have the SC-4 test write through `gdb.WithContext(innerCtx)` and also negative-test the stale ctx.
### Pitfall 8: Batch updates fire hooks on an empty model
**What goes wrong:** `db.Model(&Album{}).Where(...).Updates(...)` runs update callbacks with a zero primary key. A naive broadcaster emits `collection:0` or panics on a nil relation.
**How to avoid:** Skip broadcast and search sync when the primary key is zero. Bulk paths use explicit suppression and `Emit`.
### Pitfall 9: PHP empty array → `[]`
The allow body `info`, `generateForIdentifier` `info`, and any `$metadata = []` all encode as `[]` in PHP. Use a nil-safe `json.RawMessage("[]")` default where PHP would emit an empty array.
### Pitfall 10: Two different "configured" checks
- Token route 503 depends on `token_secret` (`isConfigured()`).
- Publishing is a silent no-op when `api_key` is empty.
- The proxy denies when `proxy_secret` is empty.
Each has its own test. Do not collapse them into one `enabled` flag.
### Pitfall 11: ws-api bucket key before auth
**What goes wrong:** PHP: `Limit::perMinute(120)->by($request->user()?->id ?: $request->ip())` (`websockets/routes.php:11`), with throttle listed before `jwt.auth`. In Go the key func runs in middleware order, so with throttle first `bouncer.User(ctx)` is empty and the bucket keys by IP.
**How to avoid:** On the token route, order `jwt.auth` before `throttle:ws-api`, or have the key func parse the bearer. The subscribe route (server-to-server) keys by IP, i.e. the Centrifugo host: one shared 120/min bucket for all subscribes, which is what PHP does too. Flag it to the user (Open Question 6).
### Pitfall 12: fonoteka.go tests break on new framework tables
**What goes wrong:** `parity/schema_diff_test.go` compares Go public tables against the PHP snapshot with an explicit Go-only allow-list. `parity/migrate_test.go:71` asserts the exact history-table list `summer_migrations_golem15_fonoteka, …_golem15_user, …_summercms_attach, …_summercms_cabana`. Adding conga migrations adds `summer_jobs`, `river_job`, `river_leader`, `river_queue`, `river_notification`, `river_migration` (River v7 final set; `river_client*` are created then dropped in 007 `[VERIFIED: migration/main/*.up.sql]`) and a new history table.
**How to avoid:** Update both lists in fonoteka.go in the same wave as the conga migration commit.
### Pitfall 13: River migrations inside gormigrate
`rivermigrate.New(riverdatabasesql.New(sqlDB), nil)`, then `MigrateTx(ctx, tx.Statement.ConnPool.(*sql.Tx), rivermigrate.DirectionUp, &rivermigrate.MigrateOpts{TargetVersion: 7})`. Pinning prevents a River bump from silently applying new DDL. For the gormigrate `Rollback`, use `MigrateTx(ctx, sqlTx, rivermigrate.DirectionDown, &rivermigrate.MigrateOpts{TargetVersion: -1})`. A down migration with `TargetVersion: 0` runs **only one step** (`case direction == DirectionDown && opts.TargetVersion == 0: maxSteps = 1`, `river_migrate.go:532-533`), and the docs say "When migrating down, TargetVersion can be set to the special value of -1 to apply all down migrations (i.e. River schema is removed completely)" (`river_migrate.go:239-241`) `[VERIFIED]`.
### Pitfall 14: Web Push has no working PHP reference in Płytarium
- `GenerateVapidKeys` uses `Minishlink\WebPush\VAPID`, which is not in `fonoteka/vendor`.
- `TestPushNotifications` imports `Golem15\Notifications\Services\PushNotificationService` and `Golem15\Notifications\Models\PushSubscription`, and that plugin does not exist in `fonoteka/plugins/golem15`.
- `vue-fonoteka-app/docs/pwa-push-seams.md`: "This phase adds zero runtime dependencies for push."
D-15 is locked, so port it, but there is no subscription store to test against (Open Question 3).
### Pitfall 15: Payload parity needs the Phase 12 serializer
The Album `created`/`updated` payload embeds `serializeAlbum(...)`. The Go `classes.SerializeAlbum` is a "minimal ... shape this phase's tests need" (`fonoteka.go/.../classes/serialize.go:5-6`). Only `deleted` (`id`, `collection_id`, `action`, `actor`, `timestamp`) and `collection.bulk_updated` (`reason`, `count`) can be proven byte-for-byte in Phase 11.
## Code Examples
### Transactional dispatch (conga)
```go
// Source: river client.go InsertTx @v0.47.0; gorm ConnPool pattern from riverqueue.com/docs/gorm
func (m *Manager) Dispatch(ctx context.Context, tx *gorm.DB, args pact.JobArgs, label string, o DispatchOpts) (int64, error) {
sqlTx, ok := tx.Statement.ConnPool.(*sql.Tx)
if !ok { return 0, errors.New("conga: Dispatch requires a transaction") }
row := JobRecord{Label: label, Status: StatusInProgress, ProgressMax: o.Count, Metadata: encodePHP(o.Metadata)} // status 1, "" when nil
if p, ok := bouncer.User(ctx); ok { row.UserID, row.IsAdmin = ptr(p.ID), p.Backend }
if err := tx.Create(&row).Error; err != nil { return 0, err }
opts := &river.InsertOpts{Queue: o.Queue, Metadata: mustJSON(map[string]any{"summer_job_id": row.ID})}
if o.Delay > 0 { opts.ScheduledAt = time.Now().Add(o.Delay) }
res, err := m.client.InsertTx(ctx, sqlTx, args.(river.JobArgs), opts)
if err != nil { return 0, err }
return row.ID, tx.Model(&row).Update("river_job_id", res.Job.ID).Error // see Open Question 2
}
```
### Subscribe proxy deny (parity: HTTP 200)
```go
// Source: ProxyController.php deny(): response()->json([...]) → 200
func deny(w http.ResponseWriter, log *slog.Logger, reason string, attrs ...any) {
log.Warn("Subscription denied", append([]any{"reason", reason}, attrs...)...)
wire.WriteJSON(w, http.StatusOK, map[string]any{"error": map[string]any{"code": 403, "message": "Access denied"}})
}
```
### Daily periodic schedule (wall-clock)
```go
// Source: river.PeriodicSchedule interface (periodic_job.go @v0.47.0)
type Daily struct{ Hour, Minute int; Loc *time.Location }
func (d Daily) Next(now time.Time) time.Time {
n := now.In(d.Loc)
t := time.Date(n.Year(), n.Month(), n.Day(), d.Hour, d.Minute, 0, 0, d.Loc)
if !t.After(n) { t = t.AddDate(0, 0, 1) }
return t
}
```
## State of the Art
| Old Approach | Current Approach | When Changed | Impact |
|--------------|------------------|--------------|--------|
| GORM + River: poll-only `riverdatabasesql.New`, or a second full pgx client for LISTEN | `riverdatabasesql.NewWithPgxListener(sqlDB, listenerPool)` | River v0.46.0, 2026-08-29 (CHANGELOG: "Added `riverdatabasesql.NewWithPgxListener` …", PR #1366) | PITFALLS.md Pitfall 14 and STACK.md predate it. The single-client design is now the official recommendation |
| River tables `river_client`, `river_client_queue` | Dropped; `river_notification` outbox added | River migration 007 | Schema-diff allow-lists must list the v7 set |
| Centrifugo `Authorization: apikey` only | `X-API-Key` recommended; `Authorization: apikey` still accepted | Centrifugo v4+ | Keep the PHP header for byte parity |
**Deprecated/outdated:**
- `gocent` in the RT-01 requirement text: superseded by D-12 (hand-rolled client). Note it in the traceability table so the verifier does not flag it.
## Assumptions Log
| # | Claim | Section | Risk if Wrong |
|---|-------|---------|---------------|
| A1 | `NewWithPgxListener` satisfies SC-1's "separate riverpgxv5 client" wording (the listener is a riverpgxv5 driver on its own pool, but not a separate River *client*) | Summary, Pattern 1 | The verifier rejects SC-1; the fallback is the two-client pattern |
| A2 | Laravel `schedule:run` minute-matching semantics for `--once` | Pattern 6 | Wrong cron behaviour for system-cron installs |
| A3 | firebase/php-jwt header order `{"typ":"JWT","alg":"HS256"}` vs golang-jwt `{"alg":"HS256","typ":"JWT"}`; Centrifugo does not care | Pattern 7 | Only affects byte-level token comparison. Compare decoded claims instead |
| A4 | Winter `SoftDelete` fires Laravel `deleted`/`restored` events, so Scout deletes on soft delete and re-adds on restore | Pattern 10 | Restore might not re-index in PHP. Go behaviour on restore needs a decision |
| A5 | River recovers worker panics and routes them to `ErrorHandler.HandlePanic` | Pitfall 3 | The wrapper should recover anyway |
| A6 | `webpush-go` legitimacy/maintenance | Alternatives | Not recommended; irrelevant if hand-rolled |
| A7 | Production Centrifugo is v6 (inferred from `.env.example` naming `client.token.hmac_secret_key`, `http_api.key`) | Pattern 8 | Header/proxy format is identical across v4–v6 per docs; low risk |
## Open Questions
1. **SC-1 wording vs `NewWithPgxListener`.**
- What we know: River now ships an official single-driver LISTEN split, and the River GORM docs recommend it.
- What's unclear: whether the user accepts it as meeting "a separate `riverpgxv5` client".
- Recommendation: use `NewWithPgxListener`, state the mapping in the plan, and keep the timed test identical.
2. **Where the River job id lives.**
- What we know: D-04 needs `JobCancel(riverID)`. The D-01 column list has no slot for it.
- Recommendation: add an internal nullable `river_job_id BIGINT` column. It is additive; cutover rows get NULL. The alternative is `JobList().Metadata('{"summer_job_id":N}')`, which is `[VERIFIED: job_list_params.go:362]` but slower.
3. **Web Push scope (D-15).**
- What we know: there is no PHP subscription store or dependency in Płytarium, and the Nuxt side has push seams only.
- Recommendation: port the `Pusher` interface, the VAPID driver (RFC 8291/8292, tested with the RFC vector), `websockets:generate-vapid-keys` (P-256 via `crypto/ecdh`, base64url) and `websockets:health`. `websockets:test-push` takes subscriptions through an app-provided `SubscriptionSource` interface and reports "no subscription source" when none is registered. Confirm with the user.
4. **tide capture mechanism (D-10 says "subscribe").**
- Recommendation: point PHP `CENTRIFUGO_API_URL` at a tide-owned fake Centrifugo HTTP recorder on loopback. The parity README already runs PHP with `QUEUE_CONNECTION=sync`, so `BroadcastEventJob` runs inline. Record `{path, Authorization-present, body}` as goldens with `timestamp`/`actor` normalised and the apikey redacted. No WebSocket dependency, and it captures channels. The Go side diffs its Centrifugo driver requests against a `httptest` fake the same way.
5. **Schedule entry for a command that ships in Phase 14.**
- Recommendation: the scheduler logs a warning and skips unknown commands (no boot failure), and a test asserts the warning. Alternatively, land the fonoteka schedule entry in Phase 14.
6. **ws-api ordering (Pitfall 11).** Confirm `jwt.auth` before `throttle:ws-api` on the token route. It only changes throttling of invalid-token spam.
7. **Payload parity for created/updated** (Pitfall 15). Record the goldens now and assert the `album` subtree in Phase 12.
## Environment Availability
| Dependency | Required By | Available | Version | Fallback |
|------------|------------|-----------|---------|----------|
| Go toolchain | all | ✓ | go1.27.0 | — |
| Docker (testcontainers) | River/Postgres tests | ✓ | 29.7.2 | `-short` skips DB tests |
| Local Postgres container `summercms-local-pg` | dev | ✓ | postgres:16 on 127.0.0.1:55432 | testcontainers `postgres:16-alpine` (existing harness, ICU pl-PL initdb) |
| Centrifugo binary | manual smoke / tide capture | ✓ | v6.5.1 (`/usr/local/bin/centrifugo`) | httptest fake Centrifugo (tests never need the real one) |
| Typesense | manual smoke | not probed | local script uses `typesense/typesense:26.0` | httptest fake Typesense |
| PHP stack (fonoteka) | D-10 goldens | ✓ (repo present) | Laravel Scout v10.25.0, typesense-php v4.9.3 | — |
| Go module proxy | `go get` River | ✓ | River v0.47.0 fetched | — |
**Missing dependencies with no fallback:** none.
## Validation Architecture
### Test Framework
| Property | Value |
|----------|-------|
| Framework | Go `testing` + testify (assertions) + testcontainers-go v0.44.0 postgres module |
| Config file | none. Pattern: package `TestMain` starting `postgres:16-alpine` with ICU pl-PL, skipped under `-short` (`modules/lagoon/postgres_test.go:25-100`) |
| Quick run command | `go test -short ./modules/conga/... ./modules/lighthouse/... ./modules/beachcomber/...` (package names discretionary) |
| Full suite command | `go vet ./... && go test ./...` in summercms.go, then `cd ../fonoteka.go && go vet ./... && go test ./...` |
### Phase Requirements → Test Map
| Req ID | Behavior | Test Type | Automated Command | File Exists? |
|--------|----------|-----------|-------------------|-------------|
| JOBS-01 | InsertTx from insert-only client picked up by NewWithPgxListener worker < 1s with 30s poll; poll-only negative control not picked up in 2s | integration (testcontainers) | `go test ./modules/conga -run TestListenPickupLatency -count=1` | ❌ Wave 0 |
| JOBS-01 | Dispatch writes row status 1 + River job in the same tx; rollback leaves neither | integration | `go test ./modules/conga -run TestDispatchTransactional` | ❌ |
| JOBS-01 | complete / fail (only on final attempt) / skip (`{skipped:true}`) / cancel (queued never runs; running ctx cancelled; is_canceled+STOPPED) queryable via manager | integration | `go test ./modules/conga -run 'TestOutcome|TestCancel'` | ❌ |
| JOBS-01 | `queue:clear` loops JobDeleteMany until 0, ignores running | integration | `go test ./modules/conga -run TestQueueClear` | ❌ |
| CLI-06 | `queue:work --queue broadcasts` only works that queue; `serve` with `queue.work_in_serve:false` starts no producer | integration | `go test ./modules/conga -run TestQueueWork` | ❌ |
| CLI-04 | `Daily`/`Every` Next() math incl. DST; periodic job enqueues on leader; `UniqueOpts.ByPeriod` dedupe; `schedule:run --once` runs minute-matched entries via `bonfire.Call` | unit + integration | `go test ./modules/conga -run 'TestSchedule'` ; `go test ./modules/bonfire -run TestCall` | ❌ |
| RT-01 | Token claims byte-for-byte per generator (decode + compare to PHP-shaped expectations); 503 when secret empty; 401 shape; `info.name` null | unit (httptest) | `go test ./modules/lighthouse/centrifugo -run TestToken` | ❌ |
| RT-01 | Publish/broadcast request path, `Authorization: apikey` header, body shape, 5s timeout, no request when api_key empty | unit (httptest fake Centrifugo) | `go test ./modules/lighthouse/centrifugo -run TestClient` | ❌ |
| RT-02 | Proxy: missing/wrong/empty secret deny (T-11-01); empty user; `presence:presence:`; >3 segments; unknown namespace; allow `{"result":{"info":[]}}`; presence overrides; deny is HTTP 200 generic body; reason logged, not returned | unit | `go test ./modules/lighthouse/centrifugo -run TestProxy` | ❌ |
| RT-02 | Collection/wishlist authorizer ports (member allowed, non-member denied, wishlist id under collection namespace denied, presence-prefixed denied) | integration (fonoteka.go) | `cd ../fonoteka.go && go test ./plugins/golem15/fonoteka -run TestWsAuthorizer` | ❌ |
| RT-03 | Broadcast enqueued in write tx; rollback → no publish; suppression of T only (other models still broadcast); SC-4 bulk: N album creates under WithoutBroadcasting + one Emit → exactly one `collection.bulk_updated` publish | integration | `go test ./modules/lighthouse -run 'TestBroadcastTx|TestSuppression|TestBulkEmitsOnce'` | ❌ |
| RT-03 | Callbacks fire when DB is published AFTER party.Activate (production order) | integration | `go test ./modules/lagoon -run TestOnDatabaseAfterActivate` | ❌ |
| RT-03 | Deleted + bulk payload goldens diff vs PHP capture (timestamp/actor normalised) | parity | `cd ../fonoteka.go && go test ./parity -run TestBroadcastGoldens` | ❌ |
| SRCH-01 | Upsert after commit (not on rollback); document has positive `collection_id`; kill-switch off / no api key / no DB → zero Typesense requests; Typesense 500 → write succeeds, warning logged; soft delete → DELETE | integration (httptest fake Typesense) | `go test ./modules/beachcomber/... -run TestSync` ; `cd ../fonoteka.go && go test ./plugins/golem15/fonoteka -run TestAlbumSearchable` | ❌ |
### Sampling Rate
- **Per task commit:** the quick `-short` command for the touched package, plus `go vet ./...`
- **Per wave merge:** the full suite in both repos (testcontainers)
- **Phase gate:** full suite green in summercms.go and fonoteka.go before `/gsd-verify-work`
### Wave 0 Gaps
- [ ] `modules/conga/postgres_test.go`: TestMain with testcontainers, exposing both `*sql.DB` and the DSN (the listener pool needs the DSN)
- [ ] A fake Centrifugo `httptest` helper (records method, path, headers, body) shared by lighthouse tests and tide
- [ ] A fake Typesense `httptest` helper
- [ ] fonoteka.go `parity/schema_diff_test.go` + `parity/migrate_test.go` allow-list updates (Pitfall 12)
- [ ] Framework install: `go get github.com/riverqueue/river@v0.47.0 github.com/riverqueue/river/riverdriver/riverdatabasesql@v0.47.0 github.com/riverqueue/river/rivertype@v0.47.0`
## Security Domain
### Applicable ASVS Categories
| ASVS Category | Applies | Standard Control |
|---------------|---------|-----------------|
| V2 Authentication | yes | Token route behind the P7 `jwt.auth` guard; the proxy authenticates Centrifugo by shared secret (`crypto/subtle`) |
| V3 Session Management | yes | Centrifugo connection JWT TTL (`token_ttl` 3600), anonymous 300 s |
| V4 Access Control | yes | The authorizer registry re-validates on every subscribe; default deny for unknown namespaces |
| V5 Input Validation | yes | Channel parsing (presence prefix, segment count), `user` parsed as uint, the proxy request body size-limited (surf body limit) |
| V6 Cryptography | yes | HS256 via golang-jwt; VAPID ES256; RFC 8291 via stdlib primitives only |
| V7 Error handling/logging | yes | Deny reasons logged, never returned; never log `api_key`, `token_secret`, `proxy_secret`, JWTs or VAPID private key |
### Known Threat Patterns
| Pattern | STRIDE | Standard Mitigation |
|---------|--------|---------------------|
| T-11-01 Forged subscribe-proxy call bypassing Centrifugo | Spoofing | Constant-time `X-Centrifugo-Secret`; an empty configured secret denies |
| T-11-02 Channel-name tricks (`presence:presence:`, extra segments, namespace confusion, wishlist id under `collection:`) | Elevation | Port parseChannel exactly; authorizers check `kind`; tests per trick |
| T-11-03 Denial-reason leakage | Info disclosure | Generic 403 body; full context only in logs |
| T-11-04 PII in connection token `info` (WS-004) | Info disclosure | `info` limited to `name`; test asserts the exact claim set |
| T-11-05 Broadcast of rolled-back writes / wrong audience | Info disclosure | Enqueue in the write tx; channels from the model (`kind=collection` only) |
| T-11-06 Secrets in logs (api key header, tokens, River args) | Info disclosure | Log driver/debug output redacts headers; broadcast payloads never contain credentials |
| T-11-07 Cross-collection search leak via stale/mis-scoped index | Info disclosure | `collection_id` in every document (required > 0); the Phase 12 SQL re-gate |
| T-11-08 Job cancel/progress IDOR on `summer_jobs` ids | Elevation | The framework exposes no HTTP route; Phase 13 endpoints scope by the owning import (`findVisible`) |
| T-11-09 Scheduled command injection | Tampering | Commands only from compiled `HasSchedule`; no config/DB-driven command names |
| T-11-10 Subscribe/token flooding | DoS | `ws-api` 120/min bucket; proxy body limit |
| T-11-11 VAPID private key exposure | Info disclosure | `generate-vapid-keys` prints the private key only on explicit request; `--show-current` truncates as PHP does |
## Sources
### Primary (HIGH confidence)
- River source @v0.47.0 read from the module cache: `client.go` (Insert/InsertTx/JobCancel/JobDeleteMany/notify, Config, defaults), `riverdriver/riverdatabasesql/river_database_sql_driver.go` (`New`, `NewWithPgxListener`, `SupportsListenNotify`), `periodic_job.go`, `worker.go`, `delete_many_params.go`, `rivermigrate/river_migrate.go`, `rivertype`, `migration/main/*.sql`, `CHANGELOG.md`, `go.mod` files
- Go module proxy: `go list -m -versions github.com/riverqueue/river` → v0.47.0 (2026-08-31)
- GORM v1.31.2 source: `finisher_api.go` Begin, `callbacks/transaction.go`, `callbacks/callbacks.go`
- PHP source (fonoteka): apparatus `JobManager.php`, `JobStatus.php`, `create_jobs_table.php`, `QueueClearCommand.php`; websockets `routes.php`, `config/config.php`, `JwtTokenGenerator.php`, `CentrifugoClient.php`, `AuthorizerRegistry.php`, `AuthorizationResult.php`, `ProxyController.php`, `BroadcastableModel.php`, `BroadcastEventJob.php`, `CentrifugoBroadcaster.php`, console commands; fonoteka `Plugin.php`, `Album.php`, `AlbumApiController.php`, `CsvImportApiController.php`, `WishlistDigestJob.php`, `WishlistDigestQueue.php`, ws authorizers, `ReindexAlbums.php`; `config/scout.php`; vendor Scout `TypesenseEngine.php`, `ModelObserver.php`
- In-repo Go: `modules/lagoon/{connection,migrations,commands,lifecycle}.go`, `modules/pact/capabilities.go`, `modules/postcard/{drivers,mailer}.go`, `modules/surf/serve.go`, `modules/party/registry.go`, `modules/bouncer/context.go`, `modules/wire/response.go`, `modules/bonfire/command.go`, `internal/build/{build.go,stubs/*.tmpl}`, `cmd/summer/{main,runtime}.go`; fonoteka.go `plugin.go`, `routes.go`, `models/{album,settings}.go`, `classes/{album_write_service,serialize}.go`, `parity/{manifest.yaml,schema_diff_test.go,migrate_test.go}`
- https://riverqueue.com/docs/gorm: `NewWithPgxListener` recommendation, `tx.Statement.ConnPool.(*sql.Tx)`
### Secondary (MEDIUM confidence)
- https://centrifugal.dev/docs/server/server_api: `/api/<method>`, `X-API-Key`, 200 + error body
- https://raw.githubusercontent.com/centrifugal/centrifugo/master/internal/middleware/auth.go: `Authorization: apikey` still accepted
- https://centrifugal.dev/docs/server/proxy: subscribe proxy request/response, error codes 400–1999, non-200 → code 100
- https://typesense.org/docs/26.0/api/documents.html: import JSONL/upsert/200, delete, `X-TYPESENSE-API-KEY`, search params
### Tertiary (LOW confidence)
- Laravel `schedule:run` semantics, firebase/php-jwt header order, webpush-go status (training knowledge, tagged `[ASSUMED]`)
## Metadata
**Confidence breakdown:**
- Standard stack: HIGH. River version and APIs read from source and the proxy; the official docs agree.
- Architecture: MEDIUM-HIGH. The seams (`OnDatabase`, after-commit) follow from verified GORM behaviour. Exact API shapes are design recommendations.
- Pitfalls: HIGH. Nearly all are cited to file lines read this session.
- Contract values (statuses, claims, bodies): HIGH. Quoted verbatim.
**Research date:** 2026-09-29
**Valid until:** 2026-10-29. River releases often (0.44→0.47 in two weeks); re-check the River version at execution time.
## Suggested plan split (for the lean-mode plan-count checkpoint)
1. **conga core** (summercms.go + fonoteka.go test allow-lists): River dependency, `lagoon.OnDatabase` + `AfterCommit`/`Transaction` seams, River v7 + `summer_jobs` migrations, job manager, `conga.Job[T]`, worker in `serve`, `queue:work`, `queue:clear`, generated-main and `summer` delegate changes, `make:job` stub, timed LISTEN smoke test.
2. **Scheduler**: `pact.HasSchedule`, `Daily`/`Every`, periodic wiring, `schedule:run` (+`--once`), `bonfire.Call`, fonoteka schedule entry.
3. **Realtime core + Centrifugo driver**: Publisher/Registry/Route/Surface/Mount, memory/log/null drivers, Centrifugo client, token issuer, proxy, Broadcastable callbacks, suppression, `Emit`, BroadcastWorker. fonoteka.go: authorizers, `Mount`, `ws-api` bucket, Album Broadcastable, PROJECT.md D-16 note.
4. **Web Push + websockets console commands** (Open Question 3 scope).
5. **Search**: Searchable/Engine/Gate, Typesense driver, after-commit sync; fonoteka Album Searchable + settings gate.
6. **tide Centrifugo capture + broadcast goldens** (deleted, bulk; created/updated recorded for Phase 12).
7. **Unit tests** (always last): full coverage for all of the above, including the T-11-xx failing-when-broken tests.