Embed schema migrations and self-provision the schema on startup
This commit is contained in:
@@ -78,6 +78,10 @@ type Config struct {
|
||||
ReaperInterval time.Duration
|
||||
// A worker silent for longer than this is marked offline by the reaper.
|
||||
WorkerOfflineAfter time.Duration
|
||||
// Whether the binary applies its embedded schema migrations on startup.
|
||||
// On by default so a downloaded binary provisions its own database; set
|
||||
// AUTO_MIGRATE=false when an operator manages migrations out of band.
|
||||
AutoMigrate bool
|
||||
}
|
||||
|
||||
// Load reads the environment and fails fast on anything required-but-missing
|
||||
@@ -173,6 +177,14 @@ func LoadConfig() (Config, error) {
|
||||
if cfg.DefaultMaxAttempts < 1 {
|
||||
return Config{}, fmt.Errorf("DEFAULT_MAX_ATTEMPTS must be positive")
|
||||
}
|
||||
cfg.AutoMigrate = true
|
||||
if raw := os.Getenv("AUTO_MIGRATE"); raw != "" {
|
||||
parsed, err := strconv.ParseBool(raw)
|
||||
if err != nil {
|
||||
return Config{}, fmt.Errorf("AUTO_MIGRATE must be true or false")
|
||||
}
|
||||
cfg.AutoMigrate = parsed
|
||||
}
|
||||
|
||||
return cfg, nil
|
||||
}
|
||||
|
||||
@@ -23,6 +23,7 @@ import (
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
|
||||
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
|
||||
"github.com/emil28092005/SciMesh/coordinator/internal/workloads"
|
||||
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
|
||||
)
|
||||
|
||||
@@ -408,7 +409,7 @@ func TestCompleteTaskReplayIsIdempotent(t *testing.T) {
|
||||
tasks, jobs, artifacts, tx := NewTaskRepo(pool), NewJobRepo(pool), NewArtifactRepo(pool), NewTxManager(pool)
|
||||
workers, results := NewWorkerRepo(pool), NewTaskResultRepo(pool)
|
||||
clk := fixedClock{now: time.Now().UTC()}
|
||||
uc := usecase.NewCompleteTask(tasks, jobs, artifacts, workers, results, tx, clk, 2)
|
||||
uc := usecase.NewCompleteTask(tasks, jobs, artifacts, workers, results, tx, clk, 2, integrationCatalog())
|
||||
|
||||
claimed, err := tasks.ClaimNext(ctx, usecase.ClaimFilter{
|
||||
Owner: "worker-1", Now: clk.now, LeaseUntil: clk.now.Add(time.Minute),
|
||||
@@ -657,3 +658,46 @@ func TestExpireLeasesRequeuesElapsedTasks(t *testing.T) {
|
||||
t.Errorf("pending = %d, want 1 — a dead worker must not strand its task", counts[domain.TaskPending])
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateProvisionsAndIsIdempotent(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
url := os.Getenv("TEST_DATABASE_URL")
|
||||
if url == "" {
|
||||
t.Skip("TEST_DATABASE_URL is not set")
|
||||
}
|
||||
if err := Migrate(ctx, url, nil); err != nil {
|
||||
t.Fatalf("first migrate: %v", err)
|
||||
}
|
||||
if err := Migrate(ctx, url, nil); err != nil {
|
||||
t.Fatalf("second migrate (idempotent): %v", err)
|
||||
}
|
||||
pool := testPool(t)
|
||||
var count int
|
||||
if err := pool.QueryRow(ctx, "SELECT count(*) FROM schema_migrations").Scan(&count); err != nil {
|
||||
t.Fatalf("read schema_migrations: %v", err)
|
||||
}
|
||||
migrations, err := listMigrations()
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if count != len(migrations) {
|
||||
t.Errorf("schema_migrations has %d rows, want %d", count, len(migrations))
|
||||
}
|
||||
var hasJobs bool
|
||||
if err := pool.QueryRow(ctx,
|
||||
"SELECT EXISTS (SELECT 1 FROM information_schema.tables WHERE table_name = 'jobs')",
|
||||
).Scan(&hasJobs); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !hasJobs {
|
||||
t.Error("jobs table was not created by the embedded migrations")
|
||||
}
|
||||
}
|
||||
|
||||
func integrationCatalog() *workloads.Catalog {
|
||||
catalog, err := workloads.Load()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return catalog
|
||||
}
|
||||
|
||||
@@ -0,0 +1,145 @@
|
||||
package postgres
|
||||
|
||||
import (
|
||||
"context"
|
||||
"embed"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strconv"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
)
|
||||
|
||||
//go:embed migrations/*.sql
|
||||
var migrationFiles embed.FS
|
||||
|
||||
var migrationNamePattern = regexp.MustCompile(`^([0-9]+)_[a-z0-9_]+\.(up|down)\.sql$`)
|
||||
|
||||
// migration is one parsed embedded migration file.
|
||||
type migration struct {
|
||||
version int
|
||||
name string
|
||||
sql string
|
||||
}
|
||||
|
||||
// listMigrations parses and orders the embedded .up.sql files by version.
|
||||
func listMigrations() ([]migration, error) {
|
||||
entries, err := migrationFiles.ReadDir("migrations")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read embedded migrations: %w", err)
|
||||
}
|
||||
up := map[int]migration{}
|
||||
for _, entry := range entries {
|
||||
match := migrationNamePattern.FindStringSubmatch(entry.Name())
|
||||
if match == nil {
|
||||
continue
|
||||
}
|
||||
if match[2] != "up" {
|
||||
continue
|
||||
}
|
||||
version, err := strconv.Atoi(match[1])
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("migration %q has an invalid version: %w", entry.Name(), err)
|
||||
}
|
||||
if _, duplicate := up[version]; duplicate {
|
||||
return nil, fmt.Errorf("migration version %d is duplicated", version)
|
||||
}
|
||||
body, err := migrationFiles.ReadFile("migrations/" + entry.Name())
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read migration %q: %w", entry.Name(), err)
|
||||
}
|
||||
up[version] = migration{version: version, name: entry.Name(), sql: string(body)}
|
||||
}
|
||||
if len(up) == 0 {
|
||||
return nil, fmt.Errorf("no .up.sql migrations are embedded")
|
||||
}
|
||||
versions := make([]int, 0, len(up))
|
||||
for version := range up {
|
||||
versions = append(versions, version)
|
||||
}
|
||||
sort.Ints(versions)
|
||||
migrations := make([]migration, 0, len(versions))
|
||||
for _, version := range versions {
|
||||
migrations = append(migrations, up[version])
|
||||
}
|
||||
for index, item := range migrations {
|
||||
if item.version != index+1 {
|
||||
return nil, fmt.Errorf("embedded migrations are not contiguous: version %d at position %d", item.version, index+1)
|
||||
}
|
||||
}
|
||||
return migrations, nil
|
||||
}
|
||||
|
||||
// Migrate applies every embedded migration that is not yet recorded in the
|
||||
// schema_migrations table, so the binary provisions its own schema. It is
|
||||
// idempotent and safe to run concurrently: a PostgreSQL advisory lock
|
||||
// serializes migrators, and each migration file runs as its own transaction
|
||||
// (the files carry explicit BEGIN/COMMIT, matching the golang-migrate format
|
||||
// the CLI and CI still use).
|
||||
func Migrate(ctx context.Context, databaseURL string, log *slog.Logger) error {
|
||||
migrations, err := listMigrations()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
connConfig, err := pgx.ParseConfig(databaseURL)
|
||||
if err != nil {
|
||||
return fmt.Errorf("parse database url: %w", err)
|
||||
}
|
||||
// Migration files contain multiple statements (BEGIN...COMMIT), which the
|
||||
// extended query protocol rejects; run them with the simple protocol.
|
||||
connConfig.DefaultQueryExecMode = pgx.QueryExecModeSimpleProtocol
|
||||
conn, err := pgx.ConnectConfig(ctx, connConfig)
|
||||
if err != nil {
|
||||
return fmt.Errorf("connect for migration: %w", err)
|
||||
}
|
||||
defer func() { _ = conn.Close(ctx) }()
|
||||
|
||||
if _, err := conn.Exec(ctx, "SELECT pg_advisory_lock(82473911)"); err != nil {
|
||||
return fmt.Errorf("acquire migration lock: %w", err)
|
||||
}
|
||||
defer func() { _, _ = conn.Exec(ctx, "SELECT pg_advisory_unlock(82473911)") }()
|
||||
|
||||
if _, err := conn.Exec(ctx,
|
||||
"CREATE TABLE IF NOT EXISTS schema_migrations (version bigint PRIMARY KEY, applied_at timestamptz NOT NULL DEFAULT now())",
|
||||
); err != nil {
|
||||
return fmt.Errorf("ensure schema_migrations: %w", err)
|
||||
}
|
||||
|
||||
applied := map[int64]bool{}
|
||||
rows, err := conn.Query(ctx, "SELECT version FROM schema_migrations")
|
||||
if err != nil {
|
||||
return fmt.Errorf("read applied migrations: %w", err)
|
||||
}
|
||||
for rows.Next() {
|
||||
var version int64
|
||||
if err := rows.Scan(&version); err != nil {
|
||||
rows.Close()
|
||||
return fmt.Errorf("scan applied migration: %w", err)
|
||||
}
|
||||
applied[version] = true
|
||||
}
|
||||
rows.Close()
|
||||
if err := rows.Err(); err != nil {
|
||||
return fmt.Errorf("read applied migrations: %w", err)
|
||||
}
|
||||
|
||||
for _, item := range migrations {
|
||||
if applied[int64(item.version)] {
|
||||
continue
|
||||
}
|
||||
if log != nil {
|
||||
log.Info("applying migration", "version", item.version, "file", item.name)
|
||||
}
|
||||
if _, err := conn.Exec(ctx, item.sql); err != nil {
|
||||
return fmt.Errorf("apply migration %s: %w", item.name, err)
|
||||
}
|
||||
if _, err := conn.Exec(ctx,
|
||||
"INSERT INTO schema_migrations (version) VALUES ($1)", item.version,
|
||||
); err != nil {
|
||||
return fmt.Errorf("record migration %s: %w", item.name, err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,60 @@
|
||||
package postgres
|
||||
|
||||
import (
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestListMigrationsParsesAndOrdersEmbeddedFiles(t *testing.T) {
|
||||
migrations, err := listMigrations()
|
||||
if err != nil {
|
||||
t.Fatalf("list migrations: %v", err)
|
||||
}
|
||||
if len(migrations) == 0 {
|
||||
t.Fatal("no embedded migrations")
|
||||
}
|
||||
for index, item := range migrations {
|
||||
if item.version != index+1 {
|
||||
t.Errorf("migration %d has version %d, want contiguous ordering", index, item.version)
|
||||
}
|
||||
if item.name != expectedMigrationName(item.version) {
|
||||
t.Errorf("migration %d file is %q, want %q", item.version, item.name, expectedMigrationName(item.version))
|
||||
}
|
||||
if strings.TrimSpace(item.sql) == "" {
|
||||
t.Errorf("migration %d is empty", item.version)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func expectedMigrationName(version int) string {
|
||||
switch version {
|
||||
case 1:
|
||||
return "0001_init.up.sql"
|
||||
case 2:
|
||||
return "0002_workers.up.sql"
|
||||
case 3:
|
||||
return "0003_artifacts.up.sql"
|
||||
case 4:
|
||||
return "0004_result_artifact.up.sql"
|
||||
case 5:
|
||||
return "0005_uploaded_input.up.sql"
|
||||
case 6:
|
||||
return "0006_task_running_enum.up.sql"
|
||||
case 7:
|
||||
return "0007_task_running_lease.up.sql"
|
||||
case 8:
|
||||
return "0008_artifact_attempt.up.sql"
|
||||
case 9:
|
||||
return "0009_unique_partial_result_attempt.up.sql"
|
||||
case 10:
|
||||
return "0010_job_reduction.up.sql"
|
||||
case 11:
|
||||
return "0011_job_owner.up.sql"
|
||||
case 12:
|
||||
return "0012_worker_trust.up.sql"
|
||||
case 13:
|
||||
return "0013_task_results.up.sql"
|
||||
default:
|
||||
return ""
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,8 @@
|
||||
BEGIN;
|
||||
|
||||
DROP TABLE IF EXISTS tasks;
|
||||
DROP TABLE IF EXISTS jobs;
|
||||
DROP TYPE IF EXISTS task_status;
|
||||
DROP TYPE IF EXISTS job_status;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,58 @@
|
||||
BEGIN;
|
||||
|
||||
CREATE TYPE job_status AS ENUM ('pending','running','completed','failed','cancelled');
|
||||
CREATE TYPE task_status AS ENUM ('pending','leased','completed','failed','cancelled');
|
||||
|
||||
-- One user submission, possibly split into several tasks.
|
||||
CREATE TABLE jobs (
|
||||
id uuid PRIMARY KEY,
|
||||
workload text NOT NULL,
|
||||
input_uri text NOT NULL,
|
||||
parameters jsonb NOT NULL DEFAULT '{}'::jsonb,
|
||||
status job_status NOT NULL DEFAULT 'pending',
|
||||
created_at timestamptz NOT NULL DEFAULT now(),
|
||||
completed_at timestamptz
|
||||
);
|
||||
|
||||
-- One independently executable chunk.
|
||||
CREATE TABLE tasks (
|
||||
id uuid PRIMARY KEY,
|
||||
job_id uuid NOT NULL REFERENCES jobs(id) ON DELETE CASCADE,
|
||||
chunk_index integer NOT NULL,
|
||||
workload text NOT NULL,
|
||||
input_uri text NOT NULL,
|
||||
input_sha256 text NOT NULL,
|
||||
parameters jsonb NOT NULL DEFAULT '{}'::jsonb,
|
||||
status task_status NOT NULL DEFAULT 'pending',
|
||||
attempt integer NOT NULL DEFAULT 0,
|
||||
max_attempts integer NOT NULL DEFAULT 3,
|
||||
lease_owner text,
|
||||
lease_expires_at timestamptz,
|
||||
result_uri text,
|
||||
result_sha256 text,
|
||||
metrics jsonb,
|
||||
error_code text,
|
||||
error_message text,
|
||||
created_at timestamptz NOT NULL DEFAULT now(),
|
||||
started_at timestamptz,
|
||||
completed_at timestamptz,
|
||||
version integer NOT NULL DEFAULT 0,
|
||||
|
||||
CONSTRAINT uq_tasks_job_chunk UNIQUE (job_id, chunk_index),
|
||||
CONSTRAINT ck_tasks_attempt CHECK (attempt >= 0),
|
||||
CONSTRAINT ck_tasks_max_attempts CHECK (max_attempts > 0),
|
||||
-- A completed task must carry its result manifest.
|
||||
CONSTRAINT ck_tasks_completed_result CHECK (
|
||||
status <> 'completed' OR (result_uri IS NOT NULL AND result_sha256 IS NOT NULL)
|
||||
),
|
||||
-- A leased task must carry its lease.
|
||||
CONSTRAINT ck_tasks_leased_owner CHECK (
|
||||
status <> 'leased' OR (lease_owner IS NOT NULL AND lease_expires_at IS NOT NULL)
|
||||
)
|
||||
);
|
||||
|
||||
-- Claim path: find the oldest pending task fast.
|
||||
CREATE INDEX ix_tasks_claim ON tasks (status, lease_expires_at, created_at);
|
||||
CREATE INDEX ix_tasks_job ON tasks (job_id);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,6 @@
|
||||
BEGIN;
|
||||
|
||||
DROP TABLE IF EXISTS workers;
|
||||
DROP TYPE IF EXISTS worker_status;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,20 @@
|
||||
BEGIN;
|
||||
|
||||
CREATE TYPE worker_status AS ENUM ('online','busy','offline');
|
||||
|
||||
-- A registered process/machine that can claim tasks. Registration returns the
|
||||
-- id; liveness is tracked by last_heartbeat_at.
|
||||
CREATE TABLE workers (
|
||||
id uuid PRIMARY KEY,
|
||||
name text NOT NULL DEFAULT '',
|
||||
capabilities jsonb NOT NULL DEFAULT '[]'::jsonb,
|
||||
status worker_status NOT NULL DEFAULT 'online',
|
||||
last_heartbeat_at timestamptz NOT NULL DEFAULT now(),
|
||||
created_at timestamptz NOT NULL DEFAULT now(),
|
||||
updated_at timestamptz NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
-- Liveness sweep: find workers that have gone quiet.
|
||||
CREATE INDEX ix_workers_liveness ON workers (status, last_heartbeat_at);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,11 @@
|
||||
BEGIN;
|
||||
|
||||
ALTER TABLE tasks DROP COLUMN IF EXISTS input_artifact_id;
|
||||
ALTER TABLE tasks DROP COLUMN IF EXISTS result_artifact_id;
|
||||
ALTER TABLE jobs DROP COLUMN IF EXISTS input_artifact_id;
|
||||
ALTER TABLE jobs DROP COLUMN IF EXISTS result_artifact_id;
|
||||
|
||||
DROP TABLE IF EXISTS artifacts;
|
||||
DROP TYPE IF EXISTS artifact_kind;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,31 @@
|
||||
BEGIN;
|
||||
|
||||
CREATE TYPE artifact_kind AS ENUM ('input','shard','partial_result','final_result','log');
|
||||
|
||||
-- A durable file the coordinator owns: input, shard, partial/final result, log.
|
||||
-- The database is the source of truth; files are found through this metadata,
|
||||
-- never by scanning directories.
|
||||
CREATE TABLE artifacts (
|
||||
id uuid PRIMARY KEY,
|
||||
job_id uuid NOT NULL REFERENCES jobs(id) ON DELETE CASCADE,
|
||||
task_id uuid REFERENCES tasks(id) ON DELETE CASCADE, -- null for job-level inputs
|
||||
kind artifact_kind NOT NULL,
|
||||
filename text NOT NULL,
|
||||
storage_key text NOT NULL UNIQUE, -- coordinator-generated, never a client path
|
||||
content_type text NOT NULL DEFAULT 'application/octet-stream',
|
||||
size_bytes bigint NOT NULL CHECK (size_bytes >= 0),
|
||||
sha256 text NOT NULL,
|
||||
created_at timestamptz NOT NULL DEFAULT now()
|
||||
);
|
||||
|
||||
CREATE INDEX ix_artifacts_job ON artifacts (job_id);
|
||||
CREATE INDEX ix_artifacts_task ON artifacts (task_id);
|
||||
|
||||
-- Jobs and tasks reference their artifacts. Nullable during the transition from
|
||||
-- URI-based inputs/results to artifact-based ones.
|
||||
ALTER TABLE jobs ADD COLUMN input_artifact_id uuid REFERENCES artifacts(id);
|
||||
ALTER TABLE jobs ADD COLUMN result_artifact_id uuid REFERENCES artifacts(id);
|
||||
ALTER TABLE tasks ADD COLUMN input_artifact_id uuid REFERENCES artifacts(id);
|
||||
ALTER TABLE tasks ADD COLUMN result_artifact_id uuid REFERENCES artifacts(id);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,11 @@
|
||||
BEGIN;
|
||||
|
||||
ALTER TABLE tasks DROP CONSTRAINT IF EXISTS ck_tasks_completed_result;
|
||||
ALTER TABLE tasks ADD COLUMN result_uri text;
|
||||
ALTER TABLE tasks ADD COLUMN result_sha256 text;
|
||||
|
||||
ALTER TABLE tasks ADD CONSTRAINT ck_tasks_completed_result CHECK (
|
||||
status <> 'completed' OR (result_uri IS NOT NULL AND result_sha256 IS NOT NULL)
|
||||
);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,14 @@
|
||||
BEGIN;
|
||||
|
||||
-- Results are now coordinator-owned artifacts, not worker-supplied URIs.
|
||||
-- Drop the URI-based completion guard and columns, and require a completed task
|
||||
-- to reference its result artifact instead (PLAN.md §6.2).
|
||||
ALTER TABLE tasks DROP CONSTRAINT IF EXISTS ck_tasks_completed_result;
|
||||
ALTER TABLE tasks DROP COLUMN IF EXISTS result_uri;
|
||||
ALTER TABLE tasks DROP COLUMN IF EXISTS result_sha256;
|
||||
|
||||
ALTER TABLE tasks ADD CONSTRAINT ck_tasks_completed_result CHECK (
|
||||
status <> 'completed' OR result_artifact_id IS NOT NULL
|
||||
);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,9 @@
|
||||
BEGIN;
|
||||
|
||||
ALTER TABLE tasks DROP CONSTRAINT IF EXISTS ck_tasks_has_input;
|
||||
|
||||
-- Restoring NOT NULL requires the columns to be populated; safe on a fresh DB.
|
||||
ALTER TABLE tasks ALTER COLUMN input_uri SET NOT NULL;
|
||||
ALTER TABLE jobs ALTER COLUMN input_uri SET NOT NULL;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,13 @@
|
||||
BEGIN;
|
||||
|
||||
-- Inputs can now arrive as uploaded artifacts (POST /jobs/upload), not only as
|
||||
-- external URIs. Relax the URI requirement and require every task to have an
|
||||
-- input one way or the other.
|
||||
ALTER TABLE jobs ALTER COLUMN input_uri DROP NOT NULL;
|
||||
ALTER TABLE tasks ALTER COLUMN input_uri DROP NOT NULL;
|
||||
|
||||
ALTER TABLE tasks ADD CONSTRAINT ck_tasks_has_input CHECK (
|
||||
input_uri IS NOT NULL OR input_artifact_id IS NOT NULL
|
||||
);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,4 @@
|
||||
-- PostgreSQL cannot drop a single enum value without recreating the type and
|
||||
-- rewriting every dependent column. Leaving 'running' in place is harmless: no
|
||||
-- code writes it after the down of 0007 restores the leased-only transitions.
|
||||
SELECT 1;
|
||||
@@ -0,0 +1,5 @@
|
||||
-- 'running' means the worker has acknowledged start via its first heartbeat.
|
||||
-- Kept in its own migration, without an explicit transaction: an enum value
|
||||
-- added in a transaction cannot be USED in that same transaction, and the next
|
||||
-- migration references it.
|
||||
ALTER TYPE task_status ADD VALUE IF NOT EXISTS 'running';
|
||||
@@ -0,0 +1,8 @@
|
||||
BEGIN;
|
||||
|
||||
ALTER TABLE tasks DROP CONSTRAINT IF EXISTS ck_tasks_leased_owner;
|
||||
ALTER TABLE tasks ADD CONSTRAINT ck_tasks_leased_owner CHECK (
|
||||
status <> 'leased' OR (lease_owner IS NOT NULL AND lease_expires_at IS NOT NULL)
|
||||
);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,10 @@
|
||||
BEGIN;
|
||||
|
||||
-- A running task holds a lease just like a leased one, so the lease-integrity
|
||||
-- check must cover both states.
|
||||
ALTER TABLE tasks DROP CONSTRAINT IF EXISTS ck_tasks_leased_owner;
|
||||
ALTER TABLE tasks ADD CONSTRAINT ck_tasks_leased_owner CHECK (
|
||||
status NOT IN ('leased','running') OR (lease_owner IS NOT NULL AND lease_expires_at IS NOT NULL)
|
||||
);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,7 @@
|
||||
BEGIN;
|
||||
|
||||
ALTER TABLE artifacts DROP CONSTRAINT IF EXISTS ck_artifact_attempt_positive;
|
||||
ALTER TABLE artifacts DROP CONSTRAINT IF EXISTS ck_partial_result_attempt;
|
||||
ALTER TABLE artifacts DROP COLUMN IF EXISTS attempt;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,34 @@
|
||||
BEGIN;
|
||||
|
||||
-- A partial result belongs to the lease attempt that uploaded it. Without this
|
||||
-- binding a worker holding a later retry could complete a task with stale bytes
|
||||
-- uploaded by an expired attempt of that same task.
|
||||
ALTER TABLE artifacts ADD COLUMN attempt integer;
|
||||
|
||||
-- A completed task never gets a later lease, so its current attempt is also
|
||||
-- the attempt that produced the stored result.
|
||||
UPDATE artifacts AS a
|
||||
SET attempt = t.attempt
|
||||
FROM tasks AS t
|
||||
WHERE a.task_id = t.id
|
||||
AND a.kind = 'partial_result'::artifact_kind
|
||||
AND t.status = 'completed'::task_status
|
||||
AND a.attempt IS NULL;
|
||||
|
||||
-- For unfinished tasks the old schema cannot tell which attempt uploaded a
|
||||
-- partial result. Keeping it would let a later retry claim stale bytes, so the
|
||||
-- worker must upload again. Blob garbage is harmless and follows the existing
|
||||
-- coordinator-owned storage cleanup policy.
|
||||
DELETE FROM artifacts AS a
|
||||
USING tasks AS t
|
||||
WHERE a.task_id = t.id
|
||||
AND a.kind = 'partial_result'::artifact_kind
|
||||
AND t.status <> 'completed'::task_status
|
||||
AND a.attempt IS NULL;
|
||||
|
||||
ALTER TABLE artifacts ADD CONSTRAINT ck_partial_result_attempt
|
||||
CHECK (kind <> 'partial_result'::artifact_kind OR attempt IS NOT NULL);
|
||||
ALTER TABLE artifacts ADD CONSTRAINT ck_artifact_attempt_positive
|
||||
CHECK (attempt IS NULL OR attempt > 0);
|
||||
|
||||
COMMIT;
|
||||
+5
@@ -0,0 +1,5 @@
|
||||
BEGIN;
|
||||
|
||||
DROP INDEX IF EXISTS uq_partial_result_task_attempt;
|
||||
|
||||
COMMIT;
|
||||
+26
@@ -0,0 +1,26 @@
|
||||
BEGIN;
|
||||
|
||||
-- Old deployments can contain more than one partial result because earlier
|
||||
-- versions accepted repeated PUTs. Preserve the one referenced by a completed
|
||||
-- task and discard stale rows; unfinished tasks must upload again after a
|
||||
-- deploy, just as they do after a lost lease.
|
||||
DELETE FROM artifacts AS a
|
||||
USING tasks AS t
|
||||
WHERE a.task_id = t.id
|
||||
AND a.kind = 'partial_result'::artifact_kind
|
||||
AND t.status <> 'completed'::task_status;
|
||||
|
||||
DELETE FROM artifacts AS a
|
||||
USING tasks AS t
|
||||
WHERE a.task_id = t.id
|
||||
AND a.kind = 'partial_result'::artifact_kind
|
||||
AND t.status = 'completed'::task_status
|
||||
AND a.id <> t.result_artifact_id;
|
||||
|
||||
-- One lease attempt has one durable partial result. This makes an upload retry
|
||||
-- idempotent and prevents repeated uploads from accumulating orphan artifacts.
|
||||
CREATE UNIQUE INDEX uq_partial_result_task_attempt
|
||||
ON artifacts (task_id, attempt)
|
||||
WHERE kind = 'partial_result'::artifact_kind;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,7 @@
|
||||
BEGIN;
|
||||
|
||||
ALTER TABLE jobs DROP COLUMN IF EXISTS error_message;
|
||||
ALTER TABLE jobs DROP COLUMN IF EXISTS error_code;
|
||||
ALTER TABLE jobs DROP COLUMN IF EXISTS reducer_started_at;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,7 @@
|
||||
-- PostgreSQL enum values must be committed before they are used by a later
|
||||
-- transaction, so this migration intentionally has no BEGIN/COMMIT wrapper.
|
||||
ALTER TYPE job_status ADD VALUE IF NOT EXISTS 'reducing';
|
||||
|
||||
ALTER TABLE jobs ADD COLUMN IF NOT EXISTS error_code text;
|
||||
ALTER TABLE jobs ADD COLUMN IF NOT EXISTS error_message text;
|
||||
ALTER TABLE jobs ADD COLUMN IF NOT EXISTS reducer_started_at timestamptz;
|
||||
@@ -0,0 +1,6 @@
|
||||
BEGIN;
|
||||
|
||||
DROP INDEX IF EXISTS ix_jobs_owner;
|
||||
ALTER TABLE jobs DROP COLUMN IF EXISTS owner_id;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,14 @@
|
||||
BEGIN;
|
||||
|
||||
-- Who submitted this job. Equals users.id from the userservice, taken from the
|
||||
-- JWT `sub` claim. NOT a foreign key: users live in a separate service/database,
|
||||
-- so integrity is guaranteed by the signed token, not by the DB.
|
||||
--
|
||||
-- Nullable because rows created before auth existed have no owner; new inserts
|
||||
-- must supply it (enforced in the app, not the schema, during the MVP).
|
||||
ALTER TABLE jobs ADD COLUMN owner_id uuid;
|
||||
|
||||
-- "List my jobs" / "admin filters by owner" scans by owner.
|
||||
CREATE INDEX ix_jobs_owner ON jobs (owner_id);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,8 @@
|
||||
BEGIN;
|
||||
|
||||
DROP INDEX IF EXISTS ix_workers_owner;
|
||||
ALTER TABLE workers DROP COLUMN IF EXISTS trust_level;
|
||||
ALTER TABLE workers DROP COLUMN IF EXISTS owner_id;
|
||||
DROP TYPE IF EXISTS worker_trust;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,18 @@
|
||||
BEGIN;
|
||||
|
||||
-- Whether a worker's results are accepted directly or must clear quorum.
|
||||
-- 'trusted' — lab machine (shared token) or a verified/admin contributor.
|
||||
-- 'untrusted' — a plain enthusiast; results are quarantined until quorum (C2).
|
||||
CREATE TYPE worker_trust AS ENUM ('trusted', 'untrusted');
|
||||
|
||||
-- Who registered this worker (userservice user id, from the JWT sub). NULL for
|
||||
-- workers registered with the shared service token. Not a foreign key: users
|
||||
-- live in a separate service/database.
|
||||
ALTER TABLE workers ADD COLUMN owner_id uuid;
|
||||
|
||||
-- Existing rows were all shared-token lab workers, hence 'trusted'.
|
||||
ALTER TABLE workers ADD COLUMN trust_level worker_trust NOT NULL DEFAULT 'trusted';
|
||||
|
||||
CREATE INDEX ix_workers_owner ON workers (owner_id);
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,5 @@
|
||||
BEGIN;
|
||||
|
||||
DROP TABLE IF EXISTS task_results;
|
||||
|
||||
COMMIT;
|
||||
@@ -0,0 +1,23 @@
|
||||
BEGIN;
|
||||
|
||||
-- Quorum votes for a task computed by untrusted (volunteer) workers. A trusted
|
||||
-- worker's result completes the task directly and never lands here; an untrusted
|
||||
-- result is recorded as one vote, and the task is only completed once enough
|
||||
-- distinct owners submit the same result_sha256.
|
||||
--
|
||||
-- One vote per (task, owner): a single volunteer cannot stuff the ballot by
|
||||
-- running many workers under one account. A resubmission updates their vote.
|
||||
CREATE TABLE task_results (
|
||||
task_id uuid NOT NULL REFERENCES tasks(id) ON DELETE CASCADE,
|
||||
owner_id uuid NOT NULL,
|
||||
result_sha256 text NOT NULL,
|
||||
result_artifact_id uuid NOT NULL REFERENCES artifacts(id) ON DELETE CASCADE,
|
||||
created_at timestamptz NOT NULL DEFAULT now(),
|
||||
|
||||
PRIMARY KEY (task_id, owner_id)
|
||||
);
|
||||
|
||||
-- Quorum check groups a task's votes by result_sha256.
|
||||
CREATE INDEX ix_task_results_quorum ON task_results (task_id, result_sha256);
|
||||
|
||||
COMMIT;
|
||||
Reference in New Issue
Block a user