From f9c15d12a702f8bd372057b147455b71589e6e8e Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Tue, 15 Sep 2026 12:13:29 -0300 Subject: [PATCH 01/10] fix: make libuv timer and signal callbacks safe against re-entrant continuations This PR fixes use-after-free crashes when a (sync := true) continuation of a timer or signal promise cancels or stops the handle, marks timer promises multi-threaded, restores timer state when starting fails, and no longer asserts when the caller resolved the promise itself. Co-Authored-By: Claude Opus 5 (1M context) --- src/Std/Internal/UV/Signal.lean | 4 +- src/Std/Internal/UV/Timer.lean | 3 + src/runtime/uv/signal.cpp | 68 +++++++++++----- src/runtime/uv/timer.cpp | 79 +++++++++++++------ tests/compile/uv_callback_reentrancy.lean | 44 +++++++++++ .../uv_callback_reentrancy.lean.out.expected | 1 + .../compile/uv_oneshot_cancel_reentrancy.lean | 35 ++++++++ ...neshot_cancel_reentrancy.lean.out.expected | 1 + tests/compile/uv_repeating_reentrancy.lean | 49 ++++++++++++ .../uv_repeating_reentrancy.lean.out.expected | 1 + tests/elab/uv_oneshot_caller_resolved.lean | 26 ++++++ 11 files changed, 268 insertions(+), 43 deletions(-) create mode 100644 tests/compile/uv_callback_reentrancy.lean create mode 100644 tests/compile/uv_callback_reentrancy.lean.out.expected create mode 100644 tests/compile/uv_oneshot_cancel_reentrancy.lean create mode 100644 tests/compile/uv_oneshot_cancel_reentrancy.lean.out.expected create mode 100644 tests/compile/uv_repeating_reentrancy.lean create mode 100644 tests/compile/uv_repeating_reentrancy.lean.out.expected create mode 100644 tests/elab/uv_oneshot_caller_resolved.lean diff --git a/src/Std/Internal/UV/Signal.lean b/src/Std/Internal/UV/Signal.lean index d34666cf5bd1..0a6c68d4d9bc 100644 --- a/src/Std/Internal/UV/Signal.lean +++ b/src/Std/Internal/UV/Signal.lean @@ -62,7 +62,9 @@ This function has different behavior depending on the state and configuration of - if it is finished, return the last `IO.Promise` created by `next`. Notably this could be one that never resolves if the signal handler was stopped before fulfilling the last one. -The resolved `IO.Promise` contains the signal number that was received. +A promise from `next` may also be resolved by the code holding it; the handler then treats it as +fulfilled when the signal arrives. The resolved `IO.Promise` contains the signal number that was +received. -/ @[extern "lean_uv_signal_next"] opaque next (signal : @& Signal) : IO (IO.Promise Int) diff --git a/src/Std/Internal/UV/Timer.lean b/src/Std/Internal/UV/Timer.lean index 28d483d6bb85..34a97303f820 100644 --- a/src/Std/Internal/UV/Timer.lean +++ b/src/Std/Internal/UV/Timer.lean @@ -59,6 +59,9 @@ This function has different behavior depending on the state and configuration of This ensures that the returned `IO.Promise` resolves at the next repetition of the timer. - if it is finished, return the last `IO.Promise` created by `next`. Notably this could be one that never resolves if the timer was stopped before fulfilling the last one. + +A promise from `next` may also be resolved by the code holding it; the timer then treats it as +fulfilled when it fires. -/ @[extern "lean_uv_timer_next"] opaque next (timer : @& Timer) : IO (IO.Promise Unit) diff --git a/src/runtime/uv/signal.cpp b/src/runtime/uv/signal.cpp index 6ea7a4f8e1cc..3c0df6f75548 100644 --- a/src/runtime/uv/signal.cpp +++ b/src/runtime/uv/signal.cpp @@ -53,19 +53,34 @@ void handle_signal_event(uv_signal_t* handle, int signum) { if (signal->m_repeating) { if (!signal_promise_is_finished(signal)) { - lean_object* res = lean_io_promise_resolve(lean_box(signum), signal->m_promise); + // Rule 1: a continuation may `cancel` or `stop` the signal, releasing the field's + // reference. + lean_object * promise = signal->m_promise; + lean_inc(promise); + lean_object* res = lean_io_promise_resolve(lean_box(signum), promise); lean_dec(res); + lean_dec(promise); } } else { - if (signal->m_promise != NULL) { - lean_object* res = lean_io_promise_resolve(lean_box(signum), signal->m_promise); - lean_dec(res); - } - uv_signal_stop(signal->m_uv_signal); signal->m_state = SIGNAL_STATE_FINISHED; + lean_object * promise = signal->m_promise; + if (promise != NULL) { + lean_inc(promise); + } + lean_dec(obj); + + // Rule 1: nothing below may touch the signal. Code holding the promise may have resolved + // it already. + if (promise != NULL) { + if (!promise_is_resolved(promise)) { + lean_object* res = lean_io_promise_resolve(lean_box(signum), promise); + lean_dec(res); + } + lean_dec(promise); + } } } @@ -208,9 +223,10 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_next(b_obj_arg obj) { signal->m_promise = create_promise(); } - lean_inc(signal->m_promise); + lean_object * promise = signal->m_promise; + lean_inc(promise); event_loop_unlock(&global_ev); - return lean_io_result_mk_ok(signal->m_promise); + return lean_io_result_mk_ok(promise); } case SIGNAL_STATE_FINISHED: { @@ -220,18 +236,20 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_next(b_obj_arg obj) { return lean_io_result_mk_ok(finished_promise); } - lean_inc(signal->m_promise); + lean_object * promise = signal->m_promise; + lean_inc(promise); event_loop_unlock(&global_ev); - return lean_io_result_mk_ok(signal->m_promise); + return lean_io_result_mk_ok(promise); } } } else { if (signal->m_state == SIGNAL_STATE_INITIAL) { return setup_signal(); } else if (signal->m_promise != NULL) { - lean_inc(signal->m_promise); + lean_object * promise = signal->m_promise; + lean_inc(promise); event_loop_unlock(&global_ev); - return lean_io_result_mk_ok(signal->m_promise); + return lean_io_result_mk_ok(promise); } else { lean_object* finished_promise = create_promise(); event_loop_unlock(&global_ev); @@ -276,20 +294,32 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_cancel(b_obj_arg obj) { // It's locking here to avoid changing the state during other operations. event_loop_lock(&global_ev); + lean_object * promise = NULL; + bool release_signal = false; + if (signal->m_state == SIGNAL_STATE_RUNNING && signal->m_promise != NULL) { - if (signal->m_repeating) { - lean_dec(signal->m_promise); - signal->m_promise = NULL; - } else { + promise = signal->m_promise; + signal->m_promise = NULL; + + // A repeating signal keeps listening, so the loop keeps its reference until `stop`. + if (!signal->m_repeating) { uv_signal_stop(signal->m_uv_signal); - lean_dec(signal->m_promise); - signal->m_promise = NULL; signal->m_state = SIGNAL_STATE_INITIAL; - lean_dec(obj); + release_signal = true; } } event_loop_unlock(&global_ev); + + // Rules 1 and 2: the cancellation is complete and the lock dropped before releasing. + if (promise != NULL) { + lean_dec(promise); + } + + if (release_signal) { + lean_dec(obj); + } + return lean_io_result_mk_ok(lean_box(0)); } diff --git a/src/runtime/uv/timer.cpp b/src/runtime/uv/timer.cpp index e7df467ed9d9..3c5ec39bc50d 100644 --- a/src/runtime/uv/timer.cpp +++ b/src/runtime/uv/timer.cpp @@ -58,22 +58,35 @@ void handle_timer_event(uv_timer_t* handle) { if (timer->m_repeating) { // For repeating timers, only resolves if the promise exists and is not finished if (timer->m_promise != NULL && !timer_promise_is_finished(timer)) { - lean_object* res = lean_io_promise_resolve(lean_box(0), timer->m_promise); + // Rule 1: a continuation may `cancel` or `stop` the timer, releasing the field's + // reference. + lean_object * promise = timer->m_promise; + lean_inc(promise); + lean_object* res = lean_io_promise_resolve(lean_box(0), promise); lean_dec(res); + lean_dec(promise); } } else { - // For non-repeating timers, resolves if the promise exists - if (timer->m_promise != NULL) { - lean_assert(!timer_promise_is_finished(timer)); - lean_object* res = lean_io_promise_resolve(lean_box(0), timer->m_promise); - lean_dec(res); - } - uv_timer_stop(timer->m_uv_timer); timer->m_state = TIMER_STATE_FINISHED; + lean_object * promise = timer->m_promise; + if (promise != NULL) { + lean_inc(promise); + } + // The loop does not need to keep the timer alive anymore. lean_dec(obj); + + // Rule 1: nothing below may touch the timer. Code holding the promise may have resolved it + // already. + if (promise != NULL) { + if (!promise_is_resolved(promise)) { + lean_object* res = lean_io_promise_resolve(lean_box(0), promise); + lean_dec(res); + } + lean_dec(promise); + } } } @@ -118,7 +131,10 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_timer_next(b_obj_arg obj) { lean_uv_timer_object * timer = lean_to_uv_timer(obj); auto create_promise = []() { - return lean_io_promise_new(); + lean_object * promise = lean_io_promise_new(); + // The loop thread resolves and releases it, so its refcount has to be atomic. + mark_mt(promise); + return promise; }; auto setup_timer = [create_promise, obj, timer]() { @@ -140,6 +156,12 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_timer_next(b_obj_arg obj) { ); if (result != 0) { + // A failed start must not leave the timer advertising a promise the loop will settle. + timer->m_state = TIMER_STATE_INITIAL; + timer->m_promise = NULL; + + lean_dec(promise); // The structure does not own it. + lean_dec(promise); // We are not going to return it. lean_dec(obj); event_loop_unlock(&global_ev); @@ -169,18 +191,20 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_timer_next(b_obj_arg obj) { timer->m_promise = create_promise(); } - lean_inc(timer->m_promise); + lean_object * promise = timer->m_promise; + lean_inc(promise); event_loop_unlock(&global_ev); - return lean_io_result_mk_ok(timer->m_promise); + return lean_io_result_mk_ok(promise); } case TIMER_STATE_FINISHED: { if (timer->m_promise != NULL) { - lean_inc(timer->m_promise); + lean_object * promise = timer->m_promise; + lean_inc(promise); event_loop_unlock(&global_ev); - return lean_io_result_mk_ok(timer->m_promise); + return lean_io_result_mk_ok(promise); } else { // Creates a resolved promise lean_object* finished_promise = create_promise(); @@ -272,24 +296,33 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_timer_cancel(b_obj_arg obj) { // It's locking here to avoid changing the state during other operations. event_loop_lock(&global_ev); + lean_object * promise = NULL; + bool release_timer = false; + if (timer->m_state == TIMER_STATE_RUNNING && timer->m_promise != NULL) { - if (timer->m_repeating) { - lean_dec(timer->m_promise); - timer->m_promise = NULL; - } else { - uv_timer_stop(timer->m_uv_timer); + promise = timer->m_promise; + timer->m_promise = NULL; - lean_dec(timer->m_promise); - timer->m_promise = NULL; + // A repeating timer keeps running, so the loop keeps its reference until `stop`. + if (!timer->m_repeating) { + uv_timer_stop(timer->m_uv_timer); timer->m_state = TIMER_STATE_INITIAL; - - // The loop does not need to keep the timer alive anymore. - lean_dec(obj); + release_timer = true; } } event_loop_unlock(&global_ev); + // Rules 1 and 2: the cancellation is complete and the lock dropped before releasing. + if (promise != NULL) { + lean_dec(promise); + } + + if (release_timer) { + // The loop does not need to keep the timer alive anymore. + lean_dec(obj); + } + return lean_io_result_mk_ok(lean_box(0)); } diff --git a/tests/compile/uv_callback_reentrancy.lean b/tests/compile/uv_callback_reentrancy.lean new file mode 100644 index 000000000000..bf58dc27bf85 --- /dev/null +++ b/tests/compile/uv_callback_reentrancy.lean @@ -0,0 +1,44 @@ +import Std.Internal.UV + +/-! +One-shot timer and signal callbacks resolved their promise and only then stopped the handle and +released the event loop's reference. A `(sync := true)` continuation runs inside that resolve, on +the loop thread; a `cancel` or `stop` from it released the loop's reference as well, so the callback +released it a second time and used the freed handle. +-/ + +open Std.Internal.UV + +def timerFromCallback (useStop : Bool) : IO Unit := do + for _ in [0:200] do + let timer ← Timer.mk 1 false + let fired ← timer.next + BaseIO.chainTask (sync := true) fired.result? fun _ => do + let _ ← (if useStop then timer.stop else timer.cancel : IO _).toBaseIO + let _ ← IO.wait fired.result? + +/-- +Sends SIGWINCH, whose default action is to ignore it, to a one-shot signal handler. `received` stays +referenced by the main thread, which only waits on `done`. +-/ +def signalFromCallback (useStop : Bool) : IO Unit := do + let pid ← IO.Process.getPID + for _ in [0:30] do + let signal ← Signal.mk 28 false + let received ← signal.next + let done ← IO.Promise.new + BaseIO.chainTask (sync := true) received.result? fun _ => do + let _ ← (if useStop then signal.stop else signal.cancel : IO _).toBaseIO + done.resolve () + discard <| IO.Process.output { cmd := "kill", args := #["-WINCH", toString pid] } + let _ ← IO.wait done.result? + let _ ← received.isResolved + +def main : IO Unit := do + timerFromCallback (useStop := false) + timerFromCallback (useStop := true) + -- Signals cannot be sent with `kill` on Windows. + unless System.Platform.isWindows do + signalFromCallback (useStop := false) + signalFromCallback (useStop := true) + IO.println "exiting" diff --git a/tests/compile/uv_callback_reentrancy.lean.out.expected b/tests/compile/uv_callback_reentrancy.lean.out.expected new file mode 100644 index 000000000000..26b03a5b3483 --- /dev/null +++ b/tests/compile/uv_callback_reentrancy.lean.out.expected @@ -0,0 +1 @@ +exiting diff --git a/tests/compile/uv_oneshot_cancel_reentrancy.lean b/tests/compile/uv_oneshot_cancel_reentrancy.lean new file mode 100644 index 000000000000..e1b1d01488f8 --- /dev/null +++ b/tests/compile/uv_oneshot_cancel_reentrancy.lean @@ -0,0 +1,35 @@ +import Std.Internal.UV + +/-! +`Timer.cancel` and `Signal.cancel` released the pending promise before clearing the field holding +it. When the handle held the only reference, that release resolved the promise's `result?` to `none` +and ran a `(sync := true)` continuation inline, which re-entered `cancel` and found the handle still +running with the stale field. For a one-shot handle it then released the event loop's reference to +the handle object a second time, which this test detects. +-/ + +open Std.Internal.UV + +/-- Leaves the handle holding the only reference to a pending promise. -/ +def arm (next : IO (IO.Promise α)) (cancel : IO Unit) : IO Unit := do + let pending ← next + BaseIO.chainTask (sync := true) pending.result? fun _ => do + let _ ← cancel.toBaseIO + +def timerCancel : IO Unit := do + for _ in [0:200] do + let timer ← Timer.mk 3600000 false + arm timer.next timer.cancel + timer.cancel + +/-- SIGWINCH is never sent here; the handler is only armed and cancelled. -/ +def signalCancel : IO Unit := do + for _ in [0:200] do + let signal ← Signal.mk 28 false + arm signal.next signal.cancel + signal.cancel + +def main : IO Unit := do + timerCancel + signalCancel + IO.println "exiting" diff --git a/tests/compile/uv_oneshot_cancel_reentrancy.lean.out.expected b/tests/compile/uv_oneshot_cancel_reentrancy.lean.out.expected new file mode 100644 index 000000000000..26b03a5b3483 --- /dev/null +++ b/tests/compile/uv_oneshot_cancel_reentrancy.lean.out.expected @@ -0,0 +1 @@ +exiting diff --git a/tests/compile/uv_repeating_reentrancy.lean b/tests/compile/uv_repeating_reentrancy.lean new file mode 100644 index 000000000000..5efba43cadc9 --- /dev/null +++ b/tests/compile/uv_repeating_reentrancy.lean @@ -0,0 +1,49 @@ +import Std.Internal.UV + +/-! +Regression guard for repeating timers used from re-entrant and concurrent callers. + +A `(sync := true)` continuation of a repeating timer's promise runs inside the firing callback's +resolve, and cancels, re-arms and stops the timer from there. + +`Timer.next` must return the promise it read under the event loop lock rather than reading +`m_promise` again after unlocking, when a concurrent `cancel` may have cleared it. The promises it +creates are refcounted from the loop thread as well, so they must be marked multi-threaded. +-/ + +open Std.Internal.UV + +/-- Cancels and re-arms a repeating timer from inside its own firing callback. -/ +def repeatingTimerFromCallback : IO Unit := do + for _ in [0:200] do + let timer ← Timer.mk 1 true + let fired ← timer.next + BaseIO.chainTask (sync := true) fired.result? fun _ => do + let _ ← (timer.cancel : IO _).toBaseIO + let _ ← (timer.next : IO _).toBaseIO + let _ ← (timer.stop : IO _).toBaseIO + pure () + let _ ← IO.wait fired.result? + +/-- +Races `next` against `cancel` on a shared repeating timer. The promises are deliberately dropped +rather than awaited: `cancel` orphans the outstanding one, so waiting on it would block forever. +-/ +def timerNextRacesCancel : IO Unit := do + for _ in [0:100] do + let timer ← Timer.mk 1 true + let arming ← IO.asTask do + for _ in [0:200] do + let _ ← timer.next + pure () + let cancelling ← IO.asTask do + for _ in [0:200] do + timer.cancel + IO.ofExcept (← IO.wait arming) + IO.ofExcept (← IO.wait cancelling) + timer.stop + +def main : IO Unit := do + repeatingTimerFromCallback + timerNextRacesCancel + IO.println "exiting" diff --git a/tests/compile/uv_repeating_reentrancy.lean.out.expected b/tests/compile/uv_repeating_reentrancy.lean.out.expected new file mode 100644 index 000000000000..26b03a5b3483 --- /dev/null +++ b/tests/compile/uv_repeating_reentrancy.lean.out.expected @@ -0,0 +1 @@ +exiting diff --git a/tests/elab/uv_oneshot_caller_resolved.lean b/tests/elab/uv_oneshot_caller_resolved.lean new file mode 100644 index 000000000000..8327ca2e6d00 --- /dev/null +++ b/tests/elab/uv_oneshot_caller_resolved.lean @@ -0,0 +1,26 @@ +import Std.Internal.UV + +/-! +A one-shot libuv timer whose promise the caller resolved before it fired. + +`IO.Promise.resolve` is available to whoever holds the promise returned by `next`, so the loop must +treat an already-resolved promise as settled when the timer fires. Builds with assertions enabled +used to abort in the callback instead. +-/ + +open Std.Internal.UV + +/-- Waits for a later timer; timers fire in order, so the earlier one has fired by then. -/ +def settle (ms : UInt64) : IO Unit := do + let later ← Timer.mk ms false + discard <| IO.wait (← later.next).result? + +#eval show IO Unit from do + let timer ← Timer.mk 5 false + let fired ← timer.next + fired.resolve () + settle 50 + -- The timer finished on its own; `next` still hands out the promise the caller resolved. + unless ← (← timer.next).isResolved do + throw <| IO.userError "timer promise not resolved" + timer.stop From 34efb9b31f9d85f73cc5b984e981f8e7214e878e Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Tue, 15 Sep 2026 12:13:29 -0300 Subject: [PATCH 02/10] fix: keep a repeating libuv timer with a 0 ms period ticking This PR makes a repeating libuv timer created with a 0 ms period keep ticking every millisecond instead of firing only once. Co-Authored-By: Claude Opus 5 (1M context) --- src/Std/Internal/UV/Timer.lean | 4 +++- src/runtime/uv/timer.cpp | 3 ++- tests/elab/uv_timer_zero_period.lean | 23 +++++++++++++++++++++++ 3 files changed, 28 insertions(+), 2 deletions(-) create mode 100644 tests/elab/uv_timer_zero_period.lean diff --git a/src/Std/Internal/UV/Timer.lean b/src/Std/Internal/UV/Timer.lean index 34a97303f820..32045ed69230 100644 --- a/src/Std/Internal/UV/Timer.lean +++ b/src/Std/Internal/UV/Timer.lean @@ -39,7 +39,9 @@ This creates a `Timer` in the initial state and doesn't run it yet. milliseconds, counting from when it's run. - If `repeating` is `true` this constructs a timer that resolves after multiples of `timeout` milliseconds, counting from when it's run. Note that this includes the 0th multiple right after - starting the timer. Furthermore a repeating timer will only be freed after `Timer.stop` is called. + starting the timer. A `timeout` of 0 ticks every millisecond, and `reset` then delays the next + tick by 1 millisecond. Furthermore a repeating timer will only be freed after `Timer.stop` is + called. -/ @[extern "lean_uv_timer_mk"] opaque mk (timeout : UInt64) (repeating : Bool) : IO Timer diff --git a/src/runtime/uv/timer.cpp b/src/runtime/uv/timer.cpp index 3c5ec39bc50d..f741236cec0f 100644 --- a/src/runtime/uv/timer.cpp +++ b/src/runtime/uv/timer.cpp @@ -96,7 +96,8 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_timer_mk(uint64_t timeout, uint8_t r if (timer == nullptr) { return lean_io_result_mk_error(decode_io_error(ENOMEM, nullptr)); } - timer->m_timeout = timeout; + // libuv treats a repeat period of 0 as a one-shot timer. + timer->m_timeout = repeating && timeout == 0 ? 1 : timeout; timer->m_repeating = repeating; timer->m_state = TIMER_STATE_INITIAL; timer->m_promise = NULL; diff --git a/tests/elab/uv_timer_zero_period.lean b/tests/elab/uv_timer_zero_period.lean new file mode 100644 index 000000000000..1fee8b670dbb --- /dev/null +++ b/tests/elab/uv_timer_zero_period.lean @@ -0,0 +1,23 @@ +import Std.Internal.UV + +/-! +A repeating timer with a 0 ms period keeps ticking instead of firing once. libuv treats a repeat +period of 0 as a one-shot timer, so every promise after the 0th used to stay pending forever. +-/ + +open Std.Internal.UV + +/-- Waits for `p`, failing after 5 s instead of hanging the test. -/ +def awaitBounded {α : Type} (what : String) (p : IO.Promise α) : IO Unit := do + let deadlineTimer ← Timer.mk 5000 false + let deadline ← deadlineTimer.next + let resolved ← IO.waitAny [p.result?.map (·.isSome), deadline.result?.map (fun _ => false)] + deadlineTimer.cancel + unless resolved do + throw <| IO.userError s!"{what}: not resolved within 5 s" + +#eval show IO Unit from do + let zero ← Timer.mk 0 true + for i in [0:5] do + awaitBounded s!"0 ms period, tick {i}" (← zero.next) + zero.stop From 8bd4b1dfcf9c98d9fcd76033e5083565e05756ec Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Tue, 15 Sep 2026 12:13:29 -0300 Subject: [PATCH 03/10] fix: make Signal.Waiter.selector detect a received signal This PR makes Signal.Waiter.selector report a waiter that already received its signal as ready. Co-Authored-By: Claude Opus 5 (1M context) --- src/Std/Async/Signal.lean | 4 +- tests/elab/async_signal_selector.lean | 57 +++++++++++++++++++++++++++ 2 files changed, 59 insertions(+), 2 deletions(-) create mode 100644 tests/elab/async_signal_selector.lean diff --git a/src/Std/Async/Signal.lean b/src/Std/Async/Signal.lean index b18d89cc0414..d040c051a7b2 100644 --- a/src/Std/Async/Signal.lean +++ b/src/Std/Async/Signal.lean @@ -235,8 +235,8 @@ does not start the signal waiter. def selector (s : Signal.Waiter) : Selector Unit := { tryFn := do - let signalWaiter : AsyncTask _ ← async s.wait - if ← IO.hasFinished signalWaiter then + let signalWaiter ← s.native.next + if ← signalWaiter.isResolved then return some () else s.native.cancel diff --git a/tests/elab/async_signal_selector.lean b/tests/elab/async_signal_selector.lean new file mode 100644 index 000000000000..0b1453e13e94 --- /dev/null +++ b/tests/elab/async_signal_selector.lean @@ -0,0 +1,57 @@ +import Std.Async + +/-! +`Signal.Waiter.selector` against a competing sleep: with no signal sent the sleep must win, and once +the process sends itself the signal the waiter must win. Covers one-shot and repeating waiters. Also +checks that `Selectable.tryOne` reports a one-shot waiter that has already received its signal as +ready; `tryFn` used to test whether a task it had just spawned had finished, and so returned `none`. +SIGWINCH is used because its default action is to ignore it. +-/ + +open Std Async + +def race (repeating : Bool) (sendSignal : Bool) : Async Bool := do + let waiter ← Signal.Waiter.mk .sigwinch repeating + if sendSignal then + let pid ← IO.Process.getPID + discard <| IO.asTask do + IO.sleep 100 + discard <| IO.Process.output { cmd := "kill", args := #["-WINCH", toString pid] } + let sleep ← Selector.sleep (if sendSignal then 10000 else 200) + let won ← Selectable.one #[ + .case waiter.selector (fun _ => pure true), + .case sleep (fun _ => pure false)] + waiter.stop + return won + +def test : IO (List Bool) := do + -- Signals cannot be sent with `kill` on Windows; report the expected outcome. + if System.Platform.isWindows then + return [false, false, true, true] + (do return [← race false false, ← race true false, ← race false true, ← race true true] + : Async (List Bool)).block + +/-- info: [false, false, true, true] -/ +#guard_msgs in +#eval test + +/-- Counts how many of 20 one-shot waiters `tryOne` reports as ready after their signal arrived. -/ +def tryOneAfterDelivery : IO Nat := do + if System.Platform.isWindows then + return 20 + let pid ← IO.Process.getPID + let mut ready := 0 + for _ in [0:20] do + let waiter ← Signal.Waiter.mk .sigwinch false + let delivered ← waiter.wait + discard <| IO.Process.output { cmd := "kill", args := #["-WINCH", toString pid] } + discard <| IO.wait delivered + let r ← (Selectable.tryOne #[.case waiter.selector (fun _ => pure ())]).block + waiter.stop + if r.isSome then + ready := ready + 1 + return ready + +/-- info: 20 -/ +#guard_msgs in +#eval tryOneAfterDelivery From fcef767d1c68096598e3c2c292529b7c2075428f Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Tue, 15 Sep 2026 12:13:30 -0300 Subject: [PATCH 04/10] doc: correct libuv timer, signal and socket docs This PR corrects the libuv timer, signal and socket documentation to match their current behavior. Co-Authored-By: Claude Opus 5 (1M context) --- src/Std/Async/Signal.lean | 8 +++++--- src/Std/Async/Timer.lean | 19 +++++++++++-------- src/Std/Async/UDP.lean | 8 ++++++-- src/Std/Internal/UV/Signal.lean | 20 ++++++++++++-------- src/Std/Internal/UV/TCP.lean | 9 +++++++-- src/Std/Internal/UV/Timer.lean | 8 +++++--- src/Std/Internal/UV/UDP.lean | 12 ++++++++---- src/runtime/uv/signal.h | 8 +++++--- src/runtime/uv/timer.cpp | 6 ++++-- src/runtime/uv/timer.h | 8 +++++--- src/runtime/uv/udp.cpp | 2 ++ 11 files changed, 70 insertions(+), 38 deletions(-) diff --git a/src/Std/Async/Signal.lean b/src/Std/Async/Signal.lean index d040c051a7b2..bad4d6999b21 100644 --- a/src/Std/Async/Signal.lean +++ b/src/Std/Async/Signal.lean @@ -208,7 +208,9 @@ def mk (signum : Signal) (repeating : Bool) : IO Signal.Waiter := do If: - `s` is not yet running start listening and return an `AsyncTask` that will resolve once the previously configured signal is received. -- `s` is already or not anymore running return the same `AsyncTask` as the first call to `wait`. +- `s` is already running, or finished after receiving the signal, return the same `AsyncTask` as the + first call to `wait`. +- `s` was stopped with `stop` before receiving the signal, return an `AsyncTask` that fails. The resolved `AsyncTask` contains the signal number that was received. -/ @@ -220,8 +222,8 @@ def wait (s : Signal.Waiter) : IO (AsyncTask Int) := do /-- If: - `s` is still running this stops `s` without resolving any remaining `AsyncTask`s that were created - through `wait`. Note that if another `AsyncTask` is binding on any of these it is going hang - forever without further intervention. + through `wait`. Those tasks fail once the last reference to their promise is dropped, rather than + producing a value. - `s` is not yet or not anymore running this is a no-op. -/ @[inline] diff --git a/src/Std/Async/Timer.lean b/src/Std/Async/Timer.lean index 58dc8d2d1c6a..2a99f04970ac 100644 --- a/src/Std/Async/Timer.lean +++ b/src/Std/Async/Timer.lean @@ -39,7 +39,9 @@ def mk (duration : Std.Time.Millisecond.Offset) : Async Sleep := do If: - `s` is not yet running start it and return an `Async` computation that will complete once the previously configured `duration` has elapsed. -- `s` is already or not anymore running return the same `Async` computation as the first call to `wait`. +- `s` is already running, or finished after completing, return the same + `Async` computation as the first call to `wait`. +- `s` was stopped with `stop` before completing, return an `Async` computation that fails. -/ @[inline] def wait (s : Sleep) : Async Unit := @@ -57,9 +59,9 @@ def reset (s : Sleep) : Async Unit := /-- If: -- `s` is still running this stops `s` without completing any remaining `Async` computations that were created - through `wait`. Note that if another `Async` computation is binding on any of these it will hang - forever without further intervention. +- `s` is still running this stops `s` without completing any remaining `Async` computations that + were created through `wait`. Those computations fail once the last reference to their promise is + dropped, rather than producing a value. - `s` is not yet or not anymore running this is a no-op. -/ @[inline] @@ -136,7 +138,8 @@ If: call - the tick from the last call of `i` has finished return a new `Async` computation that waits for the closest next tick from the time of calling this function. -- `i` is not running anymore this is a no-op. +- `i` is not running anymore, the returned `Async` computation fails, as `stop` dropped the promise + it would have completed. -/ @[inline] def tick (i : Interval) : Async Unit := do @@ -154,9 +157,9 @@ def reset (i : Interval) : IO Unit := /-- If: -- `i` is still running this stops `i` without completing any remaining `Async` computations that were created - through `tick`. Note that if another `Async` computation is binding on any of these it will hang - forever without further intervention. +- `i` is still running this stops `i` without completing any remaining `Async` computations that + were created through `tick`. Those computations fail once the last reference to their promise is + dropped, rather than producing a value. - `i` is not yet or not anymore running this is a no-op. -/ @[inline] diff --git a/src/Std/Async/UDP.lean b/src/Std/Async/UDP.lean index 14d7e08220af..ade9c47ea497 100644 --- a/src/Std/Async/UDP.lean +++ b/src/Std/Async/UDP.lean @@ -79,9 +79,11 @@ Receives data from an UDP socket. `size` is for the maximum bytes to receive. The promise resolves when some data is available or an error occurs. If the socket has not been previously bound with `bind`, it is automatically bound to `0.0.0.0` (all interfaces) with a random port. -If a datagram larger than `size` arrives, it is discarded in its entirety and an -`IO.Error.resourceExhausted` error is thrown. Furthermore calling this function in parallel with `recvSelector` is not supported. + +A datagram larger than `size` is discarded in its entirety, and this throws `EMSGSIZE` (an +`IO.Error.resourceExhausted`) instead of returning a truncated prefix. The socket stays usable, so a +receive loop should catch the error per datagram. -/ @[inline] def recv (s : Socket) (size : UInt64) : Async (ByteArray × Option SocketAddress) := @@ -93,6 +95,8 @@ and provides that data. If the socket has not been previously bound with `bind`, automatically bound to `0.0.0.0` (all interfaces) with a random port. Calling this function does starts the data wait, only when it's used with `Selectable.one` or `combine`. It must not be called in parallel with `recv`. + +Fails with `EMSGSIZE` if the datagram is larger than `size`, like `recv`. -/ def recvSelector (s : Socket) (size : UInt64) : Selector (ByteArray × Option SocketAddress) := { diff --git a/src/Std/Internal/UV/Signal.lean b/src/Std/Internal/UV/Signal.lean index 0a6c68d4d9bc..28f096b30bd4 100644 --- a/src/Std/Internal/UV/Signal.lean +++ b/src/Std/Internal/UV/Signal.lean @@ -40,8 +40,9 @@ This creates a `Signal` in the initial state and doesn't start listening yet. - If `repeating` is `false` this constructs a signal handler that resolves once when the specified signal `signum` is received, then automatically stops listening. - If `repeating` is `true` this constructs a signal handler that resolves each time the specified - signal `signum` is received and continues listening. A repeating signal handler will only be - freed after `Signal.stop` is called. + signal `signum` is received and continues listening. While it listens, a signal that arrives with + no promise from `next` pending is consumed without being reported. A repeating signal handler + will only be freed after `Signal.stop` is called. -/ @[extern "lean_uv_signal_mk"] opaque mk (signum : Int32) (repeating : Bool) : IO Signal @@ -51,7 +52,10 @@ This function has different behavior depending on the state and configuration of - if `repeating` is `false` and: - it is initial, start listening and return a new `IO.Promise` that is set to resolve once the signal `signum` is received. After this `IO.Promise` is resolved the `Signal` is finished. - - it is running or finished, return the same `IO.Promise` that the first call to `next` returned. + - it is running, or finished after receiving the signal, return the same `IO.Promise` that the + first call to `next` returned. + - it was stopped with `stop` before receiving the signal, return a new `IO.Promise` that is never + resolved. - if `repeating` is `true` and: - it is initial, start listening and return a new `IO.Promise` that resolves when the next signal `signum` is received. @@ -59,8 +63,8 @@ This function has different behavior depending on the state and configuration of - If it is, return a new `IO.Promise` that resolves upon receiving the next signal - If it is not, return the last `IO.Promise` This ensures that the returned `IO.Promise` resolves at the next occurrence of the signal. - - if it is finished, return the last `IO.Promise` created by `next`. Notably this could be one - that never resolves if the signal handler was stopped before fulfilling the last one. + - if it is finished, return a new `IO.Promise` that is never resolved, as `stop` dropped the last + one. A promise from `next` may also be resolved by the code holding it; the handler then treats it as fulfilled when the signal arrives. The resolved `IO.Promise` contains the signal number that was @@ -72,9 +76,9 @@ opaque next (signal : @& Signal) : IO (IO.Promise Int) /-- This function has different behavior depending on the state of the `Signal`: - If it is initial or finished this is a no-op. -- If it is running the signal handler is stopped and it is put into the finished state. - Note that if the last `IO.Promise` generated by `next` is unresolved and being waited - on this creates a memory leak and the waiting task is not going to be awoken anymore. +- If it is running the signal handler is stopped, it is put into the finished state and the last + `IO.Promise` generated by `next` is released. If that promise is still pending it is never + resolved; once nothing else references it, its `result?` resolves to `none`. -/ @[extern "lean_uv_signal_stop"] opaque stop (signal : @& Signal) : IO Unit diff --git a/src/Std/Internal/UV/TCP.lean b/src/Std/Internal/UV/TCP.lean index 4b72f2b261c0..998c2ea6661c 100644 --- a/src/Std/Internal/UV/TCP.lean +++ b/src/Std/Internal/UV/TCP.lean @@ -23,6 +23,10 @@ private opaque SocketImpl : NonemptyType.{0} /-- Represents a TCP socket. + +While a `recv?`, `waitReadable` or `accept` is pending, the event loop keeps the socket alive even +if nothing else references it. Two connected sockets that each wait to receive from the other +therefore stay open until one receive is cancelled with `cancelRecv` or the program exits. -/ def Socket : Type := SocketImpl.type @@ -68,8 +72,9 @@ opaque waitReadable (socket : @& Socket) : IO (IO.Promise (Except IO.Error Bool) /-- Cancels a receive operation in the form of `recv?` or `waitReadable` if there is currently one -pending. This resolves their returned `IO.Promise` to `none`. This function is considered dangerous, -as improper use can cause data loss, and is therefore not exposed to the top-level API. +pending. The event loop releases their returned `IO.Promise` without resolving it; once nothing else +references it, its `result?` resolves to `none`. This function is considered dangerous, as improper +use can cause data loss, and is therefore not exposed to the top-level API. Note that this function is idempotent and as such can be called multiple times on the same socket without causing errors, in particular also without a receive running in the first place. diff --git a/src/Std/Internal/UV/Timer.lean b/src/Std/Internal/UV/Timer.lean index 32045ed69230..e85b80b54d84 100644 --- a/src/Std/Internal/UV/Timer.lean +++ b/src/Std/Internal/UV/Timer.lean @@ -51,7 +51,9 @@ This function has different behavior depending on the state and configuration of - if `repeating` is `false` and: - it is initial, run it and return a new `IO.Promise` that is set to resolve once `timeout` milliseconds have elapsed. After this `IO.Promise` is resolved the `Timer` is finished. - - it is running or finished, return the same `IO.Promise` that the first call to `next` returned. + - it is running, or finished after firing, return the same `IO.Promise` + that the first call to `next` returned. + - it was stopped with `stop` before firing, return a new `IO.Promise` that is never resolved. - if `repeating` is `true` and: - it is initial, run it and return a new `IO.Promise` that resolves right away (as it is the 0th multiple of `timeout`). @@ -59,8 +61,8 @@ This function has different behavior depending on the state and configuration of - If it is, return a new `IO.Promise` that resolves upon finishing the next cycle - If it is not, return the last `IO.Promise` This ensures that the returned `IO.Promise` resolves at the next repetition of the timer. - - if it is finished, return the last `IO.Promise` created by `next`. Notably this could be one - that never resolves if the timer was stopped before fulfilling the last one. + - if it is finished, return a new `IO.Promise` that is never resolved, as `stop` dropped the last + one. A promise from `next` may also be resolved by the code holding it; the timer then treats it as fulfilled when it fires. diff --git a/src/Std/Internal/UV/UDP.lean b/src/Std/Internal/UV/UDP.lean index a21cd46606d8..19cd8e0e0654 100644 --- a/src/Std/Internal/UV/UDP.lean +++ b/src/Std/Internal/UV/UDP.lean @@ -58,9 +58,12 @@ opaque send (socket : @& Socket) (data : Array ByteArray) (addr : @& Option Sock /-- Receives data from an UDP socket. `size` is for the maximum bytes to receive. The promise -resolves when some data is available or an error occurs. If a datagram larger than `size` arrives, -it is discarded in its entirety and the promise resolves to an `EMSGSIZE` error. +resolves when some data is available or an error occurs. Furthermore calling this function in parallel with `waitReadable` is not supported. + +A datagram larger than `size` is discarded in its entirety, and the promise resolves to an +`EMSGSIZE` error (an `IO.Error.resourceExhausted`) instead of a truncated prefix. The socket stays +usable, so a receive loop should handle the error per datagram. -/ @[extern "lean_uv_udp_recv"] opaque recv (socket : @& Socket) (size : UInt64) : IO (IO.Promise (Except IO.Error (ByteArray × Option SocketAddress))) @@ -74,8 +77,9 @@ opaque waitReadable (socket : @& Socket) : IO (IO.Promise (Except IO.Error Unit) /-- Cancels a receive operation in the form of `recv` or `waitReadable` if there is currently one -pending. This resolves their returned `IO.Promise` to `none`. This function is considered dangerous, -as improper use can cause data loss, and is therefore not exposed to the top-level API. +pending. The event loop releases their returned `IO.Promise` without resolving it; once nothing else +references it, its `result?` resolves to `none`. This function is considered dangerous, as improper +use can cause data loss, and is therefore not exposed to the top-level API. Note that this function is idempotent and as such can be called multiple times on the same socket without causing errors, in particular also without a receive running in the first place. -/ diff --git a/src/runtime/uv/signal.h b/src/runtime/uv/signal.h index f899f6aadd1e..d54bc0630e58 100644 --- a/src/runtime/uv/signal.h +++ b/src/runtime/uv/signal.h @@ -34,11 +34,13 @@ typedef struct { lean_object * m_promise; // The associated promise for asynchronous results. int m_signum; // Signal number to watch for. bool m_repeating; // Flag indicating if the signal handler is repeating. - uv_signal_state m_state; // The state of the signal. Beyond the API description on the Lean - // side this state has the invariant: - // `m_state != SIGNAL_STATE_INITIAL` -> `m_promise != NULL` + uv_signal_state m_state; // The state of the signal. } lean_uv_signal_object; +// `m_promise` may be NULL in any state: `stop` leaves a FINISHED signal without one, and `cancel` +// on a repeating signal leaves it RUNNING without one. A repeating signal also keeps its last +// promise after resolving it, until `next` replaces it. + // ======================================= // Signal object manipulation functions. static inline lean_object* lean_uv_signal_new(lean_uv_signal_object * s) { return lean_alloc_external(g_uv_signal_external_class, s); } diff --git a/src/runtime/uv/timer.cpp b/src/runtime/uv/timer.cpp index f741236cec0f..169a6e1efde6 100644 --- a/src/runtime/uv/timer.cpp +++ b/src/runtime/uv/timer.cpp @@ -207,7 +207,8 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_timer_next(b_obj_arg obj) { event_loop_unlock(&global_ev); return lean_io_result_mk_ok(promise); } else { - // Creates a resolved promise + // `stop` dropped this timer's promise, so the fresh one is never + // resolved, as documented on `next`. lean_object* finished_promise = create_promise(); event_loop_unlock(&global_ev); return lean_io_result_mk_ok(finished_promise); @@ -224,7 +225,8 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_timer_next(b_obj_arg obj) { return lean_io_result_mk_ok(promise); } else { event_loop_unlock(&global_ev); - // Creates a resolved promise + // `stop` dropped this timer's promise, so the fresh one is never resolved, as + // documented on `next`. lean_object* finished_promise = create_promise(); return lean_io_result_mk_ok(finished_promise); } diff --git a/src/runtime/uv/timer.h b/src/runtime/uv/timer.h index b2fc3c64a02c..a5d283260913 100644 --- a/src/runtime/uv/timer.h +++ b/src/runtime/uv/timer.h @@ -33,11 +33,13 @@ typedef struct { lean_object * m_promise; // The associated promise for asynchronous results. uint64_t m_timeout; // Timeout duration in milliseconds. bool m_repeating; // Flag indicating if the timer is repeating. - uv_timer_state m_state; // The state of the timer. Beyond the API description on the Lean - // side this state has the invariant: - // `m_state != TIMER_STATE_INITIAL` -> `m_promise != NULL` + uv_timer_state m_state; // The state of the timer. } lean_uv_timer_object; +// `m_promise` may be NULL in any state: `stop` leaves a FINISHED timer without one, and `cancel` on +// a repeating timer leaves it RUNNING without one. A repeating timer also keeps its last promise +// after resolving it, until `next` replaces it. + // ======================================= // Timer object manipulation functions. static inline lean_object* lean_uv_timer_new(lean_uv_timer_object * s) { return lean_alloc_external(g_uv_timer_external_class, s); } diff --git a/src/runtime/uv/udp.cpp b/src/runtime/uv/udp.cpp index f82a9d34f4a5..9103a08df564 100644 --- a/src/runtime/uv/udp.cpp +++ b/src/runtime/uv/udp.cpp @@ -283,6 +283,8 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_udp_recv(b_obj_arg socket, uint64_t buf->base = (char*)lean_sarray_cptr(udp_socket->m_byte_array); buf->len = lean_sarray_capacity(udp_socket->m_byte_array); }, [](uv_udp_t *handle, ssize_t nread, const uv_buf_t *buf, const struct sockaddr *addr, unsigned flags) { + // libuv signals "nothing to read yet" as an empty read with no peer. No datagram arrived, + // so the receive stays armed instead of completing with an empty one. if (nread == 0 && addr == NULL) return; uv_udp_recv_stop(handle); From d0cea84c7725c02a0adefece72917309d85c73b1 Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Wed, 16 Sep 2026 18:03:56 -0300 Subject: [PATCH 05/10] fix: keep signal and sleep selectors listening and report Lean signal numbers A signal selector that loses a select no longer misses a signal delivered before the next select, and selecting on a `Sleep` no longer cancels it, so it keeps its deadline across selects. Signal waiters now resolve with the signal number in Lean's `Signal` numbering instead of the operating system's, and a `Sleep` or `Interval` longer than a `UInt64` of milliseconds is clamped instead of wrapping to a short duration. Co-Authored-By: Claude Opus 5 --- src/Std/Async/Signal.lean | 37 ++++++++++++++----- src/Std/Async/Timer.lean | 22 ++++++++--- src/runtime/uv/signal.cpp | 5 ++- src/runtime/uv/signal.h | 1 + .../elab/signal_selector_keeps_listening.lean | 32 ++++++++++++++++ tests/elab/sleep_selector_keeps_deadline.lean | 23 ++++++++++++ tests/elab/uv_signal_number_encoding.lean | 19 ++++++++++ tests/elab/uv_sleep_long_duration.lean | 18 +++++++++ 8 files changed, 141 insertions(+), 16 deletions(-) create mode 100644 tests/elab/signal_selector_keeps_listening.lean create mode 100644 tests/elab/sleep_selector_keeps_deadline.lean create mode 100644 tests/elab/uv_signal_number_encoding.lean create mode 100644 tests/elab/uv_sleep_long_duration.lean diff --git a/src/Std/Async/Signal.lean b/src/Std/Async/Signal.lean index bad4d6999b21..233abf6a9430 100644 --- a/src/Std/Async/Signal.lean +++ b/src/Std/Async/Signal.lean @@ -192,6 +192,12 @@ private def toInt32 : Signal → Int32 structure Waiter where private ofNative :: native : Internal.UV.Signal + /-- Whether a continuation is attached to the current promise of `native` for `selector`. -/ + private armed : IO.Ref Bool + /-- The `Waiter` of the select that is currently registered through `selector`, if any. -/ + private registered : IO.Ref (Option (Async.Waiter Unit)) + /-- Whether a signal arrived while no select was registered, so the next select wins at once. -/ + private unclaimed : IO.Ref Bool namespace Waiter @@ -202,7 +208,7 @@ This function only initializes but does not yet start listening for the signal. @[inline] def mk (signum : Signal) (repeating : Bool) : IO Signal.Waiter := do let native ← Internal.UV.Signal.mk signum.toInt32 repeating - return .ofNative native + return .ofNative native (← IO.mkRef false) (← IO.mkRef none) (← IO.mkRef false) /-- If: @@ -233,24 +239,37 @@ def stop (s : Signal.Waiter) : IO Unit := /-- Create a `Selector` that resolves once `s` has received the signal. Note that calling this function does not start the signal waiter. + +A select that `s` loses leaves it listening, and a signal that arrives before the next select is +reported by that select, so `stop` has to be called once `s` is no longer needed. -/ def selector (s : Signal.Waiter) : Selector Unit := + let claim (waiter : Async.Waiter Unit) : BaseIO Unit := + waiter.race (lose := s.unclaimed.set true) (win := fun promise => promise.resolve (.ok ())) { tryFn := do + if ← s.unclaimed.modifyGet (·, false) then + return some () let signalWaiter ← s.native.next if ← signalWaiter.isResolved then return some () else - s.native.cancel return none registerFn waiter := do - let signalWaiter ← s.wait - discard <| AsyncTask.mapIO (x := signalWaiter) fun _ => do - let lose := return () - let win promise := promise.resolve (.ok ()) - waiter.race lose win - - unregisterFn := s.native.cancel + s.registered.set (some waiter) + -- One continuation per promise, so that a signal is claimed at most once. + unless ← s.armed.modifyGet (·, true) do + let signalWaiter ← s.wait + discard <| AsyncTask.mapIO (x := signalWaiter) fun _ => do + s.armed.set false + match ← s.registered.modifyGet (·, none) with + | some registered => claim registered + | none => s.unclaimed.set true + -- A signal claimed by nobody between `tryFn` and this registration. + if ← s.unclaimed.modifyGet (·, false) then + claim waiter + + unregisterFn := s.registered.set none } diff --git a/src/Std/Async/Timer.lean b/src/Std/Async/Timer.lean index 2a99f04970ac..f3c35870e194 100644 --- a/src/Std/Async/Timer.lean +++ b/src/Std/Async/Timer.lean @@ -24,6 +24,13 @@ structure Sleep where private ofNative :: native : Internal.UV.Timer +/-- +The timeout passed to libuv. Longer durations are clamped, since they cannot elapse anyway. +-/ +private def timeoutOf (duration : Std.Time.Millisecond.Offset) : UInt64 := + let ms := duration.toInt.toNat + if ms < UInt64.size then ms.toUInt64 else (UInt64.size - 1).toUInt64 + namespace Sleep /-- @@ -32,7 +39,7 @@ This function only initializes but does not yet start the timer. -/ @[inline] def mk (duration : Std.Time.Millisecond.Offset) : Async Sleep := do - let native ← Internal.UV.Timer.mk duration.toInt.toNat.toUInt64 false + let native ← Internal.UV.Timer.mk (timeoutOf duration) false return ofNative native /-- @@ -69,7 +76,9 @@ def stop (s : Sleep) : IO Unit := s.native.stop /-- -Create a `Selector` that resolves once `s` has finished. `s` only starts when it runs inside of a Selectable. +Create a `Selector` that resolves once `s` has finished. `s` only starts when it runs inside of a +Selectable, and a select that it loses leaves it running, so a `Sleep` selected on repeatedly works +as a deadline that elapses `duration` after the first of those selects. -/ def selector (s : Sleep) : Selector Unit := { @@ -78,7 +87,6 @@ def selector (s : Sleep) : Selector Unit := if ← sleepWaiter.isResolved then return some () else - s.native.cancel return none registerFn waiter := do @@ -91,7 +99,7 @@ def selector (s : Sleep) : Selector Unit := let win promise := promise.resolve (.ok ()) waiter.race lose win - unregisterFn := s.native.cancel + unregisterFn := pure () } end Sleep @@ -108,7 +116,9 @@ Return a `Selector` that completes after `duration`. -/ def Selector.sleep (duration : Std.Time.Millisecond.Offset) : Async (Selector Unit) := do let sleeper ← Sleep.mk duration - return sleeper.selector + -- A fresh timer cannot have elapsed yet, and nothing else refers to `sleeper`, so it is started + -- only on registration and stopped once the select is over. + return { sleeper.selector with tryFn := pure none, unregisterFn := sleeper.native.cancel } /-- `Interval` can be used to repeatedly wait for some duration like a clock. @@ -126,7 +136,7 @@ This function only initializes but does not yet start the timer. -/ @[inline] def mk (duration : Std.Time.Millisecond.Offset) (_ : 0 < duration := by decide) : IO Interval := do - let native ← Internal.UV.Timer.mk duration.toInt.toNat.toUInt64 true + let native ← Internal.UV.Timer.mk (timeoutOf duration) true return ofNative native /-- diff --git a/src/runtime/uv/signal.cpp b/src/runtime/uv/signal.cpp index 3c0df6f75548..ba20b6103995 100644 --- a/src/runtime/uv/signal.cpp +++ b/src/runtime/uv/signal.cpp @@ -45,9 +45,11 @@ static bool signal_promise_is_finished(lean_uv_signal_object * signal) { return signal->m_promise == NULL || promise_is_resolved(signal->m_promise); } -void handle_signal_event(uv_signal_t* handle, int signum) { +void handle_signal_event(uv_signal_t* handle, int) { lean_object * obj = (lean_object*)handle->data; lean_uv_signal_object * signal = lean_to_uv_signal(obj); + // Read before `lean_dec(obj)` below may free the signal. + int const signum = signal->m_lean_signum; lean_assert(signal->m_state == SIGNAL_STATE_RUNNING); @@ -122,6 +124,7 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_mk(uint32_t signum_obj, uint8 return lean_io_result_mk_error(decode_io_error(ENOMEM, nullptr)); } signal->m_signum = signum; + signal->m_lean_signum = (int)(int32_t)signum_obj; signal->m_repeating = repeating; signal->m_state = SIGNAL_STATE_INITIAL; signal->m_promise = NULL; diff --git a/src/runtime/uv/signal.h b/src/runtime/uv/signal.h index d54bc0630e58..b24a3b17ec95 100644 --- a/src/runtime/uv/signal.h +++ b/src/runtime/uv/signal.h @@ -33,6 +33,7 @@ typedef struct { uv_signal_t * m_uv_signal; // LibUV signal handle. lean_object * m_promise; // The associated promise for asynchronous results. int m_signum; // Signal number to watch for. + int m_lean_signum; // `m_signum` in the encoding of `Signal.toInt32`, reported to waiters. bool m_repeating; // Flag indicating if the signal handler is repeating. uv_signal_state m_state; // The state of the signal. } lean_uv_signal_object; diff --git a/tests/elab/signal_selector_keeps_listening.lean b/tests/elab/signal_selector_keeps_listening.lean new file mode 100644 index 000000000000..32fa785cc871 --- /dev/null +++ b/tests/elab/signal_selector_keeps_listening.lean @@ -0,0 +1,32 @@ +import Std.Async + +/-! +A `Signal.Waiter` whose selector lost a `Selectable.one` must keep listening, so that a signal that +arrives before the next select is still delivered to it. Covers one-shot and repeating waiters. +SIGWINCH is used because its default action is to ignore it. +-/ + +open Std Async + +def receivesBetweenSelects (repeating : Bool) : Async Bool := do + let waiter ← Signal.Waiter.mk .sigwinch repeating + let first ← Selectable.one #[ + .case waiter.selector (fun _ => pure true), + .case (← Selector.sleep 50) (fun _ => pure false)] + if first then + throw <| IO.userError "received a signal before one was sent" + let pid ← IO.Process.getPID + discard <| IO.Process.output { cmd := "kill", args := #["-WINCH", toString pid] } + Async.sleep 200 + let second ← Selectable.one #[ + .case waiter.selector (fun _ => pure true), + .case (← Selector.sleep 1000) (fun _ => pure false)] + waiter.stop + return second + +#eval show IO Unit from do + -- Signals cannot be sent with `kill` on Windows. + unless System.Platform.isWindows do + for repeating in [false, true] do + unless ← (receivesBetweenSelects repeating).block do + throw <| IO.userError s!"signal sent between selects was lost (repeating := {repeating})" diff --git a/tests/elab/sleep_selector_keeps_deadline.lean b/tests/elab/sleep_selector_keeps_deadline.lean new file mode 100644 index 000000000000..6b8cce4fa225 --- /dev/null +++ b/tests/elab/sleep_selector_keeps_deadline.lean @@ -0,0 +1,23 @@ +import Std.Async + +/-! +Selecting on a `Sleep` must not restart it: a `Sleep` used as a deadline across several +`Selectable.one` calls fires once its duration has elapsed since it started. +-/ + +open Std Async + +#eval show IO Unit from do + let won ← (do + let deadline ← Sleep.mk 300 + let mut won := false + for _ in [0:20] do + let tick ← Selector.sleep 50 + if ← Selectable.one #[ + .case deadline.selector (fun _ => pure true), + .case tick (fun _ => pure false)] then + won := true + break + return won : Async Bool).block + unless won do + throw <| IO.userError "the deadline never fired while it was being selected on" diff --git a/tests/elab/uv_signal_number_encoding.lean b/tests/elab/uv_signal_number_encoding.lean new file mode 100644 index 000000000000..f3de6c5e6c87 --- /dev/null +++ b/tests/elab/uv_signal_number_encoding.lean @@ -0,0 +1,19 @@ +import Std.Async + +/-! +A `Signal.Waiter` completes with the signal number in the encoding of `Signal.toInt32`, which is +Linux's, also on platforms whose native number for the signal differs. +-/ + +open Std Async + +#eval show IO Unit from do + -- Signals cannot be sent with `kill` on Windows. + if System.Platform.isWindows then return + let waiter ← Signal.Waiter.mk .sigusr1 false + let received ← waiter.wait + let pid ← IO.Process.getPID + discard <| IO.Process.output { cmd := "kill", args := #["-USR1", toString pid] } + let signum ← received.block + unless signum == 10 do + throw <| IO.userError s!"received signal number {signum}, expected 10" diff --git a/tests/elab/uv_sleep_long_duration.lean b/tests/elab/uv_sleep_long_duration.lean new file mode 100644 index 000000000000..d7da826f124f --- /dev/null +++ b/tests/elab/uv_sleep_long_duration.lean @@ -0,0 +1,18 @@ +import Std.Async + +/-! +A `Sleep` longer than a `UInt64` of milliseconds is clamped rather than wrapped, so it does not fire +early. +-/ + +open Std Async + +#eval show IO Unit from do + let fired ← (do + -- 2^64 + 50 ms, which wraps to 50 ms. + let long ← Sleep.mk ⟨2 ^ 64 + 50⟩ + Selectable.one #[ + .case long.selector (fun _ => pure true), + .case (← Selector.sleep 500) (fun _ => pure false)] : Async Bool).block + if fired then + throw <| IO.userError "a sleep of more than 2^64 ms fired early" From 4db11db5f85d64250878504bd1f38806cdd1897d Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Wed, 16 Sep 2026 21:33:04 -0300 Subject: [PATCH 06/10] refactor: keep signals between selects in the signal handler `Signal.cancel` now only drops the pending promise and keeps the handler listening, and a signal that arrives with no promise pending is reported by the next `next` instead of being consumed. `Signal.Waiter.selector` goes back to cancelling on unregister, without the extra state it kept to avoid losing those signals. Co-Authored-By: Claude Opus 5 --- src/Std/Async/Signal.lean | 34 +++++------------- src/Std/Internal/UV/Signal.lean | 12 ++++--- src/runtime/uv/signal.cpp | 64 ++++++++++++++++++--------------- src/runtime/uv/signal.h | 3 +- 4 files changed, 54 insertions(+), 59 deletions(-) diff --git a/src/Std/Async/Signal.lean b/src/Std/Async/Signal.lean index 233abf6a9430..8238c39c173f 100644 --- a/src/Std/Async/Signal.lean +++ b/src/Std/Async/Signal.lean @@ -192,12 +192,6 @@ private def toInt32 : Signal → Int32 structure Waiter where private ofNative :: native : Internal.UV.Signal - /-- Whether a continuation is attached to the current promise of `native` for `selector`. -/ - private armed : IO.Ref Bool - /-- The `Waiter` of the select that is currently registered through `selector`, if any. -/ - private registered : IO.Ref (Option (Async.Waiter Unit)) - /-- Whether a signal arrived while no select was registered, so the next select wins at once. -/ - private unclaimed : IO.Ref Bool namespace Waiter @@ -208,7 +202,7 @@ This function only initializes but does not yet start listening for the signal. @[inline] def mk (signum : Signal) (repeating : Bool) : IO Signal.Waiter := do let native ← Internal.UV.Signal.mk signum.toInt32 repeating - return .ofNative native (← IO.mkRef false) (← IO.mkRef none) (← IO.mkRef false) + return .ofNative native /-- If: @@ -244,32 +238,22 @@ A select that `s` loses leaves it listening, and a signal that arrives before th reported by that select, so `stop` has to be called once `s` is no longer needed. -/ def selector (s : Signal.Waiter) : Selector Unit := - let claim (waiter : Async.Waiter Unit) : BaseIO Unit := - waiter.race (lose := s.unclaimed.set true) (win := fun promise => promise.resolve (.ok ())) { tryFn := do - if ← s.unclaimed.modifyGet (·, false) then - return some () let signalWaiter ← s.native.next if ← signalWaiter.isResolved then return some () else + s.native.cancel return none registerFn waiter := do - s.registered.set (some waiter) - -- One continuation per promise, so that a signal is claimed at most once. - unless ← s.armed.modifyGet (·, true) do - let signalWaiter ← s.wait - discard <| AsyncTask.mapIO (x := signalWaiter) fun _ => do - s.armed.set false - match ← s.registered.modifyGet (·, none) with - | some registered => claim registered - | none => s.unclaimed.set true - -- A signal claimed by nobody between `tryFn` and this registration. - if ← s.unclaimed.modifyGet (·, false) then - claim waiter - - unregisterFn := s.registered.set none + let signalWaiter ← s.wait + discard <| AsyncTask.mapIO (x := signalWaiter) fun _ => do + let lose := return () + let win promise := promise.resolve (.ok ()) + waiter.race lose win + + unregisterFn := s.native.cancel } diff --git a/src/Std/Internal/UV/Signal.lean b/src/Std/Internal/UV/Signal.lean index 28f096b30bd4..e2fd385f460d 100644 --- a/src/Std/Internal/UV/Signal.lean +++ b/src/Std/Internal/UV/Signal.lean @@ -41,7 +41,7 @@ This creates a `Signal` in the initial state and doesn't start listening yet. signal `signum` is received, then automatically stops listening. - If `repeating` is `true` this constructs a signal handler that resolves each time the specified signal `signum` is received and continues listening. While it listens, a signal that arrives with - no promise from `next` pending is consumed without being reported. A repeating signal handler + no promise from `next` pending resolves the promise of the following `next`. A repeating signal handler will only be freed after `Signal.stop` is called. -/ @[extern "lean_uv_signal_mk"] @@ -53,14 +53,16 @@ This function has different behavior depending on the state and configuration of - it is initial, start listening and return a new `IO.Promise` that is set to resolve once the signal `signum` is received. After this `IO.Promise` is resolved the `Signal` is finished. - it is running, or finished after receiving the signal, return the same `IO.Promise` that the - first call to `next` returned. + last call to `next` returned, or a new one if `cancel` dropped it. A signal received after + `cancel` resolves that promise. - it was stopped with `stop` before receiving the signal, return a new `IO.Promise` that is never resolved. - if `repeating` is `true` and: - it is initial, start listening and return a new `IO.Promise` that resolves when the next signal `signum` is received. - it is running, check whether the last returned `IO.Promise` is already resolved: - - If it is, return a new `IO.Promise` that resolves upon receiving the next signal + - If it is, return a new `IO.Promise` that resolves upon receiving the next signal, or that is + already resolved if a signal arrived while no promise was pending - If it is not, return the last `IO.Promise` This ensures that the returned `IO.Promise` resolves at the next occurrence of the signal. - if it is finished, return a new `IO.Promise` that is never resolved, as `stop` dropped the last @@ -86,8 +88,8 @@ opaque stop (signal : @& Signal) : IO Unit /-- This function has different behavior depending on the state of the `Signal`: - If it is initial or finished this is a no-op. -- If it's running then it drops the accept promise and if it's not repeatable it sets - the signal handler to the initial state. +- If it's running then it drops the last `IO.Promise` generated by `next`. The handler keeps + listening, so a signal that arrives afterwards is reported by the next call to `next`. -/ @[extern "lean_uv_signal_cancel"] opaque cancel (signal : @& Signal) : IO Unit diff --git a/src/runtime/uv/signal.cpp b/src/runtime/uv/signal.cpp index ba20b6103995..71d4e01604d1 100644 --- a/src/runtime/uv/signal.cpp +++ b/src/runtime/uv/signal.cpp @@ -41,6 +41,13 @@ void initialize_libuv_signal() { }); } +static lean_object * create_signal_promise() { + lean_object * promise = lean_io_promise_new(); + // The loop thread resolves and releases it, so its refcount has to be atomic. + mark_mt(promise); + return promise; +} + static bool signal_promise_is_finished(lean_uv_signal_object * signal) { return signal->m_promise == NULL || promise_is_resolved(signal->m_promise); } @@ -54,7 +61,10 @@ void handle_signal_event(uv_signal_t* handle, int) { lean_assert(signal->m_state == SIGNAL_STATE_RUNNING); if (signal->m_repeating) { - if (!signal_promise_is_finished(signal)) { + if (signal_promise_is_finished(signal)) { + // Kept for the next `next`, so that a signal between two waits is not lost. + signal->m_received = true; + } else { // Rule 1: a continuation may `cancel` or `stop` the signal, releasing the field's // reference. lean_object * promise = signal->m_promise; @@ -67,22 +77,23 @@ void handle_signal_event(uv_signal_t* handle, int) { uv_signal_stop(signal->m_uv_signal); signal->m_state = SIGNAL_STATE_FINISHED; - lean_object * promise = signal->m_promise; - if (promise != NULL) { - lean_inc(promise); + if (signal->m_promise == NULL) { + // Kept for the next `next`, so that a signal after a `cancel` is not lost. + signal->m_promise = create_signal_promise(); } + lean_object * promise = signal->m_promise; + lean_inc(promise); + lean_dec(obj); // Rule 1: nothing below may touch the signal. Code holding the promise may have resolved // it already. - if (promise != NULL) { - if (!promise_is_resolved(promise)) { - lean_object* res = lean_io_promise_resolve(lean_box(signum), promise); - lean_dec(res); - } - lean_dec(promise); + if (!promise_is_resolved(promise)) { + lean_object* res = lean_io_promise_resolve(lean_box(signum), promise); + lean_dec(res); } + lean_dec(promise); } } @@ -126,6 +137,7 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_mk(uint32_t signum_obj, uint8 signal->m_signum = signum; signal->m_lean_signum = (int)(int32_t)signum_obj; signal->m_repeating = repeating; + signal->m_received = false; signal->m_state = SIGNAL_STATE_INITIAL; signal->m_promise = NULL; @@ -158,12 +170,7 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_mk(uint32_t signum_obj, uint8 extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_next(b_obj_arg obj) { lean_uv_signal_object * signal = lean_to_uv_signal(obj); - auto create_promise = []() { - lean_object * promise = lean_io_promise_new(); - // The loop thread resolves and releases it, so its refcount has to be atomic. - mark_mt(promise); - return promise; - }; + auto create_promise = create_signal_promise; auto setup_signal = [create_promise, obj, signal]() { lean_assert(signal->m_promise == NULL); @@ -224,6 +231,11 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_next(b_obj_arg obj) { } signal->m_promise = create_promise(); + + if (signal->m_received) { + signal->m_received = false; + lean_dec(lean_io_promise_resolve(lean_box(signal->m_lean_signum), signal->m_promise)); + } } lean_object * promise = signal->m_promise; @@ -248,6 +260,13 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_next(b_obj_arg obj) { } else { if (signal->m_state == SIGNAL_STATE_INITIAL) { return setup_signal(); + } else if (signal->m_state == SIGNAL_STATE_RUNNING && signal->m_promise == NULL) { + // Still listening after a `cancel`. + lean_object * promise = create_promise(); + signal->m_promise = promise; + lean_inc(promise); + event_loop_unlock(&global_ev); + return lean_io_result_mk_ok(promise); } else if (signal->m_promise != NULL) { lean_object * promise = signal->m_promise; lean_inc(promise); @@ -298,18 +317,11 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_cancel(b_obj_arg obj) { event_loop_lock(&global_ev); lean_object * promise = NULL; - bool release_signal = false; + // The signal keeps listening, so one that arrives before the next `next` is not lost. if (signal->m_state == SIGNAL_STATE_RUNNING && signal->m_promise != NULL) { promise = signal->m_promise; signal->m_promise = NULL; - - // A repeating signal keeps listening, so the loop keeps its reference until `stop`. - if (!signal->m_repeating) { - uv_signal_stop(signal->m_uv_signal); - signal->m_state = SIGNAL_STATE_INITIAL; - release_signal = true; - } } event_loop_unlock(&global_ev); @@ -319,10 +331,6 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_signal_cancel(b_obj_arg obj) { lean_dec(promise); } - if (release_signal) { - lean_dec(obj); - } - return lean_io_result_mk_ok(lean_box(0)); } diff --git a/src/runtime/uv/signal.h b/src/runtime/uv/signal.h index b24a3b17ec95..24ce8e19e788 100644 --- a/src/runtime/uv/signal.h +++ b/src/runtime/uv/signal.h @@ -35,11 +35,12 @@ typedef struct { int m_signum; // Signal number to watch for. int m_lean_signum; // `m_signum` in the encoding of `Signal.toInt32`, reported to waiters. bool m_repeating; // Flag indicating if the signal handler is repeating. + bool m_received; // Whether a repeating signal arrived while no promise was pending. uv_signal_state m_state; // The state of the signal. } lean_uv_signal_object; // `m_promise` may be NULL in any state: `stop` leaves a FINISHED signal without one, and `cancel` -// on a repeating signal leaves it RUNNING without one. A repeating signal also keeps its last +// leaves it RUNNING without one. A repeating signal also keeps its last // promise after resolving it, until `next` replaces it. // ======================================= From d3342c308393bd55b20fc66ae9aa6eab2d047bf4 Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Wed, 16 Sep 2026 22:09:29 -0300 Subject: [PATCH 07/10] revert: let `Sleep.selector` cancel the timer when it loses A `Sleep` whose selector loses a select is cancelled again, so the next select starts it from its full duration; the documentation of `Sleep.selector` now says so. This drops the separate `Selector.sleep` override that stopped one-off timers. Co-Authored-By: Claude Opus 5 --- src/Std/Async/Timer.lean | 10 ++++---- tests/elab/sleep_selector_keeps_deadline.lean | 23 ------------------- 2 files changed, 4 insertions(+), 29 deletions(-) delete mode 100644 tests/elab/sleep_selector_keeps_deadline.lean diff --git a/src/Std/Async/Timer.lean b/src/Std/Async/Timer.lean index f3c35870e194..37f039dca8ca 100644 --- a/src/Std/Async/Timer.lean +++ b/src/Std/Async/Timer.lean @@ -77,8 +77,7 @@ def stop (s : Sleep) : IO Unit := /-- Create a `Selector` that resolves once `s` has finished. `s` only starts when it runs inside of a -Selectable, and a select that it loses leaves it running, so a `Sleep` selected on repeatedly works -as a deadline that elapses `duration` after the first of those selects. +Selectable, and a select that it loses cancels it, so the next select starts it again from `duration`. -/ def selector (s : Sleep) : Selector Unit := { @@ -87,6 +86,7 @@ def selector (s : Sleep) : Selector Unit := if ← sleepWaiter.isResolved then return some () else + s.native.cancel return none registerFn waiter := do @@ -99,7 +99,7 @@ def selector (s : Sleep) : Selector Unit := let win promise := promise.resolve (.ok ()) waiter.race lose win - unregisterFn := pure () + unregisterFn := s.native.cancel } end Sleep @@ -116,9 +116,7 @@ Return a `Selector` that completes after `duration`. -/ def Selector.sleep (duration : Std.Time.Millisecond.Offset) : Async (Selector Unit) := do let sleeper ← Sleep.mk duration - -- A fresh timer cannot have elapsed yet, and nothing else refers to `sleeper`, so it is started - -- only on registration and stopped once the select is over. - return { sleeper.selector with tryFn := pure none, unregisterFn := sleeper.native.cancel } + return sleeper.selector /-- `Interval` can be used to repeatedly wait for some duration like a clock. diff --git a/tests/elab/sleep_selector_keeps_deadline.lean b/tests/elab/sleep_selector_keeps_deadline.lean deleted file mode 100644 index 6b8cce4fa225..000000000000 --- a/tests/elab/sleep_selector_keeps_deadline.lean +++ /dev/null @@ -1,23 +0,0 @@ -import Std.Async - -/-! -Selecting on a `Sleep` must not restart it: a `Sleep` used as a deadline across several -`Selectable.one` calls fires once its duration has elapsed since it started. --/ - -open Std Async - -#eval show IO Unit from do - let won ← (do - let deadline ← Sleep.mk 300 - let mut won := false - for _ in [0:20] do - let tick ← Selector.sleep 50 - if ← Selectable.one #[ - .case deadline.selector (fun _ => pure true), - .case tick (fun _ => pure false)] then - won := true - break - return won : Async Bool).block - unless won do - throw <| IO.userError "the deadline never fired while it was being selected on" From 881c91b36b35f284ed3d030d4cc889cc6b4af243 Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Wed, 16 Sep 2026 22:15:10 -0300 Subject: [PATCH 08/10] doc: describe UDP and timer errors without low-level names Co-Authored-By: Claude Opus 5 --- src/Std/Async/Timer.lean | 2 +- src/Std/Async/UDP.lean | 6 +++--- src/Std/Internal/UV/UDP.lean | 2 +- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/src/Std/Async/Timer.lean b/src/Std/Async/Timer.lean index 37f039dca8ca..43920ef24374 100644 --- a/src/Std/Async/Timer.lean +++ b/src/Std/Async/Timer.lean @@ -25,7 +25,7 @@ structure Sleep where native : Internal.UV.Timer /-- -The timeout passed to libuv. Longer durations are clamped, since they cannot elapse anyway. +The timeout of the underlying timer. Longer durations are clamped, since they cannot elapse anyway. -/ private def timeoutOf (duration : Std.Time.Millisecond.Offset) : UInt64 := let ms := duration.toInt.toNat diff --git a/src/Std/Async/UDP.lean b/src/Std/Async/UDP.lean index ade9c47ea497..e95fce48c8fb 100644 --- a/src/Std/Async/UDP.lean +++ b/src/Std/Async/UDP.lean @@ -81,8 +81,8 @@ has not been previously bound with `bind`, it is automatically bound to `0.0.0.0 (all interfaces) with a random port. Furthermore calling this function in parallel with `recvSelector` is not supported. -A datagram larger than `size` is discarded in its entirety, and this throws `EMSGSIZE` (an -`IO.Error.resourceExhausted`) instead of returning a truncated prefix. The socket stays usable, so a +A datagram larger than `size` is discarded in its entirety, and this throws an +`IO.Error.resourceExhausted` instead of returning a truncated prefix. The socket stays usable, so a receive loop should catch the error per datagram. -/ @[inline] @@ -96,7 +96,7 @@ automatically bound to `0.0.0.0` (all interfaces) with a random port. Calling this function does starts the data wait, only when it's used with `Selectable.one` or `combine`. It must not be called in parallel with `recv`. -Fails with `EMSGSIZE` if the datagram is larger than `size`, like `recv`. +Fails with an `IO.Error.resourceExhausted` if the datagram is larger than `size`, like `recv`. -/ def recvSelector (s : Socket) (size : UInt64) : Selector (ByteArray × Option SocketAddress) := { diff --git a/src/Std/Internal/UV/UDP.lean b/src/Std/Internal/UV/UDP.lean index 19cd8e0e0654..7a5ac4d15110 100644 --- a/src/Std/Internal/UV/UDP.lean +++ b/src/Std/Internal/UV/UDP.lean @@ -62,7 +62,7 @@ resolves when some data is available or an error occurs. Furthermore calling this function in parallel with `waitReadable` is not supported. A datagram larger than `size` is discarded in its entirety, and the promise resolves to an -`EMSGSIZE` error (an `IO.Error.resourceExhausted`) instead of a truncated prefix. The socket stays +`IO.Error.resourceExhausted` instead of a truncated prefix. The socket stays usable, so a receive loop should handle the error per datagram. -/ @[extern "lean_uv_udp_recv"] From 6cae6745dbe7e96aaa4a6f576e26eac6e834d98c Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Wed, 16 Sep 2026 22:16:02 -0300 Subject: [PATCH 09/10] doc: drop the libuv reference from the signal list Co-Authored-By: Claude Opus 5 --- src/Std/Async/Signal.lean | 2 +- src/Std/Internal/UV/TCP.lean | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/src/Std/Async/Signal.lean b/src/Std/Async/Signal.lean index 8238c39c173f..c898f711d62c 100644 --- a/src/Std/Async/Signal.lean +++ b/src/Std/Async/Signal.lean @@ -18,7 +18,7 @@ namespace Async /-- Unix style signals for Unix and Windows. SIGKILL and SIGSTOP are missing because they cannot be caught. -SIGBUS, SIGFPE, SIGILL, and SIGSEGV are missing because they cannot be caught safely by libuv. +SIGBUS, SIGFPE, SIGILL, and SIGSEGV are missing because they cannot be caught safely. SIGPIPE is not present because the runtime ignores the signal. -/ inductive Signal diff --git a/src/Std/Internal/UV/TCP.lean b/src/Std/Internal/UV/TCP.lean index 998c2ea6661c..c40d75ace2d2 100644 --- a/src/Std/Internal/UV/TCP.lean +++ b/src/Std/Internal/UV/TCP.lean @@ -139,7 +139,7 @@ Enables the Nagle algorithm for a TCP socket. opaque noDelay (socket : @& Socket) : IO Unit /-- -Enables TCP keep-alive for a socket. If delay is less than 1 then UV_EINVAL is returned. +Enables TCP keep-alive for a socket. If delay is less than 1 then an `IO.Error.invalidArgument` is returned. -/ @[extern "lean_uv_tcp_keepalive"] opaque keepAlive (socket : @& Socket) (enable : Int8) (delay : UInt32) : IO Unit From a6a7345991131fc8b32442885fb432cf131f682c Mon Sep 17 00:00:00 2001 From: Sofia Rodrigues Date: Wed, 16 Sep 2026 22:33:33 -0300 Subject: [PATCH 10/10] chore: remove comment on empty UDP reads Co-Authored-By: Claude Opus 5 --- src/runtime/uv/udp.cpp | 2 -- 1 file changed, 2 deletions(-) diff --git a/src/runtime/uv/udp.cpp b/src/runtime/uv/udp.cpp index f3f7be3df6bc..8fc850814f3b 100644 --- a/src/runtime/uv/udp.cpp +++ b/src/runtime/uv/udp.cpp @@ -288,8 +288,6 @@ extern "C" LEAN_EXPORT lean_obj_res lean_uv_udp_recv(b_obj_arg socket, uint64_t buf->base = (char*)lean_sarray_cptr(udp_socket->m_byte_array); buf->len = lean_sarray_capacity(udp_socket->m_byte_array); }, [](uv_udp_t *handle, ssize_t nread, const uv_buf_t *buf, const struct sockaddr *addr, unsigned flags) { - // libuv signals "nothing to read yet" as an empty read with no peer. No datagram arrived, - // so the receive stays armed instead of completing with an empty one. if (nread == 0 && addr == NULL) return; uv_udp_recv_stop(handle);