Skip to content
File

Blob: src/workerd/io/io-channels.h

cpp420 lines
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 
18namespace kj {
19class HttpClient;
20class Network;
21} // namespace kj
22 
23namespace workerd {
24 
25class WorkerInterface;
26 
27// Interface for talking to the Cache API. Needs to be declared here so that IoContext can
28// contain it.
29class 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.
54class 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 
73class WorkerStubChannel;
74struct 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!
99class 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.
316struct 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.
331class 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.
343struct 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.
394class 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