// Path: Backend/internal/jobs/store.go package jobs import ( "context" "encoding/json" "strings" "time" "github.com/jackc/pgx/v5" "moku-backend/internal/database" ) const ( KindBootstrapStructureMaterialize = "bootstrap.structure.materialize" ) type Status string const ( StatusPending Status = "pending" StatusRunning Status = "running" StatusSucceeded Status = "succeeded" StatusFailed Status = "failed" ) type BootstrapStructureMaterializePayload struct { InstallationID string `json:"installationId"` } type Job struct { ID string Kind string Status Status Payload json.RawMessage Attempts int MaxAttempts int AvailableAt time.Time StartedAt *time.Time FinishedAt *time.Time LastError *string CreatedAt time.Time UpdatedAt time.Time } type EnqueueInput struct { Kind string Payload any AvailableAt time.Time MaxAttempts int } type Store struct { db *database.DB } func NewStore(db *database.DB) *Store { return &Store{db: db} } func (store *Store) Enqueue(ctx context.Context, input EnqueueInput) (Job, error) { payload := json.RawMessage([]byte(`{}`)) if input.Payload != nil { encoded, err := json.Marshal(input.Payload) if err != nil { return Job{}, err } payload = encoded } availableAt := input.AvailableAt if availableAt.IsZero() { availableAt = time.Now().UTC() } maxAttempts := input.MaxAttempts if maxAttempts < 1 { maxAttempts = 1 } return scanJob(store.db.Pool.QueryRow(ctx, ` INSERT INTO background_jobs (kind, status, payload, attempts, max_attempts, available_at) VALUES ($1, 'pending'::background_job_status, $2::jsonb, 0, $3, $4) RETURNING id::text, kind, status::text, payload, attempts, max_attempts, available_at, started_at, finished_at, last_error, created_at, updated_at; `, strings.TrimSpace(input.Kind), payload, maxAttempts, availableAt)) } func (store *Store) ClaimNext(ctx context.Context) (*Job, error) { job, err := scanJob(store.db.Pool.QueryRow(ctx, ` WITH next_job AS ( SELECT id FROM background_jobs WHERE status = 'pending'::background_job_status AND available_at <= NOW() ORDER BY created_at ASC LIMIT 1 FOR UPDATE SKIP LOCKED ) UPDATE background_jobs AS jobs SET status = 'running'::background_job_status, attempts = jobs.attempts + 1, started_at = NOW(), finished_at = NULL, last_error = NULL, updated_at = NOW() FROM next_job WHERE jobs.id = next_job.id RETURNING jobs.id::text, jobs.kind, jobs.status::text, jobs.payload, jobs.attempts, jobs.max_attempts, jobs.available_at, jobs.started_at, jobs.finished_at, jobs.last_error, jobs.created_at, jobs.updated_at; `)) if err != nil { if err == pgx.ErrNoRows { return nil, nil } return nil, err } return &job, nil } func (store *Store) MarkSucceeded(ctx context.Context, jobID string) error { _, err := store.db.Pool.Exec(ctx, ` UPDATE background_jobs SET status = 'succeeded'::background_job_status, finished_at = NOW(), last_error = NULL, updated_at = NOW() WHERE id = $1::uuid; `, strings.TrimSpace(jobID)) return err } func (store *Store) MarkFailed(ctx context.Context, jobID, failure string) error { _, err := store.db.Pool.Exec(ctx, ` UPDATE background_jobs SET status = 'failed'::background_job_status, finished_at = NOW(), last_error = $2, updated_at = NOW() WHERE id = $1::uuid; `, strings.TrimSpace(jobID), strings.TrimSpace(failure)) return err } func scanJob(row pgx.Row) (Job, error) { var job Job var status string if err := row.Scan( &job.ID, &job.Kind, &status, &job.Payload, &job.Attempts, &job.MaxAttempts, &job.AvailableAt, &job.StartedAt, &job.FinishedAt, &job.LastError, &job.CreatedAt, &job.UpdatedAt, ); err != nil { return Job{}, err } job.Status = Status(status) return job, nil }