Oban.Make.Store

oban · API reference

type t = Storage.Durable(D).t
type nonrec settings = Storage.settings
val default_settings : settings
val default_pool_size : int
type error = [ 
  | `Config of string
  | `Decode of string
  | `Driver of string
  | `Incompatible_schema
  | `Invalid_limit of int
  | `Invalid_node
  | `Pool of string
  | `Unique_conflict
  | `Unsupported_dialect of string
 ]
val connect : 
  ?pool_size:int ->
  ?settings:settings ->
  string ->
  (t, error) result

Create a store. The default pool has 10 connections and default_settings bounds each operation class at 500 jobs.

val close : t -> unit
val migrate : t -> (unit, error) result
val transaction : 
  t ->
  (D.connection -> ('a, error) result) ->
  ('a, error) result

Run a non-Eio callback on one pooled connection and roll back on Error or exceptions. Inside, use only enqueue_with_conn / enqueue_many_with_conn; other Store calls use another connection and may deadlock when pool_size is 1.

val enqueue_with_conn : 
  t ->
  D.connection ->
  Storage.enqueue_request ->
  (Storage.enqueue, error) result

Enqueue synchronously inside transaction, passing the same store.

val enqueue_many_with_conn : 
  t ->
  D.connection ->
  Storage.enqueue_request list ->
  (Storage.enqueue list, error) result

Enqueue an ordered batch bounded by settings.max_batch_size inside transaction.

val enqueue : t -> Storage.enqueue_request -> (Storage.enqueue, error) result
val enqueue_many : 
  t ->
  Storage.enqueue_request list ->
  (Storage.enqueue list, error) result

Atomically enqueue an ordered batch bounded by settings.max_batch_size.

val fetch : t -> Job.id -> (Job.packed option, error) result
val list : 
  t ->
  filter:Storage.filter ->
  after:Storage.cursor option ->
  limit:int ->
  (Storage.page, error) result

List in ascending ID order, bounded by settings.max_page_size. Every requested tag must be present.

val claim_one : 
  t ->
  queue:Job.queue ->
  now:Ptime.t ->
  node:string ->
  (Job.executing Job.t option, error) result

Rejects an empty or invalid UTF-8 node, or one exceeding Limits.max_owner_bytes.

val finish : 
  t ->
  now:Ptime.t ->
  Job.executing Job.t ->
  Storage.finish ->
  (Storage.write, error) result
val cancel : 
  t ->
  target:Storage.row_ref ->
  now:Ptime.t ->
  reason:string ->
  (Storage.write, error) result
val retry : 
  t ->
  target:Storage.row_ref ->
  now:Ptime.t ->
  (Storage.write, error) result
val cancel_many : 
  t ->
  targets:Storage.row_ref list ->
  now:Ptime.t ->
  reason:string ->
  (Storage.write list, error) result

Atomically cancel an ordered revision-aware batch bounded by settings.max_batch_size. Stale targets are values, not batch errors.

val retry_many : 
  t ->
  targets:Storage.row_ref list ->
  now:Ptime.t ->
  (Storage.write list, error) result

Atomically retry an ordered revision-aware batch bounded by settings.max_batch_size. Stale targets are values, not batch errors.

val rescue : 
  t ->
  before:Ptime.t ->
  now:Ptime.t ->
  error:string ->
  limit:int ->
  (Job.packed list, error) result

Rescue at most settings.max_maintenance_batch_size jobs.

val prune : 
  t ->
  states:Job.state list ->
  before:Ptime.t ->
  limit:int ->
  (Job.id list, error) result

Prune at most settings.max_maintenance_batch_size jobs.

val pause : t -> Job.queue -> (unit, error) result
val resume : t -> Job.queue -> (unit, error) result