diff --git a/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-CONTEXT.md b/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-CONTEXT.md new file mode 100644 index 0000000..0c86460 --- /dev/null +++ b/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-CONTEXT.md @@ -0,0 +1,190 @@ +# Phase 11: Jobs, realtime and search infrastructure - Context + +**Gathered:** 2026-09-28 +**Status:** Ready for planning + + +## 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. + + + + +## 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: , ServerToServer: , 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). + + + + +## 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 + + + + +## 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`. + + + + +## 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. + + + + +## 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. + + + +--- + +*Phase: 11-jobs-realtime-and-search-infrastructure* +*Context gathered: 2026-09-28* diff --git a/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-DISCUSSION-LOG.md b/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-DISCUSSION-LOG.md new file mode 100644 index 0000000..20acdd9 --- /dev/null +++ b/.planning/phases/11-jobs-realtime-and-search-infrastructure/11-DISCUSSION-LOG.md @@ -0,0 +1,123 @@ +# Phase 11: Jobs, realtime and search infrastructure - Discussion Log + +> **Audit trail only.** Do not use as input to planning, research, or execution agents. +> Decisions are captured in CONTEXT.md — this log preserves the alternatives considered. + +**Date:** 2026-09-28 +**Phase:** 11-jobs-realtime-and-search-infrastructure +**Areas discussed:** Job manager & job records, Broadcast delivery path, Realtime package split, Worker & scheduler model (incl. search sync) + +--- + +## Job manager & job records + +The first question set was interrupted; the user ran `/gsd-explore` on whether to port Apparatus as a plugin or fold it into the framework. Outcome: `.planning/notes/apparatus-dissolved-into-framework.md`. + +| Option | Description | Selected | +|--------|-------------|----------| +| (a) keep `golem15_apparatus_jobs` | Straight copy at cutover; framework carries a plugin prefix | | +| (b) `summer_jobs` | Framework-owned name, same columns and int ids, rows copied at cutover | ✓ | +| (c) configurable table name | Per-app name | | + +| Option | Description | Selected | +|--------|-------------|----------| +| Mirror PHP exactly | Same integers; skip = COMPLETE + {skipped:true}; ERROR only on River discard | ✓ | +| Add SKIPPED status | New integer 5 | | + +| Option | Description | Selected | +|--------|-------------|----------| +| Both | is_canceled + STOPPED and River JobCancel; long jobs still poll | ✓ | +| Poll is_canceled only | Literal PHP port | | +| River cancel + ctx only | Go-idiomatic | | + +**User's choice:** Dissolve Apparatus into the framework; `summer_jobs`; mirror PHP outcomes; cancel via both. + +--- + +## Broadcast delivery path + +| Option | Description | Selected | +|--------|-------------|----------| +| River job in the write tx | Published only after commit; `broadcasts` queue | ✓ | +| Direct publish after commit | Goroutine after-commit callback | | +| Synchronous publish | Inline in request | | + +| Option | Description | Selected | +|--------|-------------|----------| +| One attempt, best-effort | tries=1, 5 s timeout, warn on failure | ✓ | +| Small retry | River backoff, e.g. 3 attempts | | + +| Option | Description | Selected | +|--------|-------------|----------| +| Per model type, ctx-scoped | `WithoutBroadcasting[T](ctx, fn)` | ✓ | +| All broadcasts in ctx | One switch | | + +| Option | Description | Selected | +|--------|-------------|----------| +| Capture from PHP Centrifugo | tide subscribes during recorded flows; goldens | ✓ | +| Goldens derived from PHP code | Hand-written expected payloads | | + +--- + +## Realtime package split + +| Option | Description | Selected | +|--------|-------------|----------| +| Framework package + thin plugin | Generic realtime package in summercms.go | | +| All in fonoteka.go plugin | Wholesale app port | | +| Framework-bundled plugin | First-party plugin in summercms.go | | + +**User's choice:** Free text: concerned about hard dependency on Centrifugo; wants transports pluggable so a project can replace it. Claude proposed a transport-neutral realtime package with drivers (centrifugo, memory, log, null), a Centrifugo sub-package, and Web Push behind its own interface. User: "yep that's perfect." + +| Option | Description | Selected | +|--------|-------------|----------| +| Hand-rolled net/http | Stdlib-first, exact request bytes | ✓ | +| gocent/v3 | Official client | | + +| Option | Description | Selected | +|--------|-------------|----------| +| Health check only | Defer push/VAPID | | +| Port all of it | Web Push + both VAPID commands + health check | ✓ | +| None | | | + +| Option | Description | Selected | +|--------|-------------|----------| +| Only what's called | Connection + subscription tokens | | +| Full parity with PHP class | All generator variants | ✓ | + +Route mounting: Claude first suggested app-mounted handlers (wristband precedent). User asked: "won't that be a problem if we implement multiple drivers in future?" Claude refined it: drivers declare `Routes()` by surface (UserAuth / ServerToServer / Public) and the app maps surfaces to its groups and buckets through `realtime.Mount`. **User's choice:** Lock it. + +--- + +## Worker & scheduler model + +| Option | Description | Selected | +|--------|-------------|----------| +| In serve by default + queue:work | `queue.work_in_serve` toggle | ✓ | +| Only queue:work | Separate process always | | + +| Option | Description | Selected | +|--------|-------------|----------| +| River periodic jobs | HasSchedule → periodic jobs, leader-elected; `schedule:run` foreground + `--once` | ✓ | +| Laravel-style schedule:run | System cron every minute, own locking | | + +| Option | Description | Selected | +|--------|-------------|----------| +| Driver interface, Typesense driver | Searchable interface + engine driver; hand-rolled Typesense | ✓ | +| Typesense client only | Direct port, no abstraction | | + +| Option | Description | Selected | +|--------|-------------|----------| +| After commit, inline, non-fatal | Like Scout queue=false | ✓ | +| River job in the write tx | Retryable, lags the write | | + +--- + +## Claude's Discretion + +- Package names and layout; ctx-to-GORM-hook plumbing for suppression; actor capture; River client config and timed-test harness; whether the Phase 8 expiry sweep moves to a periodic job; tide Centrifugo capture mechanics. + +## Deferred Ideas + +- The admin extension point ("Phase 10.1") is not on the roadmap. It stays deferred, and the references that call it a phase need rewording. +- More realtime/search drivers; a Jobs admin screen; backend admin API tokens; guarded HTTP client and redacting slog handler (pending todos).