Oban.Make.Store
oban · API reference
type t = Storage.Durable(D).ttype nonrec settings = Storage.settingsval default_settings : settingsval default_pool_size : inttype 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) resultCreate a store. The default pool has 10 connections and default_settings bounds each operation class at 500 jobs.
val close : t -> unitval migrate : t -> (unit, error) resultval transaction :
t ->
(D.connection -> ('a, error) result) ->
('a, error) resultRun 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) resultEnqueue synchronously inside transaction, passing the same store.
val enqueue_many_with_conn :
t ->
D.connection ->
Storage.enqueue_request list ->
(Storage.enqueue list, error) resultEnqueue an ordered batch bounded by settings.max_batch_size inside transaction.
val enqueue : t -> Storage.enqueue_request -> (Storage.enqueue, error) resultval enqueue_many :
t ->
Storage.enqueue_request list ->
(Storage.enqueue list, error) resultAtomically enqueue an ordered batch bounded by settings.max_batch_size.
val fetch : t -> Job.id -> (Job.packed option, error) resultval list :
t ->
filter:Storage.filter ->
after:Storage.cursor option ->
limit:int ->
(Storage.page, error) resultList 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) resultRejects 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) resultval cancel :
t ->
target:Storage.row_ref ->
now:Ptime.t ->
reason:string ->
(Storage.write, error) resultval retry :
t ->
target:Storage.row_ref ->
now:Ptime.t ->
(Storage.write, error) resultval cancel_many :
t ->
targets:Storage.row_ref list ->
now:Ptime.t ->
reason:string ->
(Storage.write list, error) resultAtomically 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) resultAtomically 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) resultRescue at most settings.max_maintenance_batch_size jobs.
val prune :
t ->
states:Job.state list ->
before:Ptime.t ->
limit:int ->
(Job.id list, error) resultPrune at most settings.max_maintenance_batch_size jobs.
val pause : t -> Job.queue -> (unit, error) resultval resume : t -> Job.queue -> (unit, error) result