Refactor: improve code quality and worker flow
This commit is contained in:
@@ -0,0 +1,196 @@
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user