Yichus/Job
Builtin package (resource job.bosatsu).
This is the complete package source used by this build. The export
list names its public API; definitions below give the types and behavior.
New to the API? Start with Make a calculator or Build an API, then use this page to look up a definition.
package Yichus/Job
from Bosatsu/Predef import (
Int, String, Bool, List, Option, Some, None,
True, False, LT, EQ, add, cmp_Int, eq_Int, eq_String
)
export (
Job, JobState(), LeaseToken(), Claim(), JobError(),
JobResult(), ClaimResult(),
new_job,
job_id, job_stage, job_input_ref, job_input_revision,
job_dependencies, job_attempt, job_max_attempts,
job_fencing_version, job_state,
dependency_error, lease_error,
claim_transition, complete_transition, fail_transition,
permanent_fail_transition, expire_transition,
cancel_transition, operator_retry_transition
)
enum JobState:
Waiting(available_at: Int)
RetryWaiting(reason: String, available_at: Int)
Leased(worker_id: String, fencing_version: Int, expires_at: Int)
Succeeded(output_ref: String, completed_at: Int)
Failed(reason: String, failed_at: Int)
Cancelled(reason: String, cancelled_at: Int)
# Job construction is intentionally opaque. Callers can inspect a Job and
# pattern-match its state, but every Job value starts through new_job or one
# of the checked transitions below.
struct Job(
id: String,
stage: String,
input_ref: String,
input_revision: Int,
dependencies: List[String],
attempt: Int,
max_attempts: Int,
fencing_version: Int,
state: JobState
)
# A worker receives this token from a successful claim. Completion and
# failure compare every authority-bearing field with the current Job row and
# require the authenticated worker id that claimed it.
struct LeaseToken(
job_id: String,
input_revision: Int,
worker_id: String,
fencing_version: Int,
expires_at: Int
)
struct Claim(job: Job, token: LeaseToken)
enum JobError:
MissingJob(job_id: String)
DuplicateJob(job_id: String)
MissingDependency(job_id: String)
IncompleteDependency(job_id: String)
InvalidJob(reason: String)
NotAvailable(available_at: Int)
LeaseHeld(expires_at: Int)
AttemptLimitReached
TerminalJob
StaleLease
LeaseExpired(expires_at: Int)
NotFailed
FenceChanged
enum JobResult:
JobApplied(job: Job)
JobRefused(error: JobError)
enum ClaimResult:
ClaimApplied(claim: Claim)
ClaimRefused(error: JobError)
def nonnegative(value: Int) -> Bool:
match cmp_Int(value, 0):
case LT: False
case _: True
def positive(value: Int) -> Bool:
match cmp_Int(value, 0):
case LT | EQ: False
case _: True
def nonempty(value: String) -> Bool:
if eq_String(value, ""): False
else: True
def job_id(job: Job) -> String:
Job(id, _, _, _, _, _, _, _, _) = job
id
def job_stage(job: Job) -> String:
Job(_, stage, _, _, _, _, _, _, _) = job
stage
def job_input_ref(job: Job) -> String:
Job(_, _, input_ref, _, _, _, _, _, _) = job
input_ref
def job_input_revision(job: Job) -> Int:
Job(_, _, _, input_revision, _, _, _, _, _) = job
input_revision
def job_dependencies(job: Job) -> List[String]:
Job(_, _, _, _, dependencies, _, _, _, _) = job
dependencies
def job_attempt(job: Job) -> Int:
Job(_, _, _, _, _, attempt, _, _, _) = job
attempt
def job_max_attempts(job: Job) -> Int:
Job(_, _, _, _, _, _, max_attempts, _, _) = job
max_attempts
def job_fencing_version(job: Job) -> Int:
Job(_, _, _, _, _, _, _, fencing_version, _) = job
fencing_version
def job_state(job: Job) -> JobState:
Job(_, _, _, _, _, _, _, _, state) = job
state
def dependency_error(dependency_id: String, dependency: Option[Job]) -> Option[JobError]:
match dependency:
case None: Some(MissingDependency(dependency_id))
case Some(job):
match job_state(job):
case Succeeded(_, _): None
case _: Some(IncompleteDependency(dependency_id))
def new_job(
id: String,
stage: String,
input_ref: String,
input_revision: Int,
dependencies: List[String],
max_attempts: Int,
available_at: Int
) -> JobResult:
if nonempty(id):
if nonempty(stage):
if nonempty(input_ref):
if nonnegative(input_revision):
if positive(max_attempts):
if nonnegative(available_at):
JobApplied(Job(
id, stage, input_ref, input_revision, dependencies,
0, max_attempts, 0, Waiting(available_at)
))
else: JobRefused(InvalidJob("available_at must be nonnegative"))
else: JobRefused(InvalidJob("max_attempts must be positive"))
else: JobRefused(InvalidJob("input_revision must be nonnegative"))
else: JobRefused(InvalidJob("input_ref must not be empty"))
else: JobRefused(InvalidJob("stage must not be empty"))
else: JobRefused(InvalidJob("id must not be empty"))
def claim_ready(
job: Job,
worker_id: String,
now: Int,
lease_duration: Int,
available_at: Int
) -> ClaimResult:
Job(id, stage, input_ref, input_revision, dependencies, attempt, max_attempts, fencing_version, _) = job
if nonempty(worker_id):
if positive(lease_duration):
match cmp_Int(now, available_at):
case LT: ClaimRefused(NotAvailable(available_at))
case _:
match cmp_Int(attempt, max_attempts):
case LT:
next_attempt = add(attempt, 1)
next_fence = add(fencing_version, 1)
expires_at = add(now, lease_duration)
updated = Job(
id, stage, input_ref, input_revision, dependencies,
next_attempt, max_attempts, next_fence,
Leased(worker_id, next_fence, expires_at)
)
ClaimApplied(Claim(
updated,
LeaseToken(id, input_revision, worker_id, next_fence, expires_at)
))
case _: ClaimRefused(AttemptLimitReached)
else: ClaimRefused(InvalidJob("lease_duration must be positive"))
else: ClaimRefused(InvalidJob("worker_id must not be empty"))
def claim_transition(
job: Job,
worker_id: String,
now: Int,
lease_duration: Int
) -> ClaimResult:
match job_state(job):
case Waiting(available_at):
claim_ready(job, worker_id, now, lease_duration, available_at)
case RetryWaiting(_, available_at):
claim_ready(job, worker_id, now, lease_duration, available_at)
case Leased(_, _, expires_at):
match cmp_Int(now, expires_at):
case LT: ClaimRefused(LeaseHeld(expires_at))
case _: claim_ready(job, worker_id, now, lease_duration, now)
case _: ClaimRefused(TerminalJob)
def lease_error(job: Job, token: LeaseToken, worker_id: String, now: Int) -> Option[JobError]:
Job(id, _, _, input_revision, _, _, _, _, state) = job
LeaseToken(token_job_id, token_input_revision, token_worker, token_fence, token_expires) = token
if eq_String(worker_id, token_worker):
if eq_String(id, token_job_id):
if eq_Int(input_revision, token_input_revision):
match state:
case Leased(worker, fence, expires_at):
if eq_String(worker, token_worker):
if eq_Int(fence, token_fence):
if eq_Int(expires_at, token_expires):
match cmp_Int(now, expires_at):
case LT: None
case _: Some(LeaseExpired(expires_at))
else: Some(StaleLease)
else: Some(StaleLease)
else: Some(StaleLease)
case _: Some(StaleLease)
else: Some(StaleLease)
else: Some(StaleLease)
else: Some(StaleLease)
def with_valid_lease(
job: Job,
token: LeaseToken,
worker_id: String,
now: Int,
update: Job -> JobResult
) -> JobResult:
match lease_error(job, token, worker_id, now):
case Some(error): JobRefused(error)
case None: update(job)
def complete_transition(job: Job, token: LeaseToken, worker_id: String, output_ref: String, now: Int) -> JobResult:
def complete(current):
if nonempty(output_ref):
Job(id, stage, input_ref, input_revision, dependencies, attempt, max_attempts, fencing_version, _) = current
JobApplied(Job(
id, stage, input_ref, input_revision, dependencies,
attempt, max_attempts, fencing_version,
Succeeded(output_ref, now)
))
else: JobRefused(InvalidJob("output_ref must not be empty"))
with_valid_lease(job, token, worker_id, now, complete)
def retry_or_fail(job: Job, reason: String, now: Int, retry_delay: Int) -> JobResult:
Job(id, stage, input_ref, input_revision, dependencies, attempt, max_attempts, fencing_version, _) = job
next_state = match cmp_Int(attempt, max_attempts):
case LT: RetryWaiting(reason, add(now, retry_delay))
case _: Failed(reason, now)
JobApplied(Job(
id, stage, input_ref, input_revision, dependencies,
attempt, max_attempts, fencing_version, next_state
))
def fail_transition(job: Job, token: LeaseToken, worker_id: String, reason: String, now: Int, retry_delay: Int) -> JobResult:
def fail(current):
if nonnegative(retry_delay): retry_or_fail(current, reason, now, retry_delay)
else: JobRefused(InvalidJob("retry_delay must be nonnegative"))
with_valid_lease(job, token, worker_id, now, fail)
# A recovery plan supplies the observed fence and database time in the same
# transaction as the job read and write. It cannot expire a replacement lease.
def expire_transition(job: Job, expected_fence: Int, now: Int) -> JobResult:
if eq_Int(job_fencing_version(job), expected_fence):
match job_state(job):
case Leased(_, _, expires_at):
match cmp_Int(now, expires_at):
case LT: JobRefused(LeaseHeld(expires_at))
case _: retry_or_fail(job, "worker lease expired", now, 0)
case Waiting(_) | RetryWaiting(_, _): JobRefused(StaleLease)
case Succeeded(_, _) | Failed(_, _) | Cancelled(_, _): JobRefused(TerminalJob)
else: JobRefused(FenceChanged)
def permanent_fail_transition(
job: Job,
token: LeaseToken,
worker_id: String,
reason: String,
now: Int
) -> JobResult:
def fail(current):
if nonempty(reason):
Job(id, stage, input_ref, input_revision, dependencies, attempt, max_attempts, fencing_version, _) = current
JobApplied(Job(
id, stage, input_ref, input_revision, dependencies,
attempt, max_attempts, fencing_version, Failed(reason, now)
))
else: JobRefused(InvalidJob("reason must not be empty"))
with_valid_lease(job, token, worker_id, now, fail)
def cancel_transition(job: Job, expected_fence: Int, reason: String, now: Int) -> JobResult:
Job(id, stage, input_ref, input_revision, dependencies, attempt, max_attempts, fencing_version, state) = job
if eq_Int(fencing_version, expected_fence):
match state:
case Succeeded(_, _) | Failed(_, _) | Cancelled(_, _): JobRefused(TerminalJob)
case _:
JobApplied(Job(
id, stage, input_ref, input_revision, dependencies,
attempt, max_attempts, fencing_version, Cancelled(reason, now)
))
else: JobRefused(FenceChanged)
def operator_retry_transition(job: Job, expected_fence: Int, reason: String, now: Int) -> JobResult:
Job(id, stage, input_ref, input_revision, dependencies, attempt, max_attempts, fencing_version, state) = job
if eq_Int(fencing_version, expected_fence):
match state:
case Failed(_, _):
JobApplied(Job(
id, stage, input_ref, input_revision, dependencies,
attempt, add(max_attempts, 1), fencing_version,
RetryWaiting(reason, now)
))
case _: JobRefused(NotFailed)
else: JobRefused(FenceChanged)