Oban.Make.Runtime

oban · API reference

Run one structured fiber per configured concurrency slot.

type t
type queue_config = {
  queue : Job.queue;
  concurrency : int;
}
type limits = {
  max_queues : int;
  max_queue_concurrency : int;
  max_total_concurrency : int;
}

Positive fiber-allocation safety bounds.

val default_limits : limits
type config
type config_error = [ 
  | `Empty_node
  | `Invalid_node_utf8
  | `Node_contains_null
  | `Node_too_long of int
  | `Empty_queues
  | `Too_many_queues of int
  | `Too_many_slots of int
  | `Duplicate_queue of Job.queue
  | `Invalid_concurrency of int
  | `Invalid_poll_interval of float
  | `Duplicate_worker of Job.worker_name
  | `Invalid_runtime_limit of string * int
 ]
type event = 
  | Enqueued of Job.packed
  | Started of Job.executing Job.t
  | Finished of Job.packed
  | Stale of Job.packed option
type error = [ 
  | `Storage of Store.error
  | `Invalid_state of Job.state
 ]
type execution = [ 
  | `Completed
  | `Failed
  | `Cancelled
  | `Discarded
  | `Snoozed
  | `Stale
 ]
exception Loop_error of error
val configure : 
  node:string ->
  queues:queue_config list ->
  poll_interval:float ->
  (config, config_error) result

Configure a non-empty UTF-8 node without U+0000, bounded by Limits.max_owner_bytes, with default_limits.

val configure_with_limits : 
  limits:limits ->
  node:string ->
  queues:queue_config list ->
  poll_interval:float ->
  (config, config_error) result

Configure with explicit positive runtime limits.

val create : 
  ?on_event:(event -> unit) ->
  config:config ->
  store:Store.t ->
  workers:Worker.packed list ->
  unit ->
  (t, config_error) result
val enqueue : t -> Job.spec -> (Storage.enqueue, error) result
val enqueue_unique : t -> Job.spec -> (Storage.enqueue, error) result
val enqueue_many : 
  t ->
  Storage.enqueue_request list ->
  (Storage.enqueue list, error) result

Atomically enqueue an ordered batch and emit events after commit.

val run_one : 
  clock:_ Eio.Time.clock ->
  wall_now:(unit -> Ptime.t) ->
  t ->
  queue:Job.queue ->
  ([ `Idle | `Worked of execution ], error) result
val run : clock:_ Eio.Time.clock -> wall_now:(unit -> Ptime.t) -> t -> unit