Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
17 changes: 11 additions & 6 deletions src/Std/Async/Signal.lean
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
-/
Expand All @@ -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]
Expand All @@ -231,12 +233,15 @@ 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 :=
{
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
Expand Down
33 changes: 22 additions & 11 deletions src/Std/Async/Timer.lean
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,13 @@ structure Sleep where
private ofNative ::
native : Internal.UV.Timer

/--
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
if ms < UInt64.size then ms.toUInt64 else (UInt64.size - 1).toUInt64

namespace Sleep

/--
Expand All @@ -32,14 +39,16 @@ 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

/--
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 :=
Expand All @@ -57,17 +66,18 @@ 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]
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 cancels it, so the next select starts it again from `duration`.
-/
def selector (s : Sleep) : Selector Unit :=
{
Expand Down Expand Up @@ -124,7 +134,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

/--
Expand All @@ -136,7 +146,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
Expand All @@ -154,9 +165,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]
Expand Down
32 changes: 20 additions & 12 deletions src/Std/Internal/UV/Signal.lean
Original file line number Diff line number Diff line change
Expand Up @@ -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 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"]
opaque mk (signum : Int32) (repeating : Bool) : IO Signal
Expand All @@ -51,37 +52,44 @@ 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
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 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.

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)

/--
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

/--
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
Expand Down
9 changes: 7 additions & 2 deletions src/Std/Internal/UV/TCP.lean
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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.
Expand Down
15 changes: 11 additions & 4 deletions src/Std/Internal/UV/Timer.lean
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -49,16 +51,21 @@ 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`).
- it is running, check whether the last returned `IO.Promise` is already resolved:
- 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.
-/
@[extern "lean_uv_timer_next"]
opaque next (timer : @& Timer) : IO (IO.Promise Unit)
Expand Down
Loading
Loading