100 lines
2.5 KiB
SQL
100 lines
2.5 KiB
SQL
-- name: CreateWorkflowRun :one
|
|
INSERT INTO workflow_runs (
|
|
id,
|
|
project_id,
|
|
workflow_version_id,
|
|
status,
|
|
trigger_type,
|
|
input,
|
|
idempotency_key
|
|
)
|
|
VALUES ($1, $2, $3, 'pending', $4, $5, $6)
|
|
RETURNING *;
|
|
|
|
-- name: CreateStepRun :one
|
|
INSERT INTO step_runs (
|
|
id,
|
|
project_id,
|
|
workflow_run_id,
|
|
step_key,
|
|
executor_kind,
|
|
max_attempts,
|
|
priority,
|
|
input,
|
|
idempotency_key,
|
|
available_at
|
|
)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)
|
|
RETURNING *;
|
|
|
|
-- name: LeaseNextStepRun :one
|
|
WITH candidate AS (
|
|
SELECT queued.id
|
|
FROM step_runs AS queued
|
|
WHERE queued.executor_kind = sqlc.arg(requested_executor_kind)
|
|
AND queued.status IN ('pending', 'retryable_failed')
|
|
AND queued.available_at <= now()
|
|
AND queued.attempt < queued.max_attempts
|
|
AND (queued.leased_until IS NULL OR queued.leased_until < now())
|
|
ORDER BY queued.priority DESC, queued.created_at
|
|
FOR UPDATE SKIP LOCKED
|
|
LIMIT 1
|
|
)
|
|
UPDATE step_runs AS step
|
|
SET status = 'leased',
|
|
attempt = step.attempt + 1,
|
|
leased_by = sqlc.arg(worker_id),
|
|
leased_until = now() + make_interval(secs => sqlc.arg(lease_seconds)::integer),
|
|
updated_at = now()
|
|
FROM candidate
|
|
WHERE step.id = candidate.id
|
|
RETURNING step.*;
|
|
|
|
-- name: MarkStepRunRunning :one
|
|
UPDATE step_runs
|
|
SET status = 'running',
|
|
started_at = COALESCE(started_at, now()),
|
|
updated_at = now()
|
|
WHERE id = $1
|
|
AND leased_by = $2
|
|
AND status = 'leased'
|
|
RETURNING *;
|
|
|
|
-- name: CompleteStepRun :one
|
|
UPDATE step_runs
|
|
SET status = 'succeeded',
|
|
output = $3,
|
|
leased_by = NULL,
|
|
leased_until = NULL,
|
|
finished_at = now(),
|
|
updated_at = now()
|
|
WHERE id = $1
|
|
AND leased_by = $2
|
|
AND status IN ('leased', 'running')
|
|
RETURNING *;
|
|
|
|
-- name: FailStepRun :one
|
|
UPDATE step_runs
|
|
SET status = CASE
|
|
WHEN attempt < max_attempts AND sqlc.arg(retryable)::boolean
|
|
THEN 'retryable_failed'
|
|
ELSE 'terminal_failed'
|
|
END,
|
|
error = sqlc.arg(error_payload),
|
|
available_at = CASE
|
|
WHEN attempt < max_attempts AND sqlc.arg(retryable)::boolean
|
|
THEN now() + make_interval(secs => sqlc.arg(retry_delay_seconds)::integer)
|
|
ELSE available_at
|
|
END,
|
|
leased_by = NULL,
|
|
leased_until = NULL,
|
|
finished_at = CASE
|
|
WHEN attempt < max_attempts AND sqlc.arg(retryable)::boolean
|
|
THEN NULL
|
|
ELSE now()
|
|
END,
|
|
updated_at = now()
|
|
WHERE id = sqlc.arg(id)
|
|
AND leased_by = sqlc.arg(worker_id)
|
|
AND status IN ('leased', 'running')
|
|
RETURNING *;
|