Bosatsu packages

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)