oban — durable background jobs

Run durable background jobs in OCaml with oban, using typed workers, persistent queues, retries and scheduled upkeep.

API reference

oban is how slow work leaves the request path. You describe a job as a small spec — a queue, a worker name, JSON arguments — and enqueue it; the job is stored as a durable row through repodb on SQLite, PostgreSQL or MariaDB. A runtime then claims available jobs in priority and schedule order and hands each one to the registered worker, one fiber per concurrency slot. If the machine restarts mid-queue, the jobs are still there. The rest of this page walks through that loop piece by piece, with runnable examples.

Jobs are specs, then rows

Queue and worker names are parsed values rather than plain strings: non-empty UTF-8 without null bytes, bounded by Oban.Limits. A spec adds JSON arguments and metadata, up to 20 sorted-and-deduplicated tags, a priority from 0 (highest) through 9, an attempt limit, and an optional schedule time. Enqueueing inserts the spec as an available row; a unique request returns the existing active job instead of a duplicate.

examples/oban_quickstart/main.ml
(* A spec is an unpersisted job: queue, worker, JSON args, and limits.
   Enqueueing inserts it as an available row. *)
let args = Simdjsont.Json.Object [ ("user_id", Simdjsont.Json.Int 42L) ] in
let spec =
  get
    (Oban.Job.create ~now:(wall_now ()) ~queue ~worker:worker_name ~args ())
in
(match Runtime.enqueue runtime spec with
| Ok (Oban.Storage.Inserted job) ->
    Printf.printf "enqueued job #%Ld\n%!"
      (Oban.Job.id_to_int64 (Oban.Job.view job).id)
| Ok (Oban.Storage.Existing _) ->
    prerr_endline "error: unexpected existing job";
    exit 1
| Error _ ->
    prerr_endline "error: enqueue failed";
    exit 1);

Tags travel with the job through its whole lifecycle and can be used in Store.list filters, where every tag in the filter must be present. Metadata and tags do not participate in the uniqueness key.

Workers return outcomes

A worker is a module implementing Oban.Worker.S: a name, a perform function over a claimed job, a backoff for the delay after an Error, and a timeout per attempt. The outcome is a value, never an exception:

examples/oban_quickstart/main.ml
(* A worker handles a claimed job and returns its outcome. [perform]
   runs in a runtime fiber; [backoff] sets the retry delay after an
   [Error]; [timeout] bounds a single attempt. *)
module Email = struct
  let name = worker_name

  let perform job =
    let args = (Oban.Job.view job).args in
    (match args with
    | Simdjsont.Json.Object fields -> (
        match List.assoc_opt "user_id" fields with
        | Some (Simdjsont.Json.Int id) ->
            Printf.printf "sending email to user %Ld\n%!" id
        | _ -> print_endline "sending email (no user_id)")
    | _ -> print_endline "sending email (no args)");
    Oban.Worker.Ok

  let backoff = Oban.Worker.default_backoff
  let timeout = Oban.Worker.default_timeout
end

Ok completes the job. Error message schedules a retry at now plus backoff, or discards the job after its final attempt. Snooze seconds makes the job available again at the current time or later without using an attempt. Cancel reason finishes it as cancelled. Worker.default_backoff starts at 15 seconds, grows exponentially, and caps at one day; Worker.default_timeout is unlimited, and a timeout is recorded as a worker error under the normal attempt policy.

Runtimes claim and run

One backend functor binds the whole stack to a driver; the store and the runtime come out of it. Open the store once, migrate it, then configure queues and register workers:

examples/oban_quickstart/main.ml
(* A temp-file DB so every pooled connection sees the same rows
   (a pool over ":memory:" would give each connection its own DB).
   Migrate before enqueueing or running anything. *)
let db = Filename.temp_file "oban_quickstart" ".sqlite" in
Fun.protect ~finally:(fun () -> (try Sys.remove db with Sys_error _ -> ()))
@@ fun () ->
let store = get (Store.connect ~pool_size:1 db) in
Fun.protect ~finally:(fun () -> Store.close store) @@ fun () ->
get (Store.migrate store);
(* One queue, two worker fibers; the poll interval only matters for
   [Runtime.run], which sleeps between empty polls. *)
let config =
  get
    (Runtime.configure ~node:"mailer-1"
       ~queues:[ { queue; concurrency = 2 } ]
       ~poll_interval:0.01)
in
let runtime =
  get
    (Runtime.create ~config ~store
       ~workers:[ Oban.Worker.pack (module Email) ]
       ())
in

