File
Blob: src/workerd/io/worker-interface.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 | |
| 7 | #include <workerd/io/outcome.capnp.h> |
| 8 | #include <workerd/io/trace.h> |
| 9 | #include <workerd/io/worker-interface.capnp.h> |
| 10 | #include <workerd/util/http-util.h> |
| 11 | |
| 12 | #include <capnp/compat/http-over-capnp.h> |
| 13 | #include <kj/compat/http.h> |
| 14 | #include <kj/debug.h> |
| 15 | |
| 16 | namespace workerd { |
| 17 | |
| 18 | class Frankenvalue; |
| 19 | class IoContext_IncomingRequest; |
| 20 | struct Worker_VersionInfo; |
| 21 | |
| 22 | // An interface representing the services made available by a worker/pipeline to handle a |
| 23 | // request. |
| 24 | class WorkerInterface: public kj::HttpService { |
| 25 | public: |
| 26 | // Constructs a WorkerInterface where any method called will throw the given exception. |
| 27 | static kj::Own<WorkerInterface> fromException(kj::Exception&& e); |
| 28 | |
| 29 | // Make an HTTP request. (This method is inherited from HttpService, but re-declared here for |
| 30 | // visibility.) |
| 31 | kj::Promise<void> request(kj::HttpMethod method, |
| 32 | kj::StringPtr url, |
| 33 | const kj::HttpHeaders& headers, |
| 34 | kj::AsyncInputStream& requestBody, |
| 35 | kj::HttpService::Response& response) override = 0; |
| 36 | // TODO(perf): Consider changing this to return Promise<DeferredProxy>. This would allow |
| 37 | // more resources to be dropped when merely proxying a request. However, it means we would no |
| 38 | // longer be implementing kj::HttpService. But maybe that doesn't matter too much in practice. |
| 39 | |
| 40 | // This is the same as the inherited HttpService::connect(), but we override it to be |
| 41 | // pure-virtual to force all subclasses of WorkerInterface to implement it explicitly rather |
| 42 | // than get the default implementation which throws an unimplemented exception. |
| 43 | kj::Promise<void> connect(kj::StringPtr host, |
| 44 | const kj::HttpHeaders& headers, |
| 45 | kj::AsyncIoStream& connection, |
| 46 | ConnectResponse& response, |
| 47 | kj::HttpConnectSettings settings) override = 0; |
| 48 | |
| 49 | // Hints that this worker will likely be invoked in the near future, so should be warmed up now. |
| 50 | // This method should also call `prewarm()` on any subsequent pipeline stages that are expected |
| 51 | // to be invoked. |
| 52 | // |
| 53 | // If prewarm() has to do anything asynchronous, it should use "waitUntil" tasks. |
| 54 | virtual kj::Promise<void> prewarm(kj::StringPtr url) = 0; |
| 55 | |
| 56 | // keep in sync with `src/rust/worker/ffi.rs` |
| 57 | struct ScheduledResult { |
| 58 | bool retry = true; |
| 59 | EventOutcome outcome = EventOutcome::UNKNOWN; |
| 60 | }; |
| 61 | |
| 62 | // Copyable subset of AlarmResult, used by ForkedPromise for alarm deduplication in Worker::Actor. |
| 63 | // keep in sync with `src/rust/worker/ffi.rs` |
| 64 | struct AlarmOutcome { |
| 65 | bool retry = true; |
| 66 | bool retryCountsAgainstLimit = true; |
| 67 | EventOutcome outcome = EventOutcome::UNKNOWN; |
| 68 | }; |
| 69 | |
| 70 | // keep in sync with `src/rust/worker/ffi.rs` |
| 71 | struct AlarmResult { |
| 72 | bool retry = true; |
| 73 | bool retryCountsAgainstLimit = true; |
| 74 | EventOutcome outcome = EventOutcome::UNKNOWN; |
| 75 | kj::Maybe<kj::String> errorDescription; |
| 76 | |
| 77 | AlarmOutcome asOutcome() const { |
| 78 | return { |
| 79 | .retry = retry, .retryCountsAgainstLimit = retryCountsAgainstLimit, .outcome = outcome}; |
| 80 | } |
| 81 | }; |
| 82 | |
| 83 | class AlarmFulfiller { |
| 84 | public: |
| 85 | AlarmFulfiller(kj::Own<kj::PromiseFulfiller<AlarmOutcome>> fulfiller); |
| 86 | KJ_DISALLOW_COPY(AlarmFulfiller); |
| 87 | AlarmFulfiller(AlarmFulfiller&&) = default; |
| 88 | AlarmFulfiller& operator=(AlarmFulfiller&&) = default; |
| 89 | ~AlarmFulfiller() noexcept(false); |
| 90 | void fulfill(const AlarmOutcome& result); |
| 91 | void reject(const kj::Exception& e); |
| 92 | void cancel(); |
| 93 | |
| 94 | private: |
| 95 | kj::Maybe<kj::Own<kj::PromiseFulfiller<AlarmOutcome>>> maybeFulfiller; |
| 96 | kj::Maybe<kj::PromiseFulfiller<AlarmOutcome>&> getFulfiller(); |
| 97 | }; |
| 98 | |
| 99 | using ScheduleAlarmResult = kj::OneOf<AlarmOutcome, AlarmFulfiller>; |
| 100 | |
| 101 | // Trigger a scheduled event with the given scheduled (unix timestamp) time and cron string. |
| 102 | // The cron string must be valid until the returned promise completes. |
| 103 | // Async work is queued in a "waitUntil" task set. |
| 104 | virtual kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) = 0; |
| 105 | |
| 106 | // Trigger an alarm event with the given scheduled (unix timestamp) time. |
| 107 | virtual kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) = 0; |
| 108 | |
| 109 | // Called when AlarmManager has given up retrying an alarm after too many counted failures. |
| 110 | // The actor should clear its alarm state so getAlarm() reflects the deletion. |
| 111 | // Returns the actor's stored alarm time if it differs from scheduledTime (i.e. the user set a |
| 112 | // new alarm), or kj::none if the alarm was cleared or no alarm was stored. |
| 113 | // Default is a no-op so subclasses that don't host actors need not override it. |
| 114 | virtual kj::Promise<kj::Maybe<kj::Date>> abandonAlarm(kj::Date scheduledTime) { |
| 115 | return kj::Maybe<kj::Date>(kj::none); |
| 116 | } |
| 117 | |
| 118 | // Run the test handler. The returned promise resolves to true or false to indicate that the test |
| 119 | // passed or failed. In the case of a failure, information should have already been written to |
| 120 | // stderr and to the devtools; there is no need for the caller to write anything further. (If the |
| 121 | // promise rejects, this indicates a bug in the test harness itself.) |
| 122 | virtual kj::Promise<bool> test() { |
| 123 | return nullptr; |
| 124 | } |
| 125 | // TODO(someday): Produce a structured test report? |
| 126 | |
| 127 | // These two constants are shared by multiple systems that invoke alarms (the production |
| 128 | // implementation, and the preview implementation), whose code live in completely different |
| 129 | // places. We end up defining them here mostly for lack of a better option. |
| 130 | static constexpr auto ALARM_RETRY_START_SECONDS = 2; // not a duration so we can left shift it |
| 131 | static constexpr auto ALARM_RETRY_MAX_TRIES = 6; |
| 132 | |
| 133 | class CustomEvent { |
| 134 | public: |
| 135 | struct Result { |
| 136 | // Outcome for logging / metrics purposes. |
| 137 | EventOutcome outcome; |
| 138 | }; |
| 139 | |
| 140 | // Deliver the event to an isolate in this process. `incomingRequest` has been newly-allocated |
| 141 | // for this event. |
| 142 | virtual kj::Promise<Result> run(kj::Own<IoContext_IncomingRequest> incomingRequest, |
| 143 | kj::Maybe<kj::StringPtr> entrypointName, |
| 144 | kj::Maybe<Worker_VersionInfo> versionInfo, |
| 145 | Frankenvalue props, |
| 146 | kj::TaskSet& waitUntilTasks, |
| 147 | bool isDynamicDispatch = false) = 0; |
| 148 | |
| 149 | // Forward the event over RPC. |
| 150 | virtual kj::Promise<Result> sendRpc(capnp::HttpOverCapnpFactory& httpOverCapnpFactory, |
| 151 | capnp::ByteStreamFactory& byteStreamFactory, |
| 152 | rpc::EventDispatcher::Client dispatcher) = 0; |
| 153 | |
| 154 | // The event is not supported by the target, raise an appropriate error. |
| 155 | virtual kj::Promise<Result> notSupported() = 0; |
| 156 | |
| 157 | // Get the type for this event for logging / metrics purposes. This is intended for use by the |
| 158 | // RequestObserver. The RequestObserver implementation will define what numbers correspond to |
| 159 | // what types. |
| 160 | virtual uint16_t getType() = 0; |
| 161 | |
| 162 | // Get event info for tracing. |
| 163 | virtual tracing::EventInfo getEventInfo() const = 0; |
| 164 | |
| 165 | // If the CustomEvent fails before any of the other methods are called, this may be invoked |
| 166 | // to report the failure reason. |
| 167 | virtual void failed(const kj::Exception& e) {} |
| 168 | }; |
| 169 | |
| 170 | // Allows delivery of a variety of event types by implementing a callback that delivers the |
| 171 | // event to a particular isolate. If and when the event is delivered to an isolate, |
| 172 | // `callback->run()` will be called inside a fresh IoContext::IncomingRequest to begin the |
| 173 | // event. |
| 174 | // |
| 175 | // If the event needs to return some sort of result, it's the responsibility of the callback to |
| 176 | // store that result in a side object that the event's invoker can inspect after the promise has |
| 177 | // resolved. |
| 178 | // |
| 179 | // Note that it is guaranteed that if the returned promise is canceled, `event` will be dropped |
| 180 | // immediately; if its callbacks have not run yet, they will not run at all. So, a CustomEvent |
| 181 | // implementation can hold references to objects it doesn't own as long as the returned promise |
| 182 | // will be canceled before those objects go away. |
| 183 | [[nodiscard]] virtual kj::Promise<CustomEvent::Result> customEvent( |
| 184 | kj::Own<CustomEvent> event) = 0; |
| 185 | |
| 186 | private: |
| 187 | kj::Maybe<kj::Own<kj::HttpService>> adapterService; |
| 188 | }; |
| 189 | |
| 190 | // Given a Promise for a WorkerInterface, return a WorkerInterface whose methods will first wait |
| 191 | // for the promise, then invoke the destination object. |
| 192 | kj::Own<WorkerInterface> newPromisedWorkerInterface(kj::Promise<kj::Own<WorkerInterface>> promise); |
| 193 | |
| 194 | template <typename Func> |
| 195 | class LazyWorkerInterface final: public WorkerInterface { |
| 196 | public: |
| 197 | LazyWorkerInterface(Func func): func(kj::mv(func)) {} |
| 198 | |
| 199 | void ensureResolve() { |
| 200 | if (promise == kj::none) { |
| 201 | promise = KJ_ASSERT_NONNULL(func)() |
| 202 | .then([this](kj::Own<WorkerInterface> result) { worker = kj::mv(result); }) |
| 203 | .eagerlyEvaluate(nullptr) |
| 204 | .fork(); |
| 205 | func = kj::none; |
| 206 | } |
| 207 | } |
| 208 | |
| 209 | kj::Promise<void> request(kj::HttpMethod method, |
| 210 | kj::StringPtr url, |
| 211 | const kj::HttpHeaders& headers, |
| 212 | kj::AsyncInputStream& requestBody, |
| 213 | Response& response) override { |
| 214 | ensureResolve(); |
| 215 | KJ_IF_SOME(w, worker) { |
| 216 | co_await w->request(method, url, headers, requestBody, response); |
| 217 | } else { |
| 218 | co_await KJ_ASSERT_NONNULL(promise); |
| 219 | co_await KJ_ASSERT_NONNULL(worker)->request(method, url, headers, requestBody, response); |
| 220 | } |
| 221 | } |
| 222 | |
| 223 | kj::Promise<void> connect(kj::StringPtr host, |
| 224 | const kj::HttpHeaders& headers, |
| 225 | kj::AsyncIoStream& connection, |
| 226 | ConnectResponse& response, |
| 227 | kj::HttpConnectSettings settings) override { |
| 228 | ensureResolve(); |
| 229 | KJ_IF_SOME(w, worker) { |
| 230 | co_await w->connect(host, headers, connection, response, kj::mv(settings)); |
| 231 | } else { |
| 232 | co_await KJ_ASSERT_NONNULL(promise); |
| 233 | co_await KJ_ASSERT_NONNULL(worker)->connect( |
| 234 | host, headers, connection, response, kj::mv(settings)); |
| 235 | } |
| 236 | } |
| 237 | |
| 238 | kj::Promise<void> prewarm(kj::StringPtr url) override { |
| 239 | ensureResolve(); |
| 240 | KJ_IF_SOME(w, worker) { |
| 241 | co_return co_await w->prewarm(url); |
| 242 | } else { |
| 243 | co_await KJ_ASSERT_NONNULL(promise); |
| 244 | co_return co_await KJ_ASSERT_NONNULL(worker)->prewarm(url); |
| 245 | } |
| 246 | } |
| 247 | |
| 248 | kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override { |
| 249 | ensureResolve(); |
| 250 | KJ_IF_SOME(w, worker) { |
| 251 | co_return co_await w->runScheduled(scheduledTime, cron); |
| 252 | } else { |
| 253 | co_await KJ_ASSERT_NONNULL(promise); |
| 254 | co_return co_await KJ_ASSERT_NONNULL(worker)->runScheduled(scheduledTime, cron); |
| 255 | } |
| 256 | } |
| 257 | |
| 258 | kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override { |
| 259 | ensureResolve(); |
| 260 | KJ_IF_SOME(w, worker) { |
| 261 | co_return co_await w->runAlarm(scheduledTime, retryCount); |
| 262 | } else { |
| 263 | co_await KJ_ASSERT_NONNULL(promise); |
| 264 | co_return co_await KJ_ASSERT_NONNULL(worker)->runAlarm(scheduledTime, retryCount); |
| 265 | } |
| 266 | } |
| 267 | |
| 268 | kj::Promise<kj::Maybe<kj::Date>> abandonAlarm(kj::Date scheduledTime) override { |
| 269 | ensureResolve(); |
| 270 | KJ_IF_SOME(w, worker) { |
| 271 | co_return co_await w->abandonAlarm(scheduledTime); |
| 272 | } else { |
| 273 | co_await KJ_ASSERT_NONNULL(promise); |
| 274 | co_return co_await KJ_ASSERT_NONNULL(worker)->abandonAlarm(scheduledTime); |
| 275 | } |
| 276 | } |
| 277 | |
| 278 | kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override { |
| 279 | ensureResolve(); |
| 280 | KJ_IF_SOME(w, worker) { |
| 281 | co_return co_await w->customEvent(kj::mv(event)); |
| 282 | } else { |
| 283 | co_await KJ_ASSERT_NONNULL(promise); |
| 284 | co_return co_await KJ_ASSERT_NONNULL(worker)->customEvent(kj::mv(event)); |
| 285 | } |
| 286 | } |
| 287 | |
| 288 | private: |
| 289 | kj::Maybe<Func> func; |
| 290 | kj::Maybe<kj::ForkedPromise<void>> promise; |
| 291 | kj::Maybe<kj::Own<WorkerInterface>> worker; |
| 292 | }; |
| 293 | // Similar to newPromisedWorkerInterface but receives a function that returns a Promise for a |
| 294 | // WorkerInterface. This is useful when you are not sure if the worker will be used or not and |
| 295 | // you don't want it to be created in case it isn't used. If you just create a |
| 296 | // PromisedWorkerInterface then the async loop might run the promise before it is eventually |
| 297 | // destroyed even if it was never used. |
| 298 | template <typename Func> |
| 299 | kj::Own<WorkerInterface> newLazyWorkerInterface(Func func) { |
| 300 | return kj::heap<LazyWorkerInterface<Func>>(kj::mv(func)); |
| 301 | } |
| 302 | |
| 303 | // Adapts WorkerInterface to HttpClient, including taking ownership. |
| 304 | // |
| 305 | // (Use kj::newHttpClient() if you don't want to take ownership.) |
| 306 | kj::Own<kj::HttpClient> asHttpClient(kj::Own<WorkerInterface> workerInterface); |
| 307 | |
| 308 | // A WorkerInterface that cancels WebSockets when revokeProm is rejected. |
| 309 | // Currently only supports cancelling for upgrades. |
| 310 | kj::Own<WorkerInterface> newRevocableWebSocketWorkerInterface( |
| 311 | kj::Own<WorkerInterface> worker, kj::Promise<void> revokeProm); |
| 312 | |
| 313 | // Implementation of WorkerInterface on top of rpc::EventDispatcher. Since an EventDispatcher |
| 314 | // is intended to be single-use, this class is also inherently single-use (i.e. only one event |
| 315 | // can be delivered). |
| 316 | class RpcWorkerInterface final: public WorkerInterface { |
| 317 | public: |
| 318 | RpcWorkerInterface(capnp::HttpOverCapnpFactory& httpOverCapnpFactory, |
| 319 | capnp::ByteStreamFactory& byteStreamFactory, |
| 320 | rpc::EventDispatcher::Client dispatcher); |
| 321 | |
| 322 | kj::Promise<void> request(kj::HttpMethod method, |
| 323 | kj::StringPtr url, |
| 324 | const kj::HttpHeaders& headers, |
| 325 | kj::AsyncInputStream& requestBody, |
| 326 | Response& response) override; |
| 327 | |
| 328 | kj::Promise<void> connect(kj::StringPtr host, |
| 329 | const kj::HttpHeaders& headers, |
| 330 | kj::AsyncIoStream& connection, |
| 331 | ConnectResponse& tunnel, |
| 332 | kj::HttpConnectSettings settings) override; |
| 333 | |
| 334 | kj::Promise<void> prewarm(kj::StringPtr url) override; |
| 335 | kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override; |
| 336 | kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override; |
| 337 | kj::Promise<kj::Maybe<kj::Date>> abandonAlarm(kj::Date scheduledTime) override; |
| 338 | kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override; |
| 339 | |
| 340 | private: |
| 341 | capnp::HttpOverCapnpFactory& httpOverCapnpFactory; |
| 342 | capnp::ByteStreamFactory& byteStreamFactory; |
| 343 | rpc::EventDispatcher::Client dispatcher; |
| 344 | }; |
| 345 | |
| 346 | } // namespace workerd |