Files
summercms/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-CONTEXT.md
2026-09-28 00:54:36 +02:00

191 lines
17 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# Phase 11: Jobs, realtime and search infrastructure - Context
**Gathered:** 2026-09-28
**Status:** Ready for planning
<domain>
## Phase Boundary
The background infrastructure later API phases depend on:
- River on the dual-driver split, plus a job manager service with queryable outcomes.
- `summer queue:work` and `summer schedule:run`.
- Realtime: token issuing, publishing, subscribe-time channel authorization, and broadcastable models with bulk suppression.
- Typesense album sync behind the settings kill-switch.
Requirements: JOBS-01, CLI-04, CLI-06, RT-01, RT-02, RT-03, SRCH-01. Repos: summercms.go (framework packages) and fonoteka.go (wiring, authorizers, Album searchable/broadcastable bindings).
Not in this phase:
- Domain jobs (CSV import/match, wishlist digest), the `reindex` command and the SQL re-gate on search results. Those are Phases 12–14.
- Any Nuxt or fonoteka-mcp change.
</domain>
<decisions>
## Implementation 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).
</decisions>
<canonical_refs>
## Canonical References
**Downstream agents MUST read these before planning or implementing.**
### Scope and decisions
- `.planning/ROADMAP.md` § Phase 11: goal and success criteria
- `.planning/REQUIREMENTS.md`: JOBS-01, CLI-04, CLI-06, RT-01, RT-02, RT-03, SRCH-01
- `.planning/notes/apparatus-dissolved-into-framework.md`: Apparatus triage, `summer_jobs` decision
- `.planning/PROJECT.md`: constraints (two repos; the websockets-plugin placement is revised by D-16)
### Research
- `.planning/research/PITFALLS.md` Pitfall 14 (River/GORM driver split), Pitfall 15 (Typesense as pre-filter), "Centrifugo broadcast-suppression and channel-name contract", opaque-tenancy note on channel names
- `.planning/research/STACK.md` / project `CLAUDE.md` § "River + GORM: share one *sql.DB"; `riverqueue.com/docs/gorm`
- `.planning/research/ARCHITECTURE.md`: `conga` row, `HasSchedule` example, job flow diagram
- `.planning/research/SUMMARY.md`: flags on River dual-driver and gocent/typesense-go versions
### PHP reference (contract source)
- `/media/nvme/dev/golem15/fonoteka/plugins/golem15/apparatus/classes/JobManager.php`, `contracts/JobStatus.php`, `contracts/ApparatusQueueJob.php`, `updates/create_jobs_table.php`
- `/media/nvme/dev/golem15/fonoteka/plugins/golem15/websockets/`: `routes.php`, `config/config.php`, `classes/{JwtTokenGenerator,CentrifugoClient,AuthorizerRegistry,AuthorizationResult,CentrifugoBroadcaster}.php`, `http/controllers/ProxyController.php`, `traits/BroadcastableModel.php`, `jobs/BroadcastEventJob.php`, `contracts/*`, `console/*`, `tests/security/*`
- `/media/nvme/dev/golem15/fonoteka/plugins/golem15/fonoteka/Plugin.php`: `registerWsAuthorizers`, `registerSchedule`, `disableSearchSyncing` gate; `classes/ws/{Collection,Wishlist}ChannelAuthorizer.php`
- `/media/nvme/dev/golem15/fonoteka/plugins/golem15/fonoteka/models/Album.php`: `searchableAs`, `shouldBeSearchable`, `toSearchableArray`, broadcast overrides
- `/media/nvme/dev/golem15/fonoteka/plugins/golem15/fonoteka/controllers/api/AlbumApiController.php`: `withoutBroadcasting` + `collection.bulk_updated` usage
- `/media/nvme/dev/golem15/fonoteka/plugins/golem15/fonoteka/console/ReindexAlbums.php`: Typesense configured check
- `/media/nvme/dev/golem15/fonoteka/config/scout.php`: `queue=false`, `after_commit=false`
- `/media/nvme/dev/golem15/fonoteka/vue-fonoteka-app/app/composables/useCentrifugo.ts`: client-side contract
### Prior phase decisions that apply
- `.planning/phases/06-http-routing-auth-groups-and-rate-limiting/06-CONTEXT.md`: auth groups, throttle buckets, route-surface isolation
- `.planning/phases/08-oauth2-1-authorization-server/08-CONTEXT.md` D-09 (app-mounted framework handlers precedent), D-17 (expiry sweep)
- `.planning/phases/02-api-parity-harness-bootstrap/02-CONTEXT.md`: tide capture rules
</canonical_refs>
<code_context>
## Existing Code Insights
### Reusable Assets
- `lagoon/connection.go`: `Open` returns the shared `*sql.DB` (pgx stdlib) plus GORM on it. Its comment already reserves a separate listener pool, so River `riverdatabasesql` binds here.
- `postcard/mailer.go`: `driverFromApp` / `mail.driver` selection (memory, log, smtp), the template for `realtime.driver` and `search.driver`.
- `pact/capabilities.go`: `HasSchedule` is reserved for its first consumer. Optional capabilities are type-asserted by the kernel.
- `fetchguard`: dial-time private-IP guarded client. Not required for Centrifugo/Typesense (operator-configured internal hosts), but see the pending todo `fetchguard-guarded-http-client.md`.
- Phase 7 JWT group / guard: the `UserAuth` surface for the token route. The fonoteka `settings` model (P5/P9) holds `search_use_typesense`.
### Established Patterns
- Per-request state only in `context.Context`, never package globals (Phase 1).
- Framework packages are app-agnostic, and apps mount framework handlers explicitly (wristband, Phase 8).
- Bonfire commands are scaffolded the Phase 4 way. Security-relevant phases carry threat IDs and failing-when-broken tests (Phases 3, 6, 8). The subscribe proxy and token issuing warrant `T-11-xx` threats.
- Unit tests are the last plan. `go vet` and `go test ./...` stay green at every commit.
### Integration Points
- `summer serve` boot (worker + scheduler start), the `summer` CLI (`queue:work`, `schedule:run`, `queue:clear`, the websockets commands).
- fonoteka.go `routes.go` (`realtime.Mount`), Album model (Searchable + Broadcastable), fonoteka plugin boot (authorizer registration, schedule).
- The `parity/manifest.yaml` entries for `/api/realtime/token` and `/api/realtime/subscribe`.
</code_context>
<specifics>
## Specific Ideas
- The user wants transports to be pluggable, so a project can replace Centrifugo, without being locked into one vendor. Pluggability is server-side: swapping the driver also needs a matching frontend client, which is accepted (Płytarium stays on Centrifugo).
- A future realtime driver might be a Pusher protocol server (Soketi), Mercure, or a built-in Go WebSocket hub. That's why D-13 declares routes by surface rather than hardcoding Centrifugo's two routes.
</specifics>
<deferred>
## Deferred Ideas
- **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.
</deferred>
---
*Phase: 11-jobs-realtime-and-search-infrastructure*
*Context gathered: 2026-09-28*