Implementing async/await in OCaml

2026-10-02 ·

In JavaScript or Python, async and await are syntax, and the compiler turns each async function into a state machine. OCaml 5 needs no compiler support for this: an effect handler can suspend a computation in the middle of any expression and resume it later, so await is an ordinary function. This post builds a single-threaded async runtime with promises, an event loop, timers, waiting on file descriptors, exceptions, and the usual combinators. It assumes the basics from Effect Handlers in OCaml.

§ 01

Where this ends up

Three tasks that sleep for different times. They run concurrently, so the whole program takes as long as the slowest.

open Async
let start = Unix.gettimeofday ()
let fetch name delay =
async (fun () ->
Printf.printf "%-5s start\n" name;
sleep delay;
Printf.printf "%-5s done after %.1fs\n" name delay;
String.length name)
let () =
run (fun () ->
let a = fetch "alpha" 0.3 in
let b = fetch "beta" 0.1 in
let c = fetch "gamma" 0.2 in
let total = await a + await b + await c in
Printf.printf "total %d in %.1fs\n" total (Unix.gettimeofday () -. start))

Output.

alpha start
beta start
gamma start
beta done after 0.1s
gamma done after 0.2s
alpha done after 0.3s
total 14 in 0.3s

fetch is a plain function and sleep is called in direct style, in the middle of a let sequence. Nothing in the types says sleep suspends: the type of fetch is string -> float -> int promise, and await has type 'a promise -> 'a.

§ 02

Promises

A promise is a mutable cell that is pending, holds a value, or holds an exception. A pending promise keeps a list of callbacks to run when it is resolved.

async.ml: promises.

