File
Blob: src/workerd/io/io-gate.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 | // An I/O gate allows someone to "lock" a type of I/O so that other concurrent tasks trying to |
| 7 | // perform that type of I/O are blocked until the lock is released. |
| 8 | // |
| 9 | // I/O gates are used in actors to implement consistency guarantees, allowing in-memory state and |
| 10 | // storage to be synchronized. |
| 11 | // |
| 12 | // Each Actor has two main gates: |
| 13 | // - Input gate: While locked, blocks all incoming I/O events of any type from being delivered to |
| 14 | // the actor, other than the specific event or events that hold the lock. This includes |
| 15 | // blocking responses to subrequests, timer events, input streams, etc. Used when storage |
| 16 | // operations are outstanding, so that awaiting a storage operation does not risk allowing |
| 17 | // concurrent events that render the state inconsistent. |
| 18 | // - Output gate: While locked, blocks all outgoing messages from an actor that would allow the |
| 19 | // rest of the world to observe the actor's state. Held while writes that have been confirmed |
| 20 | // to the application are still being flushed to disk. If the flush fails, these messages will |
| 21 | // never be sent, so that the rest of the world cannot observe a prematurely-confirmed write. |
| 22 | |
| 23 | #include <workerd/io/trace.h> |
| 24 | |
| 25 | #include <kj/async.h> |
| 26 | #include <kj/list.h> |
| 27 | #include <kj/one-of.h> |
| 28 | |
| 29 | namespace workerd { |
| 30 | |
| 31 | using kj::uint; |
| 32 | |
| 33 | // An InputGate blocks incoming events from being delivered to an actor while the lock is held. |
| 34 | class InputGate { |
| 35 | |
| 36 | public: |
| 37 | // Hooks that can be used to customize InputGate behavior. |
| 38 | // |
| 39 | // Technically, everything implemented here could be accomplished by a class that wraps |
| 40 | // InputGate, but the part of the code that wants to implement these hooks (Worker::Actor) |
| 41 | // is far away from the part of the code that calls into the InputGate (ActorCache), and so |
| 42 | // it was more convenient to give Worker::Actor a way to inject behavior into InputGate which |
| 43 | // would kick in when ActorCache tried to use it. |
| 44 | class Hooks { |
| 45 | |
| 46 | public: |
| 47 | // Optionally track metrics. In practice these are implemented by MetricsCollector::Actor, but |
| 48 | // we don't want to depend on that class from here. |
| 49 | virtual void inputGateLocked() {} |
| 50 | virtual void inputGateReleased() {} |
| 51 | virtual void inputGateWaiterAdded() {} |
| 52 | virtual void inputGateWaiterRemoved() {} |
| 53 | |
| 54 | static const Hooks DEFAULT; |
| 55 | }; |
| 56 | |
| 57 | // Hooks has no member variables, so const_cast is acceptable. |
| 58 | InputGate(Hooks& hooks = const_cast<Hooks&>(Hooks::DEFAULT)); |
| 59 | ~InputGate() noexcept; |
| 60 | |
| 61 | class CriticalSection; |
| 62 | |
| 63 | // A lock that blocks all new events from being delivered while it exists. |
| 64 | class Lock { |
| 65 | public: |
| 66 | KJ_DISALLOW_COPY(Lock); |
| 67 | Lock(Lock&& other) noexcept |
| 68 | : gate(other.gate), |
| 69 | cs(kj::mv(other.cs)), |
| 70 | lockSpan(kj::mv(other.lockSpan)) { |
| 71 | other.gate = nullptr; |
| 72 | } |
| 73 | ~Lock() noexcept(false) { |
| 74 | if (gate != nullptr) { |
| 75 | lockSpan.setTag("waiters"_kjc, static_cast<int64_t>(gate->waiters.size())); |
| 76 | lockSpan.setTag("lock_count"_kjc, static_cast<int64_t>(gate->lockCount)); |
| 77 | gate->releaseLock(); |
| 78 | } |
| 79 | } |
| 80 | |
| 81 | // Increments the lock's refcount, returning a duplicate `Lock`. All `Lock`s must be dropped |
| 82 | // before the gate is unlocked. |
| 83 | Lock addRef(SpanParent parentSpan) { |
| 84 | return Lock(*gate, kj::mv(parentSpan)); |
| 85 | } |
| 86 | |
| 87 | // Start a new critical section from this lock. After `wait()` has been called on the returned |
| 88 | // critical section for the first time, no further Locks will be handed out by |
| 89 | // InputGate::wait() until the CriticalSection has been dropped. |
| 90 | // |
| 91 | // CriticalSections can be nested. If this Lock is itself part of a CriticalSection, the new |
| 92 | // CriticalSection will be nested within it and the outer CriticalSection's wait() won't |
| 93 | // produce a Lock again until the inner CriticalSection is dropped. |
| 94 | kj::Own<CriticalSection> startCriticalSection(); |
| 95 | |
| 96 | // If this lock was taken in a CriticalSection, return it. |
| 97 | kj::Maybe<CriticalSection&> getCriticalSection(); |
| 98 | |
| 99 | bool isFor(const InputGate& gate) const; |
| 100 | |
| 101 | inline bool operator==(const Lock& other) const { |
| 102 | return gate == other.gate; |
| 103 | } |
| 104 | |
| 105 | private: |
| 106 | // Becomes null on move. |
| 107 | InputGate* gate; |
| 108 | |
| 109 | kj::Maybe<kj::Own<CriticalSection>> cs; |
| 110 | |
| 111 | SpanBuilder lockSpan; |
| 112 | |
| 113 | Lock(InputGate& gate, SpanParent parentSpan); |
| 114 | friend class InputGate; |
| 115 | }; |
| 116 | |
| 117 | // Wait until there are no `Lock`s, then create a new one and return it. |
| 118 | // |
| 119 | // If parentSpan is provided, child spans will be created to track: |
| 120 | // - Time spent waiting for the lock (if waiting is required) |
| 121 | // - Time spent holding the lock |
| 122 | kj::Promise<Lock> wait(SpanParent parentSpan); |
| 123 | |
| 124 | // Rejects if and when calls to `wait()` become broken due to a failed critical section. The |
| 125 | // actor should be shut down in this case. This promise never resolves, only rejects. |
| 126 | kj::Promise<void> onBroken(); |
| 127 | |
| 128 | private: |
| 129 | Hooks& hooks; |
| 130 | |
| 131 | // How many instances of `Lock` currently exist? When this reaches zero, we'll release some |
| 132 | // waiters. |
| 133 | uint lockCount = 0; |
| 134 | |
| 135 | // CriticalSection inherits InputGate for implementation convenience (since much implementation |
| 136 | // is shared). |
| 137 | bool isCriticalSection = false; |
| 138 | |
| 139 | struct Waiter { |
| 140 | Waiter(kj::PromiseFulfiller<Lock>& fulfiller, |
| 141 | InputGate& gate, |
| 142 | bool isChildWaiter, |
| 143 | SpanParent parentSpan); |
| 144 | ~Waiter() noexcept(false); |
| 145 | |
| 146 | kj::PromiseFulfiller<Lock>& fulfiller; |
| 147 | InputGate* gate; |
| 148 | bool isChildWaiter; |
| 149 | kj::ListLink<Waiter> link; |
| 150 | |
| 151 | // Span tracking how long we wait for the lock. Ends when Waiter is destroyed. |
| 152 | SpanBuilder waitSpan; |
| 153 | // Parent span to pass to Lock when it's created. |
| 154 | SpanParent lockSpanParent; |
| 155 | }; |
| 156 | |
| 157 | kj::List<Waiter, &Waiter::link> waiters; |
| 158 | |
| 159 | // Waiters representing CriticalSections that are ready to start. These take priority over other |
| 160 | // waiters. |
| 161 | kj::List<Waiter, &Waiter::link> waitingChildren; |
| 162 | |
| 163 | // A fulfiller for onBroken(), or an exception if already broken. |
| 164 | kj::ForkedPromise<void> brokenPromise; |
| 165 | kj::OneOf<kj::Own<kj::PromiseFulfiller<void>>, kj::Exception> brokenState; |
| 166 | |
| 167 | void releaseLock(); |
| 168 | |
| 169 | // Called when a critical section fails. All future waiters will throw this exception. |
| 170 | void setBroken(const kj::Exception& e); |
| 171 | |
| 172 | InputGate(Hooks& hooks, kj::PromiseFulfillerPair<void> paf); |
| 173 | }; |
| 174 | |
| 175 | // A CriticalSection is a procedure that must not be interrupted by anything "external". |
| 176 | // While a CriticalSection is running, all events that were not initiated by the |
| 177 | // CriticalSection itself will be blocked from being delivered. |
| 178 | // |
| 179 | // The difference between a Lock and a CriticalSection is that a critical section may succeed |
| 180 | // or fail. A failed critical section permanently breaks the input gate. Locks, on the other |
| 181 | // hand, are simply released when dropped. |
| 182 | // |
| 183 | // A CriticalSection itself holds a Lock, which blocks the "parent scope" from continuing |
| 184 | // execution until the critical section is done. Meanwhile, the code running inside the critical |
| 185 | // section obtains nested Locks. These nested locks control concurrency of the operations |
| 186 | // initiated within the critical section in the same way that input locks normally do at the |
| 187 | // top-level scope. E.g., if a critical section initiates a storage read and a fetch() at the |
| 188 | // same time, the fetch() is prevented from returning until after the storage read has returned. |
| 189 | class InputGate::CriticalSection: private InputGate, public kj::Refcounted { |
| 190 | public: |
| 191 | CriticalSection(InputGate& parent); |
| 192 | ~CriticalSection() noexcept(false); |
| 193 | |
| 194 | // Wait for a nested lock in order to continue this CriticalSection. |
| 195 | // |
| 196 | // The first call to wait() begins the CriticalSection. After that wait completes, until the |
| 197 | // CriticalSection is done and dropped, no other locks will be allowed on this InputGate, except |
| 198 | // locks requested by calling wait() on this CriticalSection -- or one of its children. |
| 199 | kj::Promise<Lock> wait(SpanParent parentSpan); |
| 200 | |
| 201 | // Call when the critical section has completed successfully. If this is not called before the |
| 202 | // CriticalSection is dropped, then failed() is called implicitly. |
| 203 | // |
| 204 | // Returns the input lock that was held on the parent critical section. This can be used to |
| 205 | // continue execution in the parent before any other input arrives. |
| 206 | Lock succeeded(); |
| 207 | |
| 208 | // Call to indicate the CriticalSection has failed with the given exception. This immediately |
| 209 | // breaks the InputGate. |
| 210 | void failed(const kj::Exception& e); |
| 211 | |
| 212 | private: |
| 213 | enum State { |
| 214 | // wait() hasn't been called. |
| 215 | NOT_STARTED, |
| 216 | |
| 217 | // wait() has been called once, and that wait hasn't finished yet. |
| 218 | INITIAL_WAIT, |
| 219 | |
| 220 | // First lock has been obtained, waiting for success() or failed(). |
| 221 | RUNNING, |
| 222 | |
| 223 | // success() or failed() has been called. |
| 224 | REPARENTED |
| 225 | }; |
| 226 | |
| 227 | State state = NOT_STARTED; |
| 228 | |
| 229 | // Points to the parent scope, which may be another CriticalSection in the case of nesting. |
| 230 | kj::OneOf<InputGate*, kj::Own<CriticalSection>> parent; |
| 231 | |
| 232 | // A lock in the parent scope. `parentLock` becomes non-null after the first lock is obtained, |
| 233 | // and becomes null again when succeeded() is called. |
| 234 | kj::Maybe<Lock> parentLock; |
| 235 | |
| 236 | friend class InputGate; |
| 237 | |
| 238 | // Return a reference for the parent scope, skipping any reparented CriticalSections |
| 239 | InputGate& parentAsInputGate(); |
| 240 | }; |
| 241 | |
| 242 | // An OutputGate blocks outgoing messages from an Actor until writes which they might depend on |
| 243 | // are confirmed. |
| 244 | class OutputGate { |
| 245 | public: |
| 246 | // Hooks that can be used to customize OutputGate behavior. |
| 247 | // |
| 248 | // Technically, everything implemented here could be accomplished by a class that wraps |
| 249 | // OutputGate, but the part of the code that wants to implement these hooks (Worker::Actor) |
| 250 | // is far away from the part of the code that calls into the OutputGate (ActorCache), and so |
| 251 | // it was more convenient to give Worker::Actor a way to inject behavior into OutputGate which |
| 252 | // would kick in when ActorCache tried to use it. |
| 253 | class Hooks { |
| 254 | public: |
| 255 | // Optionally make a promise which should be exclusiveJoin()ed with the lock promise to |
| 256 | // implement a timeout. The returned promise should be something that throws an exception |
| 257 | // after some timeout has expired. |
| 258 | virtual kj::Promise<void> makeTimeoutPromise() { |
| 259 | return kj::NEVER_DONE; |
| 260 | } |
| 261 | |
| 262 | // Optionally track metrics. In practice these are implemented by MetricsCollector::Actor, but |
| 263 | // we don't want to depend on that class from here. |
| 264 | |
| 265 | virtual void outputGateLocked() {} |
| 266 | virtual void outputGateReleased() {} |
| 267 | virtual void outputGateWaiterAdded() {} |
| 268 | virtual void outputGateWaiterRemoved() {} |
| 269 | |
| 270 | static const Hooks DEFAULT; |
| 271 | }; |
| 272 | |
| 273 | // Hooks has no member variables, so const_cast is acceptable. |
| 274 | OutputGate(Hooks& hooks = const_cast<Hooks&>(Hooks::DEFAULT)); |
| 275 | ~OutputGate() noexcept(false); |
| 276 | |
| 277 | // Block all future `wait()` calls until `promise` completes. Returns a wrapper around `promise`. |
| 278 | // If `promise` rejects, the exception will propagate to all future `wait()`s. If the returned |
| 279 | // promise is canceled before completion, all future `wait()`s will also throw. |
| 280 | template <typename T> |
| 281 | kj::Promise<T> lockWhile(kj::Promise<T> promise, SpanParent parentSpan); |
| 282 | |
| 283 | // Wait until all preceding locks are released. The wait will not be affected by any future |
| 284 | // call to `lockWhile()`. |
| 285 | kj::Promise<void> wait(SpanParent parentSpan); |
| 286 | |
| 287 | // Rejects if and when calls to `wait()` become broken due to a failed lockWhile(). The actor |
| 288 | // should be shut down in this case. This promise never resolves, only rejects. |
| 289 | // |
| 290 | // This method can only be called once. |
| 291 | kj::Promise<void> onBroken(); |
| 292 | |
| 293 | bool isBroken(); |
| 294 | |
| 295 | private: |
| 296 | Hooks& hooks; |
| 297 | |
| 298 | kj::ForkedPromise<void> pastLocksPromise; |
| 299 | |
| 300 | // A fulfiller for onBroken(), or an exception if already broken. |
| 301 | kj::OneOf<kj::Own<kj::PromiseFulfiller<void>>, kj::Exception> brokenState; |
| 302 | |
| 303 | void setBroken(const kj::Exception& e); |
| 304 | |
| 305 | kj::Own<kj::PromiseFulfiller<void>> lock(); |
| 306 | static kj::Exception makeUnfulfilledException(); |
| 307 | }; |
| 308 | |
| 309 | // ======================================================================================= |
| 310 | // inline implementation details |
| 311 | |
| 312 | template <typename T> |
| 313 | kj::Promise<T> OutputGate::lockWhile(kj::Promise<T> promise, SpanParent parentSpan) { |
| 314 | auto fulfiller = lock(); |
| 315 | SpanBuilder lockSpan = parentSpan.newChild("output_gate_lock_hold"_kjc); |
| 316 | |
| 317 | if constexpr (std::is_void_v<T>) { |
| 318 | promise = promise.exclusiveJoin(hooks.makeTimeoutPromise()); |
| 319 | } else { |
| 320 | promise = promise.exclusiveJoin(hooks.makeTimeoutPromise().then([]() -> T { KJ_UNREACHABLE; })); |
| 321 | } |
| 322 | |
| 323 | hooks.outputGateLocked(); |
| 324 | auto rejectIfCanceled = kj::defer([this, &fulfiller]() { |
| 325 | hooks.outputGateReleased(); |
| 326 | if (fulfiller->isWaiting()) { |
| 327 | auto e = makeUnfulfilledException(); |
| 328 | setBroken(e); |
| 329 | fulfiller->reject(kj::mv(e)); |
| 330 | } |
| 331 | }); |
| 332 | |
| 333 | try { |
| 334 | if constexpr (std::is_void_v<T>) { |
| 335 | co_await promise; |
| 336 | fulfiller->fulfill(); |
| 337 | } else { |
| 338 | auto v = co_await promise; |
| 339 | fulfiller->fulfill(); |
| 340 | co_return v; |
| 341 | } |
| 342 | } catch (kj::Exception& e) { |
| 343 | setBroken(e); |
| 344 | lockSpan.setTag("error"_kjc, true); |
| 345 | fulfiller->reject(e.clone()); |
| 346 | kj::throwFatalException(e.clone()); |
| 347 | } |
| 348 | } |
| 349 | |
| 350 | } // namespace workerd |