File
Blob: src/workerd/io/io-channels.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/actor-id.h> |
| 8 | #include <workerd/io/compatibility-date.capnp.h> |
| 9 | #include <workerd/io/frankenvalue.h> |
| 10 | #include <workerd/io/io-util.h> |
| 11 | #include <workerd/io/trace.h> |
| 12 | #include <workerd/io/worker-source.h> |
| 13 | |
| 14 | #include <capnp/capability.h> // for Capability |
| 15 | #include <kj/debug.h> |
| 16 | #include <kj/string.h> |
| 17 | |
| 18 | namespace kj { |
| 19 | class HttpClient; |
| 20 | class Network; |
| 21 | } // namespace kj |
| 22 | |
| 23 | namespace workerd { |
| 24 | |
| 25 | class WorkerInterface; |
| 26 | |
| 27 | // Interface for talking to the Cache API. Needs to be declared here so that IoContext can |
| 28 | // contain it. |
| 29 | class CacheClient { |
| 30 | public: |
| 31 | struct SubrequestMetadata { |
| 32 | // The `request.cf` blob, JSON-encoded. |
| 33 | kj::Maybe<kj::String> cfBlobJson; |
| 34 | |
| 35 | // Specifies the parent span for the subrequest for tracing purposes. |
| 36 | SpanParent parentSpan; |
| 37 | |
| 38 | // Serialized JSON value to pass in ew_compat field of control header to FL. This has the same |
| 39 | // semantics as the field in IoChannelFactory::SubrequestMetadata. |
| 40 | kj::Maybe<kj::String> featureFlagsForFl; |
| 41 | }; |
| 42 | |
| 43 | // Get the default namespace, i.e. the one that fetch() will use for caching. |
| 44 | // |
| 45 | // The returned client is intended to be used for one request. |
| 46 | virtual kj::Own<kj::HttpClient> getDefault(SubrequestMetadata metadata) = 0; |
| 47 | |
| 48 | // Get an HttpClient for the given cache namespace. |
| 49 | virtual kj::Own<kj::HttpClient> getNamespace(kj::StringPtr name, SubrequestMetadata metadata) = 0; |
| 50 | }; |
| 51 | |
| 52 | // A timer instance, used to back Date.now(), setTimeout(), etc. This object may implement |
| 53 | // Spectre mitigations. |
| 54 | class TimerChannel { |
| 55 | public: |
| 56 | // Call each time control enters the isolate to set up the clock. |
| 57 | virtual void syncTime() = 0; |
| 58 | |
| 59 | // Return the current time. `nextTimeout` is the time at which the next setTimeout() callback |
| 60 | // is scheduled; implementations performing Spectre mitigations should clamp to this value so |
| 61 | // that Date.now() never goes backwards or reveals timing side channels. |
| 62 | virtual kj::Date now(kj::Maybe<kj::Date> nextTimeout = kj::none) = 0; |
| 63 | |
| 64 | // Returns a promise that resolves once `now() >= when`. |
| 65 | virtual kj::Promise<void> atTime(kj::Date when) = 0; |
| 66 | |
| 67 | // Returns a promise that resolves after some time. This is intended to be used for implementing |
| 68 | // time limits on some sort of operation, not for implementing application-driven timing, as it does |
| 69 | // not implement any Spectre mitigations. |
| 70 | virtual kj::Promise<void> afterLimitTimeout(kj::Duration t) = 0; |
| 71 | }; |
| 72 | |
| 73 | class WorkerStubChannel; |
| 74 | struct DynamicWorkerSource; |
| 75 | |
| 76 | // Each IoContext has a set of "channels" on which outgoing I/O can be initiated. All outgoing |
| 77 | // I/O occurs through these channels. Think of these kind of like file descriptors. They are |
| 78 | // often associated with bindings. |
| 79 | // |
| 80 | // For example, any call to fetch() uses a subrequest channel. The global fetch() specifically |
| 81 | // uses subrequest channel zero. Each service binding (aka worker-to-worker binding) is assigned |
| 82 | // a unique subrequest channel number, and calling `binding.fetch()` sends the request to the |
| 83 | // given channel. |
| 84 | // |
| 85 | // While most channels are SubrequestChannels, other channel types exist to handle I/O that is |
| 86 | // not subrequest-shaped. For example, a Workers Analytics Engine binding uses a logging channel. |
| 87 | // |
| 88 | // Note that each type of channel has its own number space. That is, subrequest channel 5 and |
| 89 | // logging channel 5 are not related. |
| 90 | // |
| 91 | // The reason we have channels, rather than binding API objects directly holding the I/O objects, |
| 92 | // is because binding API objects live across multiple requests, but the I/O objects may differ |
| 93 | // from request to request. |
| 94 | // |
| 95 | // This class encapsulates all outgoing I/O that a Worker can perform. It does not cover incoming |
| 96 | // I/O, i.e. the event that started the Worker. If IoChannelFactory is implemented such that |
| 97 | // all methods throw exceptions, then the Worker will be completely unable to communicate with |
| 98 | // anything in the world except for the client -- this is a useful property for sandboxing! |
| 99 | class IoChannelFactory { |
| 100 | public: |
| 101 | // Contains metadata attached to an outgoing subrequest from a worker, independent of the type |
| 102 | // of request. |
| 103 | struct SubrequestMetadata { |
| 104 | // The `request.cf` blob, JSON-encoded. |
| 105 | kj::Maybe<kj::String> cfBlobJson; |
| 106 | |
| 107 | // Specifies the parent span for the subrequest for tracing purposes. |
| 108 | SpanParent parentSpan = SpanParent(nullptr); |
| 109 | |
| 110 | // User Span Parent for trace propagation. Call toSpanContext() to serialize. |
| 111 | SpanParent userSpanParent = SpanParent(nullptr); |
| 112 | |
| 113 | // Serialized JSON value to pass in ew_compat field of control header to FL. If this subrequest |
| 114 | // does not go directly to FL, this value is ignored. Flags marked with `$neededByFl` in |
| 115 | // `compatibility-date.capnp` end up here. |
| 116 | kj::Maybe<kj::String> featureFlagsForFl; |
| 117 | |
| 118 | // Timestamp for when a subrequest is started. (ms since the Unix Epoch) |
| 119 | double startTime = dateNow(); |
| 120 | }; |
| 121 | |
| 122 | // Parameters that can influence the version of a worker that is used to serve a subrequest. |
| 123 | struct VersionRequest { |
| 124 | // Request a version within the given cohort. |
| 125 | kj::Maybe<kj::String> cohort; |
| 126 | |
| 127 | VersionRequest clone() const { |
| 128 | return { |
| 129 | .cohort = cohort.map([](const kj::String& s) { return kj::str(s); }), |
| 130 | }; |
| 131 | } |
| 132 | }; |
| 133 | |
| 134 | virtual kj::Own<WorkerInterface> startSubrequest(uint channel, SubrequestMetadata metadata) = 0; |
| 135 | |
| 136 | // Get a Cap'n Proto RPC capability. Various binding types are backed by capabilities. |
| 137 | // |
| 138 | // Note that some other channel types, like actor channels, may actually be wrappers around |
| 139 | // capability channels, and so may share the same channel number space, but this shouldn't be |
| 140 | // assumed. |
| 141 | virtual capnp::Capability::Client getCapability(uint channel) = 0; |
| 142 | |
| 143 | // Get a CacheClient, used to implement the Cache API. |
| 144 | virtual kj::Own<CacheClient> getCache() = 0; |
| 145 | |
| 146 | // Get the singleton timer instance, used to back Date.now(), setTimeout(), etc. This object |
| 147 | // may implement Spectre mitigations. |
| 148 | virtual TimerChannel& getTimer() = 0; |
| 149 | |
| 150 | // Write a log message to a logfwdr channel. Each log binding has its own channel number. |
| 151 | // |
| 152 | // The IoChannelFactory already knows which member of the overall message union is expected to |
| 153 | // be filled in for this channel. That member will be initialized as a pointer, and then |
| 154 | // `buildMessage` will be invoked to fill in the pointer's content. The callback is always |
| 155 | // executed immediately, before `writeLogfwdr()` returns a promise. |
| 156 | virtual kj::Promise<void> writeLogfwdr( |
| 157 | uint channel, kj::FunctionParam<void(capnp::AnyPointer::Builder)> buildMessage) = 0; |
| 158 | |
| 159 | enum ChannelTokenUsage { |
| 160 | // Token is to be sent over RPC and hence will be converted back into a SubrequestChannel |
| 161 | // soon. Such tokens have limited lifetime but are otherwise irrevocable. |
| 162 | RPC, |
| 163 | |
| 164 | // Token is to be stored in long-term storage. At present this must only be allowed to be |
| 165 | // used in workers that have the allow_irrevocable_stub_storage compat flag (checked by the |
| 166 | // caller). In the future the format for such tokens will change. |
| 167 | STORAGE, |
| 168 | }; |
| 169 | |
| 170 | // Object representing somehere where generic workers subrequests can be sent. Multiple requests |
| 171 | // may be sent. This is an I/O type so it is only valid within the `IoContext` where it was |
| 172 | // created. |
| 173 | class SubrequestChannel: public kj::Refcounted, public Frankenvalue::CapTableEntry { |
| 174 | public: |
| 175 | // Start a new request to this target. |
| 176 | // |
| 177 | // Note that not all `metadata` properties make sense here, but it didn't seem worth defining |
| 178 | // a new struct type. `cfBlobJson` and `parentSpan` make sense, but `featureFlagsForFl` and |
| 179 | // `dynamicDispatchTarget` do not. |
| 180 | // |
| 181 | // Note that the caller is expected to keep the SubrequestChannel alive until it is done with |
| 182 | // the returned WorkerInterface. |
| 183 | virtual kj::Own<WorkerInterface> startRequest(SubrequestMetadata metadata) = 0; |
| 184 | |
| 185 | kj::Own<CapTableEntry> clone() override final { |
| 186 | return kj::addRef(*this); |
| 187 | } |
| 188 | |
| 189 | // Throws a JSG error if a Fetcher backed by this channel should not be serialized and passed |
| 190 | // to other workers. The default implementation throws a generic error, but subclasses may |
| 191 | // specialize with better errror messages -- or override to just return in order to permit the |
| 192 | // serialization. |
| 193 | // |
| 194 | // This check is necessary especially in workerd in order to block serialization of types that, |
| 195 | // in production, would be difficult or impossible to serialize. In particular, |
| 196 | // dynamically-loaded workers cannot be serialized because the system does not know how to |
| 197 | // reconstruct a dynamically-loaded worker from scratch. |
| 198 | virtual void requireAllowsTransfer() = 0; |
| 199 | |
| 200 | // Get a token representing this SubrequestChannel which can be converted back into a |
| 201 | // SubrequestChannel using subrequestChannelFromToken(). Default implementation throws a |
| 202 | // TypeError. |
| 203 | virtual kj::Array<byte> getToken(ChannelTokenUsage usage); |
| 204 | }; |
| 205 | |
| 206 | // Obtain an object representing a particular subrequest channel. |
| 207 | // |
| 208 | // getSubrequestChannel(i).startRequest(meta) is exactly equivalent to startSubrequest(i, meta). |
| 209 | // The reason to use this instead is when the channel is not necessarily going to be used to |
| 210 | // start a subrequest immediately, but instead is going to be passed around as a capability. |
| 211 | // |
| 212 | // `props` and `versionRequest` can only be specified if this is a loopback channel (i.e. from |
| 213 | // ctx.exports). For any other channel, they will throw. |
| 214 | // |
| 215 | // TODO(cleanup): Consider getting rid of `startSubrequest()` in favor of this. |
| 216 | virtual kj::Own<SubrequestChannel> getSubrequestChannel(uint channel, |
| 217 | kj::Maybe<Frankenvalue> props = kj::none, |
| 218 | kj::Maybe<VersionRequest> versionRequest = kj::none) = 0; |
| 219 | |
| 220 | // Stub for a remote actor. Allows sending requests to the actor. |
| 221 | class ActorChannel: public SubrequestChannel { |
| 222 | public: |
| 223 | // At present there are no methods beyond what `SubrequestChannel` defines. However, it's |
| 224 | // easy to imagine that actor stubs may have more functionality than just sending requests |
| 225 | // someday, so we keep this as a separate type. |
| 226 | |
| 227 | // For now, actor stubs are not transferrable -- but we do intend to change that at some point. |
| 228 | void requireAllowsTransfer() override final; |
| 229 | }; |
| 230 | |
| 231 | // Get an actor stub from the given namespace for the actor with the given ID. |
| 232 | // |
| 233 | // `id` must have been constructed using one of the `ActorIdFactory` instances corresponding to |
| 234 | // one of the worker's bindings, however it doesn't necessarily have to be from the the correct |
| 235 | // `ActorIdFactory` -- if it's from some other factory, the method will throw an appropriate |
| 236 | // exception. |
| 237 | virtual kj::Own<ActorChannel> getGlobalActor(uint channel, |
| 238 | const ActorIdFactory::ActorId& id, |
| 239 | kj::Maybe<kj::String> locationHint, |
| 240 | ActorGetMode mode, |
| 241 | bool enableReplicaRouting, |
| 242 | ActorRoutingMode routingMode, |
| 243 | SpanParent parentSpan, |
| 244 | kj::Maybe<ActorVersion> version) = 0; |
| 245 | |
| 246 | // Get an actor stub from the given namespace for the actor with the given name. |
| 247 | virtual kj::Own<ActorChannel> getColoLocalActor( |
| 248 | uint channel, kj::StringPtr id, SpanParent parentSpan) = 0; |
| 249 | |
| 250 | // ActorClassChannel is a reference to an actor class in another worker. This class acts as a |
| 251 | // token which can be passed into other interfaces that might use the actor class, particularly |
| 252 | // Worker::Actor::FacetManager. |
| 253 | class ActorClassChannel: public kj::Refcounted, public Frankenvalue::CapTableEntry { |
| 254 | public: |
| 255 | kj::Own<CapTableEntry> clone() override final { |
| 256 | return kj::addRef(*this); |
| 257 | } |
| 258 | |
| 259 | // Same as the corresponding methods on SubrequestChannel. |
| 260 | virtual void requireAllowsTransfer() = 0; |
| 261 | virtual kj::Array<byte> getToken(ChannelTokenUsage usage); |
| 262 | |
| 263 | // This class has no functional methods, since it serves as a token to be passed to other |
| 264 | // interfaces (namely the facets API). |
| 265 | }; |
| 266 | |
| 267 | // Get an actor class binding corresponding to the given channel number. |
| 268 | // |
| 269 | // `props` can only be specified if this is a loopback channel (i.e. from ctx.exports). For any |
| 270 | // other channel, it will throw. |
| 271 | virtual kj::Own<ActorClassChannel> getActorClass( |
| 272 | uint channel, kj::Maybe<Frankenvalue> props = kj::none) { |
| 273 | // TODO(cleanup): Remove this once the production runtime has implemented this. |
| 274 | KJ_UNIMPLEMENTED("This runtime doesn't support actor class channels."); |
| 275 | } |
| 276 | |
| 277 | // Aborts all actors except those in namespaces marked with `preventEviction`. |
| 278 | virtual void abortAllActors(kj::Maybe<kj::Exception&> reason) { |
| 279 | KJ_UNIMPLEMENTED("Only implemented by single-tenant workerd runtime"); |
| 280 | } |
| 281 | |
| 282 | // Aborts all actors, cancels all alarms, and deletes all underlying storage for evictable |
| 283 | // namespaces. After this, DOs can be recreated with clean state. Useful for test isolation. |
| 284 | virtual void deleteAllActors(kj::Maybe<kj::Exception&> reason) { |
| 285 | KJ_UNIMPLEMENTED("Only implemented by single-tenant workerd runtime"); |
| 286 | } |
| 287 | |
| 288 | // In workerd, the handler aborts the process (unless used on a dynamic |
| 289 | // worker). In the edge runtime it will condemn and terminate the current |
| 290 | // isolate. |
| 291 | virtual void abortIsolate(kj::StringPtr reason) = 0; |
| 292 | |
| 293 | // Use a dynamic Worker loader binding to obtain an Worker by name. If name is null, or if the named Worker doesn't already exist, the callback will be called to fetch the source code from which the Worker should be created. |
| 294 | virtual kj::Own<WorkerStubChannel> loadIsolate(uint loaderChannel, |
| 295 | kj::Maybe<kj::String> name, |
| 296 | kj::Function<kj::Promise<DynamicWorkerSource>()> fetchSource) { |
| 297 | JSG_FAIL_REQUIRE(Error, "Dynamic worker loading is not supported by this runtime."); |
| 298 | } |
| 299 | |
| 300 | // Get the network for connecting to workerd debug ports. |
| 301 | // This is used by the workerdDebugPort binding to connect to remote workerd instances. |
| 302 | virtual kj::Network& getWorkerdDebugPortNetwork() { |
| 303 | JSG_FAIL_REQUIRE(Error, "WorkerdDebugPort bindings are not supported by this runtime."); |
| 304 | } |
| 305 | |
| 306 | // Converts a token created with {SubrequestChannel,ActorClassChannel}::getToken() back into a |
| 307 | // live channel. Default implementations throw. |
| 308 | virtual kj::Own<SubrequestChannel> subrequestChannelFromToken( |
| 309 | ChannelTokenUsage usage, kj::ArrayPtr<const byte> token); |
| 310 | virtual kj::Own<ActorClassChannel> actorClassFromToken( |
| 311 | ChannelTokenUsage usage, kj::ArrayPtr<const byte> token); |
| 312 | }; |
| 313 | |
| 314 | // ResourceLimits provides a means to control the resource allocation for a worker stage via a |
| 315 | // set of optionally overridden parameters. |
| 316 | struct ResourceLimits { |
| 317 | jsg::Optional<uint32_t> cpuMs; |
| 318 | jsg::Optional<uint32_t> subRequests; |
| 319 | |
| 320 | JSG_STRUCT(cpuMs, subRequests); |
| 321 | |
| 322 | ResourceLimits clone() const { |
| 323 | return {cpuMs, subRequests}; |
| 324 | } |
| 325 | }; |
| 326 | |
| 327 | // Represents a dynamically-loaded Worker to which requests can be sent. |
| 328 | // |
| 329 | // This object is returned before the Worker actually loads, so if any errors occur while loading, |
| 330 | // any requests sent to the Worker will fail, propagating the exception. |
| 331 | class WorkerStubChannel { |
| 332 | public: |
| 333 | virtual kj::Own<IoChannelFactory::SubrequestChannel> getEntrypoint( |
| 334 | kj::Maybe<kj::String> name, Frankenvalue props, kj::Maybe<ResourceLimits> limits) = 0; |
| 335 | |
| 336 | virtual kj::Own<IoChannelFactory::ActorClassChannel> getActorClass( |
| 337 | kj::Maybe<kj::String> name, Frankenvalue props, kj::Maybe<ResourceLimits> limits) = 0; |
| 338 | |
| 339 | // TODO(someday): Allow caller to enumerate entrypoints? |
| 340 | }; |
| 341 | |
| 342 | // Source code needed to dynamically load a Worker. |
| 343 | struct DynamicWorkerSource { |
| 344 | WorkerSource source; |
| 345 | CompatibilityFlags::Reader compatibilityFlags; |
| 346 | |
| 347 | kj::Maybe<ResourceLimits> limits; |
| 348 | |
| 349 | // `env` object to pass to the loaded worker. Can contain anything that can be serialized to |
| 350 | // a `Frankenvalue` (which should eventually include all binding types, RPC stubs, etc.). |
| 351 | Frankenvalue env; |
| 352 | |
| 353 | // Where should global fetch() (and connect()) be sent? |
| 354 | kj::Maybe<kj::Own<IoChannelFactory::SubrequestChannel>> globalOutbound; |
| 355 | |
| 356 | // Tail workers that should receive tail events for invocations of the dynamic worker. |
| 357 | kj::Array<kj::Own<IoChannelFactory::SubrequestChannel>> tails; |
| 358 | kj::Array<kj::Own<IoChannelFactory::SubrequestChannel>> streamingTails; |
| 359 | |
| 360 | // Owns any data structures pointed into by the other members. (E.g. `source` contains a lot of |
| 361 | // `StringPtr`s; `ownContent` owns the backing buffer for them.) |
| 362 | kj::Own<void> ownContent; |
| 363 | |
| 364 | // Indicates whether ownContent is holding onto a Cap'n Proto RPC response. This is important |
| 365 | // to know because such an RPC response must be destroyed on the same thread where it was |
| 366 | // created, and generally should be destroyed "relatively soon", not kept around forever. If |
| 367 | // this is false, then it is perfectly safe to transfer ownership of ownContent between threads |
| 368 | // and keep it alive indefinitely long. |
| 369 | bool ownContentIsRpcResponse = true; |
| 370 | |
| 371 | // Clone the DynamicWorkerSource. Caller must provide a new reference to use as `ownContent`, |
| 372 | // which must be a refcount on the same content since the pointers will not be updated. Note |
| 373 | // that if `ownContentIsRpcResponse` is false, then `ownContent` could be passed off to other |
| 374 | // threads and as such the refcount had better be atomic. |
| 375 | DynamicWorkerSource clone(kj::Own<void> newOwnContent) { |
| 376 | return { |
| 377 | .source = source.clone(), |
| 378 | .compatibilityFlags = compatibilityFlags, |
| 379 | .limits = limits.map([](auto& limits) { return limits.clone(); }), |
| 380 | .env = env.clone(), |
| 381 | .globalOutbound = mapAddRef(globalOutbound), |
| 382 | .tails = KJ_MAP(t, tails) { return kj::addRef(*t); }, |
| 383 | .streamingTails = KJ_MAP(t, streamingTails) { return kj::addRef(*t); }, |
| 384 | .ownContent = kj::mv(newOwnContent), |
| 385 | .ownContentIsRpcResponse = ownContentIsRpcResponse, |
| 386 | }; |
| 387 | } |
| 388 | }; |
| 389 | |
| 390 | // A Frankenvalue::CapTableEntry which directly references a numbered I/O channel. This is ONLY |
| 391 | // valid to use when the `Frankenvalue` is being deserialized as the `env` object of an isolate. |
| 392 | // The caller should use frankenvalue.rewriteCaps() to rewrite the cap table entries into |
| 393 | // IoChannelCapTableEntry, building the I/O channel table as it goes. |
| 394 | class IoChannelCapTableEntry final: public Frankenvalue::CapTableEntry { |
| 395 | public: |
| 396 | enum Type { |
| 397 | SUBREQUEST, |
| 398 | ACTOR_CLASS, |
| 399 | // TODO(someday): Other channel types, maybe. |
| 400 | }; |
| 401 | |
| 402 | IoChannelCapTableEntry(Type type, uint channel): type(type), channel(channel) {} |
| 403 | |
| 404 | Type getType() const { |
| 405 | return type; |
| 406 | } |
| 407 | |
| 408 | // Throws if type doesn't match. |
| 409 | uint getChannelNumber(Type expectedType); |
| 410 | |
| 411 | kj::Own<CapTableEntry> clone() override; |
| 412 | kj::Own<CapTableEntry> threadSafeClone() const override; |
| 413 | |
| 414 | private: |
| 415 | Type type; |
| 416 | uint channel; |
| 417 | }; |
| 418 | |
| 419 | } // namespace workerd |