Store.migrate must run before anything else. Each queue gets a concurrency count and each slot claims one job immediately before running it, so a worker never holds a job it is not running yet. Claim order within a queue is priority, then schedule time, then id.

Runtime.run_one claims at most one job on a queue and runs it to its recorded outcome; Runtime.run loops that across every configured slot until its scope ends. The quickstart drains with run_one so the program terminates:

examples/oban_quickstart/main.ml
(* [run_one] claims one available job on [queue] and runs it: Idle
   when there is nothing to do, otherwise the execution outcome.
   Loop until the queue stays idle. *)
let rec drain idle =
  match Runtime.run_one ~clock ~wall_now runtime ~queue with
  | Error _ ->
      prerr_endline "error: run failed";
      exit 1
  | Ok `Idle ->
      if idle >= 50 then print_endline "queue idle: done"
      else (
        Eio.Time.sleep clock 0.01;
        drain (idle + 1))
  | Ok (`Worked outcome) ->
      Printf.printf "worked: %s\n%!"
        (match outcome with
        | `Completed -> "completed"
        | `Failed -> "failed (will retry)"
        | `Cancelled -> "cancelled"
        | `Discarded -> "discarded"
        | `Snoozed -> "snoozed"
        | `Stale -> "stale");
      drain 0
in
drain 0;
$ dune exec examples/oban_quickstart/main.exe

Retries, snoozes, cancellations

Four workers show the four fates. A failing worker with an immediate backoff retries once and then discards; a snoozer defers once and completes on its second claim without spending an attempt; the others complete and cancel outright:

examples/oban_outcomes/main.ml
(* Each outcome is a value the worker returns; the runtime records it. *)
module Welcome = struct
  let name = welcome_name

  let perform _ =
    print_endline "welcome: delivering";
    Oban.Worker.Ok

  let backoff = Oban.Worker.default_backoff
  let timeout = Oban.Worker.default_timeout
end

(* Always fails with an immediate backoff, so the retry is claimable at
   once; the second failure exhausts max_attempts and discards. *)
module Flaky = struct
  let name = flaky_name

  let perform _ =
    print_endline "flaky: boom";
    Oban.Worker.Error "boom"

  let backoff _ = 0.
  let timeout = Oban.Worker.default_timeout
end

(* Snoozes once (no attempt used), then completes on the second claim. *)
module Snoozer = struct
  let name = snoozer_name
  let snoozed = ref false

  let perform _ =
    if !snoozed then (
      print_endline "snoozer: awake, completing";
      Oban.Worker.Ok)
    else (
      snoozed := true;
      print_endline "snoozer: snoozing for 0s";
      Oban.Worker.Snooze 0.)

  let backoff = Oban.Worker.default_backoff
  let timeout = Oban.Worker.default_timeout
end

module Canceller = struct
  let name = canceller_name

  let perform _ =
    print_endline "canceller: no longer needed";
    Oban.Worker.Cancel "no longer needed"

  let backoff = Oban.Worker.default_backoff
  let timeout = Oban.Worker.default_timeout
end

One job per worker, drained until the queue stays idle:

examples/oban_outcomes/main.ml
(* One job per worker; Flaky gets two attempts so we can watch the
   retry become a discard. *)
let null = Simdjsont.Json.Null in
List.iter
  (fun (worker, max_attempts) ->
    match Runtime.enqueue runtime (spec ~worker ~max_attempts null) with
    | Ok _ -> ()
    | Error _ ->
        prerr_endline "error: enqueue failed";
        exit 1)
  [
    (welcome_name, 1);
    (flaky_name, 2);
    (snoozer_name, 3);
    (canceller_name, 1);
  ];
print_endline "enqueued: 4 jobs";
examples/oban_outcomes/main.ml
(* Drain until idle: each [run_one] claims the next available job in
   priority/schedule/id order and records its outcome. *)
let rec drain idle =
  match Runtime.run_one ~clock ~wall_now runtime ~queue with
  | Error _ ->
      prerr_endline "error: run failed";
      exit 1
  | Ok `Idle ->
      if idle >= 100 then ()
      else (
        Eio.Time.sleep clock 0.005;
        drain (idle + 1))
  | Ok (`Worked outcome) ->
      Printf.printf "worked: %s\n%!"
        (match outcome with
        | `Completed -> "completed"
        | `Failed -> "failed (will retry)"
        | `Cancelled -> "cancelled"
        | `Discarded -> "discarded"
        | `Snoozed -> "snoozed"
        | `Stale -> "stale");
      drain 0
in
drain 0;
$ dune exec examples/oban_outcomes/main.exe

Bulk and unique enqueue

