From 033e236b8c36e9a7ca2735bcf3e17c92b8f1a626 Mon Sep 17 00:00:00 2001 From: Luke Wagner Date: Mon, 14 Sep 2026 11:31:57 -0500 Subject: [PATCH] Add thread.{get,set}-task and task.drop built-ins Resolves #659 --- design/mvp/Binary.md | 3 + design/mvp/CanonicalABI.md | 227 ++-- design/mvp/Concurrency.md | 67 +- design/mvp/Explainer.md | 63 +- design/mvp/canonical-abi/definitions.py | 98 +- design/mvp/canonical-abi/run_tests.py | 90 ++ test/async/thread-set-task.wast | 1604 +++++++++++++++++++++++ test/binary/binary.wast | 7 +- test/nyi.txt | 2 + test/values/post-return.wast | 24 + 10 files changed, 2034 insertions(+), 151 deletions(-) create mode 100644 test/async/thread-set-task.wast diff --git a/design/mvp/Binary.md b/design/mvp/Binary.md index eab7dc94d..fb086d4e7 100644 --- a/design/mvp/Binary.md +++ b/design/mvp/Binary.md @@ -342,6 +342,9 @@ canon ::= 0x00 0x00 f: opts: ft: => (canon lift | 0x2b 0x00 => (canon thread.yield-then-resume (core func)) 🧵 | 0x2c 0x00 => (canon thread.suspend-then-promote (core func)) 🧵 | 0x2d 0x00 => (canon thread.yield-then-promote (core func)) 🧵 + | 0x30 => (canon thread.set-task (core func)) 🧵 + | 0x31 => (canon thread.get-task (core func)) 🧵 + | 0x32 => (canon task.drop (core func)) 🧵 | 0x40 sh?: ft: => (canon thread.spawn-ref sh? ft (core func)) 🧵② | 0x41 sh?: ft: t: => (canon thread.spawn-indirect sh? ft t (core func)) 🧵② | 0x42 sh?: => (canon thread.available-parallelism sh? (core func)) 🧵② diff --git a/design/mvp/CanonicalABI.md b/design/mvp/CanonicalABI.md index 20fd830c0..03281843b 100644 --- a/design/mvp/CanonicalABI.md +++ b/design/mvp/CanonicalABI.md @@ -66,6 +66,9 @@ specified here. * [`canon thread.yield-then-resume`](#-canon-threadyield-then-resume) 🧵 * [`canon thread.suspend-then-promote`](#-canon-threadsuspend-then-promote) 🧵 * [`canon thread.yield-then-promote`](#-canon-threadyield-then-promote) 🧵 + * [`canon thread.set-task`](#-canon-threadset-task) 🧵 + * [`canon thread.get-task`](#-canon-threadget-task) 🧵 + * [`canon task.drop`](#-canon-taskdrop) 🧵 * [`canon error-context.new`](#-canon-error-contextnew) 📝 * [`canon error-context.debug-message`](#-canon-error-contextdebug-message) 📝 * [`canon error-context.drop`](#-canon-error-contextdrop) 📝 @@ -115,12 +118,12 @@ Once a `component` has been parsed/decoded and validated, it can be loaded at runtime by *instantiating* it to produce a *component instance*. The `ComponentInstance` class tracks all the spec-internal state that is used by the definitions below to specify the Canonical ABI. For example, all `i32` handles -to resources, waitables, waitable sets, error contexts and threads will index -into the `handles` or `threads` fields of `ComponentInstance`. +to resources, waitables, waitable sets, tasks, error contexts and threads will +index into the `handles` or `threads` fields of `ComponentInstance`. ```python class ComponentInstance: store: Store - handles: Table[ResourceHandle | Waitable | WaitableSet | ErrorContext] + handles: Table[ResourceHandle | Waitable | WaitableSet | Task | ErrorContext] threads: Table[Thread] may_leave: bool backpressure: int @@ -616,8 +619,7 @@ class Task: on_resolve: OnResolve state: State num_borrows: int - implicit_thread: Optional[Thread] - threads: list[Thread] + waiting_to_enter: Optional[Thread] def __init__(self, ft, opts, inst, on_start, on_resolve): self.ft = ft @@ -627,26 +629,18 @@ class Task: self.on_resolve = on_resolve self.state = Task.State.INITIAL self.num_borrows = 0 - self.implicit_thread = None - self.threads = [] + self.waiting_to_enter = None ``` The `Task.needs_exclusive` predicate returns whether this task's implicit thread -(`Task.implicit_thread`) has *not* opted in to multiple concurrent linear memory -shadow stacks (via "stackful" lift) and thus, according to [Component Invariant] -#2, requires serialization with all the other implicit threads in the component -instance that have similarly not opted in. This question only applies to -`async`-typed functions, since synchronous functions can't block and thus can -always execute in a LIFO fashion using a single linear memory shadow stack. When -`needs_exclusive` is true, core wasm execution is gated on acquiring the -`ComponentInstance.exclusive_thread` lock. Due to cooperativity, the -`exclusive_thread` "lock" is simply a mutable field holding either `None`, when -unlocked, or, when locked, a reference to the `Task.implicit_thread` currently -holding the lock. +requires the `ComponentInstance.exclusive_thread` lock to be acquired and +released in order to enforce [Component Invariant] #2. In particular, both +flavors of `async`-typed functions released in 0.3.0 (sync ABI and async +`callback` ABI) require the exclusive lock; non-`async`-typed functions and +async non-`callback` ABI functions do not. ```python def needs_exclusive(self): - assert(self.ft.async_) - return not self.opts.async_ or self.opts.callback + return self.ft.async_ and (not self.opts.async_ or self.opts.callback) ``` The `Task.enter_implicit_thread` method implements [backpressure] between when @@ -674,30 +668,26 @@ shadow stack pointer) for components with mixed `async`- and non-`async`- typed exports. ```python def enter_implicit_thread(self): + assert(current_thread().task is self) assert(self.state == Task.State.INITIAL) - self.implicit_thread = current_thread() if self.ft.async_: def has_backpressure(): return (self.inst.backpressure > 0 or (self.needs_exclusive() and self.inst.exclusive_thread is not None)) if has_backpressure() or self.inst.num_waiting_to_enter > 0: + self.waiting_to_enter = current_thread() self.inst.num_waiting_to_enter += 1 - self.implicit_thread.wait_until(lambda: not has_backpressure()) + current_thread().wait_until(lambda: not has_backpressure()) self.inst.num_waiting_to_enter -= 1 + self.waiting_to_enter = None if self.deliver_pending_cancel(): self.cancel() return False if self.needs_exclusive(): assert(self.inst.exclusive_thread is None) - self.inst.exclusive_thread = self.implicit_thread - self.register_thread(self.implicit_thread) + self.inst.exclusive_thread = current_thread() + current_thread().index = self.inst.threads.add(current_thread()) return True - - def register_thread(self, thread): - assert(thread not in self.threads and thread.task is self) - self.threads.append(thread) - assert(thread.index is None) - thread.index = self.inst.threads.add(thread) ``` Since the order in which suspended threads are resumed is nondeterministic (see `Store.tick` below), once `Task.enter_implicit_thread` suspends the task's @@ -711,58 +701,53 @@ above definition ensures the following properties: backpressure (i.e., disabling backpressure never unleashes an unstoppable thundering herd of pending tasks). -Once a task's implicit thread has cleared the backpressure gate, it is added to -the lists of threads running inside the current task and component instance by -`Task.register_thread()` (which is also called by `thread.new-indirect`, below). +As shown above, only once a task has cleared the backpressure gate is its +implicit thread visibly added to the component-instance-wide `threads` table. Symmetrically, the `Task.exit_implicit_thread` method is called before a task's implicit thread returns to reverse the effects of `Task.enter_implicit_thread`. -In particular, if the `exclusive_thread` lock was acquired, it is released. -`Task.unregister_thread` (which is also called by `thread.new-indirect`, below) -traps if the task's last thread is unregistered and the task has not yet -returned a value to its caller. -```python - def exit_implicit_thread(self): - assert(current_thread() is self.implicit_thread) - self.unregister_thread(self.implicit_thread) - if self.ft.async_ and self.needs_exclusive(): - assert(self.inst.exclusive_thread is self.implicit_thread) - self.inst.exclusive_thread = None - - def unregister_thread(self, thread): - assert(thread in self.threads and thread.task is self) - self.threads.remove(thread) - if len(self.threads) == 0: - trap_if(self.state != Task.State.RESOLVED) - assert(self.num_borrows == 0) - assert(thread.index is not None) - self.inst.threads.remove(thread.index) +For `async`-lifted implicit threads, `current_thread().task` may have been +modified by `thread.set-task` to refer to *any* task in the same component +instance when the implicit thread exits and so the `needs_exclusive` predicate +must take care to use the *original* task's `opts` and `ft` when deciding +whether to release the `exclusive_thread` lock (i.e., the same values used by +`enter_implicit_thread` when deciding whether to acquire the lock). +```python + def exit_implicit_thread(original_task): + assert(current_instance() is original_task.inst) + inst = current_instance() + inst.threads.remove(current_thread().index) + if original_task.needs_exclusive(): + assert(inst.exclusive_thread is current_thread()) + inst.exclusive_thread = None ``` The `Task.request_cancellation` method implements the `OnCancel` callback that is called by the `subtask.cancel` built-in. If a task's implicit thread is waiting to start (in `Task.enter_implicit_thread`, defined above) due to backpressure, then it is immediately cancelled without running any guest code. -Otherwise, if the task's implicit thread is `cancellable` and `ready`, it is -resumed, allowing it to execute until returning or blocking. Currently, only -`callback` threads that have returned to their event loop are `cancellable`. -When resumed, the `callback` thread will call `Task.deliver_pending_cancel`, -which will return `True`, leading to the `TASK_CANCELLED` event code being -passed to the `callback` function. The `Thread.resume` call will return as soon -as the `callback` function returns or blocks, whether or not the subtask called +Otherwise, if the task contains one or more threads that are `cancellable` and +`ready`, one is nondeterministically picked and resumed to allow the task to +execute until returning or blocking. Currently, only `callback` threads that +have returned to their event loop are `cancellable`. When resumed, the +`callback` thread will call `Task.deliver_pending_cancel`, which will return +`True`, leading to the `TASK_CANCELLED` event code being passed to the +`callback` function. The `Thread.resume` call will return as soon as the +`callback` function returns or blocks, whether or not the subtask called `task.{cancel,return}`, and thus the caller of `OnCancel` must handle both the resolved and not-resolved cases after `OnCancel` returns. ```python def request_cancellation(self): if self.state == Task.State.INITIAL: self.state = Task.State.PENDING_CANCEL - self.implicit_thread.resume() + self.waiting_to_enter.resume() assert(self.state == Task.State.RESOLVED) else: assert(self.state == Task.State.STARTED) self.state = Task.State.PENDING_CANCEL - if self.implicit_thread.cancellable and self.implicit_thread.ready(): - self.implicit_thread.resume() + candidates = { t for t in self.inst.threads if t.task is self and t.cancellable and t.ready() } + if candidates: + random.choice(list(candidates)).resume() def has_pending_cancel(self): return self.state == Task.State.PENDING_CANCEL @@ -1080,7 +1065,7 @@ that are represented in Core WebAssembly as `i32` indices into the array. Currently, every component instance contains two tables: a `threads` table containing all the component's [threads](#threads) and a `handles` table containing everything else ([resource handles](#resource-state), -[waitables and waitable sets](#waitable-state) and +[waitables and waitable sets](#waitable-state), [tasks](#tasks) and [error contexts](#-canon-error-contextnew)). ```python class Table: @@ -3405,17 +3390,22 @@ their arguments were lowered or not. assert(types_match_values(flat_ft.params, flat_args)) ``` -If the `async` `canonopt` is *not* specified, a `lift`ed function then calls -the core wasm callee, passing the lowered arguments in core function parameters -and receiving the return value as core function results. Once the core results -are lifted according to `lift_flat_values` above, the optional `post-return` -function (specified as a `canonopt` immediate of `canon lift`) is called, -passing the same core wasm results as parameters so that the `post-return` -function can free any associated allocations. +If the `async` `canonopt` is *not* specified, a `lift`ed function then calls the +core wasm callee, passing the lowered arguments in core function parameters and +receiving the return value as core function results. Before returning, the +thread may have called `thread.set-task` to change its task membership, but on +return, the thread implicitly rejoins its original task. This allows sync +implicit threads to lift their return value using the ABI options and return +type that were statically specified in the `lift`, avoiding the additional +dynamism of `task.return`. After returning the lifted values to the caller, the +optional `post-return` function (specified as a `canonopt` immediate of `canon +lift`) is called, passing the core wasm results as parameters and allowing the +callee to free any linear memory allocations used to hold the return value. ```python if not opts.async_: flat_results = call_and_trap_on_throw(callee, flat_args) assert(types_match_values(flat_ft.results, flat_results)) + thread.task = task result = lift_flat_values(cx, MAX_FLAT_RESULTS, CoreValueIter(flat_results), ft.result_type()) task.return_(result) if opts.post_return is not None: @@ -3432,11 +3422,11 @@ functions can always be implemented by a plain synchronous function call without the need for fibers which would otherwise be necessary if the `post-return` function performed a blocking operation. -In both of the `async` cases below (with or without `callback`), the -`task.return` built-in must be called, providing the return value as core wasm -*parameters* to the `task.return` built-in (rather than as core function -results as in the synchronous case). If `task.return` is *not* called by the -time the `Task`'s last `Thread` exits, there is a trap (in `Task.unregister_thread`). +In both of the `async` cases below (with or without `callback`), the return +value is provided by calling the `task.return` built-in, passing the return +value as core wasm *parameters* (rather than as core function results as in +the synchronous case). If `task.return` never ends up being called, the +task will never complete for the caller. In the `async` non-`callback` ("stackful async") case, there is a single call to the core wasm callee which must return empty core results. Waiting for async @@ -3464,7 +3454,7 @@ function (specified as a `funcidx` immediate in `canon lift`) until the if thread.task.deliver_pending_cancel(): event = (EventCode.TASK_CANCELLED, 0, 0) else: - assert(inst.exclusive_thread is task.implicit_thread) + assert(inst.exclusive_thread is thread) inst.exclusive_thread = None thread.cancellable = True match code: @@ -3482,7 +3472,7 @@ function (specified as a `funcidx` immediate in `canon lift`) until the trap() assert(inst.exclusive_thread is None) thread.cancellable = False - inst.exclusive_thread = task.implicit_thread + inst.exclusive_thread = thread event_code, p1, p2 = event [packed] = call_and_trap_on_throw(opts.callback, [event_code, p1, p2]) code,si = unpack_callback_result(packed) @@ -4568,17 +4558,17 @@ class CoreFuncRef: callee: Callable[[list[CoreValType]], list[CoreValType]] def canon_thread_new_indirect(ft, ftbl: Table[CoreFuncRef], fi, c): - task = current_task() - trap_if(not task.inst.may_leave) + inst = current_instance() + trap_if(not inst.may_leave) f = ftbl.get(fi) assert(ft == CoreFuncType(['i32'], []) or ft == CoreFuncType(['i64'], [])) trap_if(f.t != ft) def thread_func(): [] = call_and_trap_on_throw(f.callee, [c]) - task.unregister_thread(new_thread) - new_thread = Thread(task, thread_func) + inst.threads.remove(new_thread.index) + new_thread = Thread(current_task(), thread_func) assert(new_thread.suspended()) - task.register_thread(new_thread) + new_thread.index = inst.threads.add(new_thread) return [new_thread.index] ``` The newly-created thread starts out in a "suspended" state and so, to @@ -4757,6 +4747,77 @@ def canon_thread_yield_then_promote(i): ``` +### 🧵 `canon thread.set-task` + +For a canonical definition: +```wat +(canon thread.set-task (core func $thread.set-task)) +``` +validation specifies: +* `$thread.set-task` is given type `(func (param i32))` + +Calling `$thread.set-task` invokes the following function which sets the current +thread's containing task to be the task referenced by the given `i32` index. +```python +def canon_thread_set_task(taski): + thread = current_thread() + trap_if(not thread.task.inst.may_leave) + new_task = thread.task.inst.handles.get(taski) + trap_if(not isinstance(new_task, Task)) + thread.task = new_task + return [] +``` +At the moment, the only source of task indices is `thread.get-task`. + + +### 🧵 `canon thread.get-task` + +For a canonical definition: +```wat +(canon thread.get-task (core func $thread.get-task)) +``` +validation specifies: +* `$thread.get-task` is given type `(func (result i32))` + +Calling `$thread.get-task` invokes the following function which adds a new +handle containing a reference to the current task to the current component +instance's `handles` table, returning the `i32` index of the new task handle, +which can be passed as an operand to `thread.set-task`. +```python +def canon_thread_get_task(): + task = current_task() + trap_if(not task.inst.may_leave) + taski = task.inst.handles.add(task) + return [taski] +``` +Note that each call to `thread.get-task` unconditionally allocates a new handle +and thus must be paired with a call to `task.drop` to avoid leaking handles. + + +### 🧵 `canon task.drop` + +For a canonical definition: +```wat +(canon task.drop (core func $task.drop)) +``` +validation specifies: +* `$task.drop` is given type `(func (param i32))` + +Calling `$task.drop` invokes the following function which drops the task handle +at the given index. +```python +def canon_task_drop(taski): + inst = current_instance() + trap_if(not inst.may_leave) + task = inst.handles.remove(taski) + trap_if(not isinstance(task, Task)) + return [] +``` +Note that `task.drop` only drops a *reference* to a task; it does not change the +*state* of the task or destroy the task (which may well have other handle, +thread and subtask referents). + + ### 📝 `canon error-context.new` For a canonical definition: diff --git a/design/mvp/Concurrency.md b/design/mvp/Concurrency.md index c851ef70c..84e9f568c 100644 --- a/design/mvp/Concurrency.md +++ b/design/mvp/Concurrency.md @@ -268,11 +268,16 @@ Thread where a **component store** is the top-level "thing" and analogous to a Core WebAssembly [store]. -The reason for the thread/task split is that, when one thread creates a new -thread, the new thread is contained by the task of the original thread which -creates an N:1 relationship between threads and tasks that ties N threads to -the original export call (= "task") that transitively spawned those N threads. -This relationship serves several purposes described in the following sections. +When a component export is called, one new task is created for the call; this +task contains one new *implicit* thread that executes it. The reason for the +thread/task split is that this implicit thread may then go on to spawn N more +*explicit* threads (via `thread.new-indirect`) that are initially contained by +the same task, thereby creating an N:1 relationship between threads and tasks. + +While the store:instance and instance:task relationships are immutably set on +creation, the task:thread relationship is *mutable*: guest code running inside a +component instance can change which task a thread is currently executing on +behalf of by calling the [`thread.set-task`] built-in, as described below. In the Canonical ABI explainer, threads, tasks, component instances and component stores are represented by the [`Thread`], [`Task`], @@ -311,9 +316,15 @@ supertask, they can be thought of as a single node in the async call stack. A subtask/supertask relationship is immutably established when an import is called, setting the [current task](#current-thread-and-task) as the supertask -of the new subtask created for the import call. Thus, one reason for -associating every thread with a "containing task" is to ensure that there is -always a well-defined async call stack. +of the new subtask created for the import call. Thus, one reason for associating +every thread with a "containing task" is to ensure that there is always a +well-defined async call stack. Note that guest code can call [`thread.set-task`] +to change the containing task of a thread and thus the async call stack can +change completely between two program points while executing a single thread. +For example, when a JS runtime flushes its microtask queue and encounters a JS +callback associated with a task other than the current thread's task, the JS +runtime would call `thread.set-task` so that the async call stack matches the JS +developer's expectation. The async call stack is not currently observable to running components, except that it may nondeterministically appear as part of the callstack stored in @@ -336,9 +347,9 @@ not enforcing a stricter form of Structured Concurrency at the Component Model level is that there are important use cases where forcing a supertask's thread to stay resident just to wait for subtasks to finish would waste resources without tangible benefit. Instead, we can say that once a supertask's last -thread finishes execution, the supertask semantically "tail calls" any still- -executing subtasks, staying technically-alive and on the async call stack until -they complete, but not consuming real resources. +thread exits or switches to another task, the supertask semantically "tail +calls" any still-executing subtasks, staying technically-alive and on the +async call stack until they complete, but not consuming real resources. For scenarios where one component wants to *non-cooperatively* put an upper bound on execution of a call into another component, a separate "[blast zone]" @@ -374,13 +385,23 @@ New threads are created with the [`thread.new-indirect`] built-in. As mentioned [above](#threads-and-tasks), a spawned thread inherits the task of the spawning thread which is why threads and tasks are N:1. `thread.new-indirect` adds a new thread to the component instance's threads table and returns the `i32` index of -this table entry to the Core WebAssembly caller. Like [`pthread_create`], -`thread.new-indirect` takes a Core WebAssembly function (via index into a -`funcref` table) and a "closure" parameter to pass to the function when called -on the new thread. However, unlike `pthread_create`, the new thread is -initially in a "suspended" state and must be explicitly "resumed" using one of -the following 3 thread built-ins. Once the thread is resumed, the thread can -learn its own index by calling the [`thread.index`] built-in. +this table entry to the Core WebAssembly caller. + +After creation, the implicitly-set containing task of a thread can be explicitly +overridden using the [`thread.set-task`] built-in. `thread.set-task` sets the +containing task of the current thread to a task handle that was retrieved via +[`thread.get-task`]. `thread.get-task` always allocates a fresh handle storing a +reference to the current thread's containing task. The `i32` index of this new +handle must later be explicitly dropped via [`task.drop`] to avoid leaking +the task. Tasks are thus kept alive by any or all of: contained threads, subtask +handles and task handles. + +Like [`pthread_create`], `thread.new-indirect` takes a Core WebAssembly function +(via index into a `funcref` table) and a "closure" parameter to pass to the +function when called on the new thread. However, unlike `pthread_create`, the +new thread is initially in a "suspended" state and must be explicitly "resumed" +using one of the following 3 thread built-ins. Once the thread is resumed, the +thread can learn its own index by calling the [`thread.index`] built-in. A suspended thread (identified by thread-table index) can be resumed at some nondeterministic point in future via the [`thread.resume-later`] built-in. In @@ -733,10 +754,7 @@ the "started" state. The way an `async` export returns its value using the async ABI is by calling [`task.return`], passing the core values that are to be lifted as *parameters*. When using the async ABI, *any* of the threads contained by a task can call -`task.return`; there is no "main thread" of a task. When the last thread of a -task returns, there is a trap if `task.return` has not been called. Thus, *some* -thread (either the thread created implicitly for the initial export call or some -thread transitively created by that thread) must call `task.return`. +`task.return`; there is no "main thread" of a task. Returning values by calling `task.return` allows a task to continue executing even after it has passed its initial results to the caller. This is also @@ -752,7 +770,7 @@ the readable end passed for `in`) and `stream.write`s (of the writable end it `stream.new`ed) before exiting the task. Once `task.return` is called, the task is in the "returned" state. Calling -`task.return` when not in the "started" state traps. +`task.return` when already in a resolved state traps. ### Borrows @@ -1564,6 +1582,9 @@ the concurrency story: [`thread.yield-then-resume`]: Explainer.md#-threadyield-then-resume [`thread.suspend-then-promote`]: Explainer.md#-threadsuspend-then-promote [`thread.yield-then-promote`]: Explainer.md#-threadyield-then-promote +[`thread.set-task`]: Explainer.md#-threadset-task +[`thread.get-task`]: Explainer.md#-threadget-task +[`task.drop`]: Explainer.md#-taskdrop [`{stream,future}.new`]: Explainer.md#-streamnew-and-futurenew [`{stream,future}.{read,write}`]: Explainer.md#-streamread-and-streamwrite [`stream.cancel-write`]: Explainer.md#-streamcancel-read-streamcancel-write-futurecancel-read-and-futurecancel-write diff --git a/design/mvp/Explainer.md b/design/mvp/Explainer.md index 0cd10c259..0946d6c32 100644 --- a/design/mvp/Explainer.md +++ b/design/mvp/Explainer.md @@ -1593,6 +1593,9 @@ canon ::= ... | (canon thread.yield-then-resume (core func ?)) 🧵 | (canon thread.suspend-then-promote (core func ?)) 🧵 | (canon thread.yield-then-promote (core func ?)) 🧵 + | (canon thread.set-task (core func ?)) 🧵 + | (canon thread.get-task (core func ?)) 🧵 + | (canon task.drop (core func ?)) 🧵 | (canon error-context.new * (core func ?)) 📝 | (canon error-context.debug-message * (core func ?)) 📝 | (canon error-context.drop (core func ?)) 📝 @@ -1748,7 +1751,7 @@ For details, see [Backpressure] in the concurrency explainer and | Canonical ABI signature | `[lower(FuncT.results)*] -> []` | The `task.return` built-in takes as parameters the result values of the -[current task]. One of `task.return` or `task.cancel` must be called exactly +[current task]. One of `task.return` or `task.cancel` must be called at most once from any of a task's threads. The `canon task.return` definition takes component-level return type and the @@ -2320,6 +2323,60 @@ returned `i32` is always `0` and may be removed in a future ABI revision. For details, see [Thread Built-ins] in the concurrency explainer and [`canon_thread_yield_then_promote`] in the Canonical ABI explainer. +###### 🧵 `thread.set-task` + +| Synopsis | | +| -------------------------- | --------------- | +| Approximate WIT signature | `func(t: task)` | +| Canonical ABI signature | `[t:i32] -> []` | + +The `thread.set-task` built-in allows the [current thread] to set its containing +task to the given operand, which immediately changes the [current task]. +Built-ins like `task.return` and `task.cancel` are defined in terms of "the +current task", and thus `thread.set-task` allows a thread to return a value for +a task other than the task that initially spawned the thread. + +Tasks form an [async call stack] that is consulted for debugging, observability, +and host import-to-export call attribution. Thus, changing the current task +allows a guest's concurrency runtime to control what async call stack to +associate with the current wasm execution. + +The `i32` index passed to `thread.set-task` must be a task handle which can +currently only be retrieved by calling `thread.get-task`. + +For details, see [Thread Built-ins] in the concurrency explainer and +[`canon_thread_set_task`] in the Canonical ABI explainer. + +###### 🧵 `thread.get-task` + +| Synopsis | | +| -------------------------- | ---------------- | +| Approximate WIT signature | `func() -> task` | +| Canonical ABI signature | `[] -> [i32]` | + +The `thread.get-task` built-in returns a handle referring to the [current task] +(which is the containing task of the [current thread]). This handle is currently +only used as an argument in `thread.set-task`. Each call to `thread.get-task` +returns a fresh handle which must be released by `task.drop` to avoid a handle +table leak. + +For details, see [Thread Built-ins] in the concurrency explainer and +[`canon_thread_get_task`] in the Canonical ABI explainer. + +###### 🧵 `task.drop` + +| Synopsis | | +| -------------------------- | --------------- | +| Approximate WIT signature | `func(t: task)` | +| Canonical ABI signature | `[t:i32] -> []` | + +The `task.drop` built-in drops a handle allocated by `thread.get-task`. The +state of the underlying task is not modified nor is the task destroyed, as it +may have other active referents. + +For details, see [Thread Built-ins] in the concurrency explainer and +[`canon_task_drop`] in the Canonical ABI explainer. + ###### 🧵② `thread.spawn-ref` | Synopsis | | @@ -3474,6 +3531,9 @@ For some use-case-focused, worked examples, see: [`canon_thread_yield_then_resume`]: CanonicalABI.md#-canon-threadyield-then-resume [`canon_thread_suspend_then_promote`]: CanonicalABI.md#-canon-threadsuspend-then-promote [`canon_thread_yield_then_promote`]: CanonicalABI.md#-canon-threadyield-then-promote +[`canon_thread_get_task`]: CanonicalABI.md#-canon-threadget-task +[`canon_thread_set_task`]: CanonicalABI.md#-canon-threadset-task +[`canon_task_drop`]: CanonicalABI.md#-canon-taskdrop [`canon_thread_spawn_ref`]: CanonicalABI.md#-canon-threadspawn-ref [`canon_thread_spawn_indirect`]: CanonicalABI.md#-canon-threadspawn-indirect [`canon_thread_available_parallelism`]: CanonicalABI.md#-canon-threadavailable_parallelism @@ -3483,6 +3543,7 @@ For some use-case-focused, worked examples, see: [Summary]: Concurrency.md#summary [Current Thread]: Concurrency.md#current-thread-and-task [Current Task]: Concurrency.md#current-thread-and-task +[Async Call Stack]: Concurrency.md#subtasks-and-supertasks [Thread-Local Storage]: Concurrency.md#thread-local-storage [Subtask]: Concurrency.md#subtasks-and-supertasks [Stream or Future]: Concurrency.md#streams-and-futures diff --git a/design/mvp/canonical-abi/definitions.py b/design/mvp/canonical-abi/definitions.py index d87d9de3d..c55400304 100644 --- a/design/mvp/canonical-abi/definitions.py +++ b/design/mvp/canonical-abi/definitions.py @@ -188,7 +188,7 @@ class FutureType(ValType): class ComponentInstance: store: Store - handles: Table[ResourceHandle | Waitable | WaitableSet | ErrorContext] + handles: Table[ResourceHandle | Waitable | WaitableSet | Task | ErrorContext] threads: Table[Thread] may_leave: bool backpressure: int @@ -405,8 +405,7 @@ class State(Enum): on_resolve: OnResolve state: State num_borrows: int - implicit_thread: Optional[Thread] - threads: list[Thread] + waiting_to_enter: Optional[Thread] def __init__(self, ft, opts, inst, on_start, on_resolve): self.ft = ft @@ -416,65 +415,52 @@ def __init__(self, ft, opts, inst, on_start, on_resolve): self.on_resolve = on_resolve self.state = Task.State.INITIAL self.num_borrows = 0 - self.implicit_thread = None - self.threads = [] + self.waiting_to_enter = None def needs_exclusive(self): - assert(self.ft.async_) - return not self.opts.async_ or self.opts.callback + return self.ft.async_ and (not self.opts.async_ or self.opts.callback) def enter_implicit_thread(self): + assert(current_thread().task is self) assert(self.state == Task.State.INITIAL) - self.implicit_thread = current_thread() if self.ft.async_: def has_backpressure(): return (self.inst.backpressure > 0 or (self.needs_exclusive() and self.inst.exclusive_thread is not None)) if has_backpressure() or self.inst.num_waiting_to_enter > 0: + self.waiting_to_enter = current_thread() self.inst.num_waiting_to_enter += 1 - self.implicit_thread.wait_until(lambda: not has_backpressure()) + current_thread().wait_until(lambda: not has_backpressure()) self.inst.num_waiting_to_enter -= 1 + self.waiting_to_enter = None if self.deliver_pending_cancel(): self.cancel() return False if self.needs_exclusive(): assert(self.inst.exclusive_thread is None) - self.inst.exclusive_thread = self.implicit_thread - self.register_thread(self.implicit_thread) + self.inst.exclusive_thread = current_thread() + current_thread().index = self.inst.threads.add(current_thread()) return True - def register_thread(self, thread): - assert(thread not in self.threads and thread.task is self) - self.threads.append(thread) - assert(thread.index is None) - thread.index = self.inst.threads.add(thread) - - def exit_implicit_thread(self): - assert(current_thread() is self.implicit_thread) - self.unregister_thread(self.implicit_thread) - if self.ft.async_ and self.needs_exclusive(): - assert(self.inst.exclusive_thread is self.implicit_thread) - self.inst.exclusive_thread = None - - def unregister_thread(self, thread): - assert(thread in self.threads and thread.task is self) - self.threads.remove(thread) - if len(self.threads) == 0: - trap_if(self.state != Task.State.RESOLVED) - assert(self.num_borrows == 0) - assert(thread.index is not None) - self.inst.threads.remove(thread.index) + def exit_implicit_thread(original_task): + assert(current_instance() is original_task.inst) + inst = current_instance() + inst.threads.remove(current_thread().index) + if original_task.needs_exclusive(): + assert(inst.exclusive_thread is current_thread()) + inst.exclusive_thread = None def request_cancellation(self): if self.state == Task.State.INITIAL: self.state = Task.State.PENDING_CANCEL - self.implicit_thread.resume() + self.waiting_to_enter.resume() assert(self.state == Task.State.RESOLVED) else: assert(self.state == Task.State.STARTED) self.state = Task.State.PENDING_CANCEL - if self.implicit_thread.cancellable and self.implicit_thread.ready(): - self.implicit_thread.resume() + candidates = { t for t in self.inst.threads if t.task is self and t.cancellable and t.ready() } + if candidates: + random.choice(list(candidates)).resume() def has_pending_cancel(self): return self.state == Task.State.PENDING_CANCEL @@ -2077,6 +2063,7 @@ def thread_func(): if not opts.async_: flat_results = call_and_trap_on_throw(callee, flat_args) assert(types_match_values(flat_ft.results, flat_results)) + thread.task = task result = lift_flat_values(cx, MAX_FLAT_RESULTS, CoreValueIter(flat_results), ft.result_type()) task.return_(result) if opts.post_return is not None: @@ -2099,7 +2086,7 @@ def thread_func(): if thread.task.deliver_pending_cancel(): event = (EventCode.TASK_CANCELLED, 0, 0) else: - assert(inst.exclusive_thread is task.implicit_thread) + assert(inst.exclusive_thread is thread) inst.exclusive_thread = None thread.cancellable = True match code: @@ -2117,7 +2104,7 @@ def thread_func(): trap() assert(inst.exclusive_thread is None) thread.cancellable = False - inst.exclusive_thread = task.implicit_thread + inst.exclusive_thread = thread event_code, p1, p2 = event [packed] = call_and_trap_on_throw(opts.callback, [event_code, p1, p2]) code,si = unpack_callback_result(packed) @@ -2566,17 +2553,17 @@ class CoreFuncRef: callee: Callable[[list[CoreValType]], list[CoreValType]] def canon_thread_new_indirect(ft, ftbl: Table[CoreFuncRef], fi, c): - task = current_task() - trap_if(not task.inst.may_leave) + inst = current_instance() + trap_if(not inst.may_leave) f = ftbl.get(fi) assert(ft == CoreFuncType(['i32'], []) or ft == CoreFuncType(['i64'], [])) trap_if(f.t != ft) def thread_func(): [] = call_and_trap_on_throw(f.callee, [c]) - task.unregister_thread(new_thread) - new_thread = Thread(task, thread_func) + inst.threads.remove(new_thread.index) + new_thread = Thread(current_task(), thread_func) assert(new_thread.suspended()) - task.register_thread(new_thread) + new_thread.index = inst.threads.add(new_thread) return [new_thread.index] ### 🧵 `canon thread.resume-later` @@ -2648,6 +2635,33 @@ def canon_thread_yield_then_promote(i): thread.yield_then_promote(other_thread) return [0] +### 🧵 `canon thread.set-task` + +def canon_thread_set_task(taski): + thread = current_thread() + trap_if(not thread.task.inst.may_leave) + new_task = thread.task.inst.handles.get(taski) + trap_if(not isinstance(new_task, Task)) + thread.task = new_task + return [] + +### 🧵 `canon thread.get-task` + +def canon_thread_get_task(): + task = current_task() + trap_if(not task.inst.may_leave) + taski = task.inst.handles.add(task) + return [taski] + +### 🧵 `canon task.drop` + +def canon_task_drop(taski): + inst = current_instance() + trap_if(not inst.may_leave) + task = inst.handles.remove(taski) + trap_if(not isinstance(task, Task)) + return [] + ### 📝 `canon error-context.new` @dataclass diff --git a/design/mvp/canonical-abi/run_tests.py b/design/mvp/canonical-abi/run_tests.py index a301acdda..c2ab2918c 100644 --- a/design/mvp/canonical-abi/run_tests.py +++ b/design/mvp/canonical-abi/run_tests.py @@ -526,6 +526,21 @@ def core_consumer_realloc(args): fail("thread.index must trap during realloc") except Trap: pass + try: + canon_thread_get_task() + fail("thread.get-task must trap during realloc") + except Trap: + pass + try: + canon_thread_set_task(0) + fail("thread.set-task must trap during realloc") + except Trap: + pass + try: + canon_task_drop(0) + fail("task.drop must trap during realloc") + except Trap: + pass return consumer_heap.realloc(args) consumer_opts = mk_opts(MemInst(consumer_heap.memory, 'i32'), realloc = core_consumer_realloc) @@ -3116,6 +3131,80 @@ def on_resolve(v): assert(result == 42) assert(other_result == 43) +def test_thread_set_task(): + store = Store() + inst = ComponentInstance(store) + opts = mk_opts(async_ = True) + + ftbl = Table() + ft = CoreFuncType(['i32'],[]) + + t1i = None + bi = None + b_taski = None + + def thread_func1(args): + assert(args == [201]) + task_a = current_task() + assert(task_a.state == Task.State.RESOLVED) + assert(canon_thread_index() == [t1i]) + + [a_taski] = canon_thread_get_task() + [a_taski2] = canon_thread_get_task() + assert(a_taski != a_taski2) + [] = canon_thread_set_task(a_taski2) + [] = canon_task_drop(a_taski2) + + [] = canon_thread_set_task(b_taski) + task_b = current_task() + assert(task_b is not task_a) + [] = canon_task_return([U8Type()], opts, [55]) + + [] = canon_thread_set_task(a_taski) + assert(current_task() is task_a) + [] = canon_task_drop(a_taski) + + [] = canon_thread_set_task(b_taski) + assert(current_task() is task_b) + [] = canon_task_drop(b_taski) + [] = canon_thread_resume_later(bi) + return [] + fi1 = ftbl.add(CoreFuncRef(ft, thread_func1)) + + def core_func_a(args): + assert(not args) + nonlocal t1i + [t1i] = canon_thread_new_indirect(ft, ftbl, fi1, 201) + [] = canon_thread_resume_later(t1i) + [] = canon_task_return([U8Type()], opts, [11]) + return [] + + def core_func_b(args): + assert(not args) + nonlocal bi, b_taski + [bi] = canon_thread_index() + [b_taski] = canon_thread_get_task() + assert(canon_thread_suspend() == [0]) + return [] + + a_result = None + def on_resolve_a(v): + nonlocal a_result + [a_result] = v + + b_result = None + def on_resolve_b(v): + nonlocal b_result + [b_result] = v + + caller_ft = FuncType([], [U8Type()], async_ = True) + _ = store.invoke(store.lift(core_func_a, caller_ft, opts, inst), lambda:[], on_resolve_a) + _ = store.invoke(store.lift(core_func_b, caller_ft, opts, inst), lambda:[], on_resolve_b) + while store.waiting: + store.tick() + assert(a_result == 11) + assert(b_result == 55) + test_roundtrips() test_trap_propagation() test_cross_component_realloc() @@ -3146,5 +3235,6 @@ def on_resolve(v): test_async_flat_params() test_threads() test_sync_threads() +test_thread_set_task() print("All tests passed") diff --git a/test/async/thread-set-task.wast b/test/async/thread-set-task.wast new file mode 100644 index 000000000..6a0390ef1 --- /dev/null +++ b/test/async/thread-set-task.wast @@ -0,0 +1,1604 @@ +;; Test thread.{get,set}-task and task.drop + +;; Task handles: `thread.get-task` hands out a fresh handle each time and +;; `task.drop` invalidates one without touching the task, so both built-ins +;; reject anything that is not a live task handle. No second task is needed for +;; any of this; each case runs on its own instance, where the first handle +;; `thread.get-task` or `waitable-set.new` hands out has index 1. +(component definition $Handles + (component $C + (core module $Core + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + (import "" "task.drop" (func $task.drop (param i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + + (func (export "unknown-set") (result i32) + (call $thread.set-task (i32.const 100)) + unreachable) + (func (export "unknown-drop") (result i32) + (call $task.drop (i32.const 100)) + unreachable) + + ;; tasks share an index space with the other handle types, so a handle + ;; that exists can still be the wrong kind + (func (export "wrong-set") (result i32) + (call $thread.set-task (call $waitable-set.new)) + unreachable) + (func (export "wrong-drop") (result i32) + (call $task.drop (call $waitable-set.new)) + unreachable) + + ;; `task.drop` invalidates the handle without touching the task, so both + ;; built-ins reject a dropped handle + (func (export "dropped-set") (result i32) + (local $t i32) + (local.set $t (call $thread.get-task)) + (call $task.drop (local.get $t)) + (call $thread.set-task (local.get $t)) + unreachable) + (func (export "double-drop") (result i32) + (local $t i32) + (local.set $t (call $thread.get-task)) + (call $task.drop (local.get $t)) + (call $task.drop (local.get $t)) + unreachable) + + ;; index 0 is permanently reserved and so is never a task handle + (func (export "zero-set") (result i32) + (call $thread.set-task (i32.const 0)) + unreachable) + (func (export "zero-drop") (result i32) + (call $task.drop (i32.const 0)) + unreachable) + + ;; a dropped index goes back on the free list and is handed out again by + ;; the next `thread.get-task`; the recycled handle is a normal handle, + ;; and moving to a handle for one's own task is a no-op + (func (export "reuse") (result i32) + (local $first i32) + (local $second i32) + (local.set $first (call $thread.get-task)) + (call $task.drop (local.get $first)) + (local.set $second (call $thread.get-task)) + (if (i32.ne (local.get $first) (local.get $second)) + (then unreachable)) + (call $thread.set-task (local.get $second)) + (call $task.drop (local.get $second)) + (i32.const 42)) + ) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (canon task.drop (core func $task.drop)) + (canon waitable-set.new (core func $waitable-set.new)) + (core instance $core (instantiate $Core (with "" (instance + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + (export "task.drop" (func $task.drop)) + (export "waitable-set.new" (func $waitable-set.new)) + )))) + (func (export "unknown-set") (result u32) + (canon lift (core func $core "unknown-set"))) + (func (export "unknown-drop") (result u32) + (canon lift (core func $core "unknown-drop"))) + (func (export "wrong-set") (result u32) + (canon lift (core func $core "wrong-set"))) + (func (export "wrong-drop") (result u32) + (canon lift (core func $core "wrong-drop"))) + (func (export "dropped-set") (result u32) + (canon lift (core func $core "dropped-set"))) + (func (export "double-drop") (result u32) + (canon lift (core func $core "double-drop"))) + (func (export "zero-set") (result u32) + (canon lift (core func $core "zero-set"))) + (func (export "zero-drop") (result u32) + (canon lift (core func $core "zero-drop"))) + (func (export "reuse") (result u32) + (canon lift (core func $core "reuse"))) + ) + (instance $c (instantiate $C)) + (func (export "unknown-set") (alias export $c "unknown-set")) + (func (export "unknown-drop") (alias export $c "unknown-drop")) + (func (export "wrong-set") (alias export $c "wrong-set")) + (func (export "wrong-drop") (alias export $c "wrong-drop")) + (func (export "dropped-set") (alias export $c "dropped-set")) + (func (export "double-drop") (alias export $c "double-drop")) + (func (export "zero-set") (alias export $c "zero-set")) + (func (export "zero-drop") (alias export $c "zero-drop")) + (func (export "reuse") (alias export $c "reuse")) +) +(component instance $h1 $Handles) +(assert_trap (invoke "unknown-set") "unknown handle index 100") +(component instance $h2 $Handles) +(assert_trap (invoke "unknown-drop") "unknown handle index 100") +(component instance $h3 $Handles) +(assert_trap (invoke "wrong-set") "handle is not a task") +(component instance $h4 $Handles) +(assert_trap (invoke "wrong-drop") "handle is not a task") +(component instance $h5 $Handles) +(assert_trap (invoke "dropped-set") "unknown handle index 1") +(component instance $h6 $Handles) +(assert_trap (invoke "double-drop") "unknown handle index 1") +(component instance $h7 $Handles) +(assert_trap (invoke "zero-set") "unknown handle index 0") +(component instance $h8 $Handles) +(assert_trap (invoke "zero-drop") "unknown handle index 0") +(component instance $h9 $Handles) +(assert_return (invoke "reuse") (u32.const 42)) + +;; Basic task switching: an explicit thread of task A joins task B, returns B's +;; value and moves back home again. `task.return`'s result type is checked +;; against the task the calling thread is in, not the one it started in, so the +;; u8-typed `task.return` in "join-bad" traps against B2's u32 result. +(component definition $Join + (component $C + (core module $Table + (table (export "__indirect_function_table") 2 funcref)) + (core instance $table (instantiate $Table)) + (core module $Core + (import "" "task.return-u8" (func $task.return-u8 (param i32))) + (import "" "task.return-u32" (func $task.return-u32 (param i32))) + (import "" "thread.new-indirect" (func $thread.new-indirect (param i32 i32) (result i32))) + (import "" "thread.index" (func $thread.index (result i32))) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + (import "" "task.drop" (func $task.drop (param i32))) + (import "" "thread.resume-later" (func $thread.resume-later (param i32))) + (import "" "thread.suspend-then-resume" (func $thread.suspend-then-resume (param i32) (result i32))) + (import "" "__indirect_function_table" (table $indirect-function-table 2 funcref)) + + (global $worker-thread (mut i32) (i32.const 0xdead)) ;; explicit thread spawned into task A + (global $worker-bad-thread (mut i32) (i32.const 0xdead)) ;; second mover, used by join-bad + (global $join-implicit (mut i32) (i32.const 0xdead)) ;; implicit thread of the current join task + (global $join-task (mut i32) (i32.const 0xdead)) ;; handle for the current join task + + ;; $worker-thread: starts in task A (already resolved), joins task B, + ;; returns for B + (func $worker (param i32) + (local $home i32) + (local $alias i32) + (local.set $home (call $thread.get-task)) + + ;; each thread.get-task hands out a fresh handle for the same task and + ;; moving to a handle for one's own task is a no-op + (local.set $alias (call $thread.get-task)) + (if (i32.eq (local.get $home) (local.get $alias)) + (then unreachable)) + (call $thread.set-task (local.get $alias)) + (call $task.drop (local.get $alias)) + + ;; task.return for task B (not for task A) + (call $thread.set-task (global.get $join-task)) + (call $task.return-u32 (i32.const 42)) + + ;; the handle saved on entry moves this thread back to task A, from + ;; which it exits after making $join-implicit runnable again + (call $thread.set-task (local.get $home)) + (call $task.drop (local.get $home)) + (call $thread.resume-later (global.get $join-implicit))) + + ;; $worker-bad-thread: joins task B2 and calls task A's u8-typed + ;; task.return: the declared result type is checked against the + ;; *current* task, which is now B2 with a u32 result, so this traps + (func $worker-bad (param i32) + (call $thread.set-task (global.get $join-task)) + (call $task.return-u8 (i32.const 33)) + unreachable) + + (elem (table $indirect-function-table) (i32.const 0) func $worker $worker-bad) + + ;; task A: spawn the two threads above, resolve, then let the implicit + ;; thread exit + (func (export "setup") (result i32) + (global.set $worker-thread (call $thread.new-indirect (i32.const 0) (i32.const 0))) + (global.set $worker-bad-thread (call $thread.new-indirect (i32.const 1) (i32.const 0))) + (call $task.return-u8 (i32.const 1)) + (i32.const 0 (; EXIT ;))) + + ;; task B: publish a handle for itself and switch to $worker-thread, + ;; which returns 42 on B's behalf + (func (export "join") (result i32) + (global.set $join-implicit (call $thread.index)) + (global.set $join-task (call $thread.get-task)) + (drop (call $thread.suspend-then-resume (global.get $worker-thread))) + (call $task.drop (global.get $join-task)) + (i32.const 0 (; EXIT ;))) + + ;; task B2: switch to $worker-bad-thread, which joins B2 and then traps + ;; returning with the wrong task.return + (func (export "join-bad") (result i32) + (global.set $join-implicit (call $thread.index)) + (global.set $join-task (call $thread.get-task)) + (drop (call $thread.suspend-then-resume (global.get $worker-bad-thread))) + unreachable) + + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + ) + (core type $start-func-ty (func (param i32))) + (alias core export $table "__indirect_function_table" (core table $indirect-function-table)) + (core func $thread.new-indirect + (canon thread.new-indirect $start-func-ty (core table $indirect-function-table))) + (canon task.return (result u8) (core func $task.return-u8)) + (canon task.return (result u32) (core func $task.return-u32)) + (canon thread.index (core func $thread.index)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (canon task.drop (core func $task.drop)) + (canon thread.resume-later (core func $thread.resume-later)) + (canon thread.suspend-then-resume (core func $thread.suspend-then-resume)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return-u8" (func $task.return-u8)) + (export "task.return-u32" (func $task.return-u32)) + (export "thread.new-indirect" (func $thread.new-indirect)) + (export "thread.index" (func $thread.index)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + (export "task.drop" (func $task.drop)) + (export "thread.resume-later" (func $thread.resume-later)) + (export "thread.suspend-then-resume" (func $thread.suspend-then-resume)) + (export "__indirect_function_table" (table $indirect-function-table)) + )))) + (func (export "setup") async (result u8) + (canon lift (core func $core "setup") async (callback (core func $core "never")))) + (func (export "join") async (result u32) + (canon lift (core func $core "join") async (callback (core func $core "never")))) + (func (export "join-bad") async (result u32) + (canon lift (core func $core "join-bad") async (callback (core func $core "never")))) + ) + (instance $c (instantiate $C)) + (func (export "setup") (alias export $c "setup")) + (func (export "join") (alias export $c "join")) + (func (export "join-bad") (alias export $c "join-bad")) +) +(component instance $j1 $Join) +(assert_return (invoke "setup") (u8.const 1)) +(assert_return (invoke "join") (u32.const 42)) +(component instance $j2 $Join) +(assert_return (invoke "setup") (u8.const 1)) +(assert_trap (invoke "join-bad") "wasm trap: invalid `task.return` signature and/or options for current task") + +;; Any thread can leave an unresolved task, even the last one: the task is +;; simply left unresolved until (and unless) some thread rejoins it and +;; resolves it. A handle saved before the move is what brings a thread back. +(component + (component $C + (core module $Table + (table (export "__indirect_function_table") 2 funcref)) + (core instance $table (instantiate $Table)) + (core module $Core + (import "" "task.return-u8" (func $task.return-u8 (param i32))) + (import "" "task.return-u32" (func $task.return-u32 (param i32))) + (import "" "thread.new-indirect" (func $thread.new-indirect (param i32 i32) (result i32))) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + (import "" "thread.resume-later" (func $thread.resume-later (param i32))) + (import "" "thread.suspend-then-resume" (func $thread.suspend-then-resume (param i32) (result i32))) + (import "" "__indirect_function_table" (table $indirect-function-table 2 funcref)) + + (global $last-thread (mut i32) (i32.const 0xdead)) ;; second thread of task A + (global $mover-thread (mut i32) (i32.const 0xdead)) ;; thread that moves A -> B1 -> B2 -> A + (global $b1-task (mut i32) (i32.const 0xdead)) ;; handle for resolved task B1 + (global $b2-task (mut i32) (i32.const 0xdead)) ;; handle for resolved task B2 + + ;; $mover-thread: hops from its original task A to B1 and on to B2, + ;; parking there so that $last-thread has somewhere to move to as well. + ;; Resumed by $last-thread, it uses the handle saved on entry to return + ;; to A (not to the most-recently-left B1, whose u8-typed, already + ;; resolved task would trap the u32-typed task.return) and resolves A + ;; even though A contained no threads at all for a while. + (func $mover (param i32) + (local $home i32) + (local.set $home (call $thread.get-task)) + (call $thread.set-task (global.get $b1-task)) + (call $thread.set-task (global.get $b2-task)) + (drop (call $thread.suspend-then-resume (global.get $last-thread))) + (call $thread.set-task (local.get $home)) + (call $task.return-u32 (i32.const 42))) + + ;; $last-thread: after $mover-thread has left, the only thread of the + ;; unresolved task A. Joining B2 leaves A with no threads at all, which + ;; does not trap: an unresolved task is simply not required to ever + ;; resolve. + (func $last (param i32) + (call $thread.set-task (global.get $b2-task)) + (call $thread.resume-later (global.get $mover-thread))) + + (elem (table $indirect-function-table) (i32.const 0) func $mover $last) + + ;; tasks B1 and B2: publish a handle for themselves, resolve, then let + ;; the implicit thread exit + (func (export "victim1") (result i32) + (global.set $b1-task (call $thread.get-task)) + (call $task.return-u8 (i32.const 1)) + (i32.const 0 (; EXIT ;))) + (func (export "victim2") (result i32) + (global.set $b2-task (call $thread.get-task)) + (call $task.return-u8 (i32.const 1)) + (i32.const 0 (; EXIT ;))) + + ;; task A: spawn $last-thread and $mover-thread and exit the implicit + ;; thread *without* returning a value, leaving A unresolved with two + ;; threads + (func (export "unresolved") (result i32) + (global.set $last-thread (call $thread.new-indirect (i32.const 1) (i32.const 0))) + (global.set $mover-thread (call $thread.new-indirect (i32.const 0) (i32.const 0))) + (call $thread.resume-later (global.get $mover-thread)) + (i32.const 0 (; EXIT ;))) + + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + ) + (core type $start-func-ty (func (param i32))) + (alias core export $table "__indirect_function_table" (core table $indirect-function-table)) + (core func $thread.new-indirect + (canon thread.new-indirect $start-func-ty (core table $indirect-function-table))) + (canon task.return (result u8) (core func $task.return-u8)) + (canon task.return (result u32) (core func $task.return-u32)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (canon thread.resume-later (core func $thread.resume-later)) + (canon thread.suspend-then-resume (core func $thread.suspend-then-resume)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return-u8" (func $task.return-u8)) + (export "task.return-u32" (func $task.return-u32)) + (export "thread.new-indirect" (func $thread.new-indirect)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + (export "thread.resume-later" (func $thread.resume-later)) + (export "thread.suspend-then-resume" (func $thread.suspend-then-resume)) + (export "__indirect_function_table" (table $indirect-function-table)) + )))) + (func (export "victim1") async (result u8) + (canon lift (core func $core "victim1") async (callback (core func $core "never")))) + (func (export "victim2") async (result u8) + (canon lift (core func $core "victim2") async (callback (core func $core "never")))) + (func (export "unresolved") async (result u32) + (canon lift (core func $core "unresolved") async (callback (core func $core "never")))) + ) + (instance $c (instantiate $C)) + (func (export "victim1") (alias export $c "victim1")) + (func (export "victim2") (alias export $c "victim2")) + (func (export "unresolved") (alias export $c "unresolved")) +) +(assert_return (invoke "victim1") (u8.const 1)) +(assert_return (invoke "victim2") (u8.const 1)) +(assert_return (invoke "unresolved") (u32.const 42)) + +;; A thread spawned by a moved thread inherits the spawner's new task, and a +;; task's implicit thread can exit while the task is unresolved, leaving adopted +;; and spawned threads to resolve it afterwards. +(component + (component $C + (core module $Table + (table (export "__indirect_function_table") 2 funcref)) + (core instance $table (instantiate $Table)) + (core module $Core + (import "" "task.return" (func $task.return (param i32))) + (import "" "thread.new-indirect" (func $thread.new-indirect (param i32 i32) (result i32))) + (import "" "thread.index" (func $thread.index (result i32))) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + (import "" "thread.resume-later" (func $thread.resume-later (param i32))) + (import "" "thread.suspend" (func $thread.suspend (result i32))) + (import "" "thread.suspend-then-resume" (func $thread.suspend-then-resume (param i32) (result i32))) + (import "" "__indirect_function_table" (table $indirect-function-table 2 funcref)) + + (global $worker-thread (mut i32) (i32.const 0xdead)) ;; thread that moves from task A to task B + (global $child-thread (mut i32) (i32.const 0xdead)) ;; spawned by $worker-thread after its move + (global $run-implicit (mut i32) (i32.const 0xdead)) ;; implicit thread of task B + (global $run-task (mut i32) (i32.const 0xdead)) ;; handle for task B + + ;; $worker-thread: joins task B and spawns $child-thread, which inherits + ;; the spawner's *current* task (B). It then parks and, resumed by + ;; $child-thread after B has resolved, exits as the last thread of B. + (func $worker (param i32) + (call $thread.set-task (global.get $run-task)) + (global.set $child-thread (call $thread.new-indirect (i32.const 1) (i32.const 0))) + (call $thread.resume-later (global.get $run-implicit)) + (drop (call $thread.suspend))) + + ;; $child-thread: returns for the task it was spawned into. Task A has + ;; already resolved, so this only succeeds if $child-thread inherited B + ;; from $worker-thread's post-move task. + (func $child (param i32) + (call $task.return (i32.const 42)) + (call $thread.resume-later (global.get $worker-thread))) + + (elem (table $indirect-function-table) (i32.const 0) func $worker $child) + + ;; task A: spawn $worker-thread, resolve, then let the implicit thread + ;; exit + (func (export "setup") (result i32) + (global.set $worker-thread (call $thread.new-indirect (i32.const 0) (i32.const 0))) + (call $task.return (i32.const 1)) + (i32.const 0 (; EXIT ;))) + + ;; task B: switch to $worker-thread, then exit the implicit thread while + ;; B is still unresolved, leaving B to retain $worker-thread and + ;; $child-thread; $child-thread then returns 42 on B's behalf + (func (export "run") (result i32) + (global.set $run-implicit (call $thread.index)) + (global.set $run-task (call $thread.get-task)) + (drop (call $thread.suspend-then-resume (global.get $worker-thread))) + (call $thread.resume-later (global.get $child-thread)) + (i32.const 0 (; EXIT ;))) + + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + ) + (core type $start-func-ty (func (param i32))) + (alias core export $table "__indirect_function_table" (core table $indirect-function-table)) + (core func $thread.new-indirect + (canon thread.new-indirect $start-func-ty (core table $indirect-function-table))) + (canon task.return (result u32) (core func $task.return)) + (canon thread.index (core func $thread.index)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (canon thread.resume-later (core func $thread.resume-later)) + (canon thread.suspend (core func $thread.suspend)) + (canon thread.suspend-then-resume (core func $thread.suspend-then-resume)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return" (func $task.return)) + (export "thread.new-indirect" (func $thread.new-indirect)) + (export "thread.index" (func $thread.index)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + (export "thread.resume-later" (func $thread.resume-later)) + (export "thread.suspend" (func $thread.suspend)) + (export "thread.suspend-then-resume" (func $thread.suspend-then-resume)) + (export "__indirect_function_table" (table $indirect-function-table)) + )))) + (func (export "setup") async (result u32) + (canon lift (core func $core "setup") async (callback (core func $core "never")))) + (func (export "run") async (result u32) + (canon lift (core func $core "run") async (callback (core func $core "never")))) + ) + (instance $c (instantiate $C)) + (func (export "setup") (alias export $c "setup")) + (func (export "run") (alias export $c "run")) +) +(assert_return (invoke "setup") (u32.const 1)) +(assert_return (invoke "run") (u32.const 42)) + +;; A task can be resolved by an implicit thread that belongs to another task, +;; and task.return's result type is checked against the thread's *current* +;; task. These three implicit threads form a relay: each leaves its own task +;; started, unresolved and with no threads at all, and the next one joins that +;; task, returns its value and exits while still contained by it. The result +;; types alternate, so every task.return here would trap if it were checked +;; against the task whose export the calling thread entered; "hop3" performs +;; both a u8 return for its own task and a u32 return for the task it joined. +(component + (component $C + (core module $Core + (import "" "task.return-u8" (func $task.return-u8 (param i32))) + (import "" "task.return-u32" (func $task.return-u32 (param i32))) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + + (global $a-task (mut i32) (i32.const 0xdead)) ;; handle for task A (u8 result) + (global $b-task (mut i32) (i32.const 0xdead)) ;; handle for task B (u32 result) + + ;; task A: publish a handle for itself and leave at once, so that A is + ;; started, unresolved and contains no threads at all + (func (export "hop1") (result i32) + (global.set $a-task (call $thread.get-task)) + (i32.const 0 (; EXIT ;))) + + ;; task B: publish a handle for itself, then join task A and return A's + ;; u8-typed value (which would trap against B's own u32 result) before + ;; exiting while still contained by A, leaving B threadless in turn + (func (export "hop2") (result i32) + (global.set $b-task (call $thread.get-task)) + (call $thread.set-task (global.get $a-task)) + (call $task.return-u8 (i32.const 11)) + (i32.const 0 (; EXIT ;))) + + ;; task C: return C's own u8-typed value, then join task B and return its + ;; u32-typed value; the two calls differ only in which task this thread + ;; is in at the time + (func (export "hop3") (result i32) + (call $task.return-u8 (i32.const 33)) + (call $thread.set-task (global.get $b-task)) + (call $task.return-u32 (i32.const 222)) + (i32.const 0 (; EXIT ;))) + + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + ) + (canon task.return (result u8) (core func $task.return-u8)) + (canon task.return (result u32) (core func $task.return-u32)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return-u8" (func $task.return-u8)) + (export "task.return-u32" (func $task.return-u32)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + )))) + (func (export "hop1") async (result u8) + (canon lift (core func $core "hop1") async (callback (core func $core "never")))) + (func (export "hop2") async (result u32) + (canon lift (core func $core "hop2") async (callback (core func $core "never")))) + (func (export "hop3") async (result u8) + (canon lift (core func $core "hop3") async (callback (core func $core "never")))) + ) + (component $D + (import "hop1" (func $hop1 async (result u8))) + (import "hop2" (func $hop2 async (result u32))) + (import "hop3" (func $hop3 async (result u8))) + + (core module $Memory (memory (export "mem") 1)) + (core instance $memory (instantiate $Memory)) + (core module $Core + (import "" "mem" (memory 1)) + (import "" "subtask.drop" (func $subtask.drop (param i32))) + (import "" "waitable.join" (func $waitable.join (param i32 i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + (import "" "waitable-set.poll" (func $waitable-set.poll (param i32 i32) (result i32))) + (import "" "hop1" (func $hop1 (param i32) (result i32))) + (import "" "hop2" (func $hop2 (param i32) (result i32))) + (import "" "hop3" (func $hop3 (param i32) (result i32))) + + (func (export "run") (result i32) + (local $ret i32) + (local $a-subtask i32) + (local $b-subtask i32) + (local $ws i32) + (local $n i32) + + ;; start hop1, which leaves task A behind unresolved + (local.set $ret (call $hop1 (i32.const 0 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $a-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + ;; start hop2: within this call task A is resolved with 11, while task + ;; B is left behind unresolved in the same way + (local.set $ret (call $hop2 (i32.const 4 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $b-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + ;; start hop3, which resolves task B with 222 as well as its own task + ;; with 33, so this call completes eagerly + (if (i32.ne (i32.const 2 (; RETURNED ;)) (call $hop3 (i32.const 8 (; retp ;)))) + (then unreachable)) + + ;; each abandoned task was resolved by the thread that came after it, + ;; so both resolutions are already pending and a poll collects them + (local.set $ws (call $waitable-set.new)) + (call $waitable.join (local.get $a-subtask) (local.get $ws)) + (call $waitable.join (local.get $b-subtask) (local.get $ws)) + (loop $l + (if (i32.ne (i32.const 1 (; SUBTASK ;)) + (call $waitable-set.poll (local.get $ws) (i32.const 16 (; eventp ;)))) + (then unreachable)) + (if (i32.ne (i32.const 2 (; RETURNED ;)) (i32.load offset=4 (i32.const 16))) + (then unreachable)) + (call $subtask.drop (i32.load (i32.const 16))) + (local.set $n (i32.add (local.get $n) (i32.const 1))) + (br_if $l (i32.lt_u (local.get $n) (i32.const 2)))) + + (if (i32.ne (i32.const 11) (i32.load8_u (i32.const 0))) + (then unreachable)) + (if (i32.ne (i32.const 222) (i32.load (i32.const 4))) + (then unreachable)) + (if (i32.ne (i32.const 33) (i32.load8_u (i32.const 8))) + (then unreachable)) + + (i32.const 42)) + ) + (canon subtask.drop (core func $subtask.drop)) + (canon waitable.join (core func $waitable.join)) + (canon waitable-set.new (core func $waitable-set.new)) + (canon waitable-set.poll (memory (core memory $memory "mem")) (core func $waitable-set.poll)) + (canon lower (func $hop1) async (memory (core memory $memory "mem")) (core func $hop1')) + (canon lower (func $hop2) async (memory (core memory $memory "mem")) (core func $hop2')) + (canon lower (func $hop3) async (memory (core memory $memory "mem")) (core func $hop3')) + (core instance $core (instantiate $Core (with "" (instance + (export "mem" (memory $memory "mem")) + (export "subtask.drop" (func $subtask.drop)) + (export "waitable.join" (func $waitable.join)) + (export "waitable-set.new" (func $waitable-set.new)) + (export "waitable-set.poll" (func $waitable-set.poll)) + (export "hop1" (func $hop1')) + (export "hop2" (func $hop2')) + (export "hop3" (func $hop3')) + )))) + (func (export "run") async (result u32) (canon lift (core func $core "run"))) + ) + (instance $c (instantiate $C)) + (instance $d (instantiate $D + (with "hop1" (func $c "hop1")) + (with "hop2" (func $c "hop2")) + (with "hop3" (func $c "hop3")))) + (func (export "run") (alias export $d "run")) +) +(assert_return (invoke "run") (u32.const 42)) + +;; The implicit thread of a synchronously-lifted export may change tasks while +;; the core function executes and need not move back: on return it implicitly +;; rejoins the task it was spawned in, which is the task the lifted results are +;; returned for. Task A is u8-typed while all four exports below are u32-typed, +;; so these results can only come out right if they are lifted with the lift's +;; own options and result type rather than the current task's. +(component definition $VisitOrStay + (component $C + (core module $Core + (import "" "task.return" (func $task.return (param i32))) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + + (global $a-task (mut i32) (i32.const 0xdead)) ;; handle for resolved task A + + (func (export "setup") (result i32) + (global.set $a-task (call $thread.get-task)) + (call $task.return (i32.const 1)) + (i32.const 0 (; EXIT ;))) + + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + + (func (export "sync-visit") (result i32) + (local $home i32) + (local.set $home (call $thread.get-task)) + (call $thread.set-task (global.get $a-task)) + (call $thread.set-task (local.get $home)) + (i32.const 42)) + + ;; returns while still contained by task A + (func (export "sync-stay") (result i32) + (call $thread.set-task (global.get $a-task)) + (i32.const 33)) + + (func (export "async-visit") (result i32) + (local $home i32) + (local.set $home (call $thread.get-task)) + (call $thread.set-task (global.get $a-task)) + (call $thread.set-task (local.get $home)) + (i32.const 44)) + + ;; likewise, and `async`-typed rather than plain + (func (export "async-stay") (result i32) + (call $thread.set-task (global.get $a-task)) + (i32.const 55)) + ) + (canon task.return (result u8) (core func $task.return)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return" (func $task.return)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + )))) + (func (export "setup") async (result u8) + (canon lift (core func $core "setup") async (callback (core func $core "never")))) + (func (export "sync-visit") (result u32) + (canon lift (core func $core "sync-visit"))) + (func (export "sync-stay") (result u32) + (canon lift (core func $core "sync-stay"))) + (func (export "async-visit") async (result u32) + (canon lift (core func $core "async-visit"))) + (func (export "async-stay") async (result u32) + (canon lift (core func $core "async-stay"))) + ) + (instance $c (instantiate $C)) + (func (export "setup") (alias export $c "setup")) + (func (export "sync-visit") (alias export $c "sync-visit")) + (func (export "sync-stay") (alias export $c "sync-stay")) + (func (export "async-visit") (alias export $c "async-visit")) + (func (export "async-stay") (alias export $c "async-stay")) +) +(component instance $visit-or-stay $VisitOrStay) +(assert_return (invoke "setup") (u8.const 1)) +(assert_return (invoke "sync-visit") (u32.const 42)) +(assert_return (invoke "sync-stay") (u32.const 33)) +(assert_return (invoke "async-visit") (u32.const 44)) +(assert_return (invoke "async-stay") (u32.const 55)) + +;; Test switching between sync-typed and async-typed tasks. +(component definition $SyncTasks + (component $C + (core module $Core + (import "" "task.return" (func $task.return (param i32))) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + + (global $target-task (mut i32) (i32.const 0xdead)) ;; callback-lifted task + (global $sync-task (mut i32) (i32.const 0xdead)) ;; sync-lifted, not async-typed + (global $async-sync-task (mut i32) (i32.const 0xdead)) ;; sync-lifted, async-typed + + ;; callback-lifted, so this task's value can be returned by any thread + ;; that joins it; leaving the event loop keeps it unresolved + (func (export "target") (result i32) + (global.set $target-task (call $thread.get-task)) + (i32.const 0 (; EXIT ;))) + + ;; not async-typed, and so sync-lifted: this thread returns 99 for the + ;; task it joined and 42 for its own task, the latter by returning + (func (export "from-sync") (result i32) + (local $home i32) + (local.set $home (call $thread.get-task)) + (global.set $sync-task (call $thread.get-task)) + (call $thread.set-task (global.get $target-task)) + (call $task.return (i32.const 99)) + (call $thread.set-task (local.get $home)) + (i32.const 42)) + + ;; async-typed but still lifted with the sync ABI, so task.return is + ;; equally unavailable for this task + (func (export "async-sync") (result i32) + (global.set $async-sync-task (call $thread.get-task)) + (i32.const 7)) + + ;; both of these join a sync-lifted task and then try to return its + ;; value; everything but the current task's lift ABI is in order, since + ;; the same task.return resolves "target" above + (func (export "to-sync") (result i32) + (call $thread.set-task (global.get $sync-task)) + (call $task.return (i32.const 11)) + unreachable) + (func (export "to-async-sync") (result i32) + (call $thread.set-task (global.get $async-sync-task)) + (call $task.return (i32.const 11)) + unreachable) + + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + ) + (canon task.return (result u32) (core func $task.return)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return" (func $task.return)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + )))) + (func (export "target") async (result u32) + (canon lift (core func $core "target") async (callback (core func $core "never")))) + (func (export "from-sync") (result u32) + (canon lift (core func $core "from-sync"))) + (func (export "async-sync") async (result u32) + (canon lift (core func $core "async-sync"))) + (func (export "to-sync") async (result u32) + (canon lift (core func $core "to-sync") async (callback (core func $core "never")))) + (func (export "to-async-sync") async (result u32) + (canon lift (core func $core "to-async-sync") async (callback (core func $core "never")))) + ) + (component $D + (import "target" (func $target async (result u32))) + (import "from-sync" (func $from-sync (result u32))) + (import "async-sync" (func $async-sync async (result u32))) + + (core module $Memory (memory (export "mem") 1)) + (core instance $memory (instantiate $Memory)) + (core module $Core + (import "" "mem" (memory 1)) + (import "" "subtask.drop" (func $subtask.drop (param i32))) + (import "" "waitable.join" (func $waitable.join (param i32 i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + (import "" "waitable-set.poll" (func $waitable-set.poll (param i32 i32) (result i32))) + (import "" "target" (func $target (param i32) (result i32))) + (import "" "from-sync" (func $from-sync (result i32))) + (import "" "async-sync" (func $async-sync (result i32))) + + (func (export "run") (result i32) + (local $ret i32) + (local $target-subtask i32) + (local $ws i32) + + ;; start target, whose implicit thread leaves at once + (local.set $ret (call $target (i32.const 0 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $target-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + ;; the sync-lifted task's thread returns target's value on its way + ;; through, and its own by returning + (if (i32.ne (i32.const 42) (call $from-sync)) + (then unreachable)) + + ;; that task.return ran inside the call above, so target's resolution + ;; is already pending and a poll is enough to collect it + (local.set $ws (call $waitable-set.new)) + (call $waitable.join (local.get $target-subtask) (local.get $ws)) + (if (i32.ne (i32.const 1 (; SUBTASK ;)) + (call $waitable-set.poll (local.get $ws) (i32.const 16 (; eventp ;)))) + (then unreachable)) + (if (i32.ne (local.get $target-subtask) (i32.load (i32.const 16))) + (then unreachable)) + (if (i32.ne (i32.const 2 (; RETURNED ;)) (i32.load offset=4 (i32.const 16))) + (then unreachable)) + (if (i32.ne (i32.const 99) (i32.load (i32.const 0))) + (then unreachable)) + (call $subtask.drop (local.get $target-subtask)) + + ;; publish the async-typed, sync-lifted task as well + (if (i32.ne (i32.const 7) (call $async-sync)) + (then unreachable)) + + (i32.const 42)) + ) + (canon subtask.drop (core func $subtask.drop)) + (canon waitable.join (core func $waitable.join)) + (canon waitable-set.new (core func $waitable-set.new)) + (canon waitable-set.poll (memory (core memory $memory "mem")) (core func $waitable-set.poll)) + (canon lower (func $target) async (memory (core memory $memory "mem")) (core func $target')) + (canon lower (func $from-sync) (core func $from-sync')) + (canon lower (func $async-sync) (core func $async-sync')) + (core instance $core (instantiate $Core (with "" (instance + (export "mem" (memory $memory "mem")) + (export "subtask.drop" (func $subtask.drop)) + (export "waitable.join" (func $waitable.join)) + (export "waitable-set.new" (func $waitable-set.new)) + (export "waitable-set.poll" (func $waitable-set.poll)) + (export "target" (func $target')) + (export "from-sync" (func $from-sync')) + (export "async-sync" (func $async-sync')) + )))) + (func (export "run") async (result u32) (canon lift (core func $core "run"))) + ) + (instance $c (instantiate $C)) + (instance $d (instantiate $D + (with "target" (func $c "target")) + (with "from-sync" (func $c "from-sync")) + (with "async-sync" (func $c "async-sync")))) + (func (export "run") (alias export $d "run")) + (func (export "to-sync") (alias export $c "to-sync")) + (func (export "to-async-sync") (alias export $c "to-async-sync")) +) +(component instance $sync-tasks1 $SyncTasks) +(assert_return (invoke "run") (u32.const 42)) +(assert_trap (invoke "to-sync") "wasm trap: invalid `task.return` signature and/or options for current task") +(component instance $sync-tasks2 $SyncTasks) +(assert_return (invoke "run") (u32.const 42)) +(assert_trap (invoke "to-async-sync") "wasm trap: invalid `task.return` signature and/or options for current task") + +;; The implicit thread of an async callback-lifted export can change tasks too: +;; parked in its event loop after joining task B, it is not a target for its +;; own task's cancellation request (which stays pending), but receives task B's +;; cancellation as a TASK_CANCELLED event, acknowledges it on B's behalf, then +;; moves back home and resolves its own task from inside the callback before +;; exiting the event loop. +(component + (component $C + (core module $Core + (import "" "task.return" (func $task.return (param i32))) + (import "" "task.cancel" (func $task.cancel)) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + (import "" "task.drop" (func $task.drop (param i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + + (global $target-task (mut i32) (i32.const 0xdead)) ;; handle for task B + (global $cbmove-task (mut i32) (i32.const 0xdead)) ;; handle for the cbmove task + + ;; task B: publish a handle for itself and leave the event loop at once, + ;; so that B has no thread of its own and the callback thread that joins + ;; it is its only cancellable thread + (func (export "target") (result i32) + (global.set $target-task (call $thread.get-task)) + (i32.const 0 (; EXIT ;))) + + ;; the cbmove task's implicit thread: join task B, then wait on an empty + ;; waitable set in the event loop, which releases the exclusive lock and + ;; leaves this thread parked as task B's only cancellable thread + (func (export "cbmove") (result i32) + (global.set $cbmove-task (call $thread.get-task)) + (call $thread.set-task (global.get $target-task)) + (i32.or (i32.const 2 (; WAIT ;)) + (i32.shl (call $waitable-set.new) (i32.const 4)))) + + ;; the only event this thread can receive is task B's cancellation + (func (export "cbmove-cb") (param i32 i32 i32) (result i32) + (if (i32.ne (i32.const 6 (; TASK_CANCELLED ;)) (local.get 0)) + (then unreachable)) + (call $task.cancel) ;; acknowledge on B's behalf + (call $thread.set-task (global.get $cbmove-task)) ;; move back home + (call $task.drop (global.get $target-task)) + (call $task.drop (global.get $cbmove-task)) + (call $task.return (i32.const 42)) ;; resolve the cbmove task + (i32.const 0 (; EXIT ;))) + + ;; "target" leaves its event loop before any event can arrive + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + ) + (canon task.return (result u32) (core func $task.return)) + (canon task.cancel (core func $task.cancel)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (canon task.drop (core func $task.drop)) + (canon waitable-set.new (core func $waitable-set.new)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return" (func $task.return)) + (export "task.cancel" (func $task.cancel)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + (export "task.drop" (func $task.drop)) + (export "waitable-set.new" (func $waitable-set.new)) + )))) + (func (export "target") async (result u32) + (canon lift (core func $core "target") async (callback (core func $core "never")))) + (func (export "cbmove") async (result u32) + (canon lift (core func $core "cbmove") + async (callback (core func $core "cbmove-cb")))) + ) + (component $D + (import "target" (func $target async (result u32))) + (import "cbmove" (func $cbmove async (result u32))) + + (core module $Memory (memory (export "mem") 1)) + (core instance $memory (instantiate $Memory)) + (core module $Core + (import "" "mem" (memory 1)) + (import "" "subtask.cancel" (func $subtask.cancel (param i32) (result i32))) + (import "" "subtask.drop" (func $subtask.drop (param i32))) + (import "" "waitable.join" (func $waitable.join (param i32 i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + (import "" "waitable-set.poll" (func $waitable-set.poll (param i32 i32) (result i32))) + (import "" "target" (func $target (param i32) (result i32))) + (import "" "cbmove" (func $cbmove (param i32) (result i32))) + + (func (export "run") (result i32) + (local $ret i32) + (local $target-subtask i32) + (local $cbmove-subtask i32) + (local $ws i32) + + ;; start target, whose implicit thread leaves at once + (local.set $ret (call $target (i32.const 4 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $target-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + ;; start cbmove, whose implicit thread joins target's task and parks + ;; in its event loop + (local.set $ret (call $cbmove (i32.const 8 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $cbmove-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + ;; the callback thread is parked, but no longer in its own task, so + ;; this request can only be remembered as pending + (local.set $ret (call $subtask.cancel (local.get $cbmove-subtask))) + (if (i32.ne (i32.const -1 (; BLOCKED ;)) (local.get $ret)) + (then unreachable)) + + ;; the request against the joined task is delivered to the callback + ;; thread as a TASK_CANCELLED event and acknowledged there + (local.set $ret (call $subtask.cancel (local.get $target-subtask))) + (if (i32.ne (i32.const 4 (; CANCELLED_BEFORE_RETURNED ;)) (local.get $ret)) + (then unreachable)) + (call $subtask.drop (local.get $target-subtask)) + + ;; the callback moved home and resolved its own task with 42 inside the + ;; cancel above, so that resolution is already pending + (local.set $ws (call $waitable-set.new)) + (call $waitable.join (local.get $cbmove-subtask) (local.get $ws)) + (local.set $ret (call $waitable-set.poll (local.get $ws) (i32.const 16 (; eventp ;)))) + (if (i32.ne (i32.const 1 (; SUBTASK ;)) (local.get $ret)) + (then unreachable)) + (if (i32.ne (local.get $cbmove-subtask) (i32.load (i32.const 16))) + (then unreachable)) + (if (i32.ne (i32.const 2 (; RETURNED ;)) (i32.load offset=4 (i32.const 16))) + (then unreachable)) + (if (i32.ne (i32.const 42) (i32.load (i32.const 8))) + (then unreachable)) + (call $subtask.drop (local.get $cbmove-subtask)) + + (i32.const 42)) + ) + (canon subtask.cancel async (core func $subtask.cancel)) + (canon subtask.drop (core func $subtask.drop)) + (canon waitable.join (core func $waitable.join)) + (canon waitable-set.new (core func $waitable-set.new)) + (canon waitable-set.poll (memory (core memory $memory "mem")) (core func $waitable-set.poll)) + (canon lower (func $target) async (memory (core memory $memory "mem")) (core func $target')) + (canon lower (func $cbmove) async (memory (core memory $memory "mem")) (core func $cbmove')) + (core instance $core (instantiate $Core (with "" (instance + (export "mem" (memory $memory "mem")) + (export "subtask.cancel" (func $subtask.cancel)) + (export "subtask.drop" (func $subtask.drop)) + (export "waitable.join" (func $waitable.join)) + (export "waitable-set.new" (func $waitable-set.new)) + (export "waitable-set.poll" (func $waitable-set.poll)) + (export "target" (func $target')) + (export "cbmove" (func $cbmove')) + )))) + (func (export "run") async (result u32) (canon lift (core func $core "run"))) + ) + (instance $c (instantiate $C)) + (instance $d (instantiate $D + (with "target" (func $c "target")) + (with "cbmove" (func $c "cbmove")))) + (func (export "run") (alias export $d "run")) +) +(assert_return (invoke "run") (u32.const 42)) + +;; A cancellation request is delivered to exactly one of the cancelled task's +;; threads, chosen nondeterministically among those parked in a callback event +;; loop. thread.set-task makes it possible for more than one thread to +;; qualify, since a thread that joins a task becomes a candidate for that +;; task's cancellation just as it stops being one for its own. +(component definition $TwoCandidates + (component $C + (core module $Core + (import "" "task.return" (func $task.return (param i32))) + (import "" "task.cancel" (func $task.cancel)) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + + (global $joined-task (mut i32) (i32.const 0xdead)) ;; handle for the task being cancelled + (global $abandoned-task (mut i32) (i32.const 0)) ;; handle for a task left without threads + (global $deliveries (mut i32) (i32.const 0)) ;; TASK_CANCELLED events received + + ;; this task's own implicit thread leaves the event loop at once, so the + ;; task has no cancellable thread of its own and only joining threads are + ;; candidates for its cancellation + (func (export "target") (result i32) + (global.set $joined-task (call $thread.get-task)) + (i32.const 0 (; EXIT ;))) + + ;; callback-lifted: this task's own implicit thread parks in its event + ;; loop and so is a candidate alongside any thread that joins it + (func (export "owner") (result i32) + (global.set $joined-task (call $thread.get-task)) + (i32.or (i32.const 2 (; WAIT ;)) + (i32.shl (call $waitable-set.new) (i32.const 4)))) + + ;; callback-lifted: resolve this thread's own task before joining the task + ;; published above, then park in the event loop as one of its candidates + (func (export "joiner") (result i32) + (call $task.return (i32.const 42)) + (call $thread.set-task (global.get $joined-task)) + (i32.or (i32.const 2 (; WAIT ;)) + (i32.shl (call $waitable-set.new) (i32.const 4)))) + + ;; callback-lifted: keep a handle for this thread's own task, join the + ;; task published above and leave the event loop at once. That abandons + ;; this thread's own task, started and unresolved with no threads at all; + ;; the saved handle is all that is left to reach it by. + (func (export "exiter") (result i32) + (global.set $abandoned-task (call $thread.get-task)) + (call $thread.set-task (global.get $joined-task)) + (i32.const 0 (; EXIT ;))) + + ;; shared by every export above, though "target" and "exiter" leave + ;; before any event can reach them: the only event that can arrive is the + ;; cancellation of the task this thread is currently in, which it + ;; acknowledges on that task's behalf. An abandoned task has no thread of + ;; its own left to resolve it, so this thread returns for it through the + ;; saved handle on its way out. + (func (export "cb") (param i32 i32 i32) (result i32) + (if (i32.ne (i32.const 6 (; TASK_CANCELLED ;)) (local.get 0)) + (then unreachable)) + (global.set $deliveries (i32.add (global.get $deliveries) (i32.const 1))) + (call $task.cancel) + (if (global.get $abandoned-task) ;; index 0 is never a task handle + (then + (call $thread.set-task (global.get $abandoned-task)) + (call $task.return (i32.const 42)))) + (i32.const 0 (; EXIT ;))) + + (func (export "deliveries") (result i32) + (global.get $deliveries)) + ) + (canon task.return (result u32) (core func $task.return)) + (canon task.cancel (core func $task.cancel)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (canon waitable-set.new (core func $waitable-set.new)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return" (func $task.return)) + (export "task.cancel" (func $task.cancel)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + (export "waitable-set.new" (func $waitable-set.new)) + )))) + (func (export "target") async (result u32) + (canon lift (core func $core "target") async (callback (core func $core "cb")))) + (func (export "owner") async (result u32) + (canon lift (core func $core "owner") async (callback (core func $core "cb")))) + (func (export "joiner") async (result u32) + (canon lift (core func $core "joiner") async (callback (core func $core "cb")))) + (func (export "exiter") async (result u32) + (canon lift (core func $core "exiter") async (callback (core func $core "cb")))) + (func (export "deliveries") (result u32) + (canon lift (core func $core "deliveries"))) + ) + (component $D + (import "target" (func $target async (result u32))) + (import "owner" (func $owner async (result u32))) + (import "joiner" (func $joiner async (result u32))) + (import "exiter" (func $exiter async (result u32))) + (import "deliveries" (func $deliveries (result u32))) + + (core module $Memory (memory (export "mem") 1)) + (core instance $memory (instantiate $Memory)) + (core module $Core + (import "" "mem" (memory 1)) + (import "" "subtask.cancel" (func $subtask.cancel (param i32) (result i32))) + (import "" "subtask.drop" (func $subtask.drop (param i32))) + (import "" "waitable.join" (func $waitable.join (param i32 i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + (import "" "waitable-set.poll" (func $waitable-set.poll (param i32 i32) (result i32))) + (import "" "target" (func $target (param i32) (result i32))) + (import "" "owner" (func $owner (param i32) (result i32))) + (import "" "joiner" (func $joiner (param i32) (result i32))) + (import "" "exiter" (func $exiter (param i32) (result i32))) + (import "" "deliveries" (func $deliveries (result i32))) + + ;; start a joiner, which resolves eagerly with 42 and leaves its implicit + ;; thread parked in the joined task's event loop; since it has already + ;; resolved, there is no subtask index in the packed result to keep + (func $start-joiner (param $retp i32) + (if (i32.ne (i32.const 2 (; RETURNED ;)) (call $joiner (local.get $retp))) + (then unreachable)) + (if (i32.ne (i32.const 42) (i32.load (local.get $retp))) + (then unreachable))) + + ;; both candidates joined the cancelled task from elsewhere + (func (export "cancel-two-joined") (result i32) + (local $ret i32) + (local $target-subtask i32) + + ;; start target, whose implicit thread leaves at once + (local.set $ret (call $target (i32.const 0 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $target-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + (call $start-joiner (i32.const 4 (; retp ;))) + (call $start-joiner (i32.const 8 (; retp ;))) + + ;; target's task now contains two cancellable threads and neither of + ;; them is its own, so the request goes to one of the two joining + ;; threads, which acknowledges it from its event loop + (if (i32.ne (i32.const 4 (; CANCELLED_BEFORE_RETURNED ;)) + (call $subtask.cancel (local.get $target-subtask))) + (then unreachable)) + (call $subtask.drop (local.get $target-subtask)) + + ;; exactly one of the two threads received the event + (if (i32.ne (i32.const 1) (call $deliveries)) + (then unreachable)) + + (i32.const 42)) + + ;; the cancelled task's own implicit thread is a candidate too + (func (export "cancel-own-and-joined") (result i32) + (local $ret i32) + (local $owner-subtask i32) + + ;; start owner, which parks its implicit thread in its own event loop + (local.set $ret (call $owner (i32.const 0 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $owner-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + (call $start-joiner (i32.const 4 (; retp ;))) + + ;; owner's task now contains two cancellable threads: its own implicit + ;; thread and the joining one. Both acknowledge the request the same + ;; way, so the outcome does not depend on which is chosen. + (if (i32.ne (i32.const 4 (; CANCELLED_BEFORE_RETURNED ;)) + (call $subtask.cancel (local.get $owner-subtask))) + (then unreachable)) + (call $subtask.drop (local.get $owner-subtask)) + + (if (i32.ne (i32.const 1) (call $deliveries)) + (then unreachable)) + + (i32.const 42)) + + ;; the same two candidates, plus a third joining thread that leaves its + ;; event loop immediately: the chosen candidate resolves that abandoned + ;; task as well, through the handle the thread left behind + (func (export "cancel-two-joined-and-abandoned") (result i32) + (local $ret i32) + (local $target-subtask i32) + (local $exiter-subtask i32) + (local $ws i32) + + ;; start target, whose implicit thread leaves at once + (local.set $ret (call $target (i32.const 0 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $target-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + ;; start exiter, which joins target's task and exits, leaving its own + ;; task started-but-unresolved and threadless + (local.set $ret (call $exiter (i32.const 12 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (local.set $exiter-subtask (i32.shr_u (local.get $ret) (i32.const 4))) + + (call $start-joiner (i32.const 4 (; retp ;))) + (call $start-joiner (i32.const 8 (; retp ;))) + + ;; exiter's thread is gone and so is not a candidate: the request goes + ;; to one of the two threads still parked in target's task + (if (i32.ne (i32.const 4 (; CANCELLED_BEFORE_RETURNED ;)) + (call $subtask.cancel (local.get $target-subtask))) + (then unreachable)) + (call $subtask.drop (local.get $target-subtask)) + (if (i32.ne (i32.const 1) (call $deliveries)) + (then unreachable)) + + ;; whichever thread that was, it returned 42 for the abandoned task + ;; inside the cancel above, so that resolution is already pending + (local.set $ws (call $waitable-set.new)) + (call $waitable.join (local.get $exiter-subtask) (local.get $ws)) + (if (i32.ne (i32.const 1 (; SUBTASK ;)) + (call $waitable-set.poll (local.get $ws) (i32.const 16 (; eventp ;)))) + (then unreachable)) + (if (i32.ne (local.get $exiter-subtask) (i32.load (i32.const 16))) + (then unreachable)) + (if (i32.ne (i32.const 2 (; RETURNED ;)) (i32.load offset=4 (i32.const 16))) + (then unreachable)) + (if (i32.ne (i32.const 42) (i32.load (i32.const 12))) + (then unreachable)) + (call $subtask.drop (local.get $exiter-subtask)) + + (i32.const 42)) + ) + (canon subtask.cancel async (core func $subtask.cancel)) + (canon subtask.drop (core func $subtask.drop)) + (canon waitable.join (core func $waitable.join)) + (canon waitable-set.new (core func $waitable-set.new)) + (canon waitable-set.poll (memory (core memory $memory "mem")) (core func $waitable-set.poll)) + (canon lower (func $target) async (memory (core memory $memory "mem")) (core func $target')) + (canon lower (func $owner) async (memory (core memory $memory "mem")) (core func $owner')) + (canon lower (func $joiner) async (memory (core memory $memory "mem")) (core func $joiner')) + (canon lower (func $exiter) async (memory (core memory $memory "mem")) (core func $exiter')) + (canon lower (func $deliveries) (core func $deliveries')) + (core instance $core (instantiate $Core (with "" (instance + (export "mem" (memory $memory "mem")) + (export "subtask.cancel" (func $subtask.cancel)) + (export "subtask.drop" (func $subtask.drop)) + (export "waitable.join" (func $waitable.join)) + (export "waitable-set.new" (func $waitable-set.new)) + (export "waitable-set.poll" (func $waitable-set.poll)) + (export "target" (func $target')) + (export "owner" (func $owner')) + (export "joiner" (func $joiner')) + (export "exiter" (func $exiter')) + (export "deliveries" (func $deliveries')) + )))) + (func (export "cancel-two-joined") async (result u32) + (canon lift (core func $core "cancel-two-joined"))) + (func (export "cancel-own-and-joined") async (result u32) + (canon lift (core func $core "cancel-own-and-joined"))) + (func (export "cancel-two-joined-and-abandoned") async (result u32) + (canon lift (core func $core "cancel-two-joined-and-abandoned"))) + ) + (instance $c (instantiate $C)) + (instance $d (instantiate $D + (with "target" (func $c "target")) + (with "owner" (func $c "owner")) + (with "joiner" (func $c "joiner")) + (with "exiter" (func $c "exiter")) + (with "deliveries" (func $c "deliveries")))) + (func (export "cancel-two-joined") (alias export $d "cancel-two-joined")) + (func (export "cancel-own-and-joined") (alias export $d "cancel-own-and-joined")) + (func (export "cancel-two-joined-and-abandoned") + (alias export $d "cancel-two-joined-and-abandoned")) +) +(component instance $tc1 $TwoCandidates) +(assert_return (invoke "cancel-two-joined") (u32.const 42)) +(component instance $tc2 $TwoCandidates) +(assert_return (invoke "cancel-own-and-joined") (u32.const 42)) +(component instance $tc3 $TwoCandidates) +(assert_return (invoke "cancel-two-joined-and-abandoned") (u32.const 42)) + +;; The exclusive lock belongs to the task that acquired it, not to +;; whatever task its thread is currently executing on behalf of. A +;; synchronously-lifted async-typed export holds the lock for the whole core +;; call, so while its implicit thread is parked *inside another task*, a second +;; sync-lifted task must still be held at the implicit-backpressure gate. +;; Non-async-typed exports ignore backpressure entirely and can still barge +;; in, which is what releases the parked thread here. +(component + (component $C + (core module $Core + (import "" "task.return-u32" (func $task.return-u32 (param i32))) + (import "" "thread.index" (func $thread.index (result i32))) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + (import "" "task.drop" (func $task.drop (param i32))) + (import "" "thread.resume-later" (func $thread.resume-later (param i32))) + (import "" "thread.suspend" (func $thread.suspend (result i32))) + + (global $target-task (mut i32) (i32.const 0xdead)) ;; handle for task B + (global $holder-implicit (mut i32) (i32.const 0xdead)) ;; implicit thread of $holder + + ;; task B: its implicit thread publishes a handle and leaves the event + ;; loop at once, so B holds no exclusive lock and has no thread of its + ;; own; B is resolved by $holder's thread. + (func (export "target") (result i32) + (global.set $target-task (call $thread.get-task)) + (i32.const 0 (; EXIT ;))) + + ;; "target" leaves its event loop before any event can arrive + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + + ;; sync-lifted: acquires the exclusive lock on entry and releases it only + ;; when this core function returns + (func (export "holder") (result i32) + (local $home i32) + (local.set $home (call $thread.get-task)) + (global.set $holder-implicit (call $thread.index)) + + ;; resolve task B from inside B, still holding the lock + (call $thread.set-task (global.get $target-task)) + (call $task.return-u32 (i32.const 99)) + + ;; park while still contained by B: the lock stays held + (drop (call $thread.suspend)) + + ;; no move back home: the sync lift rejoins its original task for us + (call $task.drop (local.get $home)) + (call $task.drop (global.get $target-task)) + (i32.const 42)) + + ;; also sync-lifted, and so also gated on the exclusive lock + (func (export "blocked") (result i32) + (i32.const 11)) + + ;; not async-typed, so this ignores both the exclusive lock and the + ;; task queued behind it + (func (export "release") (result i32) + (call $thread.resume-later (global.get $holder-implicit)) + (i32.const 0)) + ) + (canon task.return (result u32) (core func $task.return-u32)) + (canon thread.index (core func $thread.index)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (canon task.drop (core func $task.drop)) + (canon thread.resume-later (core func $thread.resume-later)) + (canon thread.suspend (core func $thread.suspend)) + (core instance $core (instantiate $Core (with "" (instance + (export "task.return-u32" (func $task.return-u32)) + (export "thread.index" (func $thread.index)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + (export "task.drop" (func $task.drop)) + (export "thread.resume-later" (func $thread.resume-later)) + (export "thread.suspend" (func $thread.suspend)) + )))) + (func (export "target") async (result u32) + (canon lift (core func $core "target") async (callback (core func $core "never")))) + (func (export "holder") async (result u32) + (canon lift (core func $core "holder"))) + (func (export "blocked") async (result u32) + (canon lift (core func $core "blocked"))) + (func (export "release") (result u32) + (canon lift (core func $core "release"))) + ) + (component $D + (import "target" (func $target async (result u32))) + (import "holder" (func $holder async (result u32))) + (import "blocked" (func $blocked async (result u32))) + (import "release" (func $release (result u32))) + + (core module $Memory (memory (export "mem") 1)) + (core instance $memory (instantiate $Memory)) + (core module $Core + (import "" "mem" (memory 1)) + (import "" "subtask.drop" (func $subtask.drop (param i32))) + (import "" "waitable.join" (func $waitable.join (param i32 i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + (import "" "waitable-set.wait" (func $waitable-set.wait (param i32 i32) (result i32))) + (import "" "target" (func $target (param i32) (result i32))) + (import "" "holder" (func $holder (param i32) (result i32))) + (import "" "blocked" (func $blocked (param i32) (result i32))) + (import "" "release" (func $release (result i32))) + + (func (export "run") (result i32) + (local $ret i32) + (local $ws i32) + (local $n i32) + + (local.set $ws (call $waitable-set.new)) + + ;; start target, whose implicit thread leaves at once + (local.set $ret (call $target (i32.const 0 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (call $waitable.join (i32.shr_u (local.get $ret) (i32.const 4)) (local.get $ws)) + + ;; start holder: takes the lock, resolves target from inside target's + ;; task, then parks there with the lock still held + (local.set $ret (call $holder (i32.const 4 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (call $waitable.join (i32.shr_u (local.get $ret) (i32.const 4)) (local.get $ws)) + + ;; the lock is still held even though its owner is running as another + ;; task, so this one cannot even start + (local.set $ret (call $blocked (i32.const 8 (; retp ;)))) + (if (i32.ne (i32.const 0 (; STARTING ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + (call $waitable.join (i32.shr_u (local.get $ret) (i32.const 4)) (local.get $ws)) + + ;; barge past both the lock and the task queued behind it + (if (i32.ne (i32.const 0) (call $release)) + (then unreachable)) + + ;; holder returns 42 and drops the lock, which lets blocked start and + ;; return 11; target was resolved with 99 by holder's thread + (loop $l + (if (i32.ne (i32.const 1 (; SUBTASK ;)) + (call $waitable-set.wait (local.get $ws) (i32.const 16 (; eventp ;)))) + (then unreachable)) + (if (i32.ne (i32.const 2 (; RETURNED ;)) (i32.load offset=4 (i32.const 16))) + (then unreachable)) + (call $waitable.join (i32.load (i32.const 16)) (i32.const 0)) + (call $subtask.drop (i32.load (i32.const 16))) + (local.set $n (i32.add (local.get $n) (i32.const 1))) + (br_if $l (i32.lt_u (local.get $n) (i32.const 3)))) + + (if (i32.ne (i32.const 99) (i32.load (i32.const 0))) + (then unreachable)) + (if (i32.ne (i32.const 42) (i32.load (i32.const 4))) + (then unreachable)) + (if (i32.ne (i32.const 11) (i32.load (i32.const 8))) + (then unreachable)) + (i32.const 42)) + ) + (canon subtask.drop (core func $subtask.drop)) + (canon waitable.join (core func $waitable.join)) + (canon waitable-set.new (core func $waitable-set.new)) + (canon waitable-set.wait (memory (core memory $memory "mem")) (core func $waitable-set.wait)) + (canon lower (func $target) async (memory (core memory $memory "mem")) (core func $target')) + (canon lower (func $holder) async (memory (core memory $memory "mem")) (core func $holder')) + (canon lower (func $blocked) async (memory (core memory $memory "mem")) (core func $blocked')) + (canon lower (func $release) (core func $release')) + (core instance $core (instantiate $Core (with "" (instance + (export "mem" (memory $memory "mem")) + (export "subtask.drop" (func $subtask.drop)) + (export "waitable.join" (func $waitable.join)) + (export "waitable-set.new" (func $waitable-set.new)) + (export "waitable-set.wait" (func $waitable-set.wait)) + (export "target" (func $target')) + (export "holder" (func $holder')) + (export "blocked" (func $blocked')) + (export "release" (func $release')) + )))) + (func (export "run") async (result u32) (canon lift (core func $core "run"))) + ) + (instance $c (instantiate $C)) + (instance $d (instantiate $D + (with "target" (func $c "target")) + (with "holder" (func $c "holder")) + (with "blocked" (func $c "blocked")) + (with "release" (func $c "release")))) + (func (export "run") (alias export $d "run")) +) +(assert_return (invoke "run") (u32.const 42)) + +;; A callback-lifted export also holds the exclusive lock across each turn of +;; its event loop, so exiting the loop must release the lock even when the +;; thread has moved to another task in the meantime. Doing so abandons the +;; export's own task, unresolved and with no threads at all, which is *not* a +;; trap: the caller simply never hears back, and waiting for it deadlocks once +;; nothing else in the store can run. +(component + (component $C + (core module $Core + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + + (global $target-task (mut i32) (i32.const 0xdead)) ;; handle for task B + + ;; task B: publish a handle and leave the event loop; B keeps no thread + ;; of its own and is never resolved + (func (export "target") (result i32) + (global.set $target-task (call $thread.get-task)) + (i32.const 0 (; EXIT ;))) + + ;; join task B and exit the event loop immediately, while still + ;; contained by B: this task is left unresolved with no threads + (func (export "abandon") (result i32) + (call $thread.set-task (global.get $target-task)) + (i32.const 0 (; EXIT ;))) + + ;; neither export above stays in its event loop long enough to be given + ;; an event + (func (export "never") (param i32 i32 i32) (result i32) + unreachable) + + ;; sync-lifted, and so gated on the exclusive lock: this can only + ;; complete if "abandon" released the lock on its way out + (func (export "prove") (result i32) + (i32.const 7)) + ) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (core instance $core (instantiate $Core (with "" (instance + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + )))) + (func (export "target") async (result u32) + (canon lift (core func $core "target") async (callback (core func $core "never")))) + (func (export "abandon") async (result u32) + (canon lift (core func $core "abandon") + async (callback (core func $core "never")))) + (func (export "prove") async (result u32) + (canon lift (core func $core "prove"))) + ) + (component $D + (import "target" (func $target async (result u32))) + (import "abandon" (func $abandon async (result u32))) + (import "prove" (func $prove async (result u32))) + + (core module $Memory (memory (export "mem") 1)) + (core instance $memory (instantiate $Memory)) + (core module $Core + (import "" "mem" (memory 1)) + (import "" "waitable.join" (func $waitable.join (param i32 i32))) + (import "" "waitable-set.new" (func $waitable-set.new (result i32))) + (import "" "waitable-set.wait" (func $waitable-set.wait (param i32 i32) (result i32))) + (import "" "target" (func $target (param i32) (result i32))) + (import "" "abandon" (func $abandon (param i32) (result i32))) + (import "" "prove" (func $prove (param i32) (result i32))) + + (func (export "run") (result i32) + (local $ret i32) + (local $ws i32) + + ;; start target, whose implicit thread leaves at once + (local.set $ret (call $target (i32.const 0 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + + ;; start abandon: it joins target's task and exits, leaving its own + ;; task started-but-unresolved with no threads + (local.set $ret (call $abandon (i32.const 4 (; retp ;)))) + (if (i32.ne (i32.const 1 (; STARTED ;)) (i32.and (local.get $ret) (i32.const 0xf))) + (then unreachable)) + + ;; the exclusive lock was released on the way out, so this sync-lifted + ;; export runs to completion eagerly instead of blocking on the gate + (if (i32.ne (i32.const 2 (; RETURNED ;)) (call $prove (i32.const 8 (; retp ;)))) + (then unreachable)) + (if (i32.ne (i32.const 7) (i32.load (i32.const 8))) + (then unreachable)) + + ;; abandon's task can never resolve and nothing else can run + (local.set $ws (call $waitable-set.new)) + (call $waitable.join (i32.shr_u (local.get $ret) (i32.const 4)) (local.get $ws)) + (call $waitable-set.wait (local.get $ws) (i32.const 16 (; eventp ;))) + unreachable) + ) + (canon waitable.join (core func $waitable.join)) + (canon waitable-set.new (core func $waitable-set.new)) + (canon waitable-set.wait (memory (core memory $memory "mem")) (core func $waitable-set.wait)) + (canon lower (func $target) async (memory (core memory $memory "mem")) (core func $target')) + (canon lower (func $abandon) async (memory (core memory $memory "mem")) (core func $abandon')) + (canon lower (func $prove) async (memory (core memory $memory "mem")) (core func $prove')) + (core instance $core (instantiate $Core (with "" (instance + (export "mem" (memory $memory "mem")) + (export "waitable.join" (func $waitable.join)) + (export "waitable-set.new" (func $waitable-set.new)) + (export "waitable-set.wait" (func $waitable-set.wait)) + (export "target" (func $target')) + (export "abandon" (func $abandon')) + (export "prove" (func $prove')) + )))) + (func (export "run") async (result u32) (canon lift (core func $core "run"))) + ) + (instance $c (instantiate $C)) + (instance $d (instantiate $D + (with "target" (func $c "target")) + (with "abandon" (func $c "abandon")) + (with "prove" (func $c "prove")))) + (func (export "run") (alias export $d "run")) +) +(assert_trap (invoke "run") "wasm trap: deadlock detected: event loop cannot make further progress") diff --git a/test/binary/binary.wast b/test/binary/binary.wast index 3766e0180..369c2c348 100644 --- a/test/binary/binary.wast +++ b/test/binary/binary.wast @@ -1047,8 +1047,8 @@ "\03\05" ;; core type section (5 bytes) "\01" ;; 1 core type "\60\01\7f\00" ;; core functype (i32)->() - "\08\93\01" ;; canon section (147 bytes) - "\2f" ;; 47 canons + "\08\96\01" ;; canon section (150 bytes) + "\32" ;; 50 canons "\00\00\00\00\00" ;; lift core func 0 (f), no opts, type 0 "\00\00\01\03\00\03\00\04\05\01" ;; lift core func 1 (g), utf8 + (memory 0) + (realloc 5), type 1 "\00\00\02\02\06\07\03\02" ;; lift core func 2 (run), async + (callback 3), type 2 @@ -1096,6 +1096,9 @@ "\2b\00" ;; 0x2b thread.yield-then-resume "\2c\00" ;; 0x2c thread.suspend-then-promote "\2d\00" ;; 0x2d thread.yield-then-promote + "\30" ;; 0x30 thread.set-task + "\31" ;; 0x31 thread.get-task + "\32" ;; 0x32 task.drop ) (assert_malformed diff --git a/test/nyi.txt b/test/nyi.txt index 3efe95b69..5c06b5f0e 100644 --- a/test/nyi.txt +++ b/test/nyi.txt @@ -4,3 +4,5 @@ ./values/post-return.wast ./async/big-interleaving-test.wast ./async/forward.wast +./async/cancel-not-delivered.wast +./async/thread-set-task.wast diff --git a/test/values/post-return.wast b/test/values/post-return.wast index 62c0350d4..a720e57ee 100644 --- a/test/values/post-return.wast +++ b/test/values/post-return.wast @@ -23,6 +23,9 @@ (canon task.cancel (core func $task.cancel)) (canon thread.yield (core func $yield)) (canon thread.index (core func $thread.index)) + (canon thread.get-task (core func $thread.get-task)) + (canon thread.set-task (core func $thread.set-task)) + (canon task.drop (core func $task.drop)) (canon waitable-set.new (core func $waitable-set.new)) (canon waitable-set.wait (memory (core memory $memory "mem")) (core func $waitable-set.wait)) (canon waitable-set.poll (memory (core memory $memory "mem")) (core func $waitable-set.poll)) @@ -55,6 +58,9 @@ (import "" "task.cancel" (func $task.cancel)) (import "" "yield" (func $yield (result i32))) (import "" "thread.index" (func $thread.index (result i32))) + (import "" "thread.get-task" (func $thread.get-task (result i32))) + (import "" "thread.set-task" (func $thread.set-task (param i32))) + (import "" "task.drop" (func $task.drop (param i32))) (import "" "waitable-set.new" (func $waitable-set.new (result i32))) (import "" "waitable-set.wait" (func $waitable-set.wait (param i32 i32) (result i32))) (import "" "waitable-set.poll" (func $waitable-set.poll (param i32 i32) (result i32))) @@ -87,6 +93,9 @@ (func (export "trap-calling-task-cancel") (call $task.cancel)) (func (export "trap-calling-yield") (drop (call $yield))) (func (export "trap-calling-thread-index") (drop (call $thread.index))) + (func (export "trap-calling-thread-get-task") (drop (call $thread.get-task))) + (func (export "trap-calling-thread-set-task") (call $thread.set-task (i32.const 1))) + (func (export "trap-calling-task-drop") (call $task.drop (i32.const 1))) (func (export "trap-calling-waitable-set-new") (drop (call $waitable-set.new))) (func (export "trap-calling-waitable-set-wait") (drop (call $waitable-set.wait (i32.const 0) (i32.const 0)))) (func (export "trap-calling-waitable-set-poll") (drop (call $waitable-set.poll (i32.const 0) (i32.const 0)))) @@ -121,6 +130,9 @@ (export "task.cancel" (func $task.cancel)) (export "yield" (func $yield)) (export "thread.index" (func $thread.index)) + (export "thread.get-task" (func $thread.get-task)) + (export "thread.set-task" (func $thread.set-task)) + (export "task.drop" (func $task.drop)) (export "waitable-set.new" (func $waitable-set.new)) (export "waitable-set.wait" (func $waitable-set.wait)) (export "waitable-set.poll" (func $waitable-set.poll)) @@ -152,6 +164,9 @@ (func (export "trap-calling-task-cancel") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-task-cancel")))) (func (export "trap-calling-yield") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-yield")))) (func (export "trap-calling-thread-index") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-thread-index")))) + (func (export "trap-calling-thread-get-task") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-thread-get-task")))) + (func (export "trap-calling-thread-set-task") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-thread-set-task")))) + (func (export "trap-calling-task-drop") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-task-drop")))) (func (export "trap-calling-waitable-set-new") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-waitable-set-new")))) (func (export "trap-calling-waitable-set-wait") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-waitable-set-wait")))) (func (export "trap-calling-waitable-set-poll") (canon lift (core func $dm "noop") (post-return (core func $dm "trap-calling-waitable-set-poll")))) @@ -185,6 +200,9 @@ (func (export "trap-calling-task-cancel") (alias export $d "trap-calling-task-cancel")) (func (export "trap-calling-yield") (alias export $d "trap-calling-yield")) (func (export "trap-calling-thread-index") (alias export $d "trap-calling-thread-index")) + (func (export "trap-calling-thread-get-task") (alias export $d "trap-calling-thread-get-task")) + (func (export "trap-calling-thread-set-task") (alias export $d "trap-calling-thread-set-task")) + (func (export "trap-calling-task-drop") (alias export $d "trap-calling-task-drop")) (func (export "trap-calling-waitable-set-new") (alias export $d "trap-calling-waitable-set-new")) (func (export "trap-calling-waitable-set-wait") (alias export $d "trap-calling-waitable-set-wait")) (func (export "trap-calling-waitable-set-poll") (alias export $d "trap-calling-waitable-set-poll")) @@ -270,6 +288,12 @@ (assert_trap (invoke "trap-calling-stream-forward") "cannot leave component instance") (component instance $i29 $Tester) (assert_trap (invoke "trap-calling-future-forward") "cannot leave component instance") +(component instance $i30 $Tester) +(assert_trap (invoke "trap-calling-thread-get-task") "cannot leave component instance") +(component instance $i31 $Tester) +(assert_trap (invoke "trap-calling-thread-set-task") "cannot leave component instance") +(component instance $i32 $Tester) +(assert_trap (invoke "trap-calling-task-drop") "cannot leave component instance") ;; built-ins that don't trap: