17 KiB
Phase 11: Jobs, realtime and search infrastructure - Context
Gathered: 2026-09-28 Status: Ready for planning
## Phase BoundaryThe background infrastructure later API phases depend on:
- River on the dual-driver split, plus a job manager service with queryable outcomes.
summer queue:workandsummer 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
reindexcommand and the SQL re-gate on search results. Those are Phases 12–14. - Any Nuxt or fonoteka-mcp change.
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 issummer_jobs, with the PHPgolem15_apparatus_jobscolumns: integer auto-incrementid,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_idin the CSV API expose these ids, and the Phase 15 cutover copies rows fromgolem15_apparatus_jobswith ids preserved. - D-02:
Dispatchinserts thesummer_jobsrow and enqueues the River job in the same transaction (riverdatabasesqlon the shared*sql.DB). River only executes. The row is what API responses, progress and cancellation read. The service API mirrors PHPJobManager: 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}, asWishlistDigestJobdoes. 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.
- "Skip" is COMPLETE with metadata
- D-04: Cancellation does both.
CancelJobsetsis_canceled+ STOPPED (PHP semantics the CSV API relies on) and also calls RiverJobCancel. A queued job then never starts, and a running job's ctx is cancelled.- Long jobs still call
CheckIfCanceledbetween items, as PHP does.
- D-05:
queue:clear(the ApparatusQueueClearCommandequivalent, 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
broadcastsqueue. 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 queuedBroadcastEventJob. - 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.
- MaxAttempts 1, timeout from
- 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-classwithoutBroadcasting. 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.
- The bulk pattern (suppress, then emit exactly one
- 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}, withttldefaulting to 60. - A delete snapshot is captured before deletion.
shouldBroadcast(action)and a model-supplied channels list.- Single channel →
publish, several →broadcast. broadcast_namespaceprefixing.
- Event name
- D-10: Payload parity is proven against real PHP output. Phase 2
tidetooling 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 thememorydriver or a fake Centrifugo, is diffed against them withtimestampandactornormalised. 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
postcardmail-driver pattern. It owns:- the
Publisherinterface (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 singlepresence: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. - the
-
D-12: The Centrifugo driver is a sub-package. It contains:
- A hand-rolled
net/httpHTTP API client: publish, broadcast and presence, with request bytes and theapi_keyheader matching the PHPCentrifugoClient. 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 ininfo) is kept. - The subscribe-proxy adapter, which ports
ProxyController:- constant-time
X-Centrifugo-Secretcheck; - 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"}}.
- constant-time
- A hand-rolled
-
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}, withSurface∈ {UserAuth,ServerToServer,Public}. fonoteka.gocallsrealtime.Mount(driver, realtime.Surfaces{UserAuth: <P6 JWT group>, ServerToServer: <raw group>, Middleware: throttle("ws-api")})once. Switching drivers never editsroutes.go, and the guard, group and bucket stay app-owned.- The Centrifugo driver declares
GET /api/realtime/token(UserAuth) andPOST /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. memoryandnulldeclare no routes. Płytarium and its tests always run the Centrifugo driver, pointed at a fake Centrifugo overhttptestwhere needed.
- The neutral package defines
-
D-14: The
ws-apibucket (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 PHPpush.*block (enabled,public_key,private_key,subject). -
D-16: fonoteka registers its
collectionandwishlistauthorizers (ports ofCollectionChannelAuthorizerandWishlistChannelAuthorizer) into the framework registry.- The
fonoteka.gowebsocketsplugin shrinks to config binding and theMountcall, 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.
- The
Workers and scheduler
- D-17:
summer servestarts the River work client (its ownriverpgxv5pool for LISTEN/NOTIFY) in-process by default.queue.work_in_serve: falsedisables it when a separatesummer queue:workprocess runs, with--queuefilters (e.g.broadcasts). Enqueue always goes through the shared*sql.DBdriver. 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 inpact/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:runis a foreground scheduler-only process, and--onceruns due commands immediately for system-cron setups. - fonoteka's
fonoteka:prune-notificationsdaily schedule is the first consumer. The command itself is ported in Phase 14.
- Plugins declare schedules through
Search sync
- D-19: Search is pluggable as well.
- The framework
Searchablemodel interface covers index name (searchableAs), document (toSearchableArray) andshouldBeSearchable. - The engine driver interface covers upsert, delete, flush and search-IDs.
search.driver: typesense | null. The Typesense driver is hand-rollednet/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.
- The framework
- 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 whensearch_use_typesenseis off, the Typesense config/API key is absent, or there is no DB (fresh install), matching PHPPlugin.phpdisableSearchSyncingandReindexAlbums::typesenseIsConfigured.
Claude's Discretion
- Package names and layout, following the beach-themed naming. ARCHITECTURE.md suggests
congafor 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
actormetadata 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
tideCentrifugo 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_jobsdecision.planning/PROJECT.md: constraints (two repos; the websockets-plugin placement is revised by D-16)
Research
.planning/research/PITFALLS.mdPitfall 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/ projectCLAUDE.md§ "River + GORM: share one *sql.DB";riverqueue.com/docs/gorm.planning/research/ARCHITECTURE.md:congarow,HasScheduleexample, 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,disableSearchSyncinggate;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_updatedusage/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.mdD-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:Openreturns the shared*sql.DB(pgx stdlib) plus GORM on it. Its comment already reserves a separate listener pool, so Riverriverdatabasesqlbinds here.postcard/mailer.go:driverFromApp/mail.driverselection (memory, log, smtp), the template forrealtime.driverandsearch.driver.pact/capabilities.go:HasScheduleis 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 todofetchguard-guarded-http-client.md.- Phase 7 JWT group / guard: the
UserAuthsurface for the token route. The fonotekasettingsmodel (P5/P9) holdssearch_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-xxthreats. - Unit tests are the last plan.
go vetandgo test ./...stay green at every commit.
Integration Points
summer serveboot (worker + scheduler start), thesummerCLI (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.yamlentries for/api/realtime/tokenand/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.
- 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_jobswith 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