about things notes.zzstoatzz.io
notes
notes languages ziglang io synchronization.md
7.1 kB

synchronization primitives #

Io.Mutex #

extern struct. futex-based. works from any execution context (threads AND fibers).

var mutex: Io.Mutex = Io.Mutex.init;

// cancelable lock
mutex.lock(io) catch |err| switch (err) {
    error.Canceled => return,
};
defer mutex.unlock(io);

// uncancelable lock (for use in cleanup paths)
mutex.lockUncancelable(io);
defer mutex.unlock(io);

// non-blocking try
if (mutex.tryLock()) {
    defer mutex.unlock(io);
    // ...
}
  • lock(m, io) → Cancelable!void — cancellation point
  • lockUncancelable(m, io) → void — no cancellation point
  • tryLock(m) → bool — non-blocking, no io needed
  • unlock(m, io) → void

replaces std.Thread.Mutex from 0.15.

cross-context usage — CRITICAL CONSTRAINT #

Io.Mutex is futex-based and works from any context within the same Io runtime. but it CANNOT be shared across different Io types (Threaded vs Evented):

  • Threaded futex on Evented fiber → blocks the entire Uring OS thread. that thread's io_uring instance can't process CQEs → deadlock. other fibers on that thread (including the main fiber doing CA bundle loading, accept loops, etc.) are permanently stuck.
  • Evented futex on plain std.Thread → Thread.current() is a threadlocal only set on Uring-managed threads. on a plain thread it's null. in ReleaseFast, self.? on null silently gives NULL pointer → SIGSEGV at struct field offsets (0x28, 0x30, 0x38 in our case — ready_queue, free_queue, io_uring fields of the Thread struct).

rule: all callers of a given Io.Mutex must pass the same Io instance (or at least the same Io type). if you have a data structure accessed from both a Threaded worker pool and Evented fibers, you must either:

  1. run the shared structure entirely on one Io type (e.g., all on pool_io/Threaded)
  2. use raw atomics (no Io.Mutex) for the cross-boundary synchronization
  3. use an MPSC queue with atomic CAS (no futex involvement)

this was the root cause of a SIGSEGV in zlay (our atproto relay) (frame worker threads on pool_io/Threaded calling Io.Mutex.lockUncancelable with Evented io on the Resyncer) and the subsequent deadlock (first fix attempt mixed Threaded futex with Evented fiber). see zlay commits 6674812, 439c678.

Io.Condition #

pairs with Io.Mutex:

var cond: Io.Condition = Io.Condition.init;
var mutex: Io.Mutex = Io.Mutex.init;

// waiter (must hold mutex)
mutex.lockUncancelable(io);
while (!predicate()) {
    cond.wait(&cond, io, &mutex) catch break;  // releases mutex, waits, reacquires
}
mutex.unlock(io);

// signaler
cond.signal(io);     // wake one waiter
cond.broadcast(io);  // wake all waiters
  • wait(cond, io, mutex) → Cancelable!void — releases mutex, waits, reacquires
  • waitUncancelable(cond, io, mutex) → void — same but no cancellation point
  • signal(cond, io) → void — wake one
  • broadcast(cond, io) → void — wake all

no timedWait #

Io.Condition has no timed wait variant. patterns that used timedWait must be restructured:

option 1: sleep-based polling (simplest, adds latency up to sleep interval)

while (condition_not_met and alive.load(.acquire)) {
    mutex.unlock(io);
    io.sleep(Io.Duration.fromMilliseconds(100), .awake) catch {};
    mutex.lockUncancelable(io);
}

option 2: ticker task that signals the cond (preserves immediate wakeup)

var ticker = try io.concurrent(tickerLoop, .{io, &cond});
defer _ = ticker.cancel(io);

fn tickerLoop(tick_io: Io, cond: *Io.Condition) void {
    while (true) {
        tick_io.sleep(Io.Duration.fromMilliseconds(100), .awake) catch break;
        cond.signal(tick_io);
    }
}

option 1 is fine when latency tolerance matches the poll interval. option 2 is better when immediate wake on signal AND periodic timeout are both needed.

cancellation model #

every Io function that returns Cancelable!T is a cancellation point. when future.cancel(io) is called on a task:

  1. the task is flagged for cancellation
  2. at its next cancellation point, the function returns error.Canceled
  3. the task should propagate or handle the error to exit cleanly
  4. cancel blocks until the task actually returns

what's a cancellation point? #

any Io function with Cancelable in its error set:

  • io.sleep()
  • mutex.lock() (but NOT lockUncancelable)
  • cond.wait() (but NOT waitUncancelable)
  • queue.getOne() / queue.putOne()
  • io.checkCancel() — explicit cancellation point (does nothing else)
  • io.futexWait() (but NOT futexWaitUncancelable)

recancel #

if you catch error.Canceled and want to propagate it through more cancellation points:

io.sleep(...) catch |err| switch (err) {
    error.Canceled => {
        // do some cleanup...
        io.recancel(io);  // re-arm so next cancellation point also returns Canceled
        return error.Canceled;
    },
};

recancel asserts that a prior cancellation was received. it re-arms the request so subsequent cancellation points also fire.

CancelProtection #

in rare cases, a section of code must run to completion without being interrupted by cancellation:

const old = io.swapCancelProtection(.blocked);
defer _ = io.swapCancelProtection(old);

// io operations here will NOT return error.Canceled
// even if the task has a pending cancellation request
mutex.lock(io) catch unreachable;  // lock can't fail with .blocked
defer mutex.unlock(io);
// ... critical section ...
  • .unblocked — default. cancellation points are active.
  • .blocked — no Io function returns error.Canceled.

use for cleanup code that must complete, commit-then-ack patterns, etc.

Io.Event #

simple binary signal (like a one-shot condition without a mutex):

var event: Io.Event = Io.Event.init;

// waiter
event.wait(io) catch {};  // blocks until set

// signaler
event.set(io);  // wakes all waiters

use when you just need "wait until something happens" without associated data.

Io.Semaphore #

counting semaphore:

var sem: Io.Semaphore = Io.Semaphore.init;

sem.acquire(io) catch {};  // blocks if count == 0
defer sem.release(io);

use for bounding concurrent access to a shared resource.

Io.RwLock #

reader-writer lock for read-heavy workloads:

var rwlock: Io.RwLock = Io.RwLock.init;

// readers (concurrent)
rwlock.lockShared(io) catch {};
defer rwlock.unlockShared(io);

// writer (exclusive)
rwlock.lockExclusive(io) catch {};
defer rwlock.unlockExclusive(io);

lock convoys in hot logging paths #

a single mutex around an exporter or log sink serializes every thread that logs: under load the threads convoy behind the lock and throughput collapses to the sink's speed. buffer per-thread (or per-fiber) and drain from one consumer instead; and never take an Io.Mutex from a plain OS thread that is not running on the event loop — it parks the fiber machinery, not the thread.

sources #

  • zlay — Io.Mutex SIGSEGV and deadlock, commits 6674812, 439c678