Runtime.enqueue_many — and its store twin Store.enqueue_many — inserts an ordered batch atomically, bounded by the store's max_batch_size. Unique requests key on queue, worker and args: the first insert wins while it is active, and the key is released once the job reaches a terminal state:

examples/oban_maintenance/main.ml
(* [enqueue_many] inserts an ordered batch atomically; [Unique]
   requests return the existing active job instead of a duplicate. *)
let results =
  get
    (Store.enqueue_many store
       [
         Oban.Storage.Standard (spec ());
         Oban.Storage.Standard (spec ());
         Oban.Storage.Standard (spec ());
       ])
in
Printf.printf "enqueued batch: %d jobs\n%!" (List.length results);
let unique = spec () in
let first =
  match get (Store.enqueue store (Oban.Storage.Unique unique)) with
  | Oban.Storage.Inserted job -> Oban.Job.pack job
  | Oban.Storage.Existing _ ->
      prerr_endline "error: expected a fresh unique job";
      exit 1
in
(match get (Store.enqueue store (Oban.Storage.Unique unique)) with
| Oban.Storage.Existing incumbent ->
    let same =
      Oban.Job.id_to_int64 (Oban.Job.inspect incumbent).id
      = Oban.Job.id_to_int64 (Oban.Job.inspect first).id
    in
    Printf.printf "unique: second enqueue returned existing=%b\n%!" same;
    if not same then exit 1
| Oban.Storage.Inserted _ ->
    prerr_endline "error: unique enqueue duplicated";
    exit 1);

Observing jobs

Store.list pages in ascending id order with an optional queue, worker, state set and tag filter; Store.fetch reads one job by id:

examples/oban_maintenance/main.ml
(* [list] pages in ascending id order; every tag in the filter must
   be present. [fetch] reads one job by id. *)
let filter : Oban.Storage.filter =
  { queue = Some queue; worker = None; states = []; tags = [] }
in
let page = get (Store.list store ~filter ~after:None ~limit:2) in
Printf.printf "listed: %d jobs (next=%b)\n%!" (List.length page.jobs)
  (Option.is_some page.next);
(match page.jobs with
| [] ->
    prerr_endline "error: expected listed jobs";
    exit 1
| packed :: _ -> (
    let id = (Oban.Job.inspect packed).id in
    match get (Store.fetch store id) with
    | None ->
        prerr_endline "error: fetch missed a listed job";
        exit 1
    | Some fetched ->
        Printf.printf "fetched: #%Ld state=%s\n%!"
          (Oban.Job.id_to_int64 id)
          (state_string fetched)));

Cancel and retry are revision-aware

Cancellation and retry take the observed id plus revision, so a concurrent state change comes back as Stale instead of a lost update. Only available and executing jobs cancel; only terminal jobs retry:

examples/oban_maintenance/main.ml
(* Cancel and retry are revision-aware: pass the observed id+revision
   and a concurrent change comes back as [Stale], not a lost update. *)
let target = row_ref first in
(match get (Store.cancel store ~target ~now:(wall_now ()) ~reason:"manual") with
| Oban.Storage.Applied packed ->
    Printf.printf "cancelled: state=%s\n%!" (state_string packed)
| _ ->
    prerr_endline "error: cancel did not apply";
    exit 1);
let cancelled =
  match get (Store.fetch store (Oban.Job.inspect first).id) with
  | None ->
      prerr_endline "error: cancelled job vanished";
      exit 1
  | Some packed -> packed
in
(match get (Store.retry store ~target:(row_ref cancelled) ~now:(wall_now ()))
 with
| Oban.Storage.Applied packed ->
    Printf.printf "retried: state=%s\n%!" (state_string packed)
| _ ->
    prerr_endline "error: retry did not apply";
    exit 1);

The batch twins cancel_many and retry_many apply ordered revision-aware batches in one transaction: an error rolls the batch back while individual stale targets stay ordered results.

Pause and resume

A paused queue claims nothing until it is resumed — useful when a downstream service is struggling. The claim below returns nothing while paused and a job once resumed:

examples/oban_maintenance/main.ml
(* A paused queue claims nothing until it is resumed. *)
ignore (get (Store.enqueue store (Oban.Storage.Standard (spec ~queue:paused_queue ()))));
get (Store.pause store paused_queue);
(match get (Store.claim_one store ~queue:paused_queue ~now:(wall_now ()) ~node:"maintenance-1") with
| None -> print_endline "paused: claim returned nothing, as expected"
| Some _ ->
    prerr_endline "error: paused queue claimed work";
    exit 1);
