197 lines
3.8 KiB
Go
197 lines
3.8 KiB
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
|
|
}
|