diff --git a/docs/spec/concurrency_and_lifetimes.md b/docs/spec/concurrency_and_lifetimes.md index 011e1ac7e..da6fffb83 100644 --- a/docs/spec/concurrency_and_lifetimes.md +++ b/docs/spec/concurrency_and_lifetimes.md @@ -301,13 +301,23 @@ that bounded wait into an unbounded one. Four dispositions, by site: - **`installReconnectHandler`'s reconnect callback.** Left as a `liveness()` check, deliberately not moved to `BridgeLifetime` — this is the case the first paragraph above warns about. The handler runs on the backend's - transport thread and calls `registerModelWithContext`/`registerModelShared`, - which blocks on a nested `QEventLoop` for `QtWebSocketBackend`. Gating that - span would let `~Bridge` block for the same round trip, and if the reconnect - and `~Bridge` ever land on the same thread — plausible for Qt, whose nested - loop pumps the very deferred-delete event that could run the destructor — - that is a self-deadlock, not a slow teardown. No safe mechanical fix is known - for this site; it remains open, tracked as the residual scope of issue #489. + transport thread; gating its span would let `~Bridge` block for the whole + re-registration, and if the reconnect and `~Bridge` ever land on the same + thread — plausible for Qt, whose nested loop pumps the very deferred-delete + event that could run the destructor — that is a self-deadlock, not a slow + teardown. + + morph#615 removed the worst version of that span rather than the span + itself. The handler used to call `registerModelWithContext`/ + `registerModelShared`, which for `QtWebSocketBackend` block on a nested + `QEventLoop`; it now calls `bindModel` and consults + `IBackend::bindWaitPolicy()`, so a backend that says `kCallerMustNotBlock` + is not waited for at all and the handler returns promptly. A + `kCallerMayBlock` backend is still waited out, on the transport thread, + under both bridge mutexes — a bounded round trip by that backend's own + contract, but still a span a `BridgeLifetime` gate must not cover. So the + site stays on `liveness()`, and the residual scope of issue #489 stays + open. - **The `*Async` reply callbacks** — `attachHandlerAsync`, `ensureBoundAsync` and `assignHandlerPrimary`, three of the four `IBackend` async hooks. (The fourth, `registerHandlerImpl`, is covered by the `BridgeLifetime` bullet diff --git a/docs/spec/core/backend.md b/docs/spec/core/backend.md index 33f1115c8..57b4c5778 100644 --- a/docs/spec/core/backend.md +++ b/docs/spec/core/backend.md @@ -81,7 +81,7 @@ holds a `unique_ptr` and delegates all model operations to it. | `registerModelAsync(typeId, factory, contextKey, onRegistered, onError)` | Optional non-blocking counterpart to `registerModelWithContext`. Returns `false` by default, and since morph#568 no backend overrides it, so it always does; `Bridge::registerHandler()` then falls back to `bindModel`. Removed by morph#571. See [Asynchronous registration](#asynchronous-registration--registermodelasync). | | `bindModel(request, cbExec)` | Acquires a model instance and returns a `Completion` delivered on `cbExec`. One verb covering `registerModelWithContext`, `registerModelShared` and `attachModel`, selected by the request's shape. The preferred surface — see [The structural registration surface](#the-structural-registration-surface--bindmodel-and-promotemodel). | | `promoteModel(request, cbExec)` | Files an already-live instance under a key and returns a `Completion` delivered on `cbExec`. The structural counterpart of `assignPrimary`. | -| `bindWaitPolicy()` | Whether a caller may block its own thread until a `bindModel`/`promoteModel` completion settles. `BindWait::kCallerMayBlock` by default. The only framework caller is `Bridge::registerHandlerImpl`; see [Waiting for a bind — `bindWaitPolicy`](#waiting-for-a-bind--bindwaitpolicy). | +| `bindWaitPolicy()` | Whether a caller may block its own thread until a `bindModel`/`promoteModel` completion settles. `BindWait::kCallerMayBlock` by default. The framework callers are `Bridge::registerHandlerImpl`, `Bridge::switchBackend`'s phase 1 and `Bridge::installReconnectHandler`'s handler; see [Waiting for a bind — `bindWaitPolicy`](#waiting-for-a-bind--bindwaitpolicy). | | `deregisterModel(mid)` | Removes the model identified by `mid`. | | `execute(mid, call, cbExec)` | Dispatches `call` against the model identified by `mid`. Returns a `Completion>`. | | `notifyBackendChanged()` | Called by `Bridge::switchBackend()` after all handlers are re-registered. | @@ -460,6 +460,28 @@ natively](#the-structural-registration-surface-natively). - **The executor is required.** An adapter that ran the call inline when handed nothing would be a `bindModel` that blocks on some configurations and not others — contract by configuration, which is what is being removed. +- **It cancels its own pending completions rather than only the wrapped + backend's** (morph#619). `bindModel`/`promoteModel` settle from a task on + `_control`, holding a promise the wrapped backend never sees, so the + one-line `_inner->cancelPending(exc)` this verb used to be reached none of + them: a bind cancelled by `~Bridge` or by `switchBackend` went on to resolve + **successfully** afterwards, against `IBackend::cancelPending`'s "after this + call, any later `setValue`/`setException` on those states is a no-op". The + adapter therefore keeps a `weak_ptr` to each dispatched promise and rejects + the live ones first, on the same snapshot-then-deliver shape and the same + amortised compaction as + [`LocalBackend`'s pending list](#the-pending-list-and-its-amortised-compaction); + an entry expires by itself when its strand task is destroyed, so the success + path erases nothing. + + Two limits are deliberate rather than overlooked. A task that settles first + wins, because its completion was not still pending — the same race + `LocalBackend::cancelPending` has always had. And a task already **queued** + on `_control` still runs its blocking control call against the wrapped + backend after `cancelPending` returns: the caller is told the bind was + cancelled while the registration may still go through. Stopping that is a + different change — it needs the task to check before calling `op()`, not the + promise to be settled after it — and is tracked as morph#636. - **Control calls are serialised** onto one strand, so the wrapped backend sees them one at a time, as it did when the blocking call itself serialised callers. `~SynchronousBackendAdapter` waits for any in-flight control call, so @@ -468,15 +490,19 @@ natively](#the-structural-registration-surface-natively). - **A control call issued from a reconnect handler runs on the strand**, never on the wrapped backend's transport thread — *provided the handler issues it through `bindModel`/`promoteModel`*. That proviso is load-bearing, and it is - why the adapter does **not** settle `SocketBackend`'s documented - reconnect-handler deadlock hazard (see the "runs reconnect handlers on a - dedicated thread" row under [Design decisions](#design-decisions)): the - adapter forwards `setReconnectHandler` to the wrapped backend unchanged, and - `Bridge::installReconnectHandler`'s handler calls the *blocking* + why the adapter did **not** settle `SocketBackend`'s documented + reconnect-handler deadlock hazard: the adapter forwards + `setReconnectHandler` to the wrapped backend unchanged, and + `Bridge::installReconnectHandler`'s handler used to call the *blocking* `registerModelShared`/`registerModelWithContext`, which the adapter also - forwards unchanged. A wrapped `SocketBackend` would therefore run its - reconnect control calls exactly where it runs them today. See morph#569's - answer under [`SocketBackend`](#socketbackend--socketserver--raw-socket-websocket-transport). + forwards unchanged. + + morph#615 changed the second half of that: the handler now calls + `bindModel`, so a wrapped backend's reconnect control calls *do* reach the + adapter's strand. What the proviso still rules out is the first half — the + handler itself is invoked by the wrapped backend, on whatever thread that + backend chooses, and the adapter does not move it. See morph#569's answer + under [`SocketBackend`](#socketbackend--socketserver--raw-socket-websocket-transport). ### Backends with a genuinely non-blocking path @@ -603,6 +629,13 @@ morph#567 introduced had removed the only signal that told them apart (the answer at that call site instead hung five `tests/qt/` cases on the nested `QEventLoop`. Neither fixed setting of the call site is correct. +Since morph#615 the same question is asked at all three sites that acquire an +instance for a binding — `registerHandlerImpl`, `switchBackend`'s phase 1 and +the reconnect handler — because all three used to block a thread that a +`kCallerMustNotBlock` backend needs back before its reply can arrive. The +first was the only synchronous *entry point*; the other two are worse, because +the reconnect handler runs on the backend's own transport thread. + `IBackend::bindWaitPolicy()` restores exactly one bit: - `BindWait::kCallerMayBlock` (the default) — the completion settles without @@ -620,10 +653,11 @@ verb to call*, so every call site carried two paths and a backend could be half-migrated; this one chooses nothing. There is still exactly one verb, called unconditionally, and exactly one continuation — the only question is whether the thread that registered that continuation is allowed to stop and wait for it. -Only a synchronous entry point that must return a bound instance asks: -`registerHandlerImpl` is the sole framework caller, and the asynchronous entry -points (`attachHandlerAsync`, `ensureBoundAsync`, `assignHandlerPrimary`) never -wait and never consult it. +Only a frame that would otherwise stop and wait asks: `registerHandlerImpl` +(a synchronous entry point that must return a bound instance), +`switchBackend`'s phase 1 and the reconnect handler's re-registration loop. +The asynchronous entry points (`attachHandlerAsync`, `ensureBoundAsync`, +`assignHandlerPrimary`) never wait and never consult it. The wait is unbounded by design. A timeout would make "is this handler bound when `registerHandler` returns" depend on how fast the network was, which is the @@ -647,7 +681,8 @@ caller that wants to gate anyway, and is required for the two | morph#568 | `QtWebSocketBackend` implements the surface natively and drops all four `*Async` overrides; `Bridge`'s four dispatch sites fall back to it instead of to a synchronous verb. | Landed | | morph#569 | `SocketBackend` implements the surface natively, keeping every legacy verb on `sendSync`. | Landed | | morph#593 | Adds `IBackend::bindWaitPolicy()`, the one signal morph#567's surface left the call site without. Fixes the `"handler not bound"` regression morph#568 caused in `SocketBackend`. | Landed | -| morph#570 | The example GUIs and the WASM spike; `Bridge::installReconnectHandler` onto `bindModel`. | Open | +| morph#570 | The example GUIs and the WASM spike. | Open | +| morph#615 | `Bridge::switchBackend`'s phase 1 and `Bridge::installReconnectHandler`'s handler onto `bindModel`, both consulting `bindWaitPolicy()`; `switchBackend`'s rollback keys on a rejected `Completion`. Also settles who owned the reconnect half: this table used to assign it to morph#570, whose own body scopes itself to `examples/` and never mentions `bridge.hpp`. | Landed | | morph#571 | Removes the four `*Async` verbs; migrates `LocalBackend`, `SimulatedRemoteBackend` and the test doubles. | Open | Every existing implementor still compiles unchanged, and the default @@ -671,13 +706,23 @@ branch is the one that always runs; the eleven test doubles in `tests/test_async_registration.cpp` that do override them still exercise the first, which is what keeps that suite meaningful until morph#571 rewrites it. -`Bridge::installReconnectHandler` is the one dispatch site morph#568 did **not** -move: its handler still calls the blocking `registerModelShared`/ -`registerModelWithContext`. That is deliberate — moving it changes `Bridge`'s -locking model, not just a call — and it is the reason `SocketBackend`'s -reconnect-handler deadlock hazard survives morph#568 exactly as -[its own section](#the-structural-registration-surface-natively) says it -survives morph#569. +`Bridge::installReconnectHandler` and `Bridge::switchBackend`'s phase 1 were +the two dispatch sites morph#568 did **not** move: both still called the +blocking `registerModelShared`/`registerModelWithContext`, so a backend that +had just been given a way to say `kCallerMustNotBlock` was blocked at both of +them anyway. morph#615 moved them, and its shape is the one every other site +already had — dispatch, park an inline reply in an `AsyncDispatchHandoff`, +then `awaitHandoff` or `claimHandoff` according to the policy. What made it a +change of its own rather than part of morph#568 is what the two sites do +around the call: `switchBackend` stages for an all-or-nothing commit whose +rollback used to key on a thrown exception, and the reconnect handler runs on +the transport thread holding both bridge mutexes. Both are described in +[bridge.md](bridge.md). + +`SocketBackend`'s reconnect-handler deadlock hazard is narrowed rather than +closed by that: the control call it makes is no longer the blocking verb, but +the handler is still *invoked* by the backend on whatever thread the backend +chooses, which is the half morph#569 owns. The executor those call sites name is **`exec::detail::inlineExecutor()`**, which runs the continuation on the thread that settled it. That is deliberately the @@ -2200,7 +2245,7 @@ inside the class calls `close()` — no thread it joins can be waiting on it. | `registerModelWithContext` | `virtual ModelId registerModelWithContext(const string&, function()>, string_view)` | Default: drops `contextKey`, calls `registerModel`. | | `bindModel` | `virtual Completion bindModel(BindRequest, IExecutor& cbExec)` | Default: runs `bindModelBlocking` inline and settles. See [The structural registration surface](#the-structural-registration-surface--bindmodel-and-promotemodel). | | `promoteModel` | `virtual Completion promoteModel(PromoteRequest, IExecutor& cbExec)` | Default: calls `assignPrimary` inline and settles with `request.mid`. | -| `bindWaitPolicy` | `virtual BindWait bindWaitPolicy() const noexcept` | Default: `BindWait::kCallerMayBlock`. Whether a caller may block until a `bindModel`/`promoteModel` completion settles. Read only by `Bridge::registerHandlerImpl`. | +| `bindWaitPolicy` | `virtual BindWait bindWaitPolicy() const noexcept` | Default: `BindWait::kCallerMayBlock`. Whether a caller may block until a `bindModel`/`promoteModel` completion settles. Read by `Bridge::registerHandlerImpl`, `Bridge::switchBackend` and `Bridge::installReconnectHandler`'s handler. | | `bindModelBlocking` | `ModelId bindModelBlocking(BindRequest)` | Non-virtual. Routes a `BindRequest` to `registerModelWithContext` / `registerModelShared` / `attachModel` by its shape; blocks. Shared by the default `bindModel` and by `SynchronousBackendAdapter`. | | `deregisterModel` | `virtual void deregisterModel(ModelId)` | Pure virtual. | | `execute` | `virtual Completion> execute(ModelId, ActionCall, IExecutor*)` | Pure virtual. | @@ -2240,6 +2285,7 @@ inside the class calls `close()` — no thread it joins can be waiting on it. | `bindModel(request, cbExec)` | Posts `inner->bindModelBlocking(request)` onto the control strand; settles the returned `Completion` on `cbExec`. Never blocks the caller. | | `bindWaitPolicy()` | `BindWait::kCallerMustNotBlock`, always. Not forwarded: it describes the two verbs the adapter reshapes. | | `promoteModel(request, cbExec)` | Posts `inner->assignPrimary(...)` onto the control strand; resolves with `request.mid`. | +| `cancelPending(exc)` | Rejects the adapter's own still-unsettled `bindModel`/`promoteModel` promises with `exc`, **then** forwards to `inner`. Not a plain forward: those promises are settled from `_control` tasks the wrapped backend has never heard of (morph#619). | | every other `IBackend` verb | Forwarded to `inner` unchanged, including the four `*Async` twins — wrapping a backend that has a non-blocking path must not take it away. | ### Error types diff --git a/docs/spec/core/bridge.md b/docs/spec/core/bridge.md index edc8e3462..70fc19f94 100644 --- a/docs/spec/core/bridge.md +++ b/docs/spec/core/bridge.md @@ -104,12 +104,34 @@ handler so backends with recoverable transports (e.g. `QtWebSocketBackend`) re-register all live bindings on reconnection. The handler branches on `binding->shared` exactly as `switchBackend()`'s phase 1 does: an unattached shared binding (`binding->primary` empty) has no instance to recreate and is -skipped, and an attached shared binding re-registers through -`registerModelShared(typeId, modelFactory, {contextKey, primary})` rather than -`registerModelWithContext` — otherwise the reconnected registration would come -back as an ordinary non-shared instance, silently dropping the sharing a +skipped, and an attached shared binding re-registers through the +`registerModelShared` request shape (a non-empty `primary`) rather than the +`registerModelWithContext` one — otherwise the reconnected registration would +come back as an ordinary non-shared instance, silently dropping the sharing a surviving handler had before the disconnect. A non-shared binding is -unaffected and still re-registers via `registerModelWithContext`. +unaffected and still re-registers through the `registerModelWithContext` +shape. + +Since morph#615 the handler reaches the backend through +`IBackend::bindModel` — the structural registration surface — and consults +`IBackend::bindWaitPolicy()` exactly as `registerHandlerImpl` does. It matters +here more than anywhere: the handler runs on the backend's *transport* thread, +so for a backend whose reply is delivered by that same thread's event loop +(`QtWebSocketBackend` with `asyncRegistrationEnabled`) the blocking verbs this +loop used to call park the thread against itself, and on a WASM main thread +abort the page. Three consequences: + +- For a `kCallerMayBlock` backend the handler waits out each bind and behaves + exactly as it did before, blocking verb for blocking verb. +- For a `kCallerMustNotBlock` backend it does not wait. Each binding's + `currentId` is cleared — the id belonged to the connection that just + dropped, so `executeVia`'s fast fail is the honest answer and `whenBound()` + becomes gateable (see below) — and the reply publishes the new id when it + arrives. +- A failing re-registration no longer takes the rest of the loop with it. It + used to throw out of the handler and onto the transport thread; a rejected + `Completion` is reported per binding instead, the binding is left unbound, + and the loop continues. **`registerHandler()`** creates a `HandlerBinding` with the default `ModelFactory::create()` factory and registers it on the active @@ -197,8 +219,32 @@ only the two bridge-touching side effects are skipped when the token has expired. **`switchBackend(newBackend)`** replaces the active backend atomically: the -switch either fully succeeds or leaves everything exactly as it was. It has -two overloads: +switch either fully succeeds or leaves everything exactly as it was. Since +morph#615 phase 1 acquires each instance through `IBackend::bindModel` and +consults `IBackend::bindWaitPolicy()`, and **atomicity is exactly as strong as +the wait is**: + +- `kCallerMayBlock` (every backend in the tree that is not a + `SynchronousBackendAdapter` or a WASM-configured `QtWebSocketBackend`): the + frame waits out each bind, so every outcome is known before anything is + published. The guarantee is unchanged; the rollback simply keys on a + rejected `Completion` now rather than on a thrown exception, which is what + the structural surface reports failure through. +- `kCallerMustNotBlock`: it cannot be. "Did every re-registration succeed" is + not knowable without waiting, and waiting on such a backend is the deadlock + the policy exists to prevent. Those bindings are *deferred*: the swap + happens, their `currentId` is set to 0, `registrationInFlight` is set so + `whenBound()` can gate on them, and each reply binds its own. A deferred + bind that fails after the swap is logged and rejects that binding's + `whenBound()` waiters; it cannot roll the switch back. An instance created + by a deferred bind whose switch later threw is not deregistered either — + nothing in the rollback ever learns its id. + +Waiters are always resolved after `_mtx`/`_attachMtx` are released, because +resolving one runs consumer code that is free to re-enter the `Bridge`; the +staging exception is rethrown after that, for the same reason. + +It has two overloads: - `switchBackend(shared_ptr)` — the caller keeps its own reference, so the same backend instance can be re-installed later (e.g. switching back @@ -575,10 +621,14 @@ here and nowhere else. See issue #347. Scope limits worth knowing, because each is a question `whenBound()` looks like it answers and does not: -- **It tracks the initial registration only.** `registerHandlerImpl` is the - sole writer of `registrationInFlight`. Re-registration by `switchBackend()` - and by the reconnect handler is synchronous and never sets it, so - `whenBound()` says nothing about a backend swap or a reconnect in progress. +- **It tracks the initial registration, plus a re-registration that could not + be waited for.** `registerHandlerImpl` sets `registrationInFlight` on every + path. `switchBackend()` and the reconnect handler set it only when the new + backend answers `kCallerMustNotBlock` (morph#615) — the one case where they + leave a binding unbound with a reply still to come, which is exactly the + state `whenBound()` exists to describe. Against a `kCallerMayBlock` backend + both still settle every outcome inside their own frame and set nothing, so + `whenBound()` says nothing about a swap or reconnect there, as before. - **It does not track the shared attach path.** `attachHandler`/`ensureBound` and their async counterparts bind a shared handler without going through `registerHandlerImpl`, so for an `AllowShared` handler `whenBound()` is only @@ -900,7 +950,7 @@ make teardown order-independent.) | dtor | `~Bridge()` | Clears the active backend's reconnect handler, then cancels all pending completions with `BridgeDestroyedError`. | | `registerHandler` | `shared_ptr registerHandler()` | Default factory. Prefers `IBackend::registerModelAsync`; see `backend.md`. | | `registerHandler(binding)` | `void registerHandler(const shared_ptr&)` | Pre-built binding. Same async-preferring behavior. | -| `switchBackend` | `void switchBackend(unique_ptr)` / `void switchBackend(shared_ptr)` | Pushes the current default session onto the new backend via `setSession` before staging. Atomic: stages all re-registrations on the new backend, commits (publishes new ids + swaps) only if all succeed, else rolls back and rethrows leaving old backend + `currentId`s intact. Cancels old backend's pending ops with `BackendChangedError`. Holds both `_mtx` and `_attachMtx` for its duration. The `unique_ptr` overload is a template on the concrete backend type and delegates to the `shared_ptr` one — see below. | +| `switchBackend` | `void switchBackend(unique_ptr)` / `void switchBackend(shared_ptr)` | Pushes the current default session onto the new backend via `setSession` before staging. Stages all re-registrations through `bindModel` on the new backend, commits (publishes new ids + swaps) only if all succeed, else rolls back and rethrows leaving old backend + `currentId`s intact. Atomic exactly when the new backend answers `kCallerMayBlock`; a `kCallerMustNotBlock` backend's binds are deferred and the switch is not all-or-nothing (see above). Cancels old backend's pending ops with `BackendChangedError`. Holds both `_mtx` and `_attachMtx` for its staging and commit, and resolves `whenBound()` waiters after releasing them. The `unique_ptr` overload is a template on the concrete backend type and delegates to the `shared_ptr` one — see below. | | `deregisterHandler` | `void deregisterHandler(const shared_ptr&)` | Deregisters from active backend (if bound), resets `currentId` to 0, removes from tracking. | | `executeVia` | `Completion executeVia(const shared_ptr&, Action, IExecutor*)` | Lock-free dispatch. Attaches default session. On `LocalBackend`, rejects an action whose `ActionValidator::ready` returns `false` with `morph::model::ValidationError` via `onError`, before `Model::execute` runs. Records a journal `LogEntry` for loggable actions on both success (`Outcome::Succeeded`) and a throwing `Model::execute` (`Outcome::Failed`, rethrown unchanged). Value-forwarding into the typed `Completion` is `try`/`catch`-guarded — a throwing result move/copy resolves the completion via `onError` instead of hanging or terminating. The bridge-touching side effects (`onResult`, `hasSubscribers()`/`publishResult`, the `pendingCalls()` decrement, and the execute-deadline disarm) are gated on the bridge's `CallbackToken`, checked before any runs, so a completion resolving after `~Bridge()` skips them instead of touching the dangling `Bridge`. Increments `pendingCalls()` once per call before dispatch (never for the synchronous "handler not bound" early return); decrements it exactly once, from whichever of the two mutually-exclusive resolution continuations actually fires. Arms the client-side execute deadline when one is installed (see `setExecuteDeadline`); the fast-fail "handler not bound" path returns before that and arms nothing. | | `setDefaultSession` | `void setDefaultSession(session::Context)` | Installs default session context; also pushes it to the active backend via `IBackend::setSession` so control envelopes (register/attach/assign/deregister) carry it too, not only `execute`. | diff --git a/docs/spec/core/completion.md b/docs/spec/core/completion.md index 29271ec4e..d3b16ec61 100644 --- a/docs/spec/core/completion.md +++ b/docs/spec/core/completion.md @@ -401,10 +401,17 @@ Emscripten's non-pthread `pthread_create` stub, so constructing a `std::thread` throws `std::system_error` at runtime — which would have made `setExecuteDeadline` unusable from a browser tab, and with it `examples/common/gui/event_poller.hpp`, whose constructor calls it -unconditionally. Two behavioural differences, both documented in +unconditionally. Three behavioural differences, all documented in `timeout_scheduler.hpp`'s own `@file` comment: callbacks are never concurrent -with the caller, and `cancel()` releases the callback immediately but leaves the -underlying browser timer to elapse harmlessly rather than clearing it. **This +with the caller; `cancel()` releases the callback immediately but leaves the +underlying browser timer to elapse harmlessly rather than clearing it; and +`cancel()` there really does mean "no callback runs after this returns", +whereas the threaded build's `cancel()` returns while an *already-started* +callback goes on running on the scheduler thread (morph#620). A caller that +must work in both builds gets the weaker of the two: every scheduled callback +has to stay safe to run after its own `cancel()`, which the deadline callback +here does by settling a write-once `CompletionState` it holds a `shared_ptr` +to. **This build has never been compiled or run in this repository** — no Emscripten toolchain is available here; its only verification is the `ladder-wasm` CI compile gate. diff --git a/include/morph/core/backend.hpp b/include/morph/core/backend.hpp index d2d3a485a..dea0ed6cf 100644 --- a/include/morph/core/backend.hpp +++ b/include/morph/core/backend.hpp @@ -646,10 +646,14 @@ struct IBackend { /// /// Answers the one question `Completion` cannot: an unsettled `Completion` /// looks the same whether the reply is coming from a thread the caller does - /// not own or from the caller's own event loop. Only a synchronous entry - /// point that must hand back a *bound* instance asks it — - /// `Bridge::registerHandlerImpl` is the single caller in the framework; the - /// asynchronous entry points never wait and never consult it. + /// not own or from the caller's own event loop. Only a frame that would + /// otherwise stop and wait asks it — `Bridge::registerHandlerImpl` (a + /// synchronous entry point that must hand back a *bound* instance), + /// `Bridge::switchBackend`'s staging phase, and + /// `Bridge::installReconnectHandler`'s handler, which is the one that runs + /// on the backend's own transport thread and so is the one a wrong answer + /// deadlocks outright (morph#615). The asynchronous entry points never wait + /// and never consult it. /// /// A backend that returns `kCallerMayBlock` (the default) commits to /// settling every `Completion` it returns exactly once without any further @@ -925,6 +929,12 @@ class SynchronousBackendAdapter : public detail::IBackend { } } + /// @brief The producer side of a `bindModel`/`promoteModel` completion. + /// + /// Named because `cancelPending` has to hold these weakly; see + /// `trackPending`. + using BindPromise = ::morph::async::Completion<::morph::exec::detail::ModelId>::Promise; + /// @brief The wrapped backend. /// @return Reference to the backend passed at construction; never null. [[nodiscard]] detail::IBackend& wrapped() const noexcept { return *_inner; } @@ -1128,9 +1138,47 @@ class SynchronousBackendAdapter : public detail::IBackend { /// @brief Forwards to the wrapped backend. void notifyBackendChanged() override { _inner->notifyBackendChanged(); } - /// @brief Forwards to the wrapped backend. - /// @param exc Exception delivered to every still-pending completion. - void cancelPending(const std::exception_ptr& exc) override { _inner->cancelPending(exc); } + /// @brief Rejects the completions *this adapter* produced, then forwards. + /// + /// Not a plain forward, unlike everything else in this block, and for the + /// same reason `bindWaitPolicy()` is not: the two verbs this adapter + /// reshapes produce completions the wrapped backend has never heard of. + /// A `bindModel` here settles from a task on `_control`, so + /// `_inner->cancelPending` reaches nothing of it — before morph#619 a bind + /// dispatched through this adapter went on to resolve **successfully** + /// after cancellation, which is the exact opposite of what + /// `IBackend::cancelPending` promises ("after this call, any later + /// `setValue`/`setException` on those states is a no-op"). + /// + /// The adapter's own promises are rejected first and the forward happens + /// second: the strand task is still running `op()` while this executes, and + /// rejecting before handing control to the wrapped backend keeps that + /// window as short as the caller's own call. A task that finished first + /// wins — its completion was not "still pending" — and a task that finishes + /// after finds the promise settled, which `CompletionState::setValue`'s + /// `if (ready) return;` makes a no-op. + /// + /// **What this does not do:** a task already queued on `_control` still + /// runs its blocking control call against the wrapped backend after this + /// returns. The caller is told the bind was cancelled, but the registration + /// may still happen — morph#636, since stopping it needs the task to check + /// before calling `op()`, not the promise to be settled after it. + /// @param exc Exception delivered to every still-pending completion, this + /// adapter's own and then the wrapped backend's. + void cancelPending(const std::exception_ptr& exc) override { + std::vector> snapshot; + { + std::scoped_lock const lock{_pendingMtx}; + snapshot.swap(_pending); + _compactAt = kPendingCompactFloor; + } + for (auto& weak : snapshot) { + if (auto promise = weak.lock()) { + promise->reject(exc); + } + } + _inner->cancelPending(exc); + } /// @brief Forwards to the wrapped backend. /// @param handler Callable invoked after a successful reconnect; `nullptr` clears. @@ -1163,6 +1211,12 @@ class SynchronousBackendAdapter : public detail::IBackend { using Settled = ::morph::async::Completion<::morph::exec::detail::ModelId>; auto [completion, promise] = Settled::makeSettleable(&cbExec); auto shared = std::make_shared(std::move(promise)); + // Tracked *before* the post, not after: a `cancelPending` that lands in + // between would otherwise find an empty list and leave a completion + // that is genuinely pending uncancelled. Rejecting a promise whose task + // has not started yet is safe — the task's own `resolve` then finds the + // state ready and returns (morph#619). + trackPending(shared); _control.post(kControlStrand, [shared, op = std::move(op)]() mutable { try { shared->resolve(op()); @@ -1173,13 +1227,49 @@ class SynchronousBackendAdapter : public detail::IBackend { return std::move(completion); } + /// @brief Records @p promise as cancellable until its task settles it. + /// + /// The strand task holds the only `shared_ptr` to the promise, so an entry + /// here expires exactly when that task is destroyed — "still pending" needs + /// no separate bookkeeping and no erase on the success path. + /// + /// Swept on the same amortised schedule as `LocalBackend::trackPending` + /// (morph#528): dead entries are reclaimed only when the list reaches + /// `_compactAt`, which each sweep re-arms at twice the surviving count, so + /// the per-dispatch cost is O(1) and the list stays bounded at twice the + /// live count plus the floor. Control calls are serialised onto one strand, + /// so in practice the live count is one and the floor is never reached; + /// without the sweep the list would still grow without bound on an adapter + /// whose `cancelPending` is never called. + /// @param promise Promise to reject if `cancelPending` runs before its task settles it. + void trackPending(const std::shared_ptr& promise) { + std::scoped_lock const lock{_pendingMtx}; + if (_pending.size() >= _compactAt) { + std::erase_if(_pending, [](const auto& weak) { return weak.expired(); }); + _compactAt = std::max(kPendingCompactFloor, _pending.size() * 2); + } + _pending.emplace_back(promise); + } + /// @brief The single strand key every control call shares, so they run one /// at a time. Not a real model id: this `StrandExecutor` is private /// to the adapter and shares no key space with any backend's own. static constexpr ::morph::exec::detail::ModelId kControlStrand{1}; + /// @brief Smallest size at which `trackPending` sweeps; see `LocalBackend`'s. + static constexpr std::size_t kPendingCompactFloor = 32; + std::shared_ptr _inner; ::morph::exec::detail::StrandExecutor _control; + mutable std::mutex _pendingMtx; + // Every `bindModel`/`promoteModel` promise handed to a `_control` task and + // not yet settled by it. Weak, so a settled task's promise drops out on its + // own; guarded by `_pendingMtx`, because `cancelPending` is called from + // `Bridge`'s thread while `dispatch` runs on whichever thread called it. + std::vector> _pending; + // Size at which `trackPending` next sweeps `_pending`; re-armed at twice + // the surviving count. Guarded by `_pendingMtx` with `_pending` itself. + std::size_t _compactAt = kPendingCompactFloor; }; /// @brief In-process backend that executes model actions on a thread pool strand. diff --git a/include/morph/core/bridge.hpp b/include/morph/core/bridge.hpp index 518dade79..86a73aec2 100644 --- a/include/morph/core/bridge.hpp +++ b/include/morph/core/bridge.hpp @@ -260,6 +260,30 @@ struct ParkedOutcome { std::exception_ptr failure; }; +/// @brief What `Bridge::switchBackend`'s staging phase accumulated before it +/// committed or rolled back. +/// +/// Its own type rather than four parallel locals because the rollback and the +/// commit each need a different subset of it, and a subset silently missed is +/// how a waiter ends up hanging on a registration that will never report. +struct StagedRebinds { + /// @brief Bindings still live, in `_handlers` order; replaces `_handlers` on commit. + std::vector> live; + + /// @brief `(binding, newId)` for every bind that settled successfully in + /// the staging frame. Published on commit, deregistered on rollback. + std::vector, std::uint64_t>> staged; + + /// @brief Bindings whose bind is still in flight; published as unbound and + /// bound later by their own reply. Only ever non-empty for a + /// `kCallerMustNotBlock` backend. + std::vector> deferred; + + /// @brief Every binding this staging phase marked `registrationInFlight`. + /// Settled whichever way the phase ends. + std::vector> armed; +}; + /// @brief Handoff slot between an async attach/bind dispatch and its callback. /// /// `Bridge::attachHandlerAsync`/`ensureBoundAsync` dispatch to the backend while @@ -1361,6 +1385,21 @@ class Bridge { newShared->setSession(_defaultSession); } std::shared_ptr<::morph::backend::detail::IBackend> previous; + // Settled after the mutexes are released, never under them: resolving a + // waiter runs consumer code, which is free to call back into this + // `Bridge`. + std::vector> settleBound; + std::vector> settleFailed; + std::exception_ptr staging; + // Whether this frame may stop and wait for each bind to settle. A + // backend that answers `kCallerMustNotBlock` delivers its replies + // through the calling thread's own event loop, so waiting here is a + // deadlock rather than a delay — on a WASM main thread, a page abort. + // Until morph#615 this site did not ask at all: it called the blocking + // `registerModelShared`/`registerModelWithContext` directly, so a + // backend that had just been given a way to say "do not do this to me" + // was blocked here anyway. See `IBackend::bindWaitPolicy`. + bool const mayBlock = newShared->bindWaitPolicy() == ::morph::backend::detail::BindWait::kCallerMayBlock; { // Both mutexes: this phase reads/writes every live binding's // `primary`/`contextKey` (via `_attachMtx`'s ownership of those @@ -1369,58 +1408,28 @@ class Bridge { // assignHandlerPrimary() call on any one of them. std::scoped_lock const lock{_mtx, _attachMtx}; - // Phase 1 — register every live binding on the new backend WITHOUT - // mutating any `currentId` yet, staging (binding, newId) pairs. If a - // registration throws partway (a plausible remote/transport failure), - // roll back the ones already registered and rethrow, leaving the old - // backend and every `currentId` untouched — so the switch is atomic: - // it either fully succeeds or is a no-op. - std::vector> live; - std::vector, uint64_t>> staged; - try { - for (auto& weak : _handlers) { - auto binding = weak.lock(); - if (!binding) { - continue; - } - // A shared binding that never attached has no instance to - // re-create: it stays live and unbound, and acquires one on - // the new backend the first time it is attached. - if (binding->shared && binding->primary.empty()) { - live.push_back(weak); - continue; - } - auto newId = binding->shared - ? newShared->registerModelShared( - binding->typeId, binding->modelFactory, - {.contextKey = binding->contextKey, .primary = binding->primary}) - : newShared->registerModelWithContext(binding->typeId, binding->modelFactory, - binding->contextKey); - staged.emplace_back(binding, newId.v); - live.push_back(weak); - } - } catch (...) { - for (const auto& [binding, newId] : staged) { - try { - newShared->deregisterModel(::morph::exec::detail::ModelId{newId}); - } catch (const std::exception& exc) { - ::morph::log::logError(std::string{"[switchBackend] rollback deregister failed: "} + - exc.what()); - } - } - throw; - } - - // Phase 2 — commit. Every registration succeeded, so it is now safe to - // publish the new ids and swap the backend in. - for (const auto& [binding, newId] : staged) { - binding->currentId.store(newId); + detail::StagedRebinds rebinds; + staging = stageRebinds(newShared, mayBlock, rebinds); + if (staging) { + rollbackStaged(*newShared, rebinds.staged); + settleFailed = std::move(rebinds.armed); + } else { + previous = commitRebinds(newShared, mayBlock, rebinds, settleBound); } - _handlers = std::move(live); - - previous = exchangeBackend(newShared); - newShared->notifyBackendChanged(); - installReconnectHandler(newShared); + } + for (const auto& binding : settleFailed) { + resolveRegistrationWaiters(*binding, /*ok=*/false, staging); + } + for (const auto& binding : settleBound) { + resolveRegistrationWaiters(*binding, /*ok=*/true, nullptr); + } + if (staging) { + // Rethrown here rather than from inside the locked block, so the + // waiters above are settled outside `_mtx`/`_attachMtx`. Everything + // below stays unreached on this path, exactly as before: a staging + // failure leaves the outgoing backend's reconnect handler and its + // pending completions alone, because the switch did not happen. + std::rethrow_exception(staging); } if (previous && previous != newShared) { previous->setReconnectHandler(nullptr); @@ -1542,7 +1551,17 @@ class Bridge { if (_executeDeadline.count() > 0 && _timeoutScheduler) { schedulerRef = _timeoutScheduler; // The callback captures `typedState` alone -- never `this` -- so - // it stays safe to fire even while ~Bridge() is running. + // it stays safe to fire even while ~Bridge() is running, and + // equally safe to fire *after* its own `cancel()`: + // `TimeoutScheduler::cancel` stops a callback that has not + // started but returns without waiting for one that already has + // (see that function's comment), so the disarms in the two + // continuations below are best-effort by contract and not only + // when `cancel()` throws. A deadline callback that is already + // mid-flight when the real reply lands still runs its + // `setException`, which `CompletionState` discards on an + // already-ready state -- first result wins. That, not the + // disarm, is what makes the race harmless (morph#620). deadlineHandle = schedulerRef->schedule(_executeDeadline, [typedState] { typedState->setException(std::make_exception_ptr(::morph::backend::ClientTimeoutError{})); }); @@ -1869,9 +1888,6 @@ class Bridge { binding->registrationInFlight = true; } - std::weak_ptr<::morph::backend::detail::IBackend> const weakBackend{backend}; - std::weak_ptr const weakBinding{binding}; - // `binding->contextKey` is read here without `_attachMtx`, and that is a // deliberate carve-out from the "read only under `_attachMtx`" rule the // member comment states -- not an oversight. Taking the lock here is @@ -1891,57 +1907,14 @@ class Bridge { // **set `contextKey` before calling `registerHandler()`, and do not // mutate it concurrently with that call.** After registration returns, // every access goes under `_attachMtx` as documented. morph#505. - auto onRegistered = [this, weakBackend, weakBinding, - lifetime = _lifetime](::morph::exec::detail::ModelId newId) { - auto strongBinding = weakBinding.lock(); - bool applied = false; - if (strongBinding) { - // `lifetime`'s gate held across the whole touch of `this` - // below (`_mtx`, `loadBackend()`), not just at entry -- - // `CallbackToken::active()` is advisory and cannot carry - // this weight (see `detail::BridgeLifetime`'s own doc - // comment). Safe to hold across this span, unlike - // `installReconnectHandler`'s reconnect callback - // (morph#489, site 4): nothing inside is a call into - // unbounded/consumer-supplied code, only a mutex and a - // backend-pointer comparison. - std::shared_lock const gate{lifetime->mtx}; - if (lifetime->alive) { - std::scoped_lock const lock{_mtx}; - auto pinned = weakBackend.lock(); - if (pinned && pinned == loadBackend()) { - // A switchBackend() already moved past this registration - // (see this backend's own doc comment on the class) and - // its own re-registration loop already gave `binding` a - // fresh id on the *new* backend -- applying this stale - // one now would overwrite that with a dangling id from a - // backend nothing uses any more. - strongBinding->currentId.store(newId.v); - applied = true; - } - } - } - // Resolve whenBound() waiters regardless of whether the id was - // actually applied above: either way this binding's initial - // registration attempt has settled (a stale reply ignored here - // means switchBackend's own synchronous re-registration already - // bound it), so nothing should still be described as "in - // flight". Runs even when the Bridge/binding is gone -- both - // weak locks above are only guards on touching `this`/the - // binding's other fields, not on this bookkeeping, which reads - // no Bridge state. - if (strongBinding) { - resolveRegistrationWaiters(*strongBinding, /*ok=*/applied || isBound(strongBinding), nullptr); - } - }; - auto onFailed = [weakBinding, typeId = binding->typeId](const std::string& message) { - ::morph::log::logError("[registerHandler] async registration of '" + typeId + "' failed: " + message); - if (auto strongBinding = weakBinding.lock()) { - resolveRegistrationWaiters( - *strongBinding, /*ok=*/false, - std::make_exception_ptr(std::runtime_error("registration failed: " + message))); - } - }; + // + // Both continuations come from `makeBindCallbacks`, shared with + // `switchBackend`'s phase 1 and the reconnect handler: a stale reply + // is discarded rather than allowed to overwrite the id one of those + // two just installed, and `whenBound()` waiters are settled either + // way, because a discarded reply still means this binding's own + // registration attempt is over. + auto [onRegistered, onFailed] = makeBindCallbacks(binding, backend, "[registerHandler]"); bool const started = backend->registerModelAsync(binding->typeId, binding->modelFactory, binding->contextKey, onRegistered, onFailed); @@ -2051,6 +2024,319 @@ class Bridge { } } + /// @brief Builds the two continuations every model acquisition needs when + /// its reply may arrive after the dispatching frame has gone. + /// + /// One body for the three sites that acquire an instance for a binding — + /// `registerHandlerImpl`, `switchBackend`'s phase 1 and + /// `installReconnectHandler`'s handler — because all three have the same + /// two problems, and had three chances to solve one of them differently: + /// the reply may land on a thread that does not own the `Bridge`, and it + /// may be stale by the time it lands. + /// + /// The success continuation holds `lifetime`'s gate across the whole touch + /// of `this` (`_mtx`, `loadBackend()`) rather than only checking at entry — + /// `CallbackToken::active()` is advisory and cannot carry that weight, see + /// `detail::BridgeLifetime`. Holding it across this span is safe for the + /// reason morph#489 gives: nothing inside is a call into unbounded or + /// consumer-supplied code, only a mutex and a backend-pointer comparison. + /// The id is published only while @p backend is still the active one, so a + /// reply from a backend `switchBackend()` has already replaced cannot + /// overwrite the id that switch just installed. + /// + /// Waiters are resolved either way, and outside `_mtx`: whether the id was + /// applied or discarded, this binding's registration attempt has settled, + /// and nothing should still be described as in flight. + /// @param binding Binding the reply belongs to; held weakly by both callbacks. + /// @param backend Backend the call was issued to; held weakly, and compared + /// against the active one before anything is published. + /// @param site Bracketed tag naming the caller in the failure log line. + /// @return `{onRegistered, onFailed}` — the first takes the new `ModelId`, + /// the second a diagnostic string. + [[nodiscard]] std::pair, + std::function> + makeBindCallbacks(const std::shared_ptr& binding, + const std::shared_ptr<::morph::backend::detail::IBackend>& backend, std::string_view site) { + std::weak_ptr<::morph::backend::detail::IBackend> const weakBackend{backend}; + std::weak_ptr const weakBinding{binding}; + std::function onRegistered = + [this, weakBackend, weakBinding, lifetime = _lifetime](::morph::exec::detail::ModelId newId) { + auto strongBinding = weakBinding.lock(); + bool applied = false; + if (strongBinding) { + std::shared_lock const gate{lifetime->mtx}; + if (lifetime->alive) { + std::scoped_lock const lock{_mtx}; + auto pinned = weakBackend.lock(); + if (pinned && pinned == loadBackend()) { + strongBinding->currentId.store(newId.v); + applied = true; + } + } + } + if (strongBinding) { + resolveRegistrationWaiters(*strongBinding, /*ok=*/applied || isBound(strongBinding), nullptr); + } + }; + std::function onFailed = [weakBinding, tag = std::string{site}, + typeId = binding->typeId](const std::string& message) { + ::morph::log::logError(tag + " registration of '" + typeId + "' failed: " + message); + if (auto strongBinding = weakBinding.lock()) { + resolveRegistrationWaiters( + *strongBinding, /*ok=*/false, + std::make_exception_ptr(std::runtime_error("registration failed: " + message))); + } + }; + return {std::move(onRegistered), std::move(onFailed)}; + } + + /// @brief Acquires an instance for @p binding on @p backend through + /// `IBackend::bindModel`, waiting only if @p mayBlock says it may. + /// + /// The one dispatch shape `switchBackend`'s staging phase and the reconnect + /// handler share (morph#615). A non-empty `primary` with a zero `current` + /// is the request shape that means `registerModelShared`, an empty one the + /// shape that means `registerModelWithContext` — the two blocking verbs + /// both sites called directly until morph#615 — so a backend with only the + /// default `bindModel` runs exactly the call it always did and settles + /// before this returns. `inlineExecutor()` because that is where these + /// continuations ran before; see `attachHandlerAsync`. + /// + /// A reply that lands inside this frame is parked rather than acted on, + /// because both callers hold `_mtx`/`_attachMtx` here and the continuation + /// takes `_mtx` itself. + /// @param binding Binding to acquire an instance for. + /// @param backend Backend to acquire it from. + /// @param mayBlock `true` to wait the bind out (`bindWaitPolicy()` said the + /// caller may), `false` to take only what has already landed. + /// @param site Bracketed tag naming the caller in a failure log line. + /// @return The outcome if it settled in this frame; `std::nullopt` if the + /// reply is still to come, which only @p mayBlock `== false` allows. + [[nodiscard]] std::optional rebindThroughSurface( + const std::shared_ptr& binding, + const std::shared_ptr<::morph::backend::detail::IBackend>& backend, bool mayBlock, std::string_view site) { + auto handoff = std::make_shared(); + auto [onBound, onFailed] = makeBindCallbacks(binding, backend, site); + auto completion = backend->bindModel( + ::morph::backend::detail::BindRequest{.typeId = binding->typeId, + .factory = binding->modelFactory, + .contextKey = binding->contextKey, + .primary = binding->shared ? binding->primary : std::string{}, + .current = {}}, + ::morph::exec::detail::inlineExecutor()); + completion + .then([onBound, handoff](::morph::exec::detail::ModelId newId) { + if (detail::parkIfInFrame(*handoff, true, newId, nullptr)) { + return; // Settled inline: the dispatching frame owns the outcome. + } + onBound(newId); + }) + .onError([onFailed, handoff](const std::exception_ptr& failure) { + if (detail::parkIfInFrame(*handoff, false, {}, failure)) { + return; + } + onFailed(detail::describeFailure(failure)); + }); + return mayBlock ? detail::awaitHandoff(*handoff) : detail::claimHandoff(*handoff); + } + + /// @brief Marks @p binding's registration in flight, so `whenBound()` can + /// gate on the window a deferred re-registration opens. + /// + /// Set *before* the dispatch, not after: the reply can land the moment + /// `bindModel` is called, and a `whenBound()` racing in from another thread + /// must see "in flight" for the whole window the bind can resolve in. + /// @param binding Binding to arm. + static void armRegistration(detail::HandlerBinding& binding) { + std::scoped_lock const guard{binding.registrationMtx}; + binding.registrationInFlight = true; + } + + /// @brief `switchBackend`'s phase 1: acquire an instance for every live + /// binding on @p newBackend without publishing anything yet. + /// + /// Caller holds `_mtx` and `_attachMtx`. Reports failure by returning the + /// exception rather than throwing it, so the caller keeps everything + /// @p out accumulated and can roll it back. + /// + /// Atomicity is exactly as strong as the wait is. For a `kCallerMayBlock` + /// backend — everything in the tree that is not a + /// `SynchronousBackendAdapter` or a WASM-configured `QtWebSocketBackend` — + /// this frame waits out each bind, so every outcome is known before + /// anything is published and the guarantee is what it always was. For a + /// `kCallerMustNotBlock` backend it cannot be: "did every re-registration + /// succeed" is not knowable without waiting, and waiting is the deadlock + /// the policy exists to prevent. Those binds land in `out.deferred`. + /// @param newBackend Backend being switched to. + /// @param mayBlock Whether this frame may wait for each bind to settle. + /// @param out Accumulator; filled incrementally, valid on failure too. + /// @return Null on success, or the failure that ended staging. + [[nodiscard]] std::exception_ptr stageRebinds( + const std::shared_ptr<::morph::backend::detail::IBackend>& newBackend, bool mayBlock, + detail::StagedRebinds& out) { + try { + for (auto& weak : _handlers) { + auto binding = weak.lock(); + if (!binding) { + continue; + } + // A shared binding that never attached has no instance to + // re-create: it stays live and unbound, and acquires one on the + // new backend the first time it is attached. + if (binding->shared && binding->primary.empty()) { + out.live.push_back(weak); + continue; + } + // Only armed on the path that can actually defer — a waiting + // frame settles every outcome itself, so `registrationInFlight` + // stays what it was before morph#615 for every backend that + // lets this frame wait. + if (!mayBlock) { + armRegistration(*binding); + out.armed.push_back(binding); + } + auto parked = rebindThroughSurface(binding, newBackend, mayBlock, "[switchBackend]"); + if (!parked) { + out.deferred.push_back(binding); + out.live.push_back(weak); + continue; + } + if (!parked->succeeded) { + // A rejected `Completion`, where this loop used to catch a + // throw: the structural surface reports failure through the + // completion, and a bare `catch (...)` can no longer see + // it. The rollback keys on this (morph#615). + return parked->failure; + } + out.staged.emplace_back(binding, parked->modelId.v); + out.live.push_back(weak); + } + } catch (...) { + // `bindModel` rejects rather than throws by contract, so this + // guards only what the contract does not cover — an allocation in + // this loop, or a backend that throws out of `bindModel` anyway. + return std::current_exception(); + } + return nullptr; + } + + /// @brief Undoes every instance `stageRebinds` created, after it failed. + /// + /// A deregister that itself throws is logged and the sweep continues: the + /// staging failure is the error the caller must see, not this one. + /// @param backend Backend the staged instances were created on. + /// @param staged Staged `(binding, newId)` pairs to release. + static void rollbackStaged( + ::morph::backend::detail::IBackend& backend, + const std::vector, std::uint64_t>>& staged) { + for (const auto& [binding, newId] : staged) { + try { + backend.deregisterModel(::morph::exec::detail::ModelId{newId}); + } catch (const std::exception& exc) { + ::morph::log::logError(std::string{"[switchBackend] rollback deregister failed: "} + exc.what()); + } + } + } + + /// @brief `switchBackend`'s phase 2: publish the staged ids and swap the + /// backend in. Caller holds `_mtx` and `_attachMtx`. + /// + /// A deferred binding is published as **unbound**: the id it held belonged + /// to the backend being replaced and means nothing on the new one, while + /// the new id is not here yet. Unbound is the honest state, and it is the + /// one `executeVia` fails fast on and `whenBound()` gates on — strictly + /// better than a dangling id from a backend nothing uses any more, which + /// `executeVia` would happily dispatch against. + /// @param newBackend Backend to install. + /// @param mayBlock Whether staging was allowed to wait; when it was, no + /// binding was armed and none needs settling. + /// @param rebinds What staging accumulated. + /// @param settleBound Out: bindings whose waiters the caller must resolve + /// with `true`, once it has released the mutexes. + /// @return The backend just replaced, for the caller to retire. + [[nodiscard]] std::shared_ptr<::morph::backend::detail::IBackend> commitRebinds( + const std::shared_ptr<::morph::backend::detail::IBackend>& newBackend, bool mayBlock, + detail::StagedRebinds& rebinds, std::vector>& settleBound) { + for (const auto& [binding, newId] : rebinds.staged) { + binding->currentId.store(newId); + } + for (const auto& binding : rebinds.deferred) { + binding->currentId.store(0); + } + _handlers = std::move(rebinds.live); + + auto previous = exchangeBackend(newBackend); + newBackend->notifyBackendChanged(); + installReconnectHandler(newBackend); + if (!mayBlock) { + settleBound.reserve(rebinds.staged.size()); + for (const auto& [binding, newId] : rebinds.staged) { + settleBound.push_back(binding); + } + } + return previous; + } + + /// @brief Re-acquires an instance for every live binding on a backend that + /// has just reconnected. Caller holds `_mtx` and `_attachMtx`. + /// + /// The reconnect counterpart of `stageRebinds`, and deliberately not the + /// same function: there is nothing to stage here. The connection those ids + /// belonged to is gone either way, so each binding is published as it + /// settles and a failure costs only that binding. + /// @param pinned The reconnected backend, already checked to be the active one. + /// @param mayBlock Whether this frame may wait for each bind to settle. + /// @return `(binding, failure)` pairs whose waiters the caller must resolve + /// once it has released the mutexes; a null failure means success. + [[nodiscard]] std::vector, std::exception_ptr>> reregisterLive( + const std::shared_ptr<::morph::backend::detail::IBackend>& pinned, bool mayBlock) { + std::vector, std::exception_ptr>> settle; + for (auto& weak : _handlers) { + auto binding = weak.lock(); + if (!binding) { + continue; + } + // Same branch as switchBackend()'s staging phase: a shared binding + // that never attached has no instance to re-create on the + // reconnected backend either, and a shared binding that did attach + // must come back through the `registerModelShared` request shape (a + // non-empty `primary`), or it re-registers as a non-shared instance + // and silently drops the sharing. + if (binding->shared && binding->primary.empty()) { + continue; + } + if (!mayBlock) { + armRegistration(*binding); + } + auto parked = rebindThroughSurface(binding, pinned, mayBlock, "[reconnect]"); + if (!parked) { + // The reply is still to come and will publish itself. The id + // this binding holds belongs to the connection that just + // dropped, so it is cleared rather than left to be dispatched + // against: `executeVia` fails fast on an unbound handler, and + // `whenBound()` — which this path is the first re-registration + // ever to arm — lets a caller wait the reconnect out instead. + binding->currentId.store(0); + continue; + } + if (!parked->succeeded) { + // Before morph#615 a failing re-registration threw out of the + // handler and onto the transport thread, taking every binding + // after it with it. A rejected `Completion` is reported per + // binding instead: this one is left unbound and the loop + // carries on. + binding->currentId.store(0); + settle.emplace_back(binding, parked->failure); + continue; + } + binding->currentId.store(parked->modelId.v); + if (!mayBlock) { + settle.emplace_back(binding, nullptr); + } + } + return settle; + } + std::shared_ptr<::morph::backend::detail::IBackend> exchangeBackend( std::shared_ptr<::morph::backend::detail::IBackend> next) { std::scoped_lock const lock{_backendMtx}; @@ -2077,33 +2363,34 @@ class Bridge { return; // The Bridge is gone; do not touch `this`. } auto pinned = weakBackend.lock(); - // Both mutexes: reads `_handlers` (guarded by `_mtx`) and each - // binding's `contextKey` (guarded by `_attachMtx`), same - // reasoning as switchBackend() above. - std::scoped_lock const lock{_mtx, _attachMtx}; - if (!pinned || pinned != loadBackend()) { - return; // We've moved on to a different backend; ignore. - } - for (auto& weak : _handlers) { - auto binding = weak.lock(); - if (!binding) { - continue; + // Settled after the mutexes are released, for the reason + // switchBackend() gives: resolving a waiter runs consumer code. + std::vector, std::exception_ptr>> settle; + { + // Both mutexes: reads `_handlers` (guarded by `_mtx`) and each + // binding's `contextKey` (guarded by `_attachMtx`), same + // reasoning as switchBackend() above. + std::scoped_lock const lock{_mtx, _attachMtx}; + if (!pinned || pinned != loadBackend()) { + return; // We've moved on to a different backend; ignore. } - // Same branch as switchBackend()'s phase 1: a shared binding - // that never attached has no instance to re-create on the - // reconnected backend either, and a shared binding that did - // attach must come back through registerModelShared (not - // registerModelWithContext), or it re-registers as a - // non-shared instance and silently drops the sharing. - if (binding->shared && binding->primary.empty()) { - continue; + // The same question switchBackend() now asks, for the same + // reason, and it matters more here: this handler runs on the + // backend's *transport* thread, and for a `QtWebSocketBackend` + // with `asyncRegistrationEnabled` that thread is the one whose + // event loop has to deliver the reply. Calling the blocking + // verb here — which is what this loop did until morph#615 — + // parks it against itself, and on a WASM main thread aborts + // the page. See `IBackend::bindWaitPolicy` and morph#568. + settle = reregisterLive( + pinned, pinned->bindWaitPolicy() == ::morph::backend::detail::BindWait::kCallerMayBlock); + } + for (const auto& [binding, failure] : settle) { + if (failure) { + ::morph::log::logError("[reconnect] registration of '" + binding->typeId + + "' failed: " + detail::describeFailure(failure)); } - auto newId = binding->shared ? pinned->registerModelShared( - binding->typeId, binding->modelFactory, - {.contextKey = binding->contextKey, .primary = binding->primary}) - : pinned->registerModelWithContext(binding->typeId, binding->modelFactory, - binding->contextKey); - binding->currentId.store(newId.v); + resolveRegistrationWaiters(*binding, /*ok=*/failure == nullptr, failure); } }); } diff --git a/include/morph/core/remote.hpp b/include/morph/core/remote.hpp index 51ba059d6..ec8b42358 100644 --- a/include/morph/core/remote.hpp +++ b/include/morph/core/remote.hpp @@ -1502,6 +1502,15 @@ class RemoteServer : public std::enable_shared_from_this { // typo-drift risk a shared constant guards against doesn't // apply here the way it does for the consumer-side comparison // in `core/detail/reply_router.hpp`. + // + // Safe to fire after its own `cancel()`, which the strand task + // below issues on both its success and its failure arm: + // `TimeoutScheduler::cancel` stops a callback that has not + // started but returns without waiting for one that already has + // (see that function's comment). A timeout callback already + // mid-flight when the dispatch finishes therefore still calls + // `complete`, and `complete`'s reply-exactly-once flag drops + // it rather than double-answering the call (morph#620). timeoutHandle = _timeoutScheduler->schedule(limits.executeTimeout, [complete, callId]() mutable { complete(::morph::wire::encode(::morph::wire::makeErr("timeout", callId))); }); diff --git a/include/morph/core/timeout_scheduler.hpp b/include/morph/core/timeout_scheduler.hpp index 550d3d6ff..ecbb47360 100644 --- a/include/morph/core/timeout_scheduler.hpp +++ b/include/morph/core/timeout_scheduler.hpp @@ -43,6 +43,15 @@ /// the underlying browser timer is not itself cleared — it still fires at /// its original deadline and finds nothing to do. Only a small ticket /// allocation outlives `cancel()`, until that point. +/// - **Cancelling a callback that has *already started*.** Threaded build: +/// `cancel()` cannot stop it. `run()` erases the entry before invoking the +/// callback and drops `_mtx` across the invocation, so a `cancel()` racing a +/// firing callback takes the same not-found branch as one for a handle that +/// already finished, and returns while that callback is still executing on +/// the scheduler thread. Browser build: the case cannot arise — the timer +/// callback and `cancel()` run on the same single thread, so "no callback +/// will start after `cancel()` returns" holds there and only there. See +/// `cancel()`'s own comment for what this asks of a caller. /// - **Destruction.** Threaded build: the destructor joins its thread, so no /// callback can be in flight afterwards. Browser build: nothing to join; /// pending browser timers observe an expired `std::weak_ptr` to the @@ -122,12 +131,38 @@ class TimeoutScheduler { return handle; } - /// @brief Cancels a previously scheduled callback immediately. + /// @brief Cancels a previously scheduled callback: stops one that has not + /// started, and returns without waiting for one that has. + /// + /// Two cases, and telling them apart is the caller's business because the + /// scheduler cannot: + /// + /// - **@p handle has not started.** Its entry — and anything its callback + /// captured — is erased right away, the callback never runs, and the + /// caller does not have to wait for the original deadline for that memory + /// to be released. + /// - **@p handle is already running.** `run()` erases the entry *before* it + /// invokes the callback, so this call finds nothing, takes the same + /// no-op branch as a handle that already finished, and **returns while + /// the callback is still executing** on the scheduler thread. The + /// callback is neither interrupted nor waited for. /// - /// If @p handle has not fired yet, its entry (and anything its callback - /// captured) is erased right away — the caller does not have to wait for - /// the original deadline for that memory to be released. A no-op if - /// @p handle already fired or was already cancelled. + /// So `cancel()` returning does **not** mean "no callback is in flight". + /// The only thing in this class that means that is `~TimeoutScheduler`, + /// which joins the scheduler thread. A caller must therefore keep every + /// scheduled callback safe to run *after* its `cancel()`: both callbacks in + /// this repository (`Bridge::executeVia`'s deadline and `RemoteServer`'s + /// `LimitPolicy::executeTimeout`) capture a `shared_ptr` to the state they + /// settle and settle it write-once, so a late run is an ignored duplicate + /// rather than a use-after-free. + /// + /// Blocking here until the callback finished would be the wrong contract + /// rather than a missing feature: a callback that posts back to the + /// cancelling thread would deadlock it — the hazard + /// `docs/spec/concurrency_and_lifetimes.md` names, and the same reason + /// `CallbackScope` deliberately offers no block-until-drained. + /// + /// A no-op if @p handle already fired, is firing, or was already cancelled. /// @param handle Handle returned by a prior `schedule()` call. void cancel(Handle handle) { std::scoped_lock const lock{_mtx}; @@ -238,6 +273,13 @@ class TimeoutScheduler { /// build. The browser timer itself is left to elapse and find nothing — /// see the `@file` comment. A no-op if @p handle already fired or was /// already cancelled. + /// + /// Unlike the threaded build, "already fired" here can only mean + /// *finished*: `fire()` and this function run on the same single thread, so + /// a callback cannot be mid-flight while `cancel()` is called. This build + /// therefore does give the guarantee the threaded one does not — no + /// callback runs after `cancel()` returns — and a caller that must work in + /// both builds still cannot rely on it. /// @param handle Handle returned by a prior `schedule()` call. void cancel(Handle handle) { _state->pending.erase(handle); } diff --git a/scripts/branch_partial_allowlist.json b/scripts/branch_partial_allowlist.json index db4d16423..1d48d1048 100644 --- a/scripts/branch_partial_allowlist.json +++ b/scripts/branch_partial_allowlist.json @@ -80,7 +80,7 @@ }, { "file": "include/morph/core/backend.hpp", - "line": 1324, + "line": 1414, "source": "if (const auto* inst = _instances.find(modelId)) {", "reason": "Unreachable by construction given the `_changeAware`/`_instances` invariant (core audit finding BK2). `_changeAware` is an index over the instance directory: an id enters it in `createHolder` (this file, when the holder answers `isBackendChangeAware()`) in the same `_regMtx`-held critical section that files the instance, and leaves it in `deregisterModel` only when `InstanceDirectory::release` reports the instance actually destroyed. `notifyBackendChanged()` (this function) holds the same `_regMtx` while walking `_changeAware` and looking each id up at this line, so every id it walks is still live -- the null arm cannot occur without a code change that breaks that subset invariant. Formerly keyed on `_models`, the map morph#523 replaced with the directory; the invariant and its reason are unchanged." }, @@ -98,15 +98,15 @@ }, { "file": "include/morph/core/bridge.hpp", - "line": 1542, + "line": 1551, "source": "if (_executeDeadline.count() > 0 && _timeoutScheduler) {", "reason": "Unreachable by construction (core audit finding B6). `setExecuteDeadline` (this file) is the only writer of both `_executeDeadline` and `_timeoutScheduler`, and always creates `_timeoutScheduler` in the same call that sets `_executeDeadline` positive (`_executeDeadline = deadline; if (_executeDeadline.count() > 0 && !_timeoutScheduler) { _timeoutScheduler = std::make_shared<...>(); }`, both under `_executeDeadlineMtx`); nothing anywhere resets `_timeoutScheduler` back to null -- the class's own doc comment on `setExecuteDeadline` says so explicitly (\"setting the deadline back to 0 stops new calls from arming it but does not tear the thread down\"). So `_executeDeadline > 0 && !_timeoutScheduler` cannot happen at this line once any positive deadline has ever been set." }, { "file": "include/morph/core/bridge.hpp", - "line": 1654, + "line": 1673, "source": "if (deadlineHandle && schedulerRef) {", - "reason": "Unreachable by construction, same joint-assignment shape as B6 above (core audit finding B11, reclassified (a)->(b) on review). `deadlineHandle` and `schedulerRef` are assigned together, a few lines above this one in `executeVia`, only inside `if (_executeDeadline.count() > 0 && _timeoutScheduler) { schedulerRef = _timeoutScheduler; ... }` (see the bridge.hpp:1542 entry above) -- there is no path that sets `deadlineHandle` without also having set `schedulerRef` from the same non-null `_timeoutScheduler` in the same conditional. So `schedulerRef` null while `deadlineHandle` is non-null cannot occur; the only theoretically-open arm this compound condition has is structurally impossible. This entry is specifically the exception-path use of the guard (the `catch` block that undoes `_pendingCalls` and cancels the deadline before rethrowing). Two more textually-identical `if (deadlineHandle && schedulerRef)` guards exist further down, in the `.then()`/`.onError()` continuations (lines 1678, 1767) -- present on master too, this PR only shifted their line numbers -- which are a different guard on a different, reachable arm (schedulerRef going null between capture and a callback that can run after `~Bridge()`) and are not covered by this disposition." + "reason": "Unreachable by construction, same joint-assignment shape as B6 above (core audit finding B11, reclassified (a)->(b) on review). `deadlineHandle` and `schedulerRef` are assigned together, a few lines above this one in `executeVia`, only inside `if (_executeDeadline.count() > 0 && _timeoutScheduler) { schedulerRef = _timeoutScheduler; ... }` (see the bridge.hpp:1551 entry above) -- there is no path that sets `deadlineHandle` without also having set `schedulerRef` from the same non-null `_timeoutScheduler` in the same conditional. So `schedulerRef` null while `deadlineHandle` is non-null cannot occur; the only theoretically-open arm this compound condition has is structurally impossible. This entry is specifically the exception-path use of the guard (the `catch` block that undoes `_pendingCalls` and cancels the deadline before rethrowing). Two more textually-identical `if (deadlineHandle && schedulerRef)` guards exist further down, in the `.then()`/`.onError()` continuations (lines 1697, 1786) -- present on master too, this PR only shifted their line numbers -- which are a different guard on a different, reachable arm (schedulerRef going null between capture and a callback that can run after `~Bridge()`) and are not covered by this disposition." }, { "file": "include/morph/core/remote.hpp", diff --git a/scripts/mutation_survivors.json b/scripts/mutation_survivors.json index 9a54616b0..e659f1d77 100644 --- a/scripts/mutation_survivors.json +++ b/scripts/mutation_survivors.json @@ -81,7 +81,7 @@ }, { "file": "include/morph/core/backend.hpp", - "line": 1322, + "line": 1412, "mutants": 1, "mutator": "cxx_replace_scalar_call", "source": "aware.reserve(_changeAware.size());", @@ -147,13 +147,13 @@ "representative_sites": [ { "file": "include/morph/core/backend.hpp", - "line": 1215, + "line": 1305, "source": "::morph::observe::detail::emitMetric(::morph::observe::Metric::registerCount, 1.0);", "reason": "The registerCount emission on LocalBackend::registerModel. Note for whoever refreshes this hint: the same statement appears character-for-character on the registerModelShared arm as well, so if this `line` ever drifts the gate will report the citation as ambiguous rather than printing a corrected line. That is the right outcome -- which of the two arms is meant is a question for a reader, not for a resolver -- and it is recorded here so the message is not a surprise." }, { "file": "include/morph/core/backend.hpp", - "line": 1379, + "line": 1469, "source": "::morph::observe::detail::emitMetric(::morph::observe::Metric::executeInFlight,", "reason": "The executeInFlight emission on the increment side of an execute. Its decrement twin inside the posted task is the same text after stripping, so the ambiguity note on the registerCount entry above applies here too." } diff --git a/tests/test_backend_registration_surface.cpp b/tests/test_backend_registration_surface.cpp index eba4df775..4085ba4e2 100644 --- a/tests/test_backend_registration_surface.cpp +++ b/tests/test_backend_registration_surface.cpp @@ -679,3 +679,133 @@ TEST_CASE("morph::bridge::Bridge: registerHandler does not wait for a kCallerMus REQUIRE(binding->currentId.load() == 99U); REQUIRE(backend->settledOffCallerThread.load()); } + +// ── cancelPending and the completions the adapter itself produced (#619) ───── +// +// `SynchronousBackendAdapter::cancelPending` used to be a one-line forward to +// the wrapped backend. The two verbs the adapter *reshapes* settle from a task +// on its own private control strand, so the forward reached none of them: a +// bind cancelled mid-flight went on to resolve **successfully**, which is the +// opposite of `IBackend::cancelPending`'s contract. The case below holds the +// wrapped call open so the completion is provably still pending when +// `cancelPending` runs, and fails if the adapter goes back to forwarding only +// -- the success continuation runs instead of the error one. + +namespace { + +/// @brief A wrapped backend whose blocking control calls do not return until +/// the test lets them, so a cancellation can land while one is in +/// flight. +/// +/// Every field the test thread reads is atomic: `RecordingBackend::calls` is a +/// plain vector appended from the strand thread, so this double counts through +/// atomics instead and leaves that vector to the strand alone. +struct GatedBackend : RecordingBackend { + std::mutex mtx; + std::condition_variable gate; + bool released = false; + /// @brief Control calls that have entered the blocking region. + std::atomic entered{0}; + /// @brief Control calls that have come back out of it. + std::atomic finished{0}; + /// @brief Times the adapter forwarded a cancellation here. + std::atomic cancels{0}; + + /// @brief Blocks the calling (strand) thread until `letGo()`. + void hold() { + entered.fetch_add(1); + std::unique_lock lock{mtx}; + (void)gate.wait_for(lock, morph::testing::kDefaultWaitBudget, [this] { return released; }); + } + + /// @brief Releases every held control call. + void letGo() { + { + std::scoped_lock const lock{mtx}; + released = true; + } + gate.notify_all(); + } + + ModelId registerModelWithContext(const std::string& /*typeId*/, + std::function()> /*factory*/, + std::string_view /*contextKey*/) override { + hold(); + finished.fetch_add(1); + return ModelId{2}; + } + + void assignPrimary(ModelId /*mid*/, const std::string& /*typeId*/, std::string_view /*primary*/) override { + hold(); + finished.fetch_add(1); + } + + void cancelPending(const std::exception_ptr& /*exc*/) override { cancels.fetch_add(1); } +}; + +} // namespace + +TEST_CASE("morph::backend::SynchronousBackendAdapter: cancelPending rejects the completions it produced itself", + "[backend][registration-surface][threading]") { + morph::exec::ThreadPoolExecutor pool{1}; + morph::exec::MainThreadExecutor callerExec; + auto inner = std::make_shared(); + SynchronousBackendAdapter adapter{inner, pool}; + + std::atomic okRan{0}; + std::atomic errRan{0}; + std::string message; + + auto attach = [&](ModelCompletion completion) { + completion.then([&](ModelId /*mid*/) { okRan.fetch_add(1); }).onError([&](const std::exception_ptr& exc) { + try { + std::rethrow_exception(exc); + } catch (const std::exception& err) { + message = err.what(); + } + errRan.fetch_add(1); + }); + }; + + SECTION("bind") { + attach(adapter.bindModel( + BindRequest{.typeId = std::string{kTypeId}, .factory = makeHolder, .contextKey = "ctx", .primary = {}}, + callerExec)); + } + SECTION("promote") { + attach(adapter.promoteModel( + PromoteRequest{.mid = ModelId{7}, .typeId = std::string{kTypeId}, .primary = "key"}, callerExec)); + } + + // The wrapped control call is provably inside its blocking region, so the + // completion is genuinely still pending -- not merely unobserved. + REQUIRE(morph::testing::waitUntil([&] { return inner->entered.load() == 1; })); + REQUIRE(inner->finished.load() == 0); + REQUIRE(okRan.load() == 0); + REQUIRE(errRan.load() == 0); + + adapter.cancelPending(std::make_exception_ptr(morph::backend::BridgeDestroyedError{})); + + // Delivered on the caller's executor, like every other continuation from + // this adapter, and delivered *before* the wrapped call has returned. + REQUIRE(morph::testing::waitUntil([&] { + callerExec.runOnce(); + return errRan.load() == 1; + })); + REQUIRE(inner->finished.load() == 0); + REQUIRE(okRan.load() == 0); + REQUIRE(message == std::string{morph::backend::BridgeDestroyedError{}.what()}); + // ...and the wrapped backend still saw the cancellation it always saw. + REQUIRE(inner->cancels.load() == 1); + + // Letting the wrapped call finish must not resurrect the cancelled + // completion: its `resolve` finds the state already ready. + inner->letGo(); + REQUIRE(morph::testing::waitUntil([&] { + callerExec.runOnce(); + return inner->finished.load() == 1; + })); + callerExec.runOnce(); + REQUIRE(okRan.load() == 0); + REQUIRE(errRan.load() == 1); +} diff --git a/tests/test_switch_backend.cpp b/tests/test_switch_backend.cpp index 05ebef456..4e65f6117 100644 --- a/tests/test_switch_backend.cpp +++ b/tests/test_switch_backend.cpp @@ -3,6 +3,7 @@ #include #include #include +#include #include #include #include @@ -10,6 +11,7 @@ #include #include #include +#include #include #include "test_support.hpp" @@ -839,3 +841,276 @@ TEST_CASE( (void)midB; } + +// ── The two re-registration sites and `bindWaitPolicy` (morph#615) ─────────── +// +// `switchBackend`'s phase 1 and the reconnect handler both called the blocking +// `registerModelShared`/`registerModelWithContext` directly, with no policy +// check at all -- so a backend that answers `kCallerMustNotBlock`, which by +// definition delivers its reply through the calling thread's own event loop, +// was blocked at both of them anyway (on a WASM main thread: a page abort). +// Both now go through `bindModel` and ask, exactly as `registerHandlerImpl` +// does. +// +// The double below is what that costs to test honestly: its `bindModel` never +// blocks and never settles on its own, while its legacy verbs park the calling +// thread until the test releases them. So a site that went back to a blocking +// verb does not merely fail an assertion -- it fails to return, which is the +// real symptom. Each case therefore runs the site on its own thread, records +// whether it came back within the polling budget, releases the latch so the +// thread is joinable either way, and only then asserts. A regression fails the +// case; it does not wedge the suite. + +namespace { + +/// @brief A backend with a genuinely non-blocking bind and blocking legacy +/// verbs -- `QtWebSocketBackend` under `asyncRegistrationEnabled`, as +/// far as `Bridge` can tell. +class DeferredBindBackend : public morph::backend::LocalBackend { +public: + explicit DeferredBindBackend(morph::exec::IExecutor& pool MORPH_LIFETIMEBOUND) : LocalBackend{pool} {} + + /// @brief Keeps the request and its promise; settles neither. + morph::async::Completion bindModel(morph::backend::detail::BindRequest request, + morph::exec::IExecutor& cbExec) override { + auto [completion, promise] = morph::async::Completion::makeSettleable(&cbExec); + std::scoped_lock const lock{_mtx}; + _pending.emplace_back(std::make_shared(std::move(promise)), std::move(request)); + return std::move(completion); + } + + [[nodiscard]] morph::backend::detail::BindWait bindWaitPolicy() const noexcept override { + return morph::backend::detail::BindWait::kCallerMustNotBlock; + } + + // The trap. Nothing in `Bridge` should reach either of these any more; a + // site that does parks here until `letGo()`. + morph::exec::detail::ModelId registerModelWithContext( + const std::string& typeId, std::function()> factory, + std::string_view contextKey) override { + hold(); + return LocalBackend::registerModelWithContext(typeId, std::move(factory), contextKey); + } + morph::exec::detail::ModelId registerModelShared( + const std::string& typeId, std::function()> factory, + morph::backend::detail::InstanceIdentity identity) override { + hold(); + return LocalBackend::registerModelShared(typeId, std::move(factory), identity); + } + + void setReconnectHandler(const std::function& handler) override { _handler = handler; } + void fireReconnect() const { + if (_handler) { + _handler(); + } + } + + /// @brief Number of binds waiting for a reply. + [[nodiscard]] std::size_t pendingCount() const { + std::scoped_lock const lock{_mtx}; + return _pending.size(); + } + + /// @brief Answers every waiting bind for real, through the blocking + /// dispatch its own legacy verbs would have run. + void settlePending() { + std::vector, morph::backend::detail::BindRequest>> taken; + { + std::scoped_lock const lock{_mtx}; + taken.swap(_pending); + } + for (auto& [promise, request] : taken) { + try { + promise->resolve(LocalBackend::bindModelBlocking(std::move(request))); + } catch (...) { + promise->reject(std::current_exception()); + } + } + } + + /// @brief Releases anything parked in a legacy verb. + void letGo() { + { + std::scoped_lock const lock{_gateMtx}; + _released = true; + } + _gate.notify_all(); + } + +private: + using Promise = morph::async::Completion::Promise; + + // Bounded only so a regression cannot wedge the suite for ever, and + // bounded far above `kDefaultWaitBudget` on purpose: each case measures + // "did the site come back within the polling budget", so a parked legacy + // verb must still be parked when that budget runs out. A bound near the + // budget would make the measurement a coin flip -- observed, while + // mutation-testing morph#615: with both set to two seconds the mutated + // (blocking) build passed. + void hold() { + std::unique_lock lock{_gateMtx}; + (void)_gate.wait_for(lock, std::chrono::seconds{60}, [this] { return _released; }); + } + + mutable std::mutex _mtx; + std::vector, morph::backend::detail::BindRequest>> _pending; + std::mutex _gateMtx; + std::condition_variable _gate; + bool _released = false; + std::function _handler; +}; + +/// @brief A backend whose second bind is **rejected**, never thrown. +/// +/// `switchBackend`'s rollback used to key on a `catch (...)` around the +/// blocking verbs. The structural surface reports failure through the +/// `Completion` instead, which a `catch` cannot see -- so this double fails +/// the way that surface actually fails, with no exception ever crossing the +/// call. +class RejectSecondBindBackend : public morph::backend::LocalBackend { +public: + explicit RejectSecondBindBackend(morph::exec::IExecutor& pool MORPH_LIFETIMEBOUND) : LocalBackend{pool} {} + + morph::async::Completion bindModel(morph::backend::detail::BindRequest request, + morph::exec::IExecutor& cbExec) override { + auto [completion, promise] = morph::async::Completion::makeSettleable(&cbExec); + if (++_calls >= 2) { + promise.reject(std::make_exception_ptr(std::runtime_error{"bind refused"})); + } else { + try { + promise.resolve(LocalBackend::bindModelBlocking(std::move(request))); + } catch (...) { + promise.reject(std::current_exception()); + } + } + return std::move(completion); + } + + void deregisterModel(morph::exec::detail::ModelId mid) override { + ++(*_deregisters); + LocalBackend::deregisterModel(mid); + } + + /// @brief Survives the backend, which `switchBackend` destroys as it unwinds. + [[nodiscard]] std::shared_ptr deregisterCounter() const { return _deregisters; } + +private: + int _calls = 0; + std::shared_ptr _deregisters{std::make_shared(0)}; +}; + +} // namespace + +TEST_CASE("morph::bridge::Bridge::switchBackend does not block on a kCallerMustNotBlock backend", + "[bridge][switch][registration-surface]") { + morph::exec::ThreadPoolExecutor poolA{2}; + morph::exec::ThreadPoolExecutor poolB{2}; + morph::exec::MainThreadExecutor waiterExec; + morph::bridge::Bridge bridge{std::make_unique(poolA)}; + + auto binding = bridge.registerHandler(); + REQUIRE(morph::bridge::Bridge::isBound(binding)); + + auto async = std::make_shared(poolB); + std::atomic returned{false}; + std::thread switcher{[&] { + bridge.switchBackend(std::static_pointer_cast(async)); + returned.store(true); + }}; + // The measurement: did the switch come back while its bind is still + // unanswered? On the blocking verbs it cannot -- nothing releases them + // until the line below, which runs either way so the thread is joinable. + bool const returnedPromptly = morph::testing::waitUntil([&] { return returned.load(); }); + async->letGo(); + switcher.join(); + REQUIRE(returnedPromptly); + + // The swap happened, and the binding is unbound rather than holding the id + // it had on the backend that is no longer there. + REQUIRE_FALSE(morph::bridge::Bridge::isBound(binding)); + REQUIRE(async->pendingCount() == 1); + + // ...and `whenBound()` says "in flight" rather than "nothing to wait for", + // which is what makes the unbound window usable by a caller. + bool bound = false; + bool settled = false; + bridge.whenBound(binding, &waiterExec) + .then([&](bool ok) { + bound = ok; + settled = true; + }) + .onError([&](const std::exception_ptr&) { settled = true; }); + waiterExec.runOnce(); + REQUIRE_FALSE(settled); + + async->settlePending(); + REQUIRE(morph::bridge::Bridge::isBound(binding)); + waiterExec.runOnce(); + REQUIRE(settled); + REQUIRE(bound); +} + +TEST_CASE("Bridge: the reconnect handler does not block on a kCallerMustNotBlock backend", + "[bridge][switch][reconnect][registration-surface]") { + morph::exec::ThreadPoolExecutor pool{2}; + auto owned = std::make_unique(pool); + auto* backend = owned.get(); + morph::bridge::Bridge bridge{std::move(owned)}; + + auto binding = bridge.registerHandler(); + REQUIRE_FALSE(morph::bridge::Bridge::isBound(binding)); // registerHandler already honours the policy + backend->settlePending(); + REQUIRE(morph::bridge::Bridge::isBound(binding)); + + // Fired from a thread standing in for the transport thread the real + // handler runs on -- the one a `QtWebSocketBackend` needs back in its event + // loop before the reply it is waiting for can possibly arrive. + std::atomic returned{false}; + std::thread transport{[&] { + backend->fireReconnect(); + returned.store(true); + }}; + bool const returnedPromptly = morph::testing::waitUntil([&] { return returned.load(); }); + backend->letGo(); + transport.join(); + REQUIRE(returnedPromptly); + + // The id it held belonged to the connection that just dropped, so it is + // cleared rather than left to be dispatched against. + REQUIRE_FALSE(morph::bridge::Bridge::isBound(binding)); + REQUIRE(backend->pendingCount() == 1); + + backend->settlePending(); + REQUIRE(morph::bridge::Bridge::isBound(binding)); +} + +TEST_CASE("morph::bridge::Bridge::switchBackend rolls back on a rejected bind, not only on a thrown one", + "[bridge][switch][rollback][registration-surface]") { + morph::exec::ThreadPoolExecutor pool{2}; + morph::exec::ThreadPoolExecutor pool2{2}; + SyncExec cbExec; + morph::bridge::Bridge bridge{std::make_unique(pool)}; + morph::bridge::BridgeHandler handler1{bridge, &cbExec}; + morph::bridge::BridgeHandler const handler2{bridge, &cbExec}; + + auto const idBefore1 = handler1.binding()->currentId.load(); + auto const idBefore2 = handler2.binding()->currentId.load(); + + auto rejecting = std::make_unique(pool2); + auto deregisters = rejecting->deregisterCounter(); // outlives the backend + + REQUIRE_THROWS_AS(bridge.switchBackend(std::move(rejecting)), std::runtime_error); + + // The staged first bind was rolled back, and nothing was published: the + // rollback keyed on a `Completion` that rejected without ever throwing, + // which the old `catch (...)` around the blocking verbs could not see. + REQUIRE(*deregisters == 1); + REQUIRE(handler1.binding()->currentId.load() == idBefore1); + REQUIRE(handler2.binding()->currentId.load() == idBefore2); + + // The old backend is still the active one. + std::atomic res{-1}; + handler1.execute(CountAction{5}).then([&](int val) { res.store(val); }).onError([](const std::exception_ptr&) {}); + REQUIRE(morph::testing::waitUntil([&] { return res.load() != -1; })); + REQUIRE(res.load() == 5); +} diff --git a/tests/test_timeout_scheduler.cpp b/tests/test_timeout_scheduler.cpp index 925c1331d..37551421b 100644 --- a/tests/test_timeout_scheduler.cpp +++ b/tests/test_timeout_scheduler.cpp @@ -98,3 +98,57 @@ TEST_CASE("TimeoutScheduler: cancel() before the deadline prevents the callback std::this_thread::sleep_for(80ms); REQUIRE_FALSE(fired.load()); } + +// ── What `cancel()` does about a callback that has already started ─────────── +// +// The header now states the distinction these two cases make (issue #620): +// `cancel()` stops a callback that has not started, and returns *without +// waiting* for one that has. Only `~TimeoutScheduler` means "no callback is in +// flight", because only it joins. Both halves are asserted below so the prose +// is measured rather than asserted: the first case fails if `cancel()` ever +// starts waiting, the second fails if the destructor ever stops joining. + +TEST_CASE("TimeoutScheduler: cancel() returns while the callback it names is still running", "[timeout_scheduler]") { + std::atomic entered{false}; + std::atomic release{false}; + std::atomic finished{false}; + + TimeoutScheduler scheduler; + auto const handle = scheduler.schedule(1ms, [&] { + entered = true; + // Bounded so a regression fails the case rather than wedging the suite. + (void)waitFor([&] { return release.load(); }, 5s); + finished = true; + }); + + REQUIRE(waitFor([&] { return entered.load(); })); + + // The callback is provably inside its body here, and `run()` erased the + // entry before invoking it -- so this takes cancel()'s not-found branch and + // returns at once. Nothing stops or waits for the callback. + scheduler.cancel(handle); + REQUIRE_FALSE(finished.load()); + + release = true; + REQUIRE(waitFor([&] { return finished.load(); })); +} + +TEST_CASE("TimeoutScheduler: the destructor -- unlike cancel() -- waits for a running callback", + "[timeout_scheduler]") { + std::atomic entered{false}; + std::atomic finished{false}; + + { + TimeoutScheduler scheduler; + scheduler.schedule(1ms, [&] { + entered = true; + std::this_thread::sleep_for(50ms); + finished = true; + }); + REQUIRE(waitFor([&] { return entered.load(); })); + } // ~TimeoutScheduler joins its thread here. + + // The contrast that makes the cancel() case above a real distinction rather + // than a timing accident: this one *is* "no callback in flight afterwards". + REQUIRE(finished.load()); +}