52 KiB
phase, plan, type, wave, depends_on, files_modified, autonomous, requirements, estimate, must_haves
| phase | plan | type | wave | depends_on | files_modified | autonomous | requirements | estimate | must_haves | ||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| 11-jobs-realtime-and-search-infrastructure | 01 | execute | 1 |
|
true |
|
|
|
Phase Goal
ROADMAP Phase 11 goal (verbatim; it is not in user-story form, see the planner return note): River jobs run on the correct dual-driver split, Centrifugo publishing and channel authorization match the existing server, and Typesense sync stays a re-gated pre-filter — all brought up before the API phases that depend on them.
This plan's slice: a plugin can dispatch a job inside its write transaction, a worker running in summer serve or summer queue:work picks it up through LISTEN/NOTIFY, and the job's outcome is queryable from summer_jobs (JOBS-01, CLI-06).
Purpose: CSV import (Phase 13/14), broadcasts (plan 11-03) and the scheduler (plan 11-02) all run on this. Decisions implemented: D-01, D-02, D-03, D-04, D-05, D-17; user decisions 2 (NewWithPgxListener, one client) and 4 (river_job_id column); RESEARCH Patterns 1-4 and 10, Pitfalls 1-5, 12, 13.
Output: River dependency, framework migrations (River v7 + summer_jobs), modules/conga with README, lagoon OnDatabase/Transaction/AfterCommit, worker in serve, queue:work/queue:clear, generated-main and summer delegate changes, the make:job stub, the fonoteka.go allow-list and Boot-order fixes.
Repos: summercms.go (framework) and fonoteka.go (go.mod/go.sum tidy, config, allow-lists, Boot seam adoption, regenerated main.go). Commit each repository separately; fonoteka.go changes that keep its tests green land in the same task as the framework change that needs them. Planning docs and code in separate commits. Never add co-author tags.
<execution_context>
@/.claude/gsd-core/workflows/execute-plan.md
@/.claude/gsd-core/templates/summary.md
</execution_context>
Artifacts this phase produces
(This plan's share of the phase artifacts.)
- Dependencies:
github.com/riverqueue/river v0.47.0,github.com/riverqueue/river/riverdriver/riverdatabasesql v0.47.0,github.com/riverqueue/river/rivertype v0.47.0(plus the transitive riverdriver, riverpgxv5, rivershared modules River pulls in). - Tables:
summer_jobs(with internalriver_job_id BIGINT NULL), River v7 tables (river_job,river_leader,river_queue,river_migrationand whatever else version 7 leaves; the executor lists the exact set frompg_tablesafter migrating), history tablesummer_migrations_summercms_conga. - lagoon:
QueueMigrations(sqlDB *sql.DB) []*gormigrate.Migration, constQueueHistoryID = "summercms.conga", constJobsTable = "summer_jobs", constRiverSchemaVersion = 7,OnDatabase(app, fn func(*sql.DB, *gorm.DB) error) error,Transaction(ctx, gdb, fn func(ctx context.Context, tx *gorm.DB) error) error,AfterCommit(ctx, db *gorm.DB, fn func(ctx context.Context, db *gorm.DB)), GORM callback namelagoon:after_commit. - conga:
Manager,From(app) (*Manager, error),(*Manager).Register(jobs ...pact.Job) error,Dispatch(ctx, db, args, DispatchOpts) (uint, error),Enqueue(ctx, db, args, EnqueueOpts) error,StartJob,UpdateJobState,UpdateMetadata,CompleteJob,FailJob,CancelJob,StopJob,CheckIfCanceled,GetMetadata,Get;Record(summer_jobs model);StatuswithStatusInQueue,StatusInProgress,StatusComplete,StatusError,StatusStopped;DispatchOpts{Label, Count, Metadata, Queue, Delay, MaxAttempts};EnqueueOpts{Queue, Delay, MaxAttempts};Job[T pact.JobArgs](fn func(context.Context, T) error, opts ...JobOption) pact.Job;JobOption,OnQueue,MaxAttempts,Timeout;JobID(ctx) (uint, bool);WorkerOptions{Queues []string};StartWorker(ctx, app, plugins, WorkerOptions) (*Worker, error);StartServeWorker(ctx, app, plugins) (*Worker, error);(*Worker).Stop(ctx) error;RuntimeCommands(app, plugins) []bonfire.Command; errorsErrNoDatabase,ErrRegistrationClosed,ErrNotCongaJob,ErrUnknownQueue. - CLI: app commands
queue:work [--queue name ...],queue:clear [queue];summerdelegatesqueue:work(forwarding repeatable--queue) andqueue:clear. - Config keys (app-level
queue.*):queue.work_in_serve(default true),queue.max_attempts(3),queue.job_timeout(300 seconds, int or duration string),queue.queues.<name>(MaxWorkers; default queuedefault: 4). - Files:
modules/conga/*,../fonoteka.go/config/queue.yaml.
Flagged assumptions (edge probe: unclassified)
- JOBS-01 and CLI-06 came back
unclassifiedfrom the spec-less edge probe and stayunresolved; they are surfaced here, not auto-resolved. Planner reading for manual review: (a) concurrency — two workers never run the same River job concurrently (River row locking) and two Dispatch calls produce distinct summer_jobs ids; (b) idempotency — Dispatch is not idempotent (each call is a new job, like PHP); (c) empty —queue:workwith no registered jobs still starts and idles;queue:clearon an empty queue printsCleared 0 jobs. Confirm or correct at verify time.
(2) Framework migrations (D-01, user decision 4, Pitfall 13): new modules/lagoon/queue_migrations.go with consts QueueHistoryID = "summercms.conga", JobsTable = "summer_jobs", RiverSchemaVersion = 7 and func QueueMigrations(sqlDB *sql.DB) []*gormigrate.Migration returning two migrations. 202609290001_river_schema: Migrate builds rivermigrate.New(riverdatabasesql.New(sqlDB), nil) and calls Migrate(ctx, rivermigrate.DirectionUp, &rivermigrate.MigrateOpts{TargetVersion: RiverSchemaVersion}) on the shared pool (per-migration transactions; the Tx variant is deprecated upstream because some River migrations cannot share one transaction); Rollback calls DirectionDown with TargetVersion -1. 202609290002_summer_jobs: CREATE TABLE IF NOT EXISTS summer_jobs (id SERIAL PRIMARY KEY, label VARCHAR(255) NOT NULL, status INTEGER NOT NULL DEFAULT 0, progress INTEGER NOT NULL DEFAULT 0, progress_max INTEGER NOT NULL DEFAULT 0, user_id INTEGER NULL, is_admin BOOLEAN NOT NULL DEFAULT FALSE, is_canceled BOOLEAN NOT NULL DEFAULT FALSE, metadata TEXT NOT NULL, river_job_id BIGINT NULL, created_at TIMESTAMPTZ NULL, updated_at TIMESTAMPTZ NULL) (Laravel timestamps() are nullable; no extra indexes, matching PHP); Rollback drops the table. In migrations.go Migrate runs the set with migrator(gdb, QueueHistoryID, QueueMigrations(sqlDB)) right after the cabana set, where sqlDB comes from gdb.DB(); wrap errors as lagoon: migrate queue: %w. The set lives in lagoon for the same reason BackendAdminMigrations does (lagoon cannot import conga), and every app gets it on migrate.
(3) conga core, new package modules/conga (errors prefixed conga: , logger from app.Lookup[*slog.Logger]() else slog.Default() like postcard):
- record.go:
type Status intwith the five PHP constants;type Record structmapping every summer_jobs column (ID uint, Label string, Status Status, Progress int, ProgressMax int, UserID *uint, IsAdmin bool, IsCanceled bool, Metadata string, RiverJobID *int64, CreatedAt/UpdatedAt *time.Time),TableName() string { return lagoon.JobsTable }. - job.go:
type JobOption func(*jobConfig),OnQueue(name string),MaxAttempts(n int),Timeout(d time.Duration);func Job[T pact.JobArgs](fn func(ctx context.Context, args T) error, opts ...JobOption) pact.Jobreturning an unexported*typedJob[T]that implements pact.Job (its Work type-asserts args to T and calls fn) plus unexportedkind() string,register(*river.Workers) error(callsriver.AddWorkerSafely[T]with an unexported riverWorker[T] embeddingriver.WorkerDefaults[T]) andconfig() jobConfig.JobID(ctx) (uint, bool)reads the summer_jobs id the wrapper stored in ctx. - conga.go:
type Manager structholding the app, a mutex, the registered jobs by kind, the lazily built insert-only client and the running worker client;func From(app *backpack.App) (*Manager, error)does lookup-or-publish of*Manageron the app (on a publish race, look up again);Register(jobs ...pact.Job) errorreturnsErrNotCongaJobfor a job not built by Job, an error on a duplicate kind, andErrRegistrationClosedonce a client has been built;Dispatch(ctx, db *gorm.DB, args pact.JobArgs, o DispatchOpts) (uint, error): when db is not inside a transaction (db.Statement.ConnPoolis not a *sql.Tx) wrap the work indb.WithContext(ctx).Transaction; insert the Record with Status StatusInProgress (Pitfall 1), Label (required, error when empty), ProgressMax = o.Count, Metadata = JSON of o.Metadata or the two characters""when nil (PHP json_encode of ''), UserID/IsAdmin frombouncer.User(ctx)(IsAdmin = Principal.Backend), CreatedAt = UpdatedAt = now; thenclient.InsertTx(ctx, sqlTx, args, &river.InsertOpts{Queue, MaxAttempts, ScheduledAt: now+o.Delay when Delay > 0, Metadata: {"summer_job_id": id}})using the registered job's queue/max attempts when o leaves them zero; then set river_job_id with an UpdateColumn on the same tx; return the id.Enqueue(ctx, db, args, EnqueueOpts)does only the River InsertTx (Insert when db is not in a transaction) and is what plan 11-03 uses for broadcasts.CompleteJob(ctx, id, metadata map[string]any) errorandGet(ctx, id) (Record, error)land here now (the rest in Task 2); CompleteJob reads progress_max (1 when the row is missing) and sets status 2 and progress withUpdateColumnsonTable(lagoon.JobsTable)so updated_at is never auto-touched, replacing metadata only when the map is non-empty. The Manager resolves the database per call withapp.Lookup[*gorm.DB]()/*sql.DBand returnsErrNoDatabasewhen unpublished. - client.go:
riverConfig(app, workers, queues)readsqueue.max_attempts(default 3, PHP --tries=3),queue.job_timeout(default 300s; accept "300s" or an int of seconds, the postcard timeout idiom; Pitfall 4) andqueue.queues.<name>MaxWorkers (default 4); the queue set is the configured queues plus every queue a registered job names plusdefault. The insert-only client isriver.NewClient(riverdatabasesql.New(sqlDB), cfg without Queues)built once, on first Dispatch/Enqueue/cancel when no worker client runs; the worker client is used for inserts when present. - worker.go:
type WorkerOptions struct { Queues []string }(nil = every known queue) plus unexported test knobspollInterval time.DurationandpollOnly bool;StartWorker(ctx, app, plugins []party.Plugin, o WorkerOptions) (*Worker, error): registers everypact.HasJobsjob of the plugins (a non-conga job fails with an error naming the plugin id), builds the listener pool fromlagoon.DSN(app.Config)withpgxpool.ParseConfig, MaxConns 1, MinConns 0, buildsriverdatabasesql.NewWithPgxListener(sqlDB, listener)(orriverdatabasesql.Newwhen pollOnly), filters queues (ErrUnknownQueuelisting known names), sets FetchPollInterval when the knob is set, starts the client and hands it to the Manager.(*Worker).Stop(ctx)calls Stop, falls back to StopAndCancel when ctx expires, and closes the listener pool. The riverWorker[T].Work wrapper: put the summer_jobs id from job.Metadata into ctx; call fn with panic recovery (a panic becomes an error); on error, when the row is STOPPED or is_canceled, returnriver.JobCancel(err); otherwise whenjob.Attempt >= job.MaxAttemptsset status 3 with metadata = current metadata plus keyerror(D-03); return the error so River retries non-final attempts; Timeout returns the job option or 0.
(4) Smoke tests (coverage is plan 11-07): modules/conga/postgres_test.go copies the lagoon TestMain harness (testcontainers postgres:16-alpine, ICU pl-PL, -short skips, Docker failure fails TestMain) exposing a helper that returns a dedicated *sql.DB plus its DSN migrated with lagoon.Migrate(gdb, nil). modules/conga/listen_test.go: TestListenPickupLatency (Pitfall 5: worker with pollInterval 30s, wait until SELECT count(*) FROM pg_stat_activity WHERE query ILIKE 'LISTEN%' is positive, insert through a second insert-only client with InsertTx and commit, assert the job function runs within 1s; negative control with pollOnly and the same interval asserts no run within 2s) and TestDispatchTransactional (commit path: row status 1 then the job calls CompleteJob and Get shows status 2 and progress = progress_max; rollback path: no summer_jobs row and SELECT count(*) FROM river_job unchanged).
(5) fonoteka.go allow-lists (Pitfall 12), same task so both repos stay green: add allowedDiffs entries in parity/schema_diff_test.go for summer_jobs ("Framework job record (11-01, D-01); PHP golem15_apparatus_jobs is copied at cutover") and for each River table actually created (reason "River v7 schema (11-01, JOBS-01)"); determine the names by running the migration and reading pg_tables rather than from memory. Add summer_migrations_summercms_conga to wantTables in parity/migrate_test.go (keep the ORDER BY order).
(6) Docs in the same commits: new modules/conga/README.md in the standard structure (H1 conga, the one-sentence summary "Background jobs on River over the shared Postgres pool: transactional dispatch, a summer_jobs progress record, in-process or dedicated workers, and a wall-clock scheduler.", import line, Overview, Features, Usage, API reference, Configuration, CLI commands, Dependencies (River v0.47.0 and why), Testing), mentioning the PgBouncer session-pooling requirement for the listener pool; a root README.md modules-table row with the same sentence; modules/lagoon/README.md gains the QueueMigrations set in the migrations bullet. Framework text uses neutral names (acme); verify each identifier with go doc ./modules/conga <Identifier>.
go vet ./... && go test ./modules/conga -run '^(TestListenPickupLatency|TestDispatchTransactional)$' -count=1 -v && go test ./modules/lagoon -count=1 && (cd ../fonoteka.go && go test ./parity -run '^(TestMigrateSeedsCanonicalGenres|TestSchemaMatchesPHPSnapshot)$' -count=1 -v)
<fails_when>Any command exits non-zero; the verbose output lacks a "--- PASS" line for TestListenPickupLatency, TestDispatchTransactional, TestMigrateSeedsCanonicalGenres or TestSchemaMatchesPHPSnapshot, or prints "no tests to run" or "--- SKIP"; the schema test prints "Go extra table" for a River table or summer_jobs.</fails_when>
<acceptance_criteria>
- go list -m github.com/riverqueue/river prints github.com/riverqueue/river v0.47.0.
- grep -c 'NewWithPgxListener' modules/conga/worker.go prints at least 1 and grep -c 'MaxConns' modules/conga/worker.go prints at least 1.
- grep -c 'river_job_id BIGINT' modules/lagoon/queue_migrations.go prints 1 and grep -c 'TargetVersion: RiverSchemaVersion' modules/lagoon/queue_migrations.go prints at least 1.
- grep -c 'StatusInProgress' modules/conga/conga.go prints at least 1 (Dispatch writes status 1).
- grep -c 'summer_migrations_summercms_conga' ../fonoteka.go/parity/migrate_test.go prints 1 and grep -c '"summer_jobs"' ../fonoteka.go/parity/schema_diff_test.go prints 1.
- go doc ./modules/conga Manager.Dispatch, go doc ./modules/conga Job, go doc ./modules/conga StartWorker and go doc ./modules/lagoon QueueMigrations exit 0.
- grep -c '\[conga\](modules/conga/README.md)' README.md prints 1.
- TestListenPickupLatency measures pickup under 1s with a 30s poll interval and the poll-only control shows no pickup in 2s.
</acceptance_criteria>
River is a pinned dependency, migrate creates River v7 and summer_jobs, a transactional Dispatch is picked up by a NewWithPgxListener worker well under the poll interval and completes its row, rollback leaves nothing, and fonoteka.go's schema and migrate tests pass with the new tables.
(2) Worker wrapper refinement (D-03, D-04): after fn returns, if the row is status 4 treat any ctx-cancel error as river.JobCancel; a nil return with the row still at status 1 is left as is (jobs complete themselves, as in PHP). Recovered panics log the stack at Error and follow the final-attempt rule.
(3) Commands, new modules/conga/commands.go: func RuntimeCommands(app *backpack.App, plugins []party.Plugin) []bonfire.Command returning queue:work (description "Run background job workers in the foreground"; flag queue Repeatable, "Queue to work (repeatable; default all known queues)") and queue:clear (description "Clear all queued jobs, by deleting all pending jobs." as in PHP; optional arg queue, default default; the PHP connection argument has no Go counterpart because River uses the one database). Both open the DB the lagoon way (OpenFromApp + Publish, closing on return; mirror lagoon's unexported withDB inside conga). queue:work starts StartWorker with the filter, prints worker started on queues: <sorted list>, blocks on signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM) (surf/serve.go idiom) and stops with a 10s timeout. queue:clear prints Clearing queue "<q>", loops JobDeleteMany(ctx, river.NewJobDeleteManyParams().Queues(q).States(rivertype.JobStateAvailable, rivertype.JobStateScheduled, rivertype.JobStateRetryable).First(10000)) until a pass deletes zero, and prints Cleared N jobs (D-05).
(4) serve (D-17): in modules/surf/serve.go create the signal ctx before starting anything long-lived, then after Assemble call worker, err := conga.StartServeWorker(ctx, app, plugins); StartServeWorker returns nil, nil when queue.work_in_serve is explicitly false and otherwise StartWorker with all queues. On every exit path after a successful start (shutdown and ListenAndServe error) stop the worker with the same 10s shutdown ctx after srv.Shutdown. surf importing conga is allowed because conga never imports surf. Update modules/surf/README.md (serve starts workers; queue.work_in_serve).
(5) Generated main and summer delegates: internal/build/build.go imports git.golem15.com/golem15/summercms/modules/conga and appends commands = append(commands, conga.RuntimeCommands(app, plugins)...) right after lagoon.RuntimeCommands; add TestGenerateMainRegistersCongaRuntimeCommands to build_test.go asserting that exact line appears once. Regenerate examples/hello/main.go and ../fonoteka.go/main.go with the framework CLI (go run ./cmd/summer build from each app directory, or go run <path to summercms.go>/cmd/summer build from ../fonoteka.go); bin/ output stays untracked. In cmd/summer add delegateQueueWorkCommand() (declares the repeatable queue flag and forwards every value as --queue <v>, modelled on delegateRollbackCommand) and delegateCommand("queue:clear", "Clear pending queued jobs in the app binary"), register both in toolCommands, and extend the expected list in cmd/summer/main_test.go.
(6) make:job stub (RESEARCH Pattern 2): the job.go block in internal/build/stubs/artifacts.tmpl generates {{.Ident}}Args with Kind() and func {{.Func}}() pact.Job { return conga.Job(func(ctx context.Context, args {{.Ident}}Args) error { return nil }) }, importing conga and pact but never River (the command description stays "Generate a plugin job without importing River"); drop the now-unused Worker field from artifact.go's data if nothing else reads it. TestScaffoldAllArtifacts must still compile the scaffolded plugin.
(7) fonoteka.go config: new ../fonoteka.go/config/queue.yaml with commented keys work_in_serve: true, max_attempts: 3, job_timeout: 300, queues: {default: 4} (comment: set work_in_serve false when a separate fonoteka queue:work process runs; the listener needs session pooling, not PgBouncer transaction pooling).
(8) Tests in listen_test.go for the behavior list above (the Postgres ones), kept as smoke-level; plan 11-07 adds branch coverage. Update modules/conga/README.md (API reference rows for every new method, CLI commands section, Configuration section) in the same commit.
go vet ./... && go test ./modules/conga ./modules/surf ./internal/build ./cmd/summer -count=1 && go test ./internal/build -run '^(TestGenerateMainRegistersCongaRuntimeCommands|TestScaffoldAllArtifacts)$' -count=1 -v && (cd ../fonoteka.go && go vet ./... && go build ./... && SUMMER_GOLEM15__USER__JWT__SECRET=test-only-cli-secret go run . queue:clear --help)
<fails_when>Any command exits non-zero; the verbose run lacks "--- PASS" for TestGenerateMainRegistersCongaRuntimeCommands or TestScaffoldAllArtifacts or prints "no tests to run"; the help output does not contain "queue:clear".</fails_when>
<acceptance_criteria>
- grep -c 'conga.RuntimeCommands(app, plugins)' ../fonoteka.go/main.go examples/hello/main.go prints 1 for each file.
- grep -c 'conga.StartServeWorker' modules/surf/serve.go prints 1.
- grep -c 'JobStateRunning' modules/conga/commands.go prints 0 and grep -c 'JobStateAvailable' modules/conga/commands.go prints at least 1.
- grep -c 'queue:work' cmd/summer/main_test.go prints at least 1.
- grep -c 'conga.Job(' internal/build/stubs/artifacts.tmpl prints 1.
- go doc ./modules/conga Manager.CancelJob, go doc ./modules/conga Manager.StopJob and go doc ./modules/conga RuntimeCommands exit 0.
- test -f ../fonoteka.go/config/queue.yaml succeeds and grep -c 'work_in_serve' ../fonoteka.go/config/queue.yaml prints at least 1.
</acceptance_criteria>
Every PHP JobManager operation exists with PHP semantics, retries only become ERROR on the final attempt, cancellation both stops River and marks the row, and workers run in serve (unless disabled), in queue:work with queue filters, and queue:clear removes only pending jobs.
(2) modules/lagoon/transaction.go (RESEARCH Pattern 10, D-20 seam): an unexported ctx key carries an *afterCommitBuffer. func Transaction(ctx context.Context, gdb *gorm.DB, fn func(ctx context.Context, tx *gorm.DB) error) error: when ctx already carries a buffer, run a nested tx.Transaction (savepoint) with a child buffer merged into the parent only when fn succeeds; otherwise create a buffer, run gdb.WithContext(txCtx).Transaction(func(tx) error { return fn(txCtx, tx) }), and after a nil return run each buffered callback in order with gdb.Session(&gorm.Session{NewDB: true, Context: ctx}), recovering and logging (slog.Default Warn) any panic so a committed write is never reported as failed. func AfterCommit(ctx context.Context, db *gorm.DB, fn func(ctx context.Context, db *gorm.DB)): (a) buffer in ctx → append; (b) else when the statement opened its own transaction (db.InstanceGet("gorm:started_transaction")) → append to a statement-scoped buffer stored with db.InstanceSet("lagoon:after_commit", ...); (c) else run fn now with db (the PHP after_commit=false fallback; inside an explicit non-lagoon GORM transaction db is that tx). In gormFromSQL register a callback named lagoon:after_commit with .After("gorm:commit_or_rollback_transaction") on the Create, Update and Delete processors that flushes the statement buffer only when db.Error == nil, using a fresh session on the committed connection pool. Register with Register when Get(name) is nil so repeated opens of one *gorm.DB stay idempotent.
(3) fonoteka.go Boot fix (flagged in RESEARCH Pattern 4 as optional; included because it closes a documented production gap): replace the if gdb, ok := app.Lookup[*gorm.DB](); ok { classes.RegisterHooks(gdb) } block in plugins/golem15/fonoteka/plugin.go with lagoon.OnDatabase(app, func(_ *sql.DB, gdb *gorm.DB) error { return classes.RegisterHooks(gdb) }), rewrite the surrounding comment to say the hooks now register whenever the database is published, and add TestHooksRegisterWhenDatabasePublishedAfterBoot to plugin_boot_test.go (activate with no DB published, then lagoon.Publish, then assert gdb.Callback().Create().Get("fonoteka:album_artist_resolver") is non-nil; use a fresh lagoon.Use(ctx, bootSQL) handle so the shared bootGDB is not involved).
(4) Tests: modules/lagoon/ondatabase_test.go TestOnDatabaseAfterActivate and modules/lagoon/transaction_test.go TestTransactionAfterCommit covering the behavior list (use the existing lagoon Postgres harness and a throwaway table). Update modules/lagoon/README.md (Features and API reference for OnDatabase, Transaction, AfterCommit and the lagoon:after_commit callback) in the same commit; check identifiers with go doc ./modules/lagoon OnDatabase etc.
go vet ./... && go test ./modules/lagoon -run '^(TestOnDatabaseAfterActivate|TestTransactionAfterCommit)$' -count=1 -v && go test ./... && (cd ../fonoteka.go && go vet ./... ./plugins/golem15/fonoteka/... ./plugins/golem15/user/... && go test ./... ./plugins/golem15/fonoteka/... ./plugins/golem15/user/...)
<fails_when>Any command exits non-zero; the verbose run lacks "--- PASS" for TestOnDatabaseAfterActivate or TestTransactionAfterCommit, or prints "no tests to run" or "--- SKIP"; any package in either repository reports FAIL.</fails_when>
<acceptance_criteria>
- go doc ./modules/lagoon OnDatabase, go doc ./modules/lagoon Transaction and go doc ./modules/lagoon AfterCommit exit 0.
- grep -l 'lagoon:after_commit' modules/lagoon/*.go lists at least one file.
- grep -c 'lagoon.OnDatabase' ../fonoteka.go/plugins/golem15/fonoteka/plugin.go prints 1.
- (cd ../fonoteka.go && go test ./plugins/golem15/fonoteka -run '^TestHooksRegisterWhenDatabasePublishedAfterBoot$' -count=1 -v) shows "--- PASS".
- grep -c 'OnDatabase' modules/lagoon/README.md prints at least 1.
</acceptance_criteria>
Anything registered through lagoon.OnDatabase at Boot runs when serve publishes the database, the fonoteka hooks now fire in production, and lagoon.Transaction/AfterCommit give later plans a commit-safe place to run side effects; both repositories' full suites pass.
<threat_model>
Trust Boundaries
| Boundary | Description |
|---|---|
| HTTP handler → business transaction → River | Request-driven writes enqueue jobs whose rows and River entries must share the write's fate |
| Worker process → shared Postgres pool + listener pool | Workers hold a dedicated LISTEN connection and run plugin job code |
| Operator CLI → queue:clear / queue:work | Destructive and long-running commands on the job tables |
| Go module proxy → go.mod/go.sum | New third-party code (River) enters the binary |
STRIDE Threat Register
| Threat ID | Category | Component | Severity | Disposition | Mitigation Plan |
|---|---|---|---|---|---|
| T-11-08 | Elevation of Privilege | summer_jobs ids (cancel/progress IDOR) | medium | accept | Phase 11 exposes no HTTP route over summer_jobs; the Phase 13 CSV endpoints must scope ids through the owning import (findVisible). Recorded for Phase 13. |
| T-11-12 | Tampering | Dispatch/Enqueue (orphan jobs for rolled-back writes) | high | mitigate | Row insert and River InsertTx run on the caller's *sql.Tx (Dispatch opens one when the caller has none); TestDispatchTransactional asserts rollback leaves neither (Task 1). |
| T-11-13 | Denial of Service | River listener pool | medium | mitigate | Dedicated pgxpool with MaxConns 1, MinConns 0 per worker; README documents session pooling for PgBouncer (Task 1). |
| T-11-14 | Tampering | queue:clear | medium | mitigate | JobDeleteMany restricted to available, scheduled and retryable states of one queue; a running job is untouched (Task 2 behavior test). |
| T-11-15 | Information Disclosure | job failure metadata and logs | medium | mitigate | Only the error text is written to metadata key error; the wrapper never logs job args; River's logger is the app logger (Task 1-2). |
| T-11-16 | Denial of Service | panicking plugin jobs | medium | mitigate | The wrapper recovers panics into errors and applies the final-attempt ERROR rule, so a panic cannot crash the worker or strand the row (Task 2). |
| T-11-SC | Tampering | Go module installs (River v0.47.0 and sub-modules) | high | mitigate | River is named by STACK.md and RESEARCH's legitimacy audit (Go module proxy, official docs); versions pinned in go.mod with go.sum checksums verified by GOSUMDB; no other module added. |
| </threat_model> |
<success_criteria>
- River v0.47.0 runs on one client type with NewWithPgxListener; LISTEN pickup is proven not to be poll latency.
- summer_jobs mirrors PHP columns plus river_job_id; Dispatch is transactional; every PHP JobManager method exists with PHP semantics; outcomes complete/fail/skip/cancel are queryable.
- serve runs workers in-process unless
queue.work_in_serve: false;queue:work --queueandqueue:clearexist in the app binary and assummerdelegates. - lagoon.OnDatabase closes the Boot-order gap; lagoon.Transaction/AfterCommit exist for plan 11-05.
- READMEs of conga (new), lagoon and surf plus the root modules row are updated in the same commits; both repositories green. </success_criteria>