| ... | @@ -650,7 +650,8 @@ const SwitchMessage = struct { | ... | @@ -650,7 +650,8 @@ const SwitchMessage = struct { |
| 650 | mutex_wait: *Mutex.Waiter, | 650 | mutex_wait: *Mutex.Waiter, |
| 651 | futex_wait: *Futex.Waiter, | 651 | futex_wait: *Futex.Waiter, |
| 652 | futex_wake: *Futex.Waker, | 652 | futex_wake: *Futex.Waker, |
| 653 | sleep: c.dispatch.time_t, | 653 | sleep_wait: *SleepWaiter, |
| | 654 | after: c.dispatch.time_t, |
| 654 | destroy, | 655 | destroy, |
| 655 | exit, | 656 | exit, |
| 656 | }; | 657 | }; |
| ... | @@ -698,7 +699,17 @@ const SwitchMessage = struct { | ... | @@ -698,7 +699,17 @@ const SwitchMessage = struct { |
| 698 | .init(ev.queue, @alignCast(@fieldParentPtr("context", message.contexts.old))); | 699 | .init(ev.queue, @alignCast(@fieldParentPtr("context", message.contexts.old))); |
| 699 | waker.futex.queue.async(waker, &Futex.Waker.remove); | 700 | waker.futex.queue.async(waker, &Futex.Waker.remove); |
| 700 | }, | 701 | }, |
| 701 | .sleep => |when| { | 702 | .sleep_wait => |waiter| { |
| | 703 | waiter.sleeper = |
| | 704 | .init(ev.queue, @alignCast(@fieldParentPtr("context", message.contexts.old))); |
| | 705 | const queue = waiter.cancelable.queue; |
| | 706 | switch (waiter.sleeper.fiber.cancel_protection.check()) { |
| | 707 | .unblocked => {}, |
| | 708 | .blocked => waiter.cancelable = .blocked, |
| | 709 | } |
| | 710 | queue.async(waiter, &SleepWaiter.start); |
| | 711 | }, |
| | 712 | .after => |when| { |
| 702 | const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old)); | 713 | const fiber: *Fiber = @alignCast(@fieldParentPtr("context", message.contexts.old)); |
| 703 | when.after(ev.queue, fiber, &Fiber.@"resume"); | 714 | when.after(ev.queue, fiber, &Fiber.@"resume"); |
| 704 | }, | 715 | }, |
| ... | @@ -1563,7 +1574,7 @@ const Futex = struct { | ... | @@ -1563,7 +1574,7 @@ const Futex = struct { |
| 1563 | .FOREVER => {}, | 1574 | .FOREVER => {}, |
| 1564 | else => |timeout| { | 1575 | else => |timeout| { |
| 1565 | const timer = c.dispatch.source_create(.TIMER, 0, .none, futex.queue) orelse { | 1576 | const timer = c.dispatch.source_create(.TIMER, 0, .none, futex.queue) orelse { |
| 1566 | log.warn("unable to create timer for futex timeout", .{}); | 1577 | log.warn("failed to create timer for futex timeout", .{}); |
| 1567 | return error.CancelRequested; | 1578 | return error.CancelRequested; |
| 1568 | }; | 1579 | }; |
| 1569 | timer.as_object().set_context(waiter); | 1580 | timer.as_object().set_context(waiter); |
| ... | @@ -4738,9 +4749,69 @@ fn clockResolution(userdata: ?*anyopaque, clock: Io.Clock) Io.Clock.ResolutionEr | ... | @@ -4738,9 +4749,69 @@ fn clockResolution(userdata: ?*anyopaque, clock: Io.Clock) Io.Clock.ResolutionEr |
| 4738 | }; | 4749 | }; |
| 4739 | } | 4750 | } |
| 4740 | | 4751 | |
| | 4752 | const SleepWaiter = struct { |
| | 4753 | sleeper: Sleeper = undefined, |
| | 4754 | cancelable: Cancelable, |
| | 4755 | timer: c.dispatch.source_t, |
| | 4756 | started: bool = false, |
| | 4757 | |
| | 4758 | fn start(context: ?*anyopaque) callconv(.c) void { |
| | 4759 | const waiter: *SleepWaiter = @ptrCast(@alignCast(context)); |
| | 4760 | waiter.cancelable.enter(waiter.sleeper.fiber) catch |err| switch (err) { |
| | 4761 | error.CancelRequested => waiter.timer.cancel(), |
| | 4762 | }; |
| | 4763 | waiter.timer.as_object().activate(); |
| | 4764 | } |
| | 4765 | |
| | 4766 | fn timedOut(context: ?*anyopaque) callconv(.c) void { |
| | 4767 | const waiter: *SleepWaiter = @ptrCast(@alignCast(context)); |
| | 4768 | waiter.cancelable.leave(waiter.sleeper.fiber) catch |err| switch (err) { |
| | 4769 | error.CancelRequested => return, |
| | 4770 | }; |
| | 4771 | waiter.timer.cancel(); |
| | 4772 | } |
| | 4773 | |
| | 4774 | fn canceled(context: ?*anyopaque) callconv(.c) void { |
| | 4775 | const cancelable: *Cancelable = @ptrCast(@alignCast(context)); |
| | 4776 | const waiter: *SleepWaiter = @fieldParentPtr("cancelable", cancelable); |
| | 4777 | cancelable.requested(waiter.sleeper.fiber); |
| | 4778 | waiter.timer.cancel(); |
| | 4779 | } |
| | 4780 | |
| | 4781 | fn wake(context: ?*anyopaque) callconv(.c) void { |
| | 4782 | const waiter: *SleepWaiter = @ptrCast(@alignCast(context)); |
| | 4783 | var sleeper = waiter.sleeper; |
| | 4784 | waiter.* = undefined; |
| | 4785 | Sleeper.wake(&sleeper); |
| | 4786 | } |
| | 4787 | }; |
| | 4788 | |
| 4741 | fn sleep(userdata: ?*anyopaque, timeout: Io.Timeout) Io.Cancelable!void { | 4789 | fn sleep(userdata: ?*anyopaque, timeout: Io.Timeout) Io.Cancelable!void { |
| 4742 | const ev: *Evented = @ptrCast(@alignCast(userdata)); | 4790 | const ev: *Evented = @ptrCast(@alignCast(userdata)); |
| 4743 | ev.yield(.{ .sleep = ev.timeFromTimeout(timeout) }); | 4791 | const queue = c.dispatch.queue_create_with_target( |
| | 4792 | "org.ziglang.std.Io.Dispatch.sleep", |
| | 4793 | .SERIAL(), |
| | 4794 | ev.queue, |
| | 4795 | ) orelse { |
| | 4796 | log.warn("failed to create serial queue for sleep", .{}); |
| | 4797 | return ev.yield(.{ .after = ev.timeFromTimeout(timeout) }); |
| | 4798 | }; |
| | 4799 | defer queue.as_object().release(); |
| | 4800 | const timer = c.dispatch.source_create(.TIMER, 0, .none, queue) orelse { |
| | 4801 | log.warn("failed to create timer for sleep", .{}); |
| | 4802 | return ev.yield(.{ .after = ev.timeFromTimeout(timeout) }); |
| | 4803 | }; |
| | 4804 | var waiter: SleepWaiter = .{ |
| | 4805 | .cancelable = .{ .queue = queue, .cancel = &Futex.Waiter.canceled }, |
| | 4806 | .timer = timer, |
| | 4807 | }; |
| | 4808 | timer.as_object().set_context(&waiter); |
| | 4809 | timer.set_event_handler(&SleepWaiter.timedOut); |
| | 4810 | timer.set_cancel_handler(&SleepWaiter.wake); |
| | 4811 | timer.set_timer(ev.timeFromTimeout(timeout), c.dispatch.TIME_FOREVER, ev.leeway); |
| | 4812 | ev.yield(.{ .sleep_wait = &waiter }); |
| | 4813 | timer.as_object().release(); |
| | 4814 | try waiter.cancelable.acknowledge(waiter.sleeper.fiber); |
| 4744 | } | 4815 | } |
| 4745 | | 4816 | |
| 4746 | fn timeFromTimeout(ev: *Evented, timeout: Io.Timeout) c.dispatch.time_t { | 4817 | fn timeFromTimeout(ev: *Evented, timeout: Io.Timeout) c.dispatch.time_t { |