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

17 KiB
Raw Blame History

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: <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).

<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>

## 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