Skip to content
Merged
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
4 changes: 2 additions & 2 deletions crates/moonbit/docs/adr/0002-local-future-promise.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,8 +34,8 @@ component `future` endpoint, and it works for arbitrary MoonBit `T`.
and wakes the reader, completion wins a simultaneous cancellation race and the
reader receives the value.
- Explicit `Future::drop()` follows the same race rule while `get()` is pending:
dropping before settlement wakes the reader with `Cancelled`, while an
already-assigned value or error remains owned by the waiting reader.
dropping before settlement wakes the reader with `FutureReadError::Dropped`,
while an already-assigned value or error remains owned by the waiting reader.
- A local failure or close cannot settle an already-exposed component future
without a value. If such an outcome is expected across WIT, it belongs in the
payload type, for example `Future[Result[V, E]]`.
Expand Down
24 changes: 20 additions & 4 deletions crates/moonbit/docs/async-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -318,7 +318,25 @@ interface may provide it explicitly without weakening the ordinary local types.
The runtime is audited against `moonbitlang/async` main at commit `18533c8d`.
The continuation primitive, cancellation races, shielding, wake behavior,
fairness, `Task`, `TaskGroup`, `Semaphore`, `Mutex`, and `CondVar` semantics are
kept aligned.
aligned with that baseline.

