authorgravatar for jacobly@ziglang.orgJacob Young <jacobly@ziglang.org> 2026-02-14 04:49:03-05:00
committergravatar for jacobly@ziglang.orgJacob Young <jacobly@ziglang.org> 2026-02-14 05:52:59-05:00
logb7f93695f90e92b4437346bd993bbdb4d86f7b07
tree5e0a50dfadb2410c3a243c63d390765a90ebdce0
parentf996d28666e7d562b6a18e9f80fe31a3463456d0

Io.Dispatch.Mutex: fix deadlock conditions


2 files changed, 150 insertions(+), 131 deletions(-)

lib/std/Io/Dispatch.zig+147-128
......@@ -280,7 +280,7 @@ const Fiber = struct {
280280 .select => if (@atomicRmw(i32, &fiber.await_count, .Add, 1, .monotonic) == -1) {
281281 ev.queue.async(fiber, &Fiber.@"resume");
282282 },
283 _ => |awaiting| awaiting.toCancelable().canceled(),
283 _ => |awaiting| awaiting.toCancelable().async(),
284284 }
285285 }
286286
......@@ -484,7 +484,7 @@ pub fn io(ev: *Evented) Io {
484484
485485pub const InitOptions = struct {
486486 backing_allocator_needs_mutex: bool = true,
487 queue: ?c.dispatch.queue_t = null,
487 target_queue: ?c.dispatch.queue_t = .TARGET_DEFAULT,
488488 /// Upper limit on the allowable delay in processing timeouts in order to improve power
489489 /// consumption and system performance.
490490 leeway: Io.Duration = .fromMilliseconds(10),
......@@ -499,11 +499,11 @@ pub const InitOptions = struct {
499499};
500500
501501pub fn init(ev: *Evented, backing_allocator: Allocator, options: InitOptions) !void {
502 const queue = if (options.queue) |queue| queue: {
503 queue.as_object().retain();
504 break :queue queue;
505 } else c.dispatch.queue_create("org.ziglang.std.Io.Dispatch", .CONCURRENT()) orelse
506 return error.SystemResources;
502 const queue = c.dispatch.queue_create_with_target(
503 "org.ziglang.std.Io.Dispatch",
504 .CONCURRENT(),
505 options.target_queue,
506 ) orelse return error.SystemResources;
507507 errdefer queue.as_object().release();
508508 const main_loop_stack = try backing_allocator.alignedAlloc(
509509 u8,
......@@ -727,9 +727,9 @@ const Cancelable = struct {
727727
728728 const blocked: Cancelable = .{ .queue = undefined, .cancel = is_blocked };
729729
730 const AwaitError = error{CancelRequested};
730 const RequestedError = error{CancelRequested};
731731
732 fn await(cancelable: *Cancelable, fiber: *Fiber) AwaitError!void {
732 fn enter(cancelable: *Cancelable, fiber: *Fiber) RequestedError!void {
733733 const function = cancelable.cancel;
734734 assert(function != is_requested);
735735 if (function == is_blocked) {
......@@ -750,13 +750,42 @@ const Cancelable = struct {
750750 }
751751 }
752752
753 fn canceled(cancelable: *Cancelable) void {
754 assert(cancelable.cancel != is_blocked);
755 assert(cancelable.cancel != is_requested);
756 cancelable.queue.async(cancelable, cancelable.cancel);
753 fn leave(cancelable: *Cancelable, fiber: *Fiber) RequestedError!void {
754 const function = cancelable.cancel;
755 assert(function != is_requested);
756 if (function == is_blocked) {
757 @branchHint(.unlikely);
758 return;
759 }
760 const cancel_status = @atomicRmw(Fiber.CancelStatus, &fiber.cancel_status, .And, .{
761 .requested = true,
762 .awaiting = .nothing,
763 }, .monotonic);
764 assert(cancel_status.awaiting.toCancelable() == cancelable);
765 if (cancel_status.requested) return error.CancelRequested;
766 }
767
768 fn async(cancelable: *Cancelable) void {
769 const function = cancelable.cancel;
770 assert(function != is_blocked and function != is_requested);
771 cancelable.queue.async(cancelable, function);
772 }
773
774 fn requested(cancelable: *Cancelable, fiber: *Fiber) void {
775 const function = cancelable.cancel;
776 assert(function != is_blocked and function != is_requested);
777 assert(@atomicLoad(Fiber.CancelStatus, &fiber.cancel_status, .monotonic) == Fiber.CancelStatus{
778 .requested = true,
779 .awaiting = .fromCancelable(cancelable),
780 });
781 cancelable.cancel = is_requested;
782 @atomicStore(Fiber.CancelStatus, &fiber.cancel_status, .{
783 .requested = true,
784 .awaiting = .nothing,
785 }, .monotonic);
757786 }
758787
759 fn check(cancelable: *Cancelable, fiber: *Fiber) Io.Cancelable!void {
788 fn acknowledge(cancelable: *Cancelable, fiber: *Fiber) Io.Cancelable!void {
760789 if (cancelable.cancel == is_requested) {
761790 @branchHint(.unlikely);
762791 fiber.cancel_protection.acknowledge();
......@@ -784,81 +813,97 @@ const Sleeper = struct {
784813};
785814
786815const Mutex = struct {
787 /// including the locker
788 num_waiters: usize,
816 state: State,
789817 queue: c.dispatch.queue_t,
790818 waiters: std.DoublyLinkedList,
791819
820 const State = packed struct(usize) {
821 locked: bool,
822 num_waiters: NumWaiters,
823
824 const NumWaiters = @Int(.unsigned, @bitSizeOf(usize) - 1);
825 };
826
792827 const Waiter = struct {
793828 sleeper: Sleeper = undefined,
794829 cancelable: Cancelable,
795830 mutex: *Mutex,
796 node: std.DoublyLinkedList.Node = .{},
831 node: std.DoublyLinkedList.Node = undefined,
797832
798833 fn add(context: ?*anyopaque) callconv(.c) void {
799834 const waiter: *Waiter = @ptrCast(@alignCast(context));
800 waiter.tryAdd() catch |err| switch (err) {
835 waiter.cancelable.enter(waiter.sleeper.fiber) catch |err| switch (err) {
836 error.CancelRequested => return waiter.wake(),
837 };
838 var state = @atomicRmw(State, &waiter.mutex.state, .Add, .{
839 .locked = false,
840 .num_waiters = 1,
841 }, .monotonic);
842 state.num_waiters += 1;
843 while (!state.locked) {
844 @branchHint(.unlikely);
845 state = @cmpxchgWeak(State, &waiter.mutex.state, state, .{
846 .locked = true,
847 .num_waiters = state.num_waiters - 1,
848 }, .acquire, .monotonic) orelse break;
849 } else return waiter.mutex.waiters.append(&waiter.node);
850 waiter.cancelable.leave(waiter.sleeper.fiber) catch |err| switch (err) {
801851 error.CancelRequested => {
802 waiter.wake();
803 assert(@atomicRmw(usize, &waiter.mutex.num_waiters, .Sub, 1, .monotonic) >= 1);
852 waiter.node.next = &waiter.node;
853 return;
804854 },
805855 };
806 }
807
808 fn tryAdd(waiter: *Waiter) Cancelable.AwaitError!void {
809 switch (@atomicLoad(usize, &waiter.mutex.num_waiters, .acquire)) {
810 0 => unreachable,
811 1 => return waiter.wake(), // already locked exclusively
812 else => try waiter.cancelable.await(waiter.sleeper.fiber),
813 }
814 waiter.mutex.waiters.append(&waiter.node);
856 waiter.wake();
815857 }
816858
817859 fn canceled(context: ?*anyopaque) callconv(.c) void {
818860 const cancelable: *Cancelable = @ptrCast(@alignCast(context));
819 cancelable.cancel = Cancelable.is_requested;
820861 const waiter: *Waiter = @fieldParentPtr("cancelable", cancelable);
821 assert(@atomicRmw(
822 Fiber.CancelStatus,
823 &waiter.sleeper.fiber.cancel_status,
824 .Xchg,
825 .{ .requested = true, .awaiting = .nothing },
826 .monotonic,
827 ) == Fiber.CancelStatus{ .requested = true, .awaiting = .fromCancelable(cancelable) });
862 cancelable.requested(waiter.sleeper.fiber);
828863 const mutex = waiter.mutex;
829 mutex.waiters.remove(&waiter.node);
864 if (waiter.node.next != &waiter.node) {
865 @branchHint(.likely);
866 mutex.waiters.remove(&waiter.node);
867 assert(@atomicRmw(State, &mutex.state, .Sub, .{
868 .locked = false,
869 .num_waiters = 1,
870 }, .monotonic).num_waiters >= 1);
871 }
872 waiter.node = undefined;
830873 waiter.wake();
831 assert(@atomicRmw(usize, &mutex.num_waiters, .Sub, 1, .monotonic) >= 1);
832874 }
833875
834876 fn remove(context: ?*anyopaque) callconv(.c) void {
835877 const mutex: *Mutex = @ptrCast(@alignCast(context));
836 var stop_node: ?*std.DoublyLinkedList.Node = null;
837 while (mutex.waiters.first != stop_node) {
878 var state = @atomicLoad(State, &mutex.state, .monotonic);
879 while (!state.locked and state.num_waiters > 0) {
880 @branchHint(.likely);
881 state = @cmpxchgWeak(State, &mutex.state, state, .{
882 .locked = true,
883 .num_waiters = state.num_waiters - 1,
884 }, .acquire, .monotonic) orelse break;
885 } else return;
886 var num_removed: State.NumWaiters = 0;
887 while (mutex.waiters.popFirst()) |node| {
838888 @branchHint(.likely);
839 const waiter: *Waiter = @fieldParentPtr("node", mutex.waiters.popFirst().?);
840 if (waiter.cancelable.cancel != Cancelable.is_blocked) {
841 @branchHint(.likely);
842 const cancel_status = @atomicRmw(
843 Fiber.CancelStatus,
844 &waiter.sleeper.fiber.cancel_status,
845 .And,
846 .{ .requested = true, .awaiting = .nothing },
847 .monotonic,
848 );
849 assert(cancel_status.awaiting.toCancelable() == &waiter.cancelable);
850 if (cancel_status.requested) {
851 @branchHint(.unlikely);
852 // carefully place the hot potato out of the way
853 mutex.waiters.append(&waiter.node);
854 if (stop_node == null) stop_node = &waiter.node;
889 const waiter: *Waiter = @fieldParentPtr("node", node);
890 node.* = undefined;
891 waiter.cancelable.leave(waiter.sleeper.fiber) catch |err| switch (err) {
892 error.CancelRequested => {
893 num_removed += 1;
894 node.next = node;
855895 continue;
856 }
857 }
858 waiter.wake();
859 return;
896 },
897 };
898 break;
899 }
900 if (num_removed > 0) {
901 @branchHint(.unlikely);
902 assert(@atomicRmw(State, &mutex.state, .Sub, .{
903 .locked = false,
904 .num_waiters = num_removed,
905 }, .monotonic).num_waiters >= num_removed);
860906 }
861 // everyone is about to die, nobody will wake up ;-(
862907 }
863908
864909 fn wake(waiter: *Waiter) void {
......@@ -868,7 +913,7 @@ const Mutex = struct {
868913
869914 fn init(mutex: *Mutex, queue: c.dispatch.queue_t) error{SystemResources}!void {
870915 mutex.* = .{
871 .num_waiters = 0,
916 .state = .{ .locked = false, .num_waiters = 0 },
872917 .queue = c.dispatch.queue_create_with_target(
873918 "org.ziglang.std.Io.Dispatch.Mutex",
874919 .SERIAL(),
......@@ -879,56 +924,48 @@ const Mutex = struct {
879924 }
880925
881926 fn deinit(mutex: *Mutex) void {
882 assert(mutex.num_waiters == 0 and mutex.waiters.first == null and mutex.waiters.last == null);
927 assert(mutex.state == State{ .locked = false, .num_waiters = 0 });
928 assert(mutex.waiters.first == null and mutex.waiters.last == null);
883929 mutex.queue.as_object().release();
884930 mutex.* = undefined;
885931 }
886932
887933 fn tryLock(mutex: *Mutex) bool {
888 if (@cmpxchgWeak(usize, &mutex.num_waiters, 0, 1, .acquire, .monotonic) == null) {
889 @branchHint(.likely);
890 return true;
934 const state =
935 @atomicRmw(State, &mutex.state, .Or, .{ .locked = true, .num_waiters = 0 }, .acquire);
936 if (state.locked) {
937 @branchHint(.unlikely);
891938 }
892 return false;
939 return !state.locked;
893940 }
894941
895942 fn lock(mutex: *Mutex, ev: *Evented) Io.Cancelable!void {
896 switch (@atomicRmw(usize, &mutex.num_waiters, .Add, 1, .acquire)) {
897 0 => {},
898 else => {
899 @branchHint(.unlikely);
900 var waiter: Waiter = .{
901 .cancelable = .{ .queue = mutex.queue, .cancel = &Mutex.Waiter.canceled },
902 .mutex = mutex,
903 };
904 ev.yield(.{ .mutex_wait = &waiter });
905 try waiter.cancelable.check(waiter.sleeper.fiber);
906 },
907 }
943 if (mutex.tryLock()) return;
944 var waiter: Waiter = .{
945 .cancelable = .{ .queue = mutex.queue, .cancel = &Mutex.Waiter.canceled },
946 .mutex = mutex,
947 };
948 ev.yield(.{ .mutex_wait = &waiter });
949 try waiter.cancelable.acknowledge(waiter.sleeper.fiber);
908950 }
909951
910952 fn lockUncancelable(mutex: *Mutex, ev: *Evented) void {
911 switch (@atomicRmw(usize, &mutex.num_waiters, .Add, 1, .acquire)) {
912 0 => {},
913 else => {
914 @branchHint(.unlikely);
915 var waiter: Waiter = .{ .cancelable = .blocked, .mutex = mutex };
916 ev.yield(.{ .mutex_wait = &waiter });
917 waiter.cancelable.check(waiter.sleeper.fiber) catch |err| switch (err) {
918 error.Canceled => unreachable, // blocked
919 };
920 },
921 }
953 if (mutex.tryLock()) return;
954 var waiter: Waiter = .{ .cancelable = .blocked, .mutex = mutex };
955 ev.yield(.{ .mutex_wait = &waiter });
956 waiter.cancelable.acknowledge(waiter.sleeper.fiber) catch |err| switch (err) {
957 error.Canceled => unreachable, // blocked
958 };
922959 }
923960
924961 fn unlock(mutex: *Mutex) void {
925 switch (@atomicRmw(usize, &mutex.num_waiters, .Sub, 1, .release)) {
926 0 => unreachable,
927 1 => {},
928 else => {
929 @branchHint(.unlikely);
930 mutex.queue.async(mutex, &Waiter.remove);
931 },
962 const state = @atomicRmw(State, &mutex.state, .And, .{
963 .locked = false,
964 .num_waiters = std.math.maxInt(State.NumWaiters),
965 }, .release);
966 if (state.num_waiters > 0) {
967 @branchHint(.unlikely);
968 mutex.queue.async(mutex, &Waiter.remove);
932969 }
933970 }
934971};
......@@ -1517,10 +1554,10 @@ const Futex = struct {
15171554 };
15181555 }
15191556
1520 fn tryAdd(waiter: *Waiter) Cancelable.AwaitError!void {
1557 fn tryAdd(waiter: *Waiter) Cancelable.RequestedError!void {
15211558 if (@atomicLoad(u32, waiter.ptr, .monotonic) != waiter.expected)
15221559 return error.CancelRequested;
1523 try waiter.cancelable.await(waiter.sleeper.fiber);
1560 try waiter.cancelable.enter(waiter.sleeper.fiber);
15241561 const futex = waiter.futex;
15251562 switch (waiter.timeout) {
15261563 .FOREVER => {},
......@@ -1542,46 +1579,28 @@ const Futex = struct {
15421579
15431580 fn canceled(context: ?*anyopaque) callconv(.c) void {
15441581 const cancelable: *Cancelable = @ptrCast(@alignCast(context));
1545 cancelable.cancel = Cancelable.is_requested;
15461582 const waiter: *Waiter = @fieldParentPtr("cancelable", cancelable);
1547 assert(@atomicRmw(
1548 Fiber.CancelStatus,
1549 &waiter.sleeper.fiber.cancel_status,
1550 .Xchg,
1551 .{ .requested = true, .awaiting = .nothing },
1552 .monotonic,
1553 ) == Fiber.CancelStatus{ .requested = true, .awaiting = .fromCancelable(cancelable) });
1583 cancelable.requested(waiter.sleeper.fiber);
15541584 const futex = waiter.futex;
1555 waiter.removeUncancelable();
1585 waiter.remove();
15561586 assert(@atomicRmw(usize, &futex.num_waiters, .Sub, 1, .monotonic) >= 1);
15571587 }
15581588
15591589 fn timedOut(context: ?*anyopaque) callconv(.c) void {
15601590 const waiter: *Waiter = @ptrCast(@alignCast(context));
15611591 const futex = waiter.futex;
1562 waiter.remove() catch |err| switch (err) {
1592 waiter.tryRemove() catch |err| switch (err) {
15631593 error.CancelRequested => return,
15641594 };
15651595 assert(@atomicRmw(usize, &futex.num_waiters, .Sub, 1, .monotonic) >= 1);
15661596 }
15671597
1568 fn remove(waiter: *Waiter) Cancelable.AwaitError!void {
1569 if (waiter.cancelable.cancel != Cancelable.is_blocked) {
1570 @branchHint(.likely);
1571 const cancel_status = @atomicRmw(
1572 Fiber.CancelStatus,
1573 &waiter.sleeper.fiber.cancel_status,
1574 .And,
1575 .{ .requested = true, .awaiting = .nothing },
1576 .monotonic,
1577 );
1578 assert(cancel_status.awaiting.toCancelable() == &waiter.cancelable);
1579 if (cancel_status.requested) return error.CancelRequested;
1580 }
1581 waiter.removeUncancelable();
1598 fn tryRemove(waiter: *Waiter) Cancelable.RequestedError!void {
1599 try waiter.cancelable.leave(waiter.sleeper.fiber);
1600 waiter.remove();
15821601 }
15831602
1584 fn removeUncancelable(waiter: *Waiter) void {
1603 fn remove(waiter: *Waiter) void {
15851604 waiter.futex.waiters.remove(&waiter.node);
15861605 if (waiter.timer) |timer| timer.cancel() else wake(waiter);
15871606 }
......@@ -1614,7 +1633,7 @@ const Futex = struct {
16141633 @branchHint(.unlikely);
16151634 continue;
16161635 }
1617 waiter.remove() catch |err| switch (err) {
1636 waiter.tryRemove() catch |err| switch (err) {
16181637 error.CancelRequested => continue,
16191638 };
16201639 num_removed += 1;
......@@ -1720,7 +1739,7 @@ fn futexWait(
17201739 .leeway = ev.leeway,
17211740 };
17221741 ev.yield(.{ .futex_wait = &waiter });
1723 try waiter.cancelable.check(waiter.sleeper.fiber);
1742 try waiter.cancelable.acknowledge(waiter.sleeper.fiber);
17241743}
17251744
17261745fn futexWaitUncancelable(userdata: ?*anyopaque, ptr: *const u32, expected: u32) void {
......@@ -1734,7 +1753,7 @@ fn futexWaitUncancelable(userdata: ?*anyopaque, ptr: *const u32, expected: u32)
17341753 .leeway = ev.leeway,
17351754 };
17361755 ev.yield(.{ .futex_wait = &waiter });
1737 waiter.cancelable.check(waiter.sleeper.fiber) catch |err| switch (err) {
1756 waiter.cancelable.acknowledge(waiter.sleeper.fiber) catch |err| switch (err) {
17381757 error.Canceled => unreachable, // blocked
17391758 };
17401759}
lib/std/c/darwin/dispatch.zig+3-3
......@@ -39,14 +39,13 @@ pub const once_t = enum(isize) {
3939 _,
4040
4141 pub inline fn once(predicate: *once_t, context: ?*anyopaque, function: function_t) void {
42 if (@atomicLoad(once_t, predicate, .unordered) != .done) {
42 if (predicate.* != .done) {
4343 @branchHint(.unlikely);
4444 once_f(predicate, context, function);
4545 } else asm volatile ("" ::: .{ .memory = true });
4646 switch (builtin.mode) {
4747 .Debug, .ReleaseSafe => {},
48 .ReleaseFast, .ReleaseSmall => if (@atomicLoad(once_t, predicate, .unordered) != .done)
49 unreachable,
48 .ReleaseFast, .ReleaseSmall => if (predicate.* != .done) unreachable,
5049 }
5150 }
5251};
......@@ -110,6 +109,7 @@ const queue_s = opaque {
110109 pub const get_current = get_current_queue;
111110 pub const get_main = get_main_queue;
112111 pub const get_global = get_global_queue;
112 pub const TARGET_DEFAULT = TARGET_QUEUE_DEFAULT;
113113 pub const create_with_target = queue_create_with_target;
114114 pub const create = queue_create;
115115 pub const get_label = queue_get_label;