type 'a state =
| Pending of (('a, exn) result -> unit) list
| Done of 'a
| Failed of exn
type 'a promise = { mutable state : 'a state }
let resolve p r =
match p.state with
| Pending waiters ->
p.state <- (match r with Ok v -> Done v | Error e -> Failed e);
List.iter (fun w -> w r) (List.rev waiters)
| Done _ | Failed _ -> invalid_arg "resolve: promise already resolved"
let on_resolve p w =
match p.state with
| Pending ws -> p.state <- Pending (w :: ws)
| Done v -> w (Ok v)
| Failed e -> w (Error e)

Waiters are stored newest first, so resolve reverses the list to call them in the order they registered. on_resolve on a promise that is already resolved calls the callback at once. Nothing here knows about tasks or scheduling: a callback is just a function.

§ 03

Three effects

The runtime has three operations. Async f starts f as a new task and returns its promise, Await p suspends until p is resolved, and Yield lets other tasks run. Each is an effect, and the user-facing functions just perform them.

async.ml: the effects.

type _ Effect.t +=
| Async : (unit -> 'a) -> 'a promise Effect.t
| Await : 'a promise -> 'a Effect.t
| Yield : unit Effect.t
let async f = Effect.perform (Async f)
let await p = Effect.perform (Await p)
let yield () = Effect.perform Yield

The type parameter of Effect.t is what perform returns. Async returns a promise to the task that performed it; Await p returns the value inside p, which is why await is 'a promise -> 'a.

§ 04

The scheduler

A task is a function run under a handler. When it performs an effect, the handler receives a continuation k, which is the rest of the task, suspended. Resuming it with continue k v makes perform return v; discontinue k e makes perform raise e instead. The scheduler is a queue of thunks, and each handler case decides what to put on it.

async.ml: the run queue, and fork, which runs one task under the handler.

let run_queue : (unit -> unit) Queue.t = Queue.create ()
let enqueue task = Queue.push task run_queue
let rec fork : 'a. 'a promise -> (unit -> 'a) -> unit =
fun p f ->
let open Effect.Deep in
match_with f ()
{ retc = (fun v -> resolve p (Ok v));
exnc = (fun e -> resolve p (Error e));
effc =
(fun (type b) (eff : b Effect.t) ->
match eff with
| Async g ->
Some
(fun (k : (b, unit) continuation) ->
let child = { state = Pending [] } in
enqueue (fun () -> fork child g);
enqueue (fun () -> continue k child))
| Await q ->
Some
(fun (k : (b, unit) continuation) ->
on_resolve q (function
| Ok v -> enqueue (fun () -> continue k v)
| Error e -> enqueue (fun () -> discontinue k e)))
| Yield ->
Some (fun (k : (b, unit) continuation) -> enqueue (fun () -> continue k ()))
| _ -> None) }

fork p f runs f and resolves p with its result or with the exception it raised. The handler cases:

  • Async g makes a fresh promise for the child, queues the child to start, and queues the parent to resume with that promise.
  • Await q registers a callback on q. When q is resolved, the callback queues the task to resume with the value, or to have the exception raised at the await.
  • Yield queues the task to resume and returns.

No case calls continue directly. Each one queues the resumption and returns unit, and that return is what hands control back to the event loop: the thunk that ran this task, either fork or an earlier continue k v, returns as soon as the task suspends. So every await is a point where other tasks can run, even when the promise is already resolved, as in JavaScript.

Two typing details. The handler is a GADT match, so effc takes a locally abstract type b, and in the Await case b is refined to the promise's element type, so continue k v typechecks. And fork calls itself on the child at a different type, which is polymorphic recursion and needs the explicit annotation 'a. 'a promise -> .... The handler is deep, so it stays installed after each continue and catches every later effect the task performs; each child gets its own handler from its own fork.

§ 05

Sleeping and waiting for input

Timers and I/O need no new effects. sleep creates a pending promise, records when it should be resolved, and awaits it. readable fd does the same with a file descriptor.

async.ml: timers and readers.

let timers : (float * unit promise) list ref = ref []
let readers : (Unix.file_descr * unit promise) list ref = ref []
let sleep d =
let p = { state = Pending [] } in
let at = Unix.gettimeofday () +. d in
let rec insert = function
| (t, _) as x :: rest when t <= at -> x :: insert rest
| rest -> (at, p) :: rest
in
timers := insert !timers;
await p
let readable fd =
let p = { state = Pending [] } in
readers := (fd, p) :: !readers;
await p

The timers are kept sorted by deadline, which is enough for a post; a real runtime would use a heap. The event loop runs queued tasks until the queue is empty. Then, if any timers or readers are pending, it blocks in Unix.select until a descriptor is readable or the earliest timer is due, resolves what is ready, and goes round again. Resolving a promise only queues the tasks waiting on it, so they run on the next pass.

async.ml: the event loop and run.

let fire_timers () =
let now = Unix.gettimeofday () in
let due, later = List.partition (fun (t, _) -> t <= now) !timers in
timers := later;
List.iter (fun (_, p) -> resolve p (Ok ())) due
let poll_io () =
let timeout =
match !timers with
| [] -> -1.0
| (t, _) :: _ -> Float.max 0.0 (t -. Unix.gettimeofday ())
in
let fds = List.map fst !readers in
let ready, _, _ =
try Unix.select fds [] [] timeout
with Unix.Unix_error (Unix.EINTR, _, _) -> ([], [], [])
in
let woken, waiting = List.partition (fun (fd, _) -> List.mem fd ready) !readers in
readers := waiting;
List.iter (fun (_, p) -> resolve p (Ok ())) woken
let rec loop () =
match Queue.take_opt run_queue with
| Some task -> task (); loop ()
| None when !timers = [] && !readers = [] -> ()
| None -> poll_io (); fire_timers (); loop ()
let run main =
let p = { state = Pending [] } in
enqueue (fun () -> fork p main);
loop ();
match p.state with
| Done v -> v
| Failed e -> raise e
| Pending _ -> failwith "run: main is waiting on a promise nothing can resolve"

run starts main as the first task and returns when there is nothing left to do: no queued tasks, no timers and no readers. That includes tasks nobody awaited, so run waits for them too.

§ 06

Exceptions

An exception raised inside a task is caught by exnc and stored in the task's promise. await on a failed promise discontinues the waiting task with that exception, so it is raised at the await and can be caught with an ordinary match or try. If main fails, run re-raises.

Two lookups, one of which fails.

open Async
let lookup key =
async (fun () ->
sleep 0.05;
match List.assoc_opt key [ ("ocaml", 1996); ("haskell", 1990) ] with
| Some year -> year
| None -> raise Not_found)
let () =
run (fun () ->
let ok = lookup "ocaml" and bad = lookup "cobol" in
Printf.printf "ocaml: %d\n" (await ok);
match await bad with
| year -> Printf.printf "cobol: %d\n" year
| exception Not_found -> print_endline "cobol: not found");
try Printf.printf "%d\n" (run (fun () -> await (lookup "cobol")))
with Not_found -> print_endline "run re-raised Not_found"

Output.

ocaml: 1996
cobol: not found
run re-raised Not_found

A task that fails and is never awaited fails silently, since its exception sits in a promise nobody reads. Lwt has the same issue for tasks started with Lwt.async, and reports their failures through a global hook, Lwt.async_exception_hook.

§ 07

Interleaving

A writer sends three messages down a pipe, one every 100 ms, and closes it. A reader waits for each with readable, and a ticker prints in between. The reader's loop is plain recursion with a blocking-looking read_line in it.

A pipe, a writer, a reader and a ticker in one thread.

open Async
let read_line fd =
let buf = Bytes.create 64 in
readable fd;
match Unix.read fd buf 0 64 with
| 0 -> None
| n -> Some (Bytes.sub_string buf 0 n)
let () =
let r, w = Unix.pipe () in
run (fun () ->
let writer = async (fun () ->
List.iter (fun msg ->
sleep 0.1;
ignore (Unix.write_substring w msg 0 (String.length msg)))
[ "one"; "two"; "three" ];
Unix.close w) in
let ticker = async (fun () ->
sleep 0.05;
for i = 1 to 3 do print_endline ("tick " ^ string_of_int i); sleep 0.1 done) in
let rec reader () =
match read_line r with
| Some msg -> print_endline ("read " ^ msg); reader ()
| None -> print_endline "eof"
in
reader ();
await writer;
await ticker)

Output, over 0.35 s.

tick 1
read one
tick 2
read two
tick 3
read three
eof

The reader calls Unix.read only after select has reported the pipe readable, so the read returns immediately. When the writer closes its end, the pipe becomes readable and read returns 0 bytes, which is end of file.

yield shows the scheduling order directly. Two workers each print and yield three times.

Cooperative workers, and await outside run.

open Async
let worker name n =
async (fun () ->
for i = 1 to n do Printf.printf "%s%d " name i; yield () done)
let () =
run (fun () ->
let a = worker "a" 3 and b = worker "b" 3 in
await a; await b);
print_newline ();
match await (async (fun () -> 1)) with
| _ -> ()
| exception Effect.Unhandled _ -> print_endline "async outside run: Effect.Unhandled"

Output.

a1 a2 b1 a3 b2 b3
async outside run: Effect.Unhandled

The order follows the queue exactly. async queues the child before the parent, so worker a prints a1 before main has even started worker b, and a2 comes before b1 because a was queued again before b started. After that the two alternate until a runs out. Scheduling is cooperative: a task that loops without performing an effect never gives up control. Outside run no handler is installed, and performing the effect raises Effect.Unhandled.

§ 08

Combinators

wait returns a pending promise and a function that resolves it, which is how callback-based code is turned into a promise. all is one line: the tasks are already running, so awaiting them in order costs nothing extra and the result is ready when the slowest finishes. race resolves with whichever promise settles first and ignores the others. timeout races a task against a sleep.

async.ml: wait, all, race, timeout.

let wait () =
let p = { state = Pending [] } in
(p, fun v -> resolve p (Ok v))
let all ps = async (fun () -> List.map await ps)
let race ps =
let p = { state = Pending [] } in
let settle r = match p.state with Pending _ -> resolve p r | _ -> () in
List.iter (fun q -> on_resolve q settle) ps;
p
let timeout d p =
race [ async (fun () -> Some (await p)); async (fun () -> sleep d; None) ]

Using them.

open Async
let slow name d = async (fun () -> sleep d; name)
let show = function Some s -> s | None -> "timed out"
let () =
run (fun () ->
let winner = await (race [ slow "tortoise" 0.3; slow "hare" 0.1 ]) in
print_endline ("race: " ^ winner);
print_endline ("fast: " ^ show (await (timeout 0.2 (slow "quick" 0.05))));
print_endline ("slow: " ^ show (await (timeout 0.2 (slow "sluggish" 0.5))));
let names = await (all [ slow "a" 0.2; slow "b" 0.1; slow "c" 0.15 ]) in
print_endline ("all: " ^ String.concat " " names))

Output.

race: hare
fast: quick
slow: timed out
all: a b c

This program takes 0.65 s, although its last line prints at 0.55 s. race does not cancel the losers: the sluggish task keeps sleeping after its timeout fires, and run waits for it. Cancellation needs each task to carry a switch that sleep and readable check, so that a cancelled task is discontinued with an exception at its next suspension point. Eio does this: its Fiber.first cancels the fiber that loses.

§ 09

Deadlock

If main awaits a promise that nothing will ever resolve, the queue empties with no timers or readers left, the loop stops, and main's promise is still pending. run reports that instead of hanging.

Awaiting a promise nobody resolves.

open Async
let () =
match run (fun () -> let p, _resolve = wait () in await p) with
| () -> ()
| exception Failure msg -> print_endline msg

Output.

run: main is waiting on a promise nothing can resolve
§ 10

What a real runtime adds

The scheduler state is global, so this runs on one domain; running tasks on several domains needs a run queue per domain and stealing between them, as in the work-stealing scheduler. Writes, sockets and accepting connections work the same way as readable, through select's other two lists, though select itself scales badly and real runtimes use epoll, kqueue or io_uring. Eio and Miou are libraries built on exactly this pattern, with cancellation, structured concurrency and those backends.

§ 11

Download

The runtime, async.ml, and the examples d1_fetch.ml, d2_exn.ml, d3_pipe.ml, d4_race.ml, d5_deadlock.ml and d6_yield.ml. Each example builds with ocamlfind ocamlopt -package unix -linkpkg async.ml d1_fetch.ml on OCaml 5.