File
Blob: src/workerd/util/canceler.h
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #pragma once |
| 6 | #include <kj/async.h> |
| 7 | #include <kj/common.h> |
| 8 | #include <kj/debug.h> |
| 9 | #include <kj/function.h> |
| 10 | #include <kj/list.h> |
| 11 | #include <kj/refcount.h> |
| 12 | |
| 13 | namespace workerd { |
| 14 | |
| 15 | // A simple wrapper around kj::Canceler that can be safely |
| 16 | // shared by multiple objects. This is used, for instance, |
| 17 | // to support fetch() requests that use an AbortSignal. |
| 18 | // The AbortSignal (see api/basics.h) creates an instance |
| 19 | // of RefcountedCanceler then passes references to it out |
| 20 | // to various other objects that will use it to wrap their |
| 21 | // Promises. |
| 22 | class RefcountedCanceler: public kj::Refcounted { |
| 23 | public: |
| 24 | class Listener { |
| 25 | public: |
| 26 | explicit Listener(RefcountedCanceler& canceler, kj::Function<void()> fn) |
| 27 | : fn(kj::mv(fn)), |
| 28 | canceler(canceler) { |
| 29 | canceler.addListener(*this); |
| 30 | } |
| 31 | |
| 32 | ~Listener() { |
| 33 | canceler.removeListener(*this); |
| 34 | } |
| 35 | |
| 36 | private: |
| 37 | kj::Function<void()> fn; |
| 38 | RefcountedCanceler& canceler; |
| 39 | kj::ListLink<Listener> link; |
| 40 | |
| 41 | friend class RefcountedCanceler; |
| 42 | }; |
| 43 | |
| 44 | RefcountedCanceler(kj::Maybe<kj::Exception> reason = kj::none): reason(kj::mv(reason)) {} |
| 45 | |
| 46 | ~RefcountedCanceler() noexcept(false) { |
| 47 | // `listeners` has to be empty since each listener should have held a strong reference. |
| 48 | KJ_ASSERT(listeners.empty()); |
| 49 | |
| 50 | // RefcountedCanceler is used in use cases where we don't want to cancel by default if the |
| 51 | // canceler is destroyed, so release any remaining wrapped promises. |
| 52 | canceler.release(); |
| 53 | } |
| 54 | |
| 55 | KJ_DISALLOW_COPY_AND_MOVE(RefcountedCanceler); |
| 56 | |
| 57 | template <typename T> |
| 58 | kj::Promise<T> wrap(kj::Promise<T> promise) { |
| 59 | KJ_IF_SOME(ex, reason) { |
| 60 | return ex.clone(); |
| 61 | } |
| 62 | return canceler.wrap(kj::mv(promise)); |
| 63 | } |
| 64 | |
| 65 | void cancel(kj::StringPtr cancelReason) { |
| 66 | if (reason == kj::none) { |
| 67 | cancel(kj::Exception( |
| 68 | kj::Exception::Type::DISCONNECTED, __FILE__, __LINE__, kj::str(cancelReason))); |
| 69 | } |
| 70 | } |
| 71 | |
| 72 | void cancel(const kj::Exception& exception) { |
| 73 | if (reason == kj::none) { |
| 74 | reason = exception.clone(); |
| 75 | canceler.cancel(exception); |
| 76 | for (auto& listener: listeners) { |
| 77 | listener.fn(); |
| 78 | } |
| 79 | } |
| 80 | } |
| 81 | |
| 82 | bool isEmpty() const { |
| 83 | return canceler.isEmpty(); |
| 84 | } |
| 85 | |
| 86 | void throwIfCanceled() { |
| 87 | KJ_IF_SOME(ex, reason) { |
| 88 | kj::throwFatalException(ex.clone()); |
| 89 | } |
| 90 | } |
| 91 | |
| 92 | bool isCanceled() const { |
| 93 | return reason != kj::none; |
| 94 | } |
| 95 | |
| 96 | void addListener(Listener& listener) { |
| 97 | listeners.add(listener); |
| 98 | } |
| 99 | |
| 100 | void removeListener(Listener& listener) { |
| 101 | listeners.remove(listener); |
| 102 | } |
| 103 | |
| 104 | private: |
| 105 | kj::Canceler canceler; |
| 106 | kj::Maybe<kj::Exception> reason; |
| 107 | |
| 108 | kj::List<Listener, &Listener::link> listeners; |
| 109 | }; |
| 110 | |
| 111 | } // namespace workerd |