The cancellation primitives and their callers also follow upstream's
[`cancel`/`nocancel` migration](https://github.com/moonbitlang/async/commit/19f57c5988e0387f6a79d5a60852ee64fd4d33e6),
checked at upstream commit `a4cbfabbcdf4fa70ef28ad7ccc388082d92371de`:

- cancellation bypasses ordinary `catch` and runs `defer`/`errdefer`;
- `handle_cancellation` observes cancellation without clearing the task's
cancelled state;
- waiting for a cancelled task raises the ordinary `TaskCancelled` error;
- `protect_from_cancel` is `nocancel` and finishes normally; pending
cancellation is delivered at the next unshielded cancellation point;
- task-group defer callbacks must be `nocancel`.

Component cleanup must finish before propagating cancellation. Generated
future/stream reads and subtask cancellation wait for terminal events before
releasing canonical buffers and handles. An already completed endpoint copy
still takes precedence over simultaneous task cancellation, so ownership of
transferred values is preserved.

Intentional differences are limited to the component environment:

Expand All @@ -342,9 +360,7 @@ Implemented and covered by composed runtime tests:
payload cleanup;
- concurrent component-task isolation;
- `wasi:cli@0.3.0` stream output;
- `wasi:http@0.3.0` body, trailers, and post-response background work;
- deletion guards proving endpoint-free sync generation is unchanged without
async support.
- `wasi:http@0.3.0` body, trailers, and post-response background work.

Known limits:

Expand Down
8 changes: 8 additions & 0 deletions crates/moonbit/src/async/async_abi.mbt
Original file line number Diff line number Diff line change
Expand Up @@ -104,6 +104,14 @@ extern "wasm" fn int_array2ptr(array : FixedArray[Int]) -> Int =
///|
struct WaitableSet(Int) derive(Eq, Hash)

///|
#deprecated
pub extend WaitableSet with Eq::{not_equal, equal}

///|
#deprecated
pub extend WaitableSet with Hash::{hash, hash_combine}

///|
fn WaitableSet::new() -> WaitableSet {
WaitableSet(waitable_set_new())
Expand Down
18 changes: 14 additions & 4 deletions crates/moonbit/src/async/async_primitive.mbt
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,19 @@
// limitations under the License.

///|
async fn[T, E : Error] async_suspend(
cb : ((T) -> Unit, (E) -> Unit) -> Unit,
) -> T raise E = "%async.suspend"
async fn[T] async_suspend(cb : ((T) -> Unit) -> Unit) -> T noraise + nocancel = "%async.suspend"

///|
fn run_async(f : async () -> Unit noraise) = "%async.run"
fn run_async(f : async () -> Unit noraise + nocancel) = "%async.run"

///|
/// Deliver cancellation from the runtime, or resume it after boundary cleanup.
#internal(wit_bindgen, "runtime and generated binding code only")
#doc(hidden)
pub async fn[X] raise_cancellation_signal() -> X noraise = "%async.raise_cancellation_signal"

///|
async fn[X] capture_cancellation(
f : async () -> X,
handler~ : async () -> X nocancel,
) -> X nocancel = "%async.capture_cancellation"
11 changes: 6 additions & 5 deletions crates/moonbit/src/async/cond_var.mbt
Original file line number Diff line number Diff line change
Expand Up @@ -35,14 +35,15 @@ pub fn CondVar::CondVar() -> CondVar {

///|
/// Wait until the condition variable is signaled.
pub async fn CondVar::wait(self : CondVar) -> Unit {
pub async fn CondVar::wait(self : CondVar) -> Unit noraise {
let waiter = { woken: false, coro: Some(current_coroutine()) }
self.waiters.push_back(waiter)
suspend() catch {
_ if waiter.woken => ()
err => {
match suspend_check_cancel() {
Continue => ()
Cancelled if waiter.woken => ()
Cancelled => {
waiter.coro = None
raise err
raise_cancellation_signal()
}
}
}
Expand Down
118 changes: 66 additions & 52 deletions crates/moonbit/src/async/coroutine.mbt
Original file line number Diff line number Diff line change
Expand Up @@ -12,12 +12,19 @@
// See the License for the specific language governing permissions and
// limitations under the License.

///|
priv enum SuspendResult {
Continue
Cancelled
}

///|
priv enum State {
Done
Fail(Error)
Cancelled
Running
Suspend(ok_cont~ : (Unit) -> Unit, err_cont~ : (Error) -> Unit)
Suspend((SuspendResult) -> Unit)
}

///|
Expand All @@ -39,7 +46,7 @@ impl Eq for Coroutine with fn equal(c1, c2) {

///|
impl Hash for Coroutine with fn hash_combine(self, hasher) {
self.coro_id.hash_combine(hasher)
Hash::hash_combine(self.coro_id, hasher)
}

///|
Expand All @@ -58,19 +65,16 @@ pub fn is_being_cancelled() -> Bool {
}

///|
pub fn check_cancellation() -> Unit raise {
pub async fn check_cancellation() -> Unit noraise {
if is_being_cancelled() {
raise Cancelled::Cancelled
raise_cancellation_signal()
}
}

///|
pub(all) suberror Cancelled derive(Debug)

///|
fn Coroutine::cancel(self : Coroutine) -> Unit {
match self.state {
Done | Fail(_) => return
Done | Fail(_) | Cancelled => return
Running | Suspend(_) => ()
}
self.cancelled = true
Expand All @@ -80,21 +84,29 @@ fn Coroutine::cancel(self : Coroutine) -> Unit {
}

///|
async fn suspend() -> Unit {
guard scheduler.curr_coro is Some(coro)
async fn suspend_check_cancel() -> SuspendResult noraise {
guard! scheduler.curr_coro is Some(coro)
if coro.cancelled && !coro.shielded {
raise Cancelled::Cancelled
return Cancelled
}
coro.lane.blocking += 1
defer {
coro.lane.blocking -= 1
}
async_suspend(fn(ok_cont, err_cont) {
guard coro.state is Running
coro.state = Suspend(ok_cont~, err_cont~)
async_suspend(fn(cont) {
guard! coro.state is Running
coro.state = Suspend(cont)
})
}

///|
async fn suspend() -> Unit noraise {
match suspend_check_cancel() {
Continue => ()
Cancelled => raise_cancellation_signal()
}
}

///|
fn spawn_owned(
f : async () -> Unit,
Expand All @@ -105,7 +117,7 @@ fn spawn_owned(
let coro = {
state: Running,
ready: true,
shielded: true,
shielded: false,
downstream: Set([]),
coro_id: scheduler.coro_id,
lane,
Expand All @@ -114,11 +126,13 @@ fn spawn_owned(
}
fn run(_) {
run_async(() => {
coro.shielded = false
try f() catch {
// Even a task cancelled before its first turn enters its body, so entry
// defers run. Cancellation is delivered at its first suspension/check.
try capture_cancellation(f, handler=() => coro.state = Cancelled) catch {
err => coro.state = Fail(err)
} noraise {
_ => coro.state = Done
_ if coro.state is Running => coro.state = Done
_ => ()
}
for coro in coro.downstream {
coro.wake()
Expand All @@ -127,7 +141,7 @@ fn spawn_owned(
})
}

coro.state = Suspend(ok_cont=run, err_cont=_ => ())
coro.state = Suspend(run)
enqueue(coro)
coro
}
Expand All @@ -145,60 +159,60 @@ fn spawn(f : async () -> Unit) -> Coroutine {
}

///|
fn Coroutine::unwrap(self : Coroutine) -> Unit raise {
fn Coroutine::unwrap(self : Coroutine) -> SuspendResult raise {
match self.state {
Done => ()
Done => Continue
Cancelled => Cancelled
Fail(err) => raise err
Running | Suspend(_) => panic()
}
}

///|
async fn Coroutine::wait(target : Coroutine) -> Unit {
guard scheduler.curr_coro is Some(coro)
guard !physical_equal(coro, target)
async fn Coroutine::wait(target : Coroutine) -> SuspendResult {
guard! scheduler.curr_coro is Some(coro)
guard! !physical_equal(coro, target)
match target.state {
Done => return
Done => return Continue
Cancelled => return Cancelled
Fail(err) => raise err
Running | Suspend(_) => ()
}
target.downstream.add(coro)
try suspend() catch {
err => {
target.downstream.remove(coro)
raise err
}
} noraise {
_ => target.unwrap()
}
defer target.downstream.remove(coro)
suspend()
target.unwrap()
}

///|
fn Coroutine::check_error(coro : Coroutine) -> Unit raise {
match coro.state {
Fail(err) => raise err
Done | Running | Suspend(_) => ()
Done | Cancelled | Running | Suspend(_) => ()
}
}

///|
pub async fn[X] protect_from_cancel(
f : async () -> X,
resume_on_cancel? : Bool = false,
) -> X {
guard scheduler.curr_coro is Some(coro)
if coro.shielded {
// already in a shield, do nothing
f()
} else {
coro.shielded = true
defer {
coro.shielded = false
}
let result = f()
if !resume_on_cancel && coro.cancelled {
raise Cancelled::Cancelled
pub async fn[X] protect_from_cancel(f : async () -> X) -> X nocancel {
guard! scheduler.curr_coro is Some(coro)
// The shield prevents delivery; capture_cancellation expresses that invariant
// in the effect type without changing normal error propagation.
capture_cancellation(handler=() => panic(), () => {
if coro.shielded {
f()
} else {
coro.shielded = true
defer {
coro.shielded = false
}
f()
}
result
}
})
}

///|
/// Observe cancellation without clearing the task's sticky cancelled state.
/// Returns None on cancellation, Some on success, and propagates normal errors.
pub async fn[X] handle_cancellation(f : async () -> X) -> X? {
capture_cancellation(handler=() => None, () => Some(f()))
}
Loading
Loading