get (Store.resume store paused_queue);
(match get (Store.claim_one store ~queue:paused_queue ~now:(wall_now ()) ~node:"maintenance-1") with
| None ->
    prerr_endline "error: resumed queue claimed nothing";
    exit 1
| Some job ->
    Printf.printf "resumed: claimed #%Ld\n%!"
      (Oban.Job.id_to_int64 (Oban.Job.view job).id);
    (match
       get (Store.finish store ~now:(wall_now ()) job Oban.Storage.Complete)
     with
    | Oban.Storage.Applied _ -> print_endline "resumed: finished as completed"
    | _ ->
        prerr_endline "error: finish did not apply";
        exit 1));

Rescue and prune are application-scheduled

Jobs send no heartbeats while a worker runs, so a rescue threshold should exceed the longest expected worker duration. Rescue requeues executing jobs stuck since a cutoff; prune deletes old terminal rows. Both take a limit bounded by max_maintenance_batch_size:

examples/oban_maintenance/main.ml
(* Rescue requeues executing jobs stuck since [before]; prune deletes
   old terminal rows. Both are application-scheduled and bounded. *)
ignore (get (Store.enqueue store (Oban.Storage.Standard (spec ()))));
let stuck =
  match
    get (Store.claim_one store ~queue ~now:(wall_now ()) ~node:"worker-1")
  with
  | None ->
      prerr_endline "error: nothing to rescue";
      exit 1
  | Some job -> job
in
let rescued =
  get
    (Store.rescue store ~before:(wall_now ()) ~now:(wall_now ())
       ~error:"stuck worker" ~limit:10)
in
Printf.printf "rescued: %d job(s)\n%!" (List.length rescued);
(match
   get
     (Store.finish store ~now:(wall_now ()) stuck Oban.Storage.Complete)
 with
| Oban.Storage.Stale _ ->
    print_endline "rescue: old claim is stale, as expected"
| _ ->
    prerr_endline "error: rescued claim still completed";
    exit 1);
let done_job =
  match
    get (Store.claim_one store ~queue ~now:(wall_now ()) ~node:"worker-2")
  with
  | None ->
      prerr_endline "error: nothing to finish";
      exit 1
  | Some job -> job
in
ignore
  (get (Store.finish store ~now:(wall_now ()) done_job Oban.Storage.Complete));
let before =
  match Ptime.add_span (wall_now ()) (Ptime.Span.of_int_s 3600) with
  | Some t -> t
  | None -> wall_now ()
in
let pruned =
  get (Store.prune store ~states:[ Oban.Job.Completed ] ~before ~limit:10)
in
Printf.printf "pruned: %d completed job(s)\n%!" (List.length pruned);

Execution is at least once: a worker may perform an external effect and stop before its completion is stored, or rescued work may overlap an older worker. Keep worker-side effects idempotent.

Cron computes times; you enqueue

Oban.Cron parses five-field UTC cron expressions and tests or finds times with Ptime.t. It never enqueues anything itself — the application ticks, computes the next time, and enqueues the associated job:

examples/oban_maintenance/main.ml
(* Cron only computes times; the application decides what to enqueue. *)
let cron = get (Oban.Cron.parse "*/15 * * * *") in
let now = wall_now () in
Printf.printf "cron: matches now=%b\n%!" (Oban.Cron.matches cron now);
(match Oban.Cron.next_after cron now with
| None ->
    prerr_endline "error: cron has no next time";
    exit 1
| Some next ->
    let mins =
      Option.value (Ptime.Span.to_int_s (Ptime.diff next now)) ~default:0 / 60
    in
    Printf.printf "cron: next run in %d minute(s)\n%!" mins);

Enqueue inside your own transaction

Store.transaction runs a synchronous callback on one pooled connection. Inside it, the _with_conn variants join job inserts to application writes in the same database transaction, so the job commits or rolls back with the data it describes:

examples/oban_maintenance/main.ml
(* Join job enqueueing to application writes in one database
   transaction: inside, use only the [_with_conn] variants. *)
let tx_specs =
  [ Oban.Storage.Standard (spec ()); Oban.Storage.Standard (spec ()) ]
in
(match
   Store.transaction store (fun conn ->
       Store.enqueue_many_with_conn store conn tx_specs)
 with
| Ok inserted ->
    Printf.printf "transaction: enqueued %d jobs atomically\n%!"
      (List.length inserted)
| Error _ ->
    prerr_endline "error: transaction failed";
    exit 1)
$ dune exec examples/oban_maintenance/main.exe

Where this fits in a full application: every store and runtime call lives in a model or a context — controllers and views never touch the queue. Each context function that writes owns one transaction and enqueues post-commit effects after it returns. The API reference linked above documents every module and function in full.