Belay.Storage behaviour (Belay v1.0.0)

Copy Markdown View Source

The storage contract: coarse, semantic, individually-atomic operations.

Engine semantics that must be transactional (admission control inside claim, workflow settlement inside ack) live inside single operations, so adapters guarantee atomicity per call. Two adapters ship: Belay.Storage.Memory (serialized GenServer, deterministic — the test/simulation reference) and Belay.Storage.Postgres.

All time-dependent operations take now explicitly; adapters never read the wall clock.

Summary

Types

now()

@type now() :: DateTime.t()

outcome()

@type outcome() ::
  {:succeeded, binary() | nil}
  | {:retry, map(), now()}
  | {:failed, map()}
  | {:cancelled, map()}
  | {:snooze, now()}
  | {:await, String.t(), String.t(), now() | nil}

queue_spec()

@type queue_spec() :: %{
  queue: String.t(),
  local_limit: pos_integer(),
  limit_min: pos_integer() | nil,
  global_limit: pos_integer() | nil,
  rate: %{allowed: pos_integer(), period: pos_integer()} | nil,
  partition: {:input | :meta, String.t()} | nil
}

ref()

@type ref() :: term()

settle_result()

@type settle_result() :: %{
  job: Belay.Job.t(),
  released: [Belay.Job.t()],
  cancelled: [Belay.Job.t()]
}

Callbacks

ack(ref, t, outcome, now)

@callback ack(ref(), Belay.Job.t(), outcome(), now()) ::
  {:ok, settle_result()} | {:error, :stale}

append_event(ref, integer, map, now)

@callback append_event(ref(), integer(), map(), now()) :: {:ok, seq :: integer()}

child_spec(tuple)

@callback child_spec({Belay.Config.t(), keyword()}) :: Supervisor.child_spec()

children(ref, integer)

@callback children(ref(), integer()) :: {:ok, [Belay.Job.t()]}

claim(ref, queue_spec, demand, node_id, lease_ttl_ms, now)

@callback claim(
  ref(),
  queue_spec(),
  demand :: pos_integer(),
  node_id :: String.t(),
  lease_ttl_ms :: pos_integer(),
  now()
) :: {:ok, [Belay.Job.t()]}

clear_signal(ref, t, t)

@callback clear_signal(ref(), String.t(), String.t()) :: :ok

debit_rate(ref, bucket, period, amount, now)

@callback debit_rate(
  ref(),
  bucket :: String.t(),
  period :: pos_integer(),
  amount :: integer(),
  now()
) :: :ok

delete_dynamic_cron(ref, t)

@callback delete_dynamic_cron(ref(), String.t()) :: :ok

delete_dynamic_queue(ref, t)

@callback delete_dynamic_queue(ref(), String.t()) :: :ok

get_by_unique_key(ref, t)

@callback get_by_unique_key(ref(), String.t()) :: {:ok, Belay.Job.t()} | :error

get_job(ref, integer)

@callback get_job(ref(), integer()) :: {:ok, Belay.Job.t()} | :error

get_signal(ref, list, t)

@callback get_signal(ref(), [String.t()], String.t()) :: {:ok, map()} | :none

get_step(ref, integer, t)

@callback get_step(ref(), integer(), String.t()) :: {:ok, binary()} | :none

insert_jobs(ref, list, now)

@callback insert_jobs(ref(), [map()], now()) :: {:ok, [Belay.Job.t()]}

list_dynamic_crons(ref)

@callback list_dynamic_crons(ref()) :: {:ok, [map()]}

list_dynamic_queues(ref)

@callback list_dynamic_queues(ref()) :: {:ok, [{String.t(), map()}]}

list_events(ref, integer, after_seq)

@callback list_events(ref(), integer(), after_seq :: integer()) :: {:ok, [map()]}

list_jobs(ref, map)

@callback list_jobs(ref(), map()) :: {:ok, [Belay.Job.t()]}

list_steps(ref, integer)

@callback list_steps(ref(), integer()) :: {:ok, [map()]}

prune_jobs(ref, state, now, keep_seconds, limit)

@callback prune_jobs(
  ref(),
  state :: String.t(),
  now(),
  keep_seconds :: integer(),
  limit :: pos_integer()
) :: {:ok, non_neg_integer()}

prune_rate(ref, before_unix)

@callback prune_rate(ref(), before_unix :: integer()) :: :ok

prune_signals(ref, now, ttl_seconds)

@callback prune_signals(ref(), now(), ttl_seconds :: integer()) :: :ok

put_dynamic_cron(ref, map, now)

@callback put_dynamic_cron(ref(), map(), now()) :: :ok

put_dynamic_queue(ref, t, map, now)

@callback put_dynamic_queue(ref(), String.t(), map(), now()) :: :ok

put_signal(ref, t, t, map, now)

@callback put_signal(ref(), String.t(), String.t(), map(), now()) ::
  {:ok, woken :: [Belay.Job.t()]}

put_step(ref, integer, t, binary, map, now)

@callback put_step(
  ref(),
  integer(),
  String.t(),
  binary(),
  %{usd_micros: integer(), tokens: integer()},
  now()
) :: {:ok, %{spent_usd_micros: integer(), spent_tokens: integer()}}

queue_stats(ref)

@callback queue_stats(ref()) ::
  {:ok,
   [
     %{
       queue: String.t(),
       state: String.t(),
       count: integer(),
       usd_micros: integer(),
       tokens: integer()
     }
   ]}

reclaim_expired(ref, now, function)

@callback reclaim_expired(ref(), now(), (Belay.Job.t() -> now())) ::
  {:ok, %{retried: [integer()], failed: [integer()]}}

renew_leases(ref, list, t, now)

@callback renew_leases(ref(), [integer()], String.t(), now()) :: {:ok, [integer()]}

request_cancel(ref, integer, now)

@callback request_cancel(ref(), integer(), now()) ::
  {:ok,
   %{
     status: :cancelled | :requested | :noop,
     cancelled: [Belay.Job.t()],
     released: [Belay.Job.t()]
   }}

resettle_parents(ref, now)

@callback resettle_parents(ref(), now()) :: {:ok, [Belay.Job.t()]}

retry(ref, integer, now)

@callback retry(ref(), integer(), now()) ::
  {:ok, Belay.Job.t()} | {:error, :not_retryable | :not_found}

set_cron_paused(ref, t, boolean)

@callback set_cron_paused(ref(), String.t(), boolean()) :: :ok

workflow_jobs(ref, t)

@callback workflow_jobs(ref(), String.t()) :: {:ok, [Belay.Job.t()]}