feat(coordinator): scaffold task-queue service in Go

Adds the SciMesh coordinator: a durable task-queue server on PostgreSQL
that owns all database access, with workers reaching it over HTTP only.

Structured as a modular monolith following Clean Architecture:

  domain     entities and their invariants, no I/O
  usecase    business operations + repository/clock ports
  transport  HTTP handlers, DTOs, auth, error mapping
  storage    PostgreSQL repositories, transactions carried in context
  infra      config, pool, clock, server, lease reaper

Dependencies point strictly inward; domain imports nothing from the module.

Working: layer wiring, routing, shared-token auth, access logging, request
IDs, domain-error to status-code mapping, transactional boundaries,
graceful shutdown (HTTP drain -> reaper stop -> pool close), migrations,
and a Compose stack starting Postgres -> migrations -> coordinator.

The domain is complete and covered by unit tests that need no database:
lease ownership, stale attempts, idempotent result replay, retry budgets,
and lease expiry.

Repository methods are stubs returning ErrNotImplemented (HTTP 501). The
SQL for atomic claiming (FOR UPDATE SKIP LOCKED) and for lease expiry is
written and ready to wire up.

See coordinator/ARCHITECTURE.md for the layer map and a request traced
through every layer.
This commit is contained in:
Efremenko Arhip
2026-07-22 13:49:01 +03:00
parent ccf11403bf
commit bda22666d7
37 changed files with 3005 additions and 0 deletions
+18
View File
@@ -0,0 +1,18 @@
package domain
import "errors"
// Business-rule violations. They live in the innermost layer because they
// describe what the rules are, not how a transport reports them: the HTTP
// adapter maps these to status codes, and nothing here knows 409 exists.
//
// Always compare with errors.Is — outer layers may wrap these with %w.
var (
ErrJobNotFound = errors.New("job not found")
ErrTaskNotFound = errors.New("task not found")
ErrLeaseConflict = errors.New("task leased to another worker")
ErrStaleAttempt = errors.New("attempt does not match lease")
ErrResultConflict = errors.New("different result already recorded")
ErrInvalidInput = errors.New("invalid input")
ErrTaskNotLeased = errors.New("task is not currently leased")
)
+107
View File
@@ -0,0 +1,107 @@
package domain
import (
"time"
"github.com/google/uuid"
)
type JobStatus string
const (
JobPending JobStatus = "pending"
JobRunning JobStatus = "running"
JobCompleted JobStatus = "completed"
JobFailed JobStatus = "failed"
JobCancelled JobStatus = "cancelled"
)
// Job is one user submission that fans out into one or more tasks.
type Job struct {
ID uuid.UUID
Workload string
InputURI string
Parameters map[string]any
Status JobStatus
CreatedAt time.Time
CompletedAt *time.Time
}
// ChunkSpec describes one piece a job is split into. Callers build these from
// whatever chunking strategy the workload uses; the domain only validates them.
type ChunkSpec struct {
ChunkIndex int
Workload string // empty inherits the job's workload
InputURI string
InputSHA256 string
Parameters map[string]any
MaxAttempts int
}
// NewJobWithTasks builds a job together with all of its tasks, validating the
// set as a whole. Returning both from one constructor keeps the invariant
// visible: a job without tasks, or with duplicate chunk indexes, cannot exist.
func NewJobWithTasks(workload, inputURI string, params map[string]any,
chunks []ChunkSpec, now time.Time) (*Job, []*Task, error) {
if workload == "" || inputURI == "" || len(chunks) == 0 {
return nil, nil, ErrInvalidInput
}
job := &Job{
ID: uuid.New(),
Workload: workload,
InputURI: inputURI,
Parameters: params,
Status: JobPending,
CreatedAt: now,
}
seen := make(map[int]struct{}, len(chunks))
tasks := make([]*Task, 0, len(chunks))
for _, c := range chunks {
if _, dup := seen[c.ChunkIndex]; dup {
return nil, nil, ErrInvalidInput // unique (job_id, chunk_index)
}
seen[c.ChunkIndex] = struct{}{}
w := c.Workload
if w == "" {
w = workload
}
task, err := NewTask(job.ID, c.ChunkIndex, w, c.InputURI, c.InputSHA256,
c.Parameters, c.MaxAttempts, now)
if err != nil {
return nil, nil, err
}
tasks = append(tasks, task)
}
return job, tasks, nil
}
// JobProgress is the aggregate view of a job and the state of its tasks.
type JobProgress struct {
Job Job
Total int
Pending int
Leased int
Done int
Failed int
}
// DeriveStatus computes what the job's status should be from its task counts,
// so the rule lives here rather than in a SQL trigger or a handler.
func (p JobProgress) DeriveStatus() JobStatus {
switch {
case p.Total == 0:
return JobPending
case p.Done == p.Total:
return JobCompleted
case p.Failed > 0 && p.Done+p.Failed == p.Total:
return JobFailed
case p.Leased > 0 || p.Done > 0 || p.Failed > 0:
return JobRunning
default:
return JobPending
}
}
+244
View File
@@ -0,0 +1,244 @@
// Package domain holds SciMesh's entities and the rules that govern them. It
// is the innermost layer: it imports nothing from this module and knows nothing
// about HTTP, SQL, or configuration. Every state transition a task can undergo
// is a method here, so the rules are unit-testable without a database.
package domain
import (
"time"
"github.com/google/uuid"
)
type TaskStatus string
const (
TaskPending TaskStatus = "pending"
TaskLeased TaskStatus = "leased"
TaskCompleted TaskStatus = "completed"
TaskFailed TaskStatus = "failed"
TaskCancelled TaskStatus = "cancelled"
)
// ErrCodeLeaseExpired marks tasks failed by the reaper rather than by a worker.
const ErrCodeLeaseExpired = "lease_expired"
// Task is one independently executable chunk of a job.
//
// Nullable columns are pointers so "no lease" stays distinguishable from
// "lease owned by the empty string" — a plain string cannot express both.
type Task struct {
ID uuid.UUID
JobID uuid.UUID
ChunkIndex int
Workload string
InputURI string
InputSHA256 string
Parameters map[string]any
Status TaskStatus
Attempt int
MaxAttempts int
LeaseOwner *string
LeaseExpiresAt *time.Time
ResultURI *string
ResultSHA256 *string
Metrics map[string]any
ErrorCode *string
ErrorMessage *string
CreatedAt time.Time
StartedAt *time.Time
CompletedAt *time.Time
Version int
}
// NewTask builds a pending task. maxAttempts <= 0 falls back to the default.
func NewTask(jobID uuid.UUID, chunkIndex int, workload, inputURI, inputSHA256 string,
params map[string]any, maxAttempts int, now time.Time) (*Task, error) {
if inputURI == "" {
return nil, ErrInvalidInput
}
if inputSHA256 == "" {
return nil, ErrInvalidInput // checksum is mandatory: workers verify inputs
}
if chunkIndex < 0 {
return nil, ErrInvalidInput
}
if maxAttempts <= 0 {
maxAttempts = DefaultMaxAttempts
}
return &Task{
ID: uuid.New(),
JobID: jobID,
ChunkIndex: chunkIndex,
Workload: workload,
InputURI: inputURI,
InputSHA256: inputSHA256,
Parameters: params,
Status: TaskPending,
Attempt: 0,
MaxAttempts: maxAttempts,
CreatedAt: now,
}, nil
}
// DefaultMaxAttempts applies when a task does not specify its own ceiling.
const DefaultMaxAttempts = 3
// CanRetry reports whether any attempts remain.
func (t *Task) CanRetry() bool { return t.Attempt < t.MaxAttempts }
// IsLeaseHeldBy reports whether worker currently holds this task at attempt.
func (t *Task) IsLeaseHeldBy(worker string, attempt int) bool {
return t.LeaseOwner != nil && *t.LeaseOwner == worker && t.Attempt == attempt
}
// AsClaimed projects the task into the trimmed view handed to a worker:
// everything needed to execute, nothing it has no business seeing.
func (t *Task) AsClaimed() ClaimedTask {
ct := ClaimedTask{
TaskID: t.ID,
JobID: t.JobID,
ChunkIndex: t.ChunkIndex,
Workload: t.Workload,
InputURI: t.InputURI,
InputSHA256: t.InputSHA256,
Parameters: t.Parameters,
Attempt: t.Attempt,
}
if t.LeaseOwner != nil {
ct.LeaseOwner = *t.LeaseOwner
}
if t.LeaseExpiresAt != nil {
ct.LeaseExpiresAt = *t.LeaseExpiresAt
}
return ct
}
// verifyLease is the guard every worker-driven transition shares: the caller
// must own the lease and reference the attempt it was granted.
func (t *Task) verifyLease(worker string, attempt int) error {
if t.Status != TaskLeased {
return ErrTaskNotLeased
}
if t.LeaseOwner == nil || *t.LeaseOwner != worker {
return ErrLeaseConflict
}
if t.Attempt != attempt {
return ErrStaleAttempt
}
return nil
}
// RenewLease extends the lease of the worker that holds it.
func (t *Task) RenewLease(worker string, attempt int, until time.Time) error {
if err := t.verifyLease(worker, attempt); err != nil {
return err
}
t.LeaseExpiresAt = &until
t.Version++
return nil
}
// CompleteWith records a successful result.
//
// Idempotency comes first deliberately: a worker whose network dropped will
// retry the same manifest, and that must succeed rather than trip the lease
// check on a task the coordinator already finished. A *different* manifest for
// an already-completed task is a genuine conflict.
func (t *Task) CompleteWith(resultURI, resultSHA256 string, metrics map[string]any,
worker string, attempt int, now time.Time) error {
if resultURI == "" || resultSHA256 == "" {
return ErrInvalidInput
}
if t.Status == TaskCompleted {
if t.Attempt == attempt && t.ResultURI != nil && *t.ResultURI == resultURI &&
t.ResultSHA256 != nil && *t.ResultSHA256 == resultSHA256 {
return nil // same attempt, same manifest — replay of a successful call
}
return ErrResultConflict
}
if err := t.verifyLease(worker, attempt); err != nil {
return err
}
t.Status = TaskCompleted
t.ResultURI = &resultURI
t.ResultSHA256 = &resultSHA256
t.Metrics = metrics
t.CompletedAt = &now
t.LeaseOwner = nil
t.LeaseExpiresAt = nil
t.ErrorCode = nil
t.ErrorMessage = nil
t.Version++
return nil
}
// Fail records a worker-reported failure. A retryable failure with attempts
// left returns the task to the queue; otherwise it terminates as failed.
func (t *Task) Fail(worker string, attempt int, code, message string, retryable bool, now time.Time) error {
if err := t.verifyLease(worker, attempt); err != nil {
return err
}
t.ErrorCode = &code
t.ErrorMessage = &message
t.LeaseOwner = nil
t.LeaseExpiresAt = nil
t.Version++
if retryable && t.CanRetry() {
t.Status = TaskPending
return nil
}
t.Status = TaskFailed
t.CompletedAt = &now
return nil
}
// ExpireLease is applied by the reaper when a lease elapses without a
// heartbeat: requeue while attempts remain, otherwise fail terminally.
func (t *Task) ExpireLease(now time.Time) {
if t.Status != TaskLeased {
return
}
t.LeaseOwner = nil
t.LeaseExpiresAt = nil
t.Version++
if t.CanRetry() {
t.Status = TaskPending
return
}
code, msg := ErrCodeLeaseExpired, "lease expired after the final attempt"
t.ErrorCode = &code
t.ErrorMessage = &msg
t.Status = TaskFailed
t.CompletedAt = &now
}
// ClaimedTask is the worker-facing projection of a leased task.
type ClaimedTask struct {
TaskID uuid.UUID
JobID uuid.UUID
ChunkIndex int
Workload string
InputURI string
InputSHA256 string
Parameters map[string]any
Attempt int
LeaseOwner string
LeaseExpiresAt time.Time
}
// ResultManifest is a completed task's output, ordered for the stitcher.
type ResultManifest struct {
TaskID uuid.UUID
ChunkIndex int
ResultURI string
ResultSHA256 string
Metrics map[string]any
}
+183
View File
@@ -0,0 +1,183 @@
package domain
import (
"errors"
"testing"
"time"
"github.com/google/uuid"
)
var (
testNow = time.Date(2026, 7, 21, 12, 0, 0, 0, time.UTC)
testLater = testNow.Add(time.Hour)
testWorker = "worker-1"
)
// leasedTask builds a task already leased to testWorker at the given attempt.
func leasedTask(attempt, maxAttempts int) *Task {
owner := testWorker
expires := testLater
return &Task{
ID: uuid.New(),
JobID: uuid.New(),
Status: TaskLeased,
Attempt: attempt,
MaxAttempts: maxAttempts,
LeaseOwner: &owner,
LeaseExpiresAt: &expires,
}
}
func TestCompleteWithRecordsResult(t *testing.T) {
task := leasedTask(1, 3)
if err := task.CompleteWith("s3://r.csv", "abc", nil, testWorker, 1, testNow); err != nil {
t.Fatalf("unexpected error: %v", err)
}
if task.Status != TaskCompleted {
t.Errorf("status = %q, want completed", task.Status)
}
if task.LeaseOwner != nil || task.LeaseExpiresAt != nil {
t.Error("lease must be released on completion")
}
if task.CompletedAt == nil || !task.CompletedAt.Equal(testNow) {
t.Error("completed_at must be stamped")
}
}
// A worker whose network dropped retries the same manifest; that must succeed
// rather than fail on the lease it has already given up.
func TestCompleteWithIsIdempotentForSameManifest(t *testing.T) {
task := leasedTask(1, 3)
if err := task.CompleteWith("s3://r.csv", "abc", nil, testWorker, 1, testNow); err != nil {
t.Fatalf("first call: %v", err)
}
versionAfterFirst := task.Version
if err := task.CompleteWith("s3://r.csv", "abc", nil, testWorker, 1, testLater); err != nil {
t.Fatalf("replay must be idempotent, got %v", err)
}
if task.Version != versionAfterFirst {
t.Error("replay must not mutate the task")
}
}
func TestCompleteWithRejectsDifferentManifest(t *testing.T) {
task := leasedTask(1, 3)
if err := task.CompleteWith("s3://r.csv", "abc", nil, testWorker, 1, testNow); err != nil {
t.Fatalf("first call: %v", err)
}
err := task.CompleteWith("s3://other.csv", "def", nil, testWorker, 1, testLater)
if !errors.Is(err, ErrResultConflict) {
t.Errorf("err = %v, want ErrResultConflict", err)
}
}
func TestCompleteWithRejectsForeignWorker(t *testing.T) {
task := leasedTask(1, 3)
err := task.CompleteWith("s3://r.csv", "abc", nil, "worker-2", 1, testNow)
if !errors.Is(err, ErrLeaseConflict) {
t.Errorf("err = %v, want ErrLeaseConflict", err)
}
}
func TestCompleteWithRejectsStaleAttempt(t *testing.T) {
task := leasedTask(2, 3) // task is on attempt 2
err := task.CompleteWith("s3://r.csv", "abc", nil, testWorker, 1, testNow) // worker thinks it is 1
if !errors.Is(err, ErrStaleAttempt) {
t.Errorf("err = %v, want ErrStaleAttempt", err)
}
}
func TestFailRequeuesWhileAttemptsRemain(t *testing.T) {
task := leasedTask(1, 3)
if err := task.Fail(testWorker, 1, "boom", "exploded", true, testNow); err != nil {
t.Fatalf("unexpected error: %v", err)
}
if task.Status != TaskPending {
t.Errorf("status = %q, want pending", task.Status)
}
if task.LeaseOwner != nil {
t.Error("lease must be released so another worker can claim it")
}
}
func TestFailTerminatesOnFinalAttempt(t *testing.T) {
task := leasedTask(3, 3) // no attempts left
if err := task.Fail(testWorker, 3, "boom", "exploded", true, testNow); err != nil {
t.Fatalf("unexpected error: %v", err)
}
if task.Status != TaskFailed {
t.Errorf("status = %q, want failed", task.Status)
}
}
func TestFailIsTerminalWhenNotRetryable(t *testing.T) {
task := leasedTask(1, 3) // attempts remain, but the error is fatal
if err := task.Fail(testWorker, 1, "bad_input", "checksum mismatch", false, testNow); err != nil {
t.Fatalf("unexpected error: %v", err)
}
if task.Status != TaskFailed {
t.Errorf("status = %q, want failed", task.Status)
}
}
// This is the MVP acceptance criterion: a dead worker must not strand its task.
func TestExpireLeaseRequeuesWhileAttemptsRemain(t *testing.T) {
task := leasedTask(1, 3)
task.ExpireLease(testNow)
if task.Status != TaskPending {
t.Errorf("status = %q, want pending", task.Status)
}
if task.LeaseOwner != nil || task.LeaseExpiresAt != nil {
t.Error("expired lease must be cleared")
}
}
func TestExpireLeaseFailsAfterFinalAttempt(t *testing.T) {
task := leasedTask(3, 3)
task.ExpireLease(testNow)
if task.Status != TaskFailed {
t.Errorf("status = %q, want failed", task.Status)
}
if task.ErrorCode == nil || *task.ErrorCode != ErrCodeLeaseExpired {
t.Error("expected a lease_expired error code")
}
}
func TestExpireLeaseIgnoresUnleasedTasks(t *testing.T) {
task := &Task{Status: TaskCompleted, Attempt: 1, MaxAttempts: 3}
task.ExpireLease(testNow)
if task.Status != TaskCompleted {
t.Errorf("status = %q, completed tasks must be untouched", task.Status)
}
}
func TestRenewLeaseExtendsOnlyForHolder(t *testing.T) {
task := leasedTask(1, 3)
until := testLater.Add(time.Hour)
if err := task.RenewLease(testWorker, 1, until); err != nil {
t.Fatalf("unexpected error: %v", err)
}
if !task.LeaseExpiresAt.Equal(until) {
t.Error("lease must be extended")
}
if err := task.RenewLease("worker-2", 1, until); !errors.Is(err, ErrLeaseConflict) {
t.Errorf("err = %v, want ErrLeaseConflict", err)
}
}
+13
View File
@@ -0,0 +1,13 @@
// Clock: the real implementation of the usecase.Clock port. It lives out here
// because reading the system clock is infrastructure; tests substitute a fixed one.
package infra
import "time"
type System struct{}
func NewClock() System { return System{} }
// Now returns UTC so every timestamp the coordinator writes is comparable
// regardless of the host's timezone.
func (System) Now() time.Time { return time.Now().UTC() }
+135
View File
@@ -0,0 +1,135 @@
// Config: coordinator settings, read only from the environment, so the same
// binary behaves identically in CI, local, and prod.
package infra
import (
"errors"
"fmt"
"io/fs"
"math"
"os"
"strconv"
"time"
"github.com/joho/godotenv"
)
// defaultEnvFile is loaded by Load unless ENV_FILE points elsewhere.
const defaultEnvFile = ".env"
type Config struct {
// HTTP listen address, e.g. ":8080".
Addr string
// PostgreSQL connection string (pgx format / libpq URL).
DatabaseURL string
// Shared bearer token workers must present. Empty disables auth (dev only).
WorkerAuthToken string
// Connection pool upper bound.
DBMaxConns int32
// Per-request context timeout applied to handlers and DB calls.
RequestTimeout time.Duration
// Default lease length handed out on claim.
LeaseDuration time.Duration
// Default attempt ceiling for newly created tasks.
DefaultMaxAttempts int
// How often the background lease-reaper runs.
ReaperInterval time.Duration
}
// Load reads the environment and fails fast on anything required-but-missing
// or malformed, so a misconfigured process never limps along half-wired.
//
// A .env file (path overridable via ENV_FILE) is loaded first as a local-dev
// convenience. It only fills variables the environment does not already define.
func Load() (Config, error) {
envFile := os.Getenv("ENV_FILE")
if envFile == "" {
envFile = defaultEnvFile
}
// godotenv.Load never overwrites variables already present in the
// environment, so an orchestrator's values always beat the file. A missing
// file is expected in production, where env vars are injected directly.
if err := godotenv.Load(envFile); err != nil && !errors.Is(err, fs.ErrNotExist) {
return Config{}, fmt.Errorf("load env file %q: %w", envFile, err)
}
cfg := Config{
Addr: getEnv("COORDINATOR_ADDR", ":8080"),
DatabaseURL: os.Getenv("DATABASE_URL"),
WorkerAuthToken: os.Getenv("WORKER_AUTH_TOKEN"),
DBMaxConns: 10,
RequestTimeout: 15 * time.Second,
LeaseDuration: 2 * time.Minute,
DefaultMaxAttempts: 3,
ReaperInterval: 30 * time.Second,
}
if cfg.DatabaseURL == "" {
return Config{}, fmt.Errorf("DATABASE_URL is required")
}
var err error
if cfg.DBMaxConns, err = getEnvInt32("DB_MAX_CONNS", cfg.DBMaxConns); err != nil {
return Config{}, err
}
if cfg.RequestTimeout, err = getEnvDuration("REQUEST_TIMEOUT", cfg.RequestTimeout); err != nil {
return Config{}, err
}
if cfg.LeaseDuration, err = getEnvDuration("LEASE_DURATION", cfg.LeaseDuration); err != nil {
return Config{}, err
}
if cfg.ReaperInterval, err = getEnvDuration("REAPER_INTERVAL", cfg.ReaperInterval); err != nil {
return Config{}, err
}
if cfg.DefaultMaxAttempts, err = getEnvInt("DEFAULT_MAX_ATTEMPTS", cfg.DefaultMaxAttempts); err != nil {
return Config{}, err
}
return cfg, nil
}
func getEnv(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}
func getEnvInt(key string, def int) (int, error) {
v := os.Getenv(key)
if v == "" {
return def, nil
}
n, err := strconv.Atoi(v)
if err != nil {
return 0, fmt.Errorf("%s: %w", key, err)
}
return n, nil
}
func getEnvInt32(key string, def int32) (int32, error) {
n, err := getEnvInt(key, int(def))
if err != nil {
return 0, err
}
// On 64-bit builds int is wider than int32, so an oversized value would
// wrap silently — DB_MAX_CONNS=2147483648 becoming a negative pool size.
if n < math.MinInt32 || n > math.MaxInt32 {
return 0, fmt.Errorf("%s: %d is out of range for int32", key, n)
}
return int32(n), nil
}
func getEnvDuration(key string, def time.Duration) (time.Duration, error) {
v := os.Getenv(key)
if v == "" {
return def, nil
}
d, err := time.ParseDuration(v)
if err != nil {
return 0, fmt.Errorf("%s: %w", key, err)
}
return d, nil
}
+30
View File
@@ -0,0 +1,30 @@
// DB: the PostgreSQL connection pool.
package infra
import (
"context"
"github.com/jackc/pgx/v5/pgxpool"
)
// NewPool builds the single shared pool. The caller owns its lifetime and must
// Close() it on shutdown.
func NewPool(ctx context.Context, cfg Config) (*pgxpool.Pool, error) {
poolCfg, err := pgxpool.ParseConfig(cfg.DatabaseURL)
if err != nil {
return nil, err
}
poolCfg.MaxConns = cfg.DBMaxConns
pool, err := pgxpool.NewWithConfig(ctx, poolCfg)
if err != nil {
return nil, err
}
// pgxpool.New is lazy, so without this ping a bad DATABASE_URL would only
// surface on the first request instead of at startup.
if err := pool.Ping(ctx); err != nil {
pool.Close()
return nil, err
}
return pool, nil
}
+7
View File
@@ -0,0 +1,7 @@
// Package infra holds the outermost layer: frameworks, drivers, and process
// wiring. It reads configuration, opens the database pool, supplies the real
// clock, and runs the HTTP server and background reaper.
//
// Nothing inward depends on this package — it is the last thing constructed and
// the first thing that would be swapped when the runtime environment changes.
package infra
+72
View File
@@ -0,0 +1,72 @@
// Server: the HTTP listener and the background lease reaper, both shut down
// cleanly on a signal.
package infra
import (
"context"
"errors"
"log/slog"
"net/http"
"time"
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
)
const shutdownGrace = 15 * time.Second
// Run serves handler until ctx is cancelled, then drains in-flight requests.
func RunServer(ctx context.Context, log *slog.Logger, addr string, handler http.Handler) error {
srv := &http.Server{
Addr: addr,
Handler: handler,
ReadHeaderTimeout: 5 * time.Second,
}
// Buffered so this goroutine can exit even when nobody reads the channel
// (the ctx.Done branch below) — an unbuffered send would leak it forever.
errCh := make(chan error, 1)
go func() {
log.Info("coordinator listening", "addr", addr)
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
errCh <- err
}
}()
select {
case err := <-errCh:
return err
case <-ctx.Done():
log.Info("shutdown signal received")
}
// A fresh context: ctx is already cancelled, and reusing it would abort the
// very requests we are trying to let finish.
shutdownCtx, cancel := context.WithTimeout(context.Background(), shutdownGrace)
defer cancel()
return srv.Shutdown(shutdownCtx)
}
// RunReaper periodically reclaims tasks whose lease elapsed, so a worker that
// died without a heartbeat cannot strand its task in 'leased' forever.
func RunReaper(ctx context.Context, log *slog.Logger, uc *usecase.ExpireLeases, interval time.Duration) {
t := time.NewTicker(interval)
defer t.Stop()
for {
select {
case <-ctx.Done():
return
case <-t.C:
n, err := uc.Execute(ctx)
if err != nil {
// Demoted to debug while the repository is still a stub;
// raise to Warn once phase 5 lands.
log.Debug("reaper skipped", "err", err)
continue
}
if n > 0 {
log.Info("reaper requeued expired leases", "count", n)
}
}
}
}
@@ -0,0 +1,40 @@
package postgres
import (
"context"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
)
// JobRepo implements usecase.JobRepository.
type JobRepo struct {
pool *pgxpool.Pool
}
func NewJobRepo(pool *pgxpool.Pool) *JobRepo {
return &JobRepo{pool: pool}
}
var _ usecase.JobRepository = (*JobRepo)(nil)
// TODO(phase 2-3): replace stubs with real pgx queries.
func (r *JobRepo) Insert(ctx context.Context, j *domain.Job) error {
// Phase 2: INSERT INTO jobs ...; runs inside the caller's transaction.
return usecase.ErrNotImplemented
}
func (r *JobRepo) Get(ctx context.Context, id uuid.UUID) (*domain.Job, error) {
// Phase 3: SELECT ... WHERE id = $1; no rows -> domain.ErrJobNotFound.
return nil, usecase.ErrNotImplemented
}
func (r *JobRepo) UpdateStatus(ctx context.Context, id uuid.UUID, status domain.JobStatus, completedAt *time.Time) error {
// Phase 3: UPDATE jobs SET status = $2, completed_at = $3 WHERE id = $1.
return usecase.ErrNotImplemented
}
@@ -0,0 +1,82 @@
package postgres
import (
"context"
"errors"
"time"
"github.com/cenkalti/backoff/v4"
"github.com/jackc/pgx/v5/pgconn"
)
// Transient PostgreSQL failures. Under concurrent claiming these are expected
// rather than exceptional: two coordinators touching neighbouring rows can
// deadlock or fail to serialize, and the correct response is to try again.
const (
codeSerializationFailure = "40001"
codeDeadlockDetected = "40P01"
codeTooManyConnections = "53300"
codeCannotConnectNow = "57P03"
)
// Retry budget: short and bounded. A worker polling for tasks would rather get
// a fast error and poll again than have its request hang for half a minute.
const (
retryInitialInterval = 50 * time.Millisecond
retryMaxInterval = 1 * time.Second
retryMaxElapsedTime = 5 * time.Second
)
// isTransient reports whether err is worth retrying.
//
// The default is *not* to retry: a constraint violation or a syntax error will
// fail identically every time, and retrying it only multiplies the damage.
func isTransient(err error) bool {
if err == nil {
return false
}
// A cancelled caller does not want another attempt.
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
return false
}
var pgErr *pgconn.PgError
if errors.As(err, &pgErr) {
switch pgErr.Code {
case codeSerializationFailure, codeDeadlockDetected,
codeTooManyConnections, codeCannotConnectNow:
return true
default:
return false
}
}
// Connection-level trouble (dropped socket, closed pool). pgconn knows
// whether the query could have been executed before the failure — retrying
// a maybe-executed write would risk duplicating it.
return pgconn.SafeToRetry(err)
}
// withRetry runs op, retrying only transient database failures with
// exponential backoff and jitter, and giving up as soon as ctx is done.
//
// Jitter matters here: without it, several coordinators that collide once will
// retry in lockstep and collide again at exactly the same moment.
func withRetry(ctx context.Context, op func(context.Context) error) error {
b := backoff.NewExponentialBackOff()
b.InitialInterval = retryInitialInterval
b.MaxInterval = retryMaxInterval
b.MaxElapsedTime = retryMaxElapsedTime
// RandomizationFactor defaults to 0.5, which is the jitter.
return backoff.Retry(func() error {
err := op(ctx)
if err == nil {
return nil
}
if !isTransient(err) {
return backoff.Permanent(err) // stop now, do not burn the budget
}
return err
}, backoff.WithContext(b, ctx))
}
@@ -0,0 +1,98 @@
package postgres
import (
"context"
"errors"
"testing"
"time"
"github.com/jackc/pgx/v5/pgconn"
)
func TestIsTransient(t *testing.T) {
tests := []struct {
name string
err error
want bool
}{
{"nil", nil, false},
{"serialization failure", &pgconn.PgError{Code: codeSerializationFailure}, true},
{"deadlock", &pgconn.PgError{Code: codeDeadlockDetected}, true},
{"too many connections", &pgconn.PgError{Code: codeTooManyConnections}, true},
// A unique-violation repeats identically forever — retrying is pointless.
{"unique violation", &pgconn.PgError{Code: "23505"}, false},
{"syntax error", &pgconn.PgError{Code: "42601"}, false},
{"context cancelled", context.Canceled, false},
{"deadline exceeded", context.DeadlineExceeded, false},
{"unknown error", errors.New("boom"), false},
// Wrapping must not hide the cause: errors.As walks the chain.
{"wrapped deadlock", errors2Wrap(&pgconn.PgError{Code: codeDeadlockDetected}), true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := isTransient(tt.err); got != tt.want {
t.Errorf("isTransient(%v) = %v, want %v", tt.err, got, tt.want)
}
})
}
}
func errors2Wrap(err error) error {
return errors.Join(errors.New("query failed"), err)
}
func TestWithRetrySucceedsAfterTransientFailures(t *testing.T) {
calls := 0
err := withRetry(context.Background(), func(context.Context) error {
calls++
if calls < 3 {
return &pgconn.PgError{Code: codeSerializationFailure}
}
return nil
})
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if calls != 3 {
t.Errorf("calls = %d, want 3", calls)
}
}
func TestWithRetryStopsOnPermanentError(t *testing.T) {
permanent := &pgconn.PgError{Code: "23505"} // unique violation
calls := 0
err := withRetry(context.Background(), func(context.Context) error {
calls++
return permanent
})
if !errors.Is(err, permanent) {
t.Errorf("err = %v, want the original error", err)
}
if calls != 1 {
t.Errorf("calls = %d, want 1 — a permanent error must not be retried", calls)
}
}
func TestWithRetryHonoursContextCancellation(t *testing.T) {
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
defer cancel()
calls := 0
start := time.Now()
err := withRetry(ctx, func(context.Context) error {
calls++
return &pgconn.PgError{Code: codeDeadlockDetected}
})
if err == nil {
t.Fatal("expected an error once the context expired")
}
// Must abort at the deadline, not run the full 5s retry budget.
if elapsed := time.Since(start); elapsed > time.Second {
t.Errorf("took %v, expected to stop at the context deadline", elapsed)
}
}
@@ -0,0 +1,114 @@
package postgres
import (
"context"
"time"
"github.com/google/uuid"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
)
// TaskRepo implements usecase.TaskRepository.
type TaskRepo struct {
pool *pgxpool.Pool
}
func NewTaskRepo(pool *pgxpool.Pool) *TaskRepo {
return &TaskRepo{pool: pool}
}
var _ usecase.TaskRepository = (*TaskRepo)(nil)
// claimNextSQL leases one task in a single statement.
//
// FOR UPDATE SKIP LOCKED is what makes concurrent coordinators safe: each
// process locks a different candidate row instead of queueing on the same one,
// so no task is ever handed to two workers and no claim blocks behind another.
// Splitting this into SELECT + UPDATE would reintroduce exactly that race.
//
//nolint:unused // wired up in phase 2; kept beside the repository it belongs to
const claimNextSQL = `
WITH candidate AS (
SELECT id
FROM tasks
WHERE status = 'pending'
AND attempt < max_attempts
AND (cardinality($1::text[]) = 0 OR workload = ANY($1))
ORDER BY created_at, chunk_index
FOR UPDATE SKIP LOCKED
LIMIT 1
)
UPDATE tasks
SET status = 'leased',
attempt = attempt + 1,
lease_owner = $2,
lease_expires_at = $3,
started_at = COALESCE(started_at, $4),
version = version + 1
FROM candidate
WHERE tasks.id = candidate.id
RETURNING tasks.id, tasks.job_id, tasks.chunk_index, tasks.workload,
tasks.input_uri, tasks.input_sha256, tasks.parameters,
tasks.status, tasks.attempt, tasks.max_attempts,
tasks.lease_owner, tasks.lease_expires_at, tasks.version;
`
// expireLeasesSQL applies the lease-expiry rule set-based, mirroring
// domain.Task.ExpireLease: requeue while attempts remain, otherwise fail.
//
//nolint:unused // wired up in phase 5; mirrors domain.Task.ExpireLease
const expireLeasesSQL = `
UPDATE tasks
SET status = CASE WHEN attempt < max_attempts THEN 'pending'::task_status
ELSE 'failed'::task_status END,
lease_owner = NULL,
lease_expires_at = NULL,
error_code = CASE WHEN attempt >= max_attempts THEN $2 ELSE error_code END,
error_message = CASE WHEN attempt >= max_attempts
THEN 'lease expired after the final attempt'
ELSE error_message END,
completed_at = CASE WHEN attempt >= max_attempts THEN $1 ELSE completed_at END,
version = version + 1
WHERE status = 'leased' AND lease_expires_at < $1;
`
// TODO(phase 2-6): replace stubs with real pgx queries; the SQL above is ready
// to wire up. Each method maps 1:1 to a roadmap phase.
func (r *TaskRepo) ClaimNext(ctx context.Context, f usecase.ClaimFilter) (*domain.Task, error) {
// Phase 2: run claimNextSQL; pgx.ErrNoRows -> (nil, nil).
return nil, usecase.ErrNotImplemented
}
func (r *TaskRepo) GetForUpdate(ctx context.Context, id uuid.UUID) (*domain.Task, error) {
// Phase 3: SELECT ... WHERE id = $1 FOR UPDATE; no rows -> domain.ErrTaskNotFound.
return nil, usecase.ErrNotImplemented
}
func (r *TaskRepo) Update(ctx context.Context, t *domain.Task) error {
// Phase 3: UPDATE ... WHERE id = $1 AND version = $2 (optimistic concurrency).
return usecase.ErrNotImplemented
}
func (r *TaskRepo) InsertBatch(ctx context.Context, tasks []*domain.Task) error {
// Phase 2: pgx.Batch or COPY; runs inside the caller's transaction.
return usecase.ErrNotImplemented
}
func (r *TaskRepo) ListCompleted(ctx context.Context, jobID uuid.UUID) ([]*domain.Task, error) {
// Phase 6: WHERE job_id = $1 AND status = 'completed' ORDER BY chunk_index.
return nil, usecase.ErrNotImplemented
}
func (r *TaskRepo) CountByStatus(ctx context.Context, jobID uuid.UUID) (map[domain.TaskStatus]int, error) {
// Phase 3: SELECT status, count(*) ... GROUP BY status.
return nil, usecase.ErrNotImplemented
}
func (r *TaskRepo) ExpireLeases(ctx context.Context, now time.Time) (int64, error) {
// Phase 5: run expireLeasesSQL, return the affected row count.
return 0, usecase.ErrNotImplemented
}
@@ -0,0 +1,69 @@
// Package postgres implements the usecase repository ports on PostgreSQL.
// SQL and pgx types never escape this package.
package postgres
import (
"context"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgconn"
"github.com/jackc/pgx/v5/pgxpool"
)
// querier is satisfied by both *pgxpool.Pool and pgx.Tx, letting every
// repository method run identically inside or outside a transaction.
//
//nolint:unused // used by repository methods once phase 2 replaces the stubs
type querier interface {
Query(ctx context.Context, sql string, args ...any) (pgx.Rows, error)
QueryRow(ctx context.Context, sql string, args ...any) pgx.Row
Exec(ctx context.Context, sql string, args ...any) (pgconn.CommandTag, error)
}
// txKey is an unexported struct type, so no other package can collide with it
// or reach the transaction we stash in the context.
type txKey struct{}
// TxManager implements usecase.TxManager.
type TxManager struct {
pool *pgxpool.Pool
}
func NewTxManager(pool *pgxpool.Pool) *TxManager {
return &TxManager{pool: pool}
}
// WithinTx runs fn inside one transaction, committing on success and rolling
// back on any error or panic.
//
// The transaction travels in the context rather than in fn's signature, which
// is what lets the usecase layer express "do these repository calls atomically"
// without its port ever mentioning pgx.
func (m *TxManager) WithinTx(ctx context.Context, fn func(ctx context.Context) error) error {
if _, ok := ctx.Value(txKey{}).(pgx.Tx); ok {
return fn(ctx) // already inside a transaction — join it, don't nest
}
tx, err := m.pool.Begin(ctx)
if err != nil {
return err
}
// Rollback after a successful Commit is a no-op, so this defer is safe and
// also covers the panic path.
defer func() { _ = tx.Rollback(ctx) }()
if err := fn(context.WithValue(ctx, txKey{}, tx)); err != nil {
return err
}
return tx.Commit(ctx)
}
// conn returns the transaction bound to ctx, or the pool when there is none.
//
//nolint:unused // every repository method will route through this in phase 2
func conn(ctx context.Context, pool *pgxpool.Pool) querier {
if tx, ok := ctx.Value(txKey{}).(pgx.Tx); ok {
return tx
}
return pool
}
+119
View File
@@ -0,0 +1,119 @@
package http
import (
"time"
"github.com/google/uuid"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
)
// Wire formats. Keeping them separate from domain entities means the API
// contract can evolve without reshaping the database, and nothing internal
// (version counters, other workers' errors) leaks by accident.
type createJobRequest struct {
Workload string `json:"workload"`
InputURI string `json:"input_uri"`
Parameters map[string]any `json:"parameters"`
Chunks []chunkDTO `json:"chunks"`
}
type chunkDTO struct {
ChunkIndex int `json:"chunk_index"`
Workload string `json:"workload"`
InputURI string `json:"input_uri"`
InputSHA256 string `json:"input_sha256"`
Parameters map[string]any `json:"parameters"`
MaxAttempts int `json:"max_attempts"`
}
type claimRequest struct {
WorkerID string `json:"worker_id"`
Workloads []string `json:"workloads"`
}
type heartbeatRequest struct {
WorkerID string `json:"worker_id"`
Attempt int `json:"attempt"`
}
type resultRequest struct {
WorkerID string `json:"worker_id"`
Attempt int `json:"attempt"`
ResultURI string `json:"result_uri"`
ResultSHA256 string `json:"result_sha256"`
Metrics map[string]any `json:"metrics"`
}
type failureRequest struct {
WorkerID string `json:"worker_id"`
Attempt int `json:"attempt"`
ErrorCode string `json:"error_code"`
ErrorMessage string `json:"error_message"`
Retryable bool `json:"retryable"`
}
type jobResponse struct {
ID uuid.UUID `json:"id"`
Status string `json:"status"`
}
type taskResponse struct {
ID uuid.UUID `json:"id"`
JobID uuid.UUID `json:"job_id"`
Status string `json:"status"`
}
type claimedTaskResponse struct {
TaskID uuid.UUID `json:"task_id"`
JobID uuid.UUID `json:"job_id"`
ChunkIndex int `json:"chunk_index"`
Workload string `json:"workload"`
InputURI string `json:"input_uri"`
InputSHA256 string `json:"input_sha256"`
Parameters map[string]any `json:"parameters"`
Attempt int `json:"attempt"`
LeaseExpiresAt time.Time `json:"lease_expires_at"`
}
type jobProgressResponse struct {
ID uuid.UUID `json:"id"`
Status string `json:"status"`
Total int `json:"total"`
Pending int `json:"pending"`
Leased int `json:"leased"`
Done int `json:"completed"`
Failed int `json:"failed"`
}
type errorResponse struct {
Error string `json:"error"`
RequestID string `json:"request_id,omitempty"`
}
func toClaimedTaskResponse(c domain.ClaimedTask) claimedTaskResponse {
return claimedTaskResponse{
TaskID: c.TaskID,
JobID: c.JobID,
ChunkIndex: c.ChunkIndex,
Workload: c.Workload,
InputURI: c.InputURI,
InputSHA256: c.InputSHA256,
Parameters: c.Parameters,
Attempt: c.Attempt,
LeaseExpiresAt: c.LeaseExpiresAt,
}
}
func toJobProgressResponse(p domain.JobProgress) jobProgressResponse {
return jobProgressResponse{
ID: p.Job.ID,
Status: string(p.DeriveStatus()),
Total: p.Total,
Pending: p.Pending,
Leased: p.Leased,
Done: p.Done,
Failed: p.Failed,
}
}
@@ -0,0 +1,55 @@
package http
import (
"encoding/json"
"errors"
"net/http"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
)
func writeJSON(w http.ResponseWriter, status int, v any) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(status)
_ = json.NewEncoder(w).Encode(v)
}
func decodeJSON(r *http.Request, dst any) error {
dec := json.NewDecoder(r.Body)
// Reject unknown fields: silently ignoring a misspelled "worker_ID" would
// surface later as a baffling validation failure.
dec.DisallowUnknownFields()
return dec.Decode(dst)
}
// writeError translates domain errors into status codes. This mapping is the
// only place in the codebase that knows HTTP status codes exist — the inner
// layers speak only in business terms.
func (s *Server) writeError(w http.ResponseWriter, r *http.Request, err error) {
reqID := requestIDFrom(r.Context())
status := http.StatusInternalServerError
switch {
case errors.Is(err, domain.ErrInvalidInput):
status = http.StatusBadRequest
case errors.Is(err, domain.ErrJobNotFound), errors.Is(err, domain.ErrTaskNotFound):
status = http.StatusNotFound
case errors.Is(err, domain.ErrLeaseConflict),
errors.Is(err, domain.ErrStaleAttempt),
errors.Is(err, domain.ErrResultConflict),
errors.Is(err, domain.ErrTaskNotLeased):
status = http.StatusConflict
case errors.Is(err, usecase.ErrNotImplemented):
status = http.StatusNotImplemented
}
if status >= 500 {
// Never echo an internal error: it can carry table names, query
// fragments, and values. The request ID is the bridge to the logs.
s.log.Error("request failed", "request_id", reqID, "path", r.URL.Path, "err", err)
writeJSON(w, status, errorResponse{Error: "internal error", RequestID: reqID})
return
}
writeJSON(w, status, errorResponse{Error: err.Error(), RequestID: reqID})
}
@@ -0,0 +1,181 @@
package http
import (
"context"
"net/http"
"github.com/google/uuid"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
)
// Every handler follows the same shape: decode, map to a use-case input,
// execute, translate. Anything resembling a rule belongs one layer inward.
func (s *Server) handleCreateJob(w http.ResponseWriter, r *http.Request) {
ctx, cancel := s.reqCtx(r)
defer cancel()
var req createJobRequest
if err := decodeJSON(r, &req); err != nil {
s.writeError(w, r, domain.ErrInvalidInput)
return
}
in := usecase.CreateJobInput{
Workload: req.Workload,
InputURI: req.InputURI,
Parameters: req.Parameters,
}
for _, c := range req.Chunks {
in.Chunks = append(in.Chunks, usecase.ChunkInput(c))
}
job, err := s.uc.CreateJob.Execute(ctx, in)
if err != nil {
s.writeError(w, r, err)
return
}
writeJSON(w, http.StatusCreated, jobResponse{ID: job.ID, Status: string(job.Status)})
}
func (s *Server) handleClaim(w http.ResponseWriter, r *http.Request) {
ctx, cancel := s.reqCtx(r)
defer cancel()
var req claimRequest
if err := decodeJSON(r, &req); err != nil {
s.writeError(w, r, domain.ErrInvalidInput)
return
}
claimed, err := s.uc.ClaimTask.Execute(ctx, usecase.ClaimTaskInput{
WorkerID: req.WorkerID,
Workloads: req.Workloads,
})
if err != nil {
s.writeError(w, r, err)
return
}
if claimed == nil {
w.WriteHeader(http.StatusNoContent) // empty queue, not an error
return
}
writeJSON(w, http.StatusOK, toClaimedTaskResponse(*claimed))
}
func (s *Server) handleHeartbeat(w http.ResponseWriter, r *http.Request) {
ctx, cancel := s.reqCtx(r)
defer cancel()
taskID, ok := s.pathUUID(w, r, "task_id")
if !ok {
return
}
var req heartbeatRequest
if err := decodeJSON(r, &req); err != nil {
s.writeError(w, r, domain.ErrInvalidInput)
return
}
claimed, err := s.uc.RenewLease.Execute(ctx, usecase.RenewLeaseInput{
TaskID: taskID,
WorkerID: req.WorkerID,
Attempt: req.Attempt,
})
if err != nil {
s.writeError(w, r, err)
return
}
writeJSON(w, http.StatusOK, toClaimedTaskResponse(*claimed))
}
func (s *Server) handleResult(w http.ResponseWriter, r *http.Request) {
ctx, cancel := s.reqCtx(r)
defer cancel()
taskID, ok := s.pathUUID(w, r, "task_id")
if !ok {
return
}
var req resultRequest
if err := decodeJSON(r, &req); err != nil {
s.writeError(w, r, domain.ErrInvalidInput)
return
}
task, err := s.uc.CompleteTask.Execute(ctx, usecase.CompleteTaskInput{
TaskID: taskID,
WorkerID: req.WorkerID,
Attempt: req.Attempt,
ResultURI: req.ResultURI,
ResultSHA256: req.ResultSHA256,
Metrics: req.Metrics,
})
if err != nil {
s.writeError(w, r, err)
return
}
writeJSON(w, http.StatusOK, taskResponse{ID: task.ID, JobID: task.JobID, Status: string(task.Status)})
}
func (s *Server) handleFailure(w http.ResponseWriter, r *http.Request) {
ctx, cancel := s.reqCtx(r)
defer cancel()
taskID, ok := s.pathUUID(w, r, "task_id")
if !ok {
return
}
var req failureRequest
if err := decodeJSON(r, &req); err != nil {
s.writeError(w, r, domain.ErrInvalidInput)
return
}
task, err := s.uc.FailTask.Execute(ctx, usecase.FailTaskInput{
TaskID: taskID,
WorkerID: req.WorkerID,
Attempt: req.Attempt,
ErrorCode: req.ErrorCode,
ErrorMessage: req.ErrorMessage,
Retryable: req.Retryable,
})
if err != nil {
s.writeError(w, r, err)
return
}
writeJSON(w, http.StatusOK, taskResponse{ID: task.ID, JobID: task.JobID, Status: string(task.Status)})
}
func (s *Server) handleGetJob(w http.ResponseWriter, r *http.Request) {
ctx, cancel := s.reqCtx(r)
defer cancel()
jobID, ok := s.pathUUID(w, r, "job_id")
if !ok {
return
}
progress, err := s.uc.GetJobStatus.Execute(ctx, jobID)
if err != nil {
s.writeError(w, r, err)
return
}
writeJSON(w, http.StatusOK, toJobProgressResponse(progress))
}
// --- helpers ---
func (s *Server) reqCtx(r *http.Request) (context.Context, context.CancelFunc) {
return context.WithTimeout(r.Context(), s.requestTimeout)
}
func (s *Server) pathUUID(w http.ResponseWriter, r *http.Request, name string) (uuid.UUID, bool) {
id, err := uuid.Parse(r.PathValue(name))
if err != nil {
s.writeError(w, r, domain.ErrInvalidInput)
return uuid.Nil, false
}
return id, true
}
@@ -0,0 +1,103 @@
package http
import (
"context"
"crypto/rand"
"crypto/subtle"
"encoding/hex"
"log/slog"
"net/http"
"strings"
"time"
)
type ctxKey string
const requestIDKey ctxKey = "request_id"
// withRequestID stamps every request with an ID for correlated logs and error
// bodies. It wraps the auth middleware rather than the other way round, so even
// a rejected request carries an ID the caller can quote in a bug report.
func withRequestID(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
id := newRequestID()
w.Header().Set("X-Request-ID", id)
next.ServeHTTP(w, r.WithContext(context.WithValue(r.Context(), requestIDKey, id)))
})
}
func requestIDFrom(ctx context.Context) string {
if v, ok := ctx.Value(requestIDKey).(string); ok {
return v
}
return ""
}
func newRequestID() string {
var b [8]byte
_, _ = rand.Read(b[:])
return hex.EncodeToString(b[:])
}
// withAuth enforces the shared bearer token every worker presents.
// An empty token disables the check (local development only).
func withAuth(token string) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if token == "" {
next.ServeHTTP(w, r)
return
}
presented := strings.TrimPrefix(r.Header.Get("Authorization"), "Bearer ")
// Constant-time compare: a byte-by-byte early exit would let an
// attacker recover the token by timing responses.
if subtle.ConstantTimeCompare([]byte(presented), []byte(token)) != 1 {
w.Header().Set("WWW-Authenticate", "Bearer")
writeJSON(w, http.StatusUnauthorized, errorResponse{
Error: "unauthorized",
RequestID: requestIDFrom(r.Context()),
})
return
}
next.ServeHTTP(w, r)
})
}
}
// statusRecorder captures the status code for the access log.
type statusRecorder struct {
http.ResponseWriter
status int
}
func (s *statusRecorder) WriteHeader(code int) {
s.status = code
s.ResponseWriter.WriteHeader(code)
}
// withAccessLog records one structured line per request — the minimum needed to
// debug a distributed system after the fact.
func withAccessLog(log *slog.Logger) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
start := time.Now()
rec := &statusRecorder{ResponseWriter: w, status: http.StatusOK}
next.ServeHTTP(rec, r)
log.Info("request",
"request_id", requestIDFrom(r.Context()),
"method", r.Method,
"path", r.URL.Path,
"status", rec.status,
"duration_ms", time.Since(start).Milliseconds(),
)
})
}
}
// chain applies middleware so that the first argument is the outermost layer.
func chain(h http.Handler, mw ...func(http.Handler) http.Handler) http.Handler {
for i := len(mw) - 1; i >= 0; i-- {
h = mw[i](h)
}
return h
}
@@ -0,0 +1,60 @@
// Package http adapts the use-case layer to HTTP. Handlers decode requests,
// map them onto use-case inputs, and translate results and errors back — no
// business rules live here.
package http
import (
"log/slog"
"net/http"
"time"
"github.com/emil28092005/SciMesh/coordinator/internal/usecase"
)
// UseCases collects everything the transport needs. Depending on concrete
// use-case types (not one fat interface) keeps each handler's dependency
// explicit and the wiring visible in the composition root.
type UseCases struct {
CreateJob *usecase.CreateJob
ClaimTask *usecase.ClaimTask
RenewLease *usecase.RenewLease
CompleteTask *usecase.CompleteTask
FailTask *usecase.FailTask
GetJobStatus *usecase.GetJobStatus
}
type Server struct {
uc UseCases
log *slog.Logger
requestTimeout time.Duration
}
func NewServer(uc UseCases, log *slog.Logger, requestTimeout time.Duration) *Server {
return &Server{uc: uc, log: log, requestTimeout: requestTimeout}
}
// Handler builds the router. Go 1.22's ServeMux matches on method and path
// wildcards, so no third-party router is needed.
func (s *Server) Handler(token string) http.Handler {
protected := http.NewServeMux()
protected.HandleFunc("POST /jobs", s.handleCreateJob)
protected.HandleFunc("GET /jobs/{job_id}", s.handleGetJob)
protected.HandleFunc("POST /tasks/claim", s.handleClaim)
protected.HandleFunc("POST /tasks/{task_id}/heartbeat", s.handleHeartbeat)
protected.HandleFunc("POST /tasks/{task_id}/result", s.handleResult)
protected.HandleFunc("POST /tasks/{task_id}/failure", s.handleFailure)
mux := http.NewServeMux()
// A more specific pattern wins, so /health stays outside the auth wall.
mux.HandleFunc("GET /health", s.handleHealth)
mux.Handle("/", chain(protected,
withRequestID, // outermost: every response gets an ID,
withAccessLog(s.log), // including the 401s below
withAuth(token),
))
return mux
}
func (s *Server) handleHealth(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]string{"status": "ok"})
}
+51
View File
@@ -0,0 +1,51 @@
package usecase
import "github.com/google/uuid"
// Use-case boundary types. Adapters map their wire formats onto these, so the
// HTTP shape can change without touching business code.
type CreateJobInput struct {
Workload string
InputURI string
Parameters map[string]any
Chunks []ChunkInput
}
type ChunkInput struct {
ChunkIndex int
Workload string
InputURI string
InputSHA256 string
Parameters map[string]any
MaxAttempts int
}
type ClaimTaskInput struct {
WorkerID string
Workloads []string
}
type RenewLeaseInput struct {
TaskID uuid.UUID
WorkerID string
Attempt int
}
type CompleteTaskInput struct {
TaskID uuid.UUID
WorkerID string
Attempt int
ResultURI string
ResultSHA256 string
Metrics map[string]any
}
type FailTaskInput struct {
TaskID uuid.UUID
WorkerID string
Attempt int
ErrorCode string
ErrorMessage string
Retryable bool
}
+174
View File
@@ -0,0 +1,174 @@
package usecase
import (
"context"
"time"
"github.com/google/uuid"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
)
// Job operations: the submitter-facing lifecycle of a whole submission.
//
// CreateJob register a job and fan it out into tasks
// GetJobStatus aggregate progress
// ListResults completed manifests, ordered for the stitcher
// StitchJob merge partial results into the final artifact
// --- CreateJob -----------------------------------------------------------
type CreateJob struct {
jobs JobRepository
tasks TaskRepository
tx TxManager
clock Clock
}
func NewCreateJob(jobs JobRepository, tasks TaskRepository, tx TxManager, clock Clock) *CreateJob {
return &CreateJob{jobs: jobs, tasks: tasks, tx: tx, clock: clock}
}
// Execute builds the job and its tasks, then writes them in one transaction.
// The all-or-none guarantee comes from TxManager: a half-created job would
// leave chunks no worker could ever complete.
func (uc *CreateJob) Execute(ctx context.Context, in CreateJobInput) (*domain.Job, error) {
chunks := make([]domain.ChunkSpec, 0, len(in.Chunks))
for _, c := range in.Chunks {
chunks = append(chunks, domain.ChunkSpec(c))
}
job, tasks, err := domain.NewJobWithTasks(in.Workload, in.InputURI, in.Parameters, chunks, uc.clock.Now())
if err != nil {
return nil, err
}
err = uc.tx.WithinTx(ctx, func(ctx context.Context) error {
if err := uc.jobs.Insert(ctx, job); err != nil {
return err
}
return uc.tasks.InsertBatch(ctx, tasks)
})
if err != nil {
return nil, err
}
return job, nil
}
// --- GetJobStatus --------------------------------------------------------
type GetJobStatus struct {
jobs JobRepository
tasks TaskRepository
}
func NewGetJobStatus(jobs JobRepository, tasks TaskRepository) *GetJobStatus {
return &GetJobStatus{jobs: jobs, tasks: tasks}
}
func (uc *GetJobStatus) Execute(ctx context.Context, jobID uuid.UUID) (domain.JobProgress, error) {
job, err := uc.jobs.Get(ctx, jobID)
if err != nil {
return domain.JobProgress{}, err
}
counts, err := uc.tasks.CountByStatus(ctx, jobID)
if err != nil {
return domain.JobProgress{}, err
}
return progressFrom(*job, counts), nil
}
// --- ListResults ---------------------------------------------------------
type ListResults struct {
tasks TaskRepository
}
func NewListResults(tasks TaskRepository) *ListResults {
return &ListResults{tasks: tasks}
}
// Execute preserves chunk_index order: the stitcher merges these into one
// artifact, and a non-deterministic order would make the final result depend on
// which worker happened to finish first.
func (uc *ListResults) Execute(ctx context.Context, jobID uuid.UUID) ([]domain.ResultManifest, error) {
tasks, err := uc.tasks.ListCompleted(ctx, jobID)
if err != nil {
return nil, err
}
manifests := make([]domain.ResultManifest, 0, len(tasks))
for _, t := range tasks {
if t.ResultURI == nil || t.ResultSHA256 == nil {
continue // a completed task always carries both; skip defensively
}
manifests = append(manifests, domain.ResultManifest{
TaskID: t.ID,
ChunkIndex: t.ChunkIndex,
ResultURI: *t.ResultURI,
ResultSHA256: *t.ResultSHA256,
Metrics: t.Metrics,
})
}
return manifests, nil
}
// --- StitchJob -----------------------------------------------------------
// StitchJob merges every chunk's partial result into the job's final artifact.
// For similarity search that means concatenating each worker's local top-k,
// sorting by similarity, and keeping the global top-k — the distributed result
// must match what a single local run would produce.
type StitchJob struct {
results *ListResults
}
func NewStitchJob(results *ListResults) *StitchJob {
return &StitchJob{results: results}
}
// Execute returns the URI of the assembled artifact.
//
// TODO(phase 6): fetch each manifest's CSV, merge, and persist the result.
func (uc *StitchJob) Execute(ctx context.Context, jobID uuid.UUID) (string, error) {
if _, err := uc.results.Execute(ctx, jobID); err != nil {
return "", err
}
return "", ErrNotImplemented
}
// --- shared helpers ------------------------------------------------------
// progressFrom turns a status histogram into the domain's progress view.
func progressFrom(job domain.Job, counts map[domain.TaskStatus]int) domain.JobProgress {
p := domain.JobProgress{
Job: job,
Pending: counts[domain.TaskPending],
Leased: counts[domain.TaskLeased],
Done: counts[domain.TaskCompleted],
Failed: counts[domain.TaskFailed],
}
for _, n := range counts {
p.Total += n
}
return p
}
// syncJobStatus recomputes a job's status from its task counts and persists it.
// Shared by CompleteTask and FailTask so both close a job by the same rule —
// the rule itself lives in domain.JobProgress.DeriveStatus.
func syncJobStatus(ctx context.Context, jobs JobRepository, tasks TaskRepository,
jobID uuid.UUID, now time.Time) error {
counts, err := tasks.CountByStatus(ctx, jobID)
if err != nil {
return err
}
status := progressFrom(domain.Job{}, counts).DeriveStatus()
var completedAt *time.Time
if status == domain.JobCompleted || status == domain.JobFailed {
completedAt = &now
}
return jobs.UpdateStatus(ctx, jobID, status, completedAt)
}
+80
View File
@@ -0,0 +1,80 @@
// Package usecase holds the application's business operations. Each use case is
// a small type with its dependencies injected and a single Execute method.
//
// The interfaces below are *ports*: they are declared here, by the consumer,
// and implemented further out in storage/postgres. That is what keeps the
// dependency rule intact — usecase never imports storage or transport.
package usecase
import (
"context"
"errors"
"time"
"github.com/google/uuid"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
)
// ClaimFilter narrows which task a worker may be handed.
type ClaimFilter struct {
Workloads []string // workloads this worker can execute
Owner string // worker ID taking the lease
Now time.Time
LeaseUntil time.Time
}
// TaskRepository persists tasks.
//
// ClaimNext is deliberately coarse: leasing must be a single atomic statement
// (SELECT ... FOR UPDATE SKIP LOCKED + UPDATE), so it cannot be decomposed into
// Get+Update without losing the guarantee that one task goes to one worker.
type TaskRepository interface {
// ClaimNext atomically leases one matching pending task.
// Returns (nil, nil) when nothing is available.
ClaimNext(ctx context.Context, f ClaimFilter) (*domain.Task, error)
// GetForUpdate reads a task and locks its row for the enclosing
// transaction, so read-modify-write use cases stay serialized.
GetForUpdate(ctx context.Context, id uuid.UUID) (*domain.Task, error)
// Update persists a mutated task, honouring its Version for optimistic
// concurrency.
Update(ctx context.Context, t *domain.Task) error
InsertBatch(ctx context.Context, tasks []*domain.Task) error
// ListCompleted returns completed tasks ordered by chunk_index.
ListCompleted(ctx context.Context, jobID uuid.UUID) ([]*domain.Task, error)
// CountByStatus aggregates a job's tasks for progress reporting.
CountByStatus(ctx context.Context, jobID uuid.UUID) (map[domain.TaskStatus]int, error)
// ExpireLeases applies the lease-expiry rule to every elapsed task and
// reports how many were affected.
ExpireLeases(ctx context.Context, now time.Time) (int64, error)
}
// JobRepository persists jobs.
type JobRepository interface {
Insert(ctx context.Context, j *domain.Job) error
Get(ctx context.Context, id uuid.UUID) (*domain.Job, error)
UpdateStatus(ctx context.Context, id uuid.UUID, status domain.JobStatus, completedAt *time.Time) error
}
// TxManager runs a function inside one database transaction. The transaction
// travels in the context, so repositories pick it up without this port ever
// mentioning pgx.
type TxManager interface {
WithinTx(ctx context.Context, fn func(ctx context.Context) error) error
}
// Clock supplies the current time. Injecting it keeps lease and expiry rules
// testable without sleeping or freezing the system clock.
type Clock interface {
Now() time.Time
}
// ErrNotImplemented marks scaffold code with no body yet. Unlike the errors in
// domain, it describes the state of this codebase, not a business rule.
var ErrNotImplemented = errors.New("not implemented")
+206
View File
@@ -0,0 +1,206 @@
package usecase
import (
"context"
"time"
"github.com/emil28092005/SciMesh/coordinator/internal/domain"
)
// Task operations: the worker-facing lifecycle of a single chunk.
//
// ClaimTask lease the next available task
// RenewLease extend a held lease (heartbeat)
// CompleteTask record a successful result
// FailTask record a failure
// ExpireLeases reclaim leases that elapsed without a heartbeat
// --- ClaimTask -----------------------------------------------------------
type ClaimTask struct {
tasks TaskRepository
clock Clock
leaseDuration time.Duration
}
func NewClaimTask(tasks TaskRepository, clock Clock, leaseDuration time.Duration) *ClaimTask {
return &ClaimTask{tasks: tasks, clock: clock, leaseDuration: leaseDuration}
}
// Execute reclaims elapsed leases first, then hands out one task.
//
// Sweeping before claiming matters: otherwise a task abandoned by a dead worker
// stays invisible until the reaper's next tick, and a waiting worker is told the
// queue is empty while work sits idle.
//
// This use case is thin by design — the atomicity that makes claiming correct
// lives in one SQL statement behind ClaimNext, and splitting it across the layer
// boundary would break it.
func (uc *ClaimTask) Execute(ctx context.Context, in ClaimTaskInput) (*domain.ClaimedTask, error) {
if in.WorkerID == "" {
return nil, domain.ErrInvalidInput
}
now := uc.clock.Now()
if _, err := uc.tasks.ExpireLeases(ctx, now); err != nil {
return nil, err
}
task, err := uc.tasks.ClaimNext(ctx, ClaimFilter{
Workloads: in.Workloads,
Owner: in.WorkerID,
Now: now,
LeaseUntil: now.Add(uc.leaseDuration),
})
if err != nil {
return nil, err
}
if task == nil {
return nil, nil // empty queue is a normal state, not an error
}
claimed := task.AsClaimed()
return &claimed, nil
}
// --- RenewLease ----------------------------------------------------------
type RenewLease struct {
tasks TaskRepository
tx TxManager
clock Clock
leaseDuration time.Duration
}
func NewRenewLease(tasks TaskRepository, tx TxManager, clock Clock, leaseDuration time.Duration) *RenewLease {
return &RenewLease{tasks: tasks, tx: tx, clock: clock, leaseDuration: leaseDuration}
}
// Execute is a read-modify-write, so it runs inside a transaction with the row
// locked: two concurrent heartbeats must not interleave into a lost update.
// Whether the caller may renew at all is decided by the entity, not here.
func (uc *RenewLease) Execute(ctx context.Context, in RenewLeaseInput) (*domain.ClaimedTask, error) {
var claimed domain.ClaimedTask
err := uc.tx.WithinTx(ctx, func(ctx context.Context) error {
task, err := uc.tasks.GetForUpdate(ctx, in.TaskID)
if err != nil {
return err
}
if err := task.RenewLease(in.WorkerID, in.Attempt, uc.clock.Now().Add(uc.leaseDuration)); err != nil {
return err
}
if err := uc.tasks.Update(ctx, task); err != nil {
return err
}
claimed = task.AsClaimed()
return nil
})
if err != nil {
return nil, err
}
return &claimed, nil
}
// --- CompleteTask --------------------------------------------------------
type CompleteTask struct {
tasks TaskRepository
jobs JobRepository
tx TxManager
clock Clock
}
func NewCompleteTask(tasks TaskRepository, jobs JobRepository, tx TxManager, clock Clock) *CompleteTask {
return &CompleteTask{tasks: tasks, jobs: jobs, tx: tx, clock: clock}
}
// Execute applies the result and, when that was the job's last outstanding
// task, closes the job in the same transaction — so a caller who sees a
// completed task never observes its job still marked running.
//
// Lease ownership, staleness, and idempotent replays are all decided by
// Task.CompleteWith; this use case only orchestrates.
func (uc *CompleteTask) Execute(ctx context.Context, in CompleteTaskInput) (*domain.Task, error) {
var out *domain.Task
err := uc.tx.WithinTx(ctx, func(ctx context.Context) error {
task, err := uc.tasks.GetForUpdate(ctx, in.TaskID)
if err != nil {
return err
}
now := uc.clock.Now()
if err := task.CompleteWith(in.ResultURI, in.ResultSHA256, in.Metrics,
in.WorkerID, in.Attempt, now); err != nil {
return err
}
if err := uc.tasks.Update(ctx, task); err != nil {
return err
}
out = task
return syncJobStatus(ctx, uc.jobs, uc.tasks, task.JobID, now)
})
if err != nil {
return nil, err
}
return out, nil
}
// --- FailTask ------------------------------------------------------------
type FailTask struct {
tasks TaskRepository
jobs JobRepository
tx TxManager
clock Clock
}
func NewFailTask(tasks TaskRepository, jobs JobRepository, tx TxManager, clock Clock) *FailTask {
return &FailTask{tasks: tasks, jobs: jobs, tx: tx, clock: clock}
}
// Execute delegates the requeue-or-terminate decision to Task.Fail, then keeps
// the parent job's status consistent in the same transaction.
func (uc *FailTask) Execute(ctx context.Context, in FailTaskInput) (*domain.Task, error) {
var out *domain.Task
err := uc.tx.WithinTx(ctx, func(ctx context.Context) error {
task, err := uc.tasks.GetForUpdate(ctx, in.TaskID)
if err != nil {
return err
}
now := uc.clock.Now()
if err := task.Fail(in.WorkerID, in.Attempt, in.ErrorCode, in.ErrorMessage, in.Retryable, now); err != nil {
return err
}
if err := uc.tasks.Update(ctx, task); err != nil {
return err
}
out = task
return syncJobStatus(ctx, uc.jobs, uc.tasks, task.JobID, now)
})
if err != nil {
return nil, err
}
return out, nil
}
// --- ExpireLeases --------------------------------------------------------
type ExpireLeases struct {
tasks TaskRepository
clock Clock
}
func NewExpireLeases(tasks TaskRepository, clock Clock) *ExpireLeases {
return &ExpireLeases{tasks: tasks, clock: clock}
}
// Execute reports how many tasks were reclaimed.
//
// The sweep is one set-based statement rather than a load-decide-save loop:
// several coordinators run it concurrently, and a single atomic UPDATE makes
// the duplicate work harmless — the loser simply updates 0 rows.
func (uc *ExpireLeases) Execute(ctx context.Context) (int64, error) {
return uc.tasks.ExpireLeases(ctx, uc.clock.Now())
}