- Pin River v0.47.0 (riverdatabasesql, rivertype) and tidy the example modules - lagoon.Migrate runs the summercms.conga set: River schema v7 and summer_jobs with the apparatus columns plus an internal river_job_id link - conga.Manager.Dispatch writes the summer_jobs row (status IN_PROGRESS) and the River job on the caller's *sql.Tx; a rollback leaves neither - conga.Job wraps typed job functions so plugins never import River - conga.StartWorker runs one client on riverdatabasesql.NewWithPgxListener with a single-connection LISTEN pool; the final failed attempt sets ERROR - TestListenPickupLatency: 30s poll, pickup under 1s; poll-only control 2s miss
87 lines
2.9 KiB
Go
87 lines
2.9 KiB
Go
package lagoon
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"fmt"
|
|
|
|
"github.com/go-gormigrate/gormigrate/v2"
|
|
"github.com/riverqueue/river/riverdriver/riverdatabasesql"
|
|
"github.com/riverqueue/river/rivermigrate"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
const (
|
|
// QueueHistoryID is the migration history id of the framework job-queue
|
|
// set. Its history table is summer_migrations_summercms_conga.
|
|
QueueHistoryID = "summercms.conga"
|
|
// JobsTable is the framework job record table. Its columns match the
|
|
// WinterCMS apparatus jobs table plus the internal river_job_id link.
|
|
JobsTable = "summer_jobs"
|
|
// RiverSchemaVersion pins the River schema the framework migrates to, so a
|
|
// River upgrade never applies new DDL without a new framework migration.
|
|
RiverSchemaVersion = 7
|
|
)
|
|
|
|
// QueueMigrations returns the framework job-queue migration set: River's
|
|
// schema pinned at RiverSchemaVersion, then the summer_jobs record table.
|
|
// River's migrator runs each of its own steps in a separate transaction on
|
|
// sqlDB, because some River migrations cannot share one transaction.
|
|
func QueueMigrations(sqlDB *sql.DB) []*gormigrate.Migration {
|
|
return []*gormigrate.Migration{
|
|
{
|
|
ID: "202609290001_river_schema",
|
|
Migrate: func(tx *gorm.DB) error {
|
|
return migrateRiver(txContext(tx), sqlDB, rivermigrate.DirectionUp, &rivermigrate.MigrateOpts{TargetVersion: RiverSchemaVersion})
|
|
},
|
|
Rollback: func(tx *gorm.DB) error {
|
|
// TargetVersion -1 applies every down step, removing River's schema.
|
|
return migrateRiver(txContext(tx), sqlDB, rivermigrate.DirectionDown, &rivermigrate.MigrateOpts{TargetVersion: -1})
|
|
},
|
|
},
|
|
{
|
|
ID: "202609290002_summer_jobs",
|
|
Migrate: func(tx *gorm.DB) error {
|
|
return tx.Exec(`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
|
|
)`).Error
|
|
},
|
|
Rollback: func(tx *gorm.DB) error {
|
|
return tx.Exec(`DROP TABLE IF EXISTS summer_jobs`).Error
|
|
},
|
|
},
|
|
}
|
|
}
|
|
|
|
func migrateRiver(ctx context.Context, sqlDB *sql.DB, dir rivermigrate.Direction, opts *rivermigrate.MigrateOpts) error {
|
|
if sqlDB == nil {
|
|
return fmt.Errorf("lagoon: river migrate: sql db is nil")
|
|
}
|
|
m, err := rivermigrate.New(riverdatabasesql.New(sqlDB), nil)
|
|
if err != nil {
|
|
return fmt.Errorf("lagoon: river migrate: %w", err)
|
|
}
|
|
if _, err := m.Migrate(ctx, dir, opts); err != nil {
|
|
return fmt.Errorf("lagoon: river migrate %s: %w", dir, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func txContext(tx *gorm.DB) context.Context {
|
|
if tx != nil && tx.Statement != nil && tx.Statement.Context != nil {
|
|
return tx.Statement.Context
|
|
}
|
|
return context.Background()
|
|
}
|