Skip to content
File

Blob: src/workerd/api/actor.c++

12.4 KB
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#include "actor.h"
6 
7#include <workerd/io/features.h>
8 
9#include <capnp/compat/byte-stream.h>
10#include <capnp/compat/http-over-capnp.h>
11#include <capnp/message.h>
12#include <capnp/schema.h>
13#include <kj/compat/http.h>
14#include <kj/encoding.h>
15 
16namespace workerd::api {
17 
18kj::Own<WorkerInterface> LocalActorOutgoingFactory::newSingleUseClient(
19 kj::Maybe<kj::String> cfStr) {
20 auto& context = IoContext::current();
21 
22 return context.getMetrics().wrapActorSubrequestClient(context.getSubrequest(
23 [&](TraceContext& tracing, IoChannelFactory& ioChannelFactory) {
24 tracing.setTag("objectId"_kjc, actorId.asPtr());
25 
26 // Lazily initialize actorChannel
27 if (actorChannel == kj::none) {
28 actorChannel =
29 context.getColoLocalActorChannel(channelId, actorId, tracing.getInternalSpanParent());
30 }
31 
32 return KJ_REQUIRE_NONNULL(actorChannel)
33 ->startRequest({.cfBlobJson = kj::mv(cfStr),
34 .parentSpan = tracing.getInternalSpanParent(),
35 .userSpanParent = tracing.getUserSpanParent()});
36 },
37 {.inHouse = true,
38 .wrapMetrics = true,
39 .operationName = kj::ConstString("durable_object_subrequest"_kjc)}));
40}
41 
42kj::Own<WorkerInterface> GlobalActorOutgoingFactory::newSingleUseClient(
43 kj::Maybe<kj::String> cfStr) {
44 auto& context = IoContext::current();
45 
46 return context.getMetrics().wrapActorSubrequestClient(context.getSubrequest(
47 [&](TraceContext& tracing, IoChannelFactory& ioChannelFactory) {
48 tracing.setTag("objectId"_kjc, id->toString());
49 
50 // Lazily initialize actorChannel
51 if (actorChannel == kj::none) {
52 KJ_SWITCH_ONEOF(channelIdOrFactory) {
53 KJ_CASE_ONEOF(channelId, uint) {
54 actorChannel = context.getGlobalActorChannel(channelId, id->getInner(),
55 kj::mv(locationHint), mode, enableReplicaRouting, routingMode,
56 tracing.getInternalSpanParent(), kj::mv(version));
57 }
58 KJ_CASE_ONEOF(factory, kj::Own<DurableObjectNamespace::ActorChannelFactory>) {
59 actorChannel = factory->getGlobalActor(id->getInner(), kj::mv(locationHint), mode,
60 enableReplicaRouting, routingMode, tracing.getInternalSpanParent(), kj::mv(version));
61 }
62 }
63 }
64 
65 return KJ_REQUIRE_NONNULL(actorChannel)
66 ->startRequest({.cfBlobJson = kj::mv(cfStr),
67 .parentSpan = tracing.getInternalSpanParent(),
68 .userSpanParent = tracing.getUserSpanParent()});
69 },
70 {.inHouse = true,
71 .wrapMetrics = true,
72 .operationName = kj::ConstString("durable_object_subrequest"_kjc)}));
73}
74 
75kj::Own<WorkerInterface> ReplicaActorOutgoingFactory::newSingleUseClient(
76 kj::Maybe<kj::String> cfStr) {
77 auto& context = IoContext::current();
78 
79 return context.getMetrics().wrapActorSubrequestClient(context.getSubrequest(
80 [&](TraceContext& tracing, IoChannelFactory& ioChannelFactory) {
81 tracing.setTag("objectId"_kjc, actorId.asPtr());
82 
83 // Unlike in `GlobalActorOutgoingFactory`, we do not create this lazily, since our channel was
84 // already open prior to this DO starting up.
85 return actorChannel->startRequest({.cfBlobJson = kj::mv(cfStr),
86 .parentSpan = tracing.getInternalSpanParent(),
87 .userSpanParent = tracing.getUserSpanParent()});
88 },
89 {.inHouse = true,
90 .wrapMetrics = true,
91 .operationName = kj::ConstString("durable_object_subrequest"_kjc)}));
92}
93 
94jsg::Ref<Fetcher> ColoLocalActorNamespace::get(jsg::Lock& js, kj::String actorId) {
95 JSG_REQUIRE(actorId.size() > 0 && actorId.size() <= 2048, TypeError,
96 "Actor ID length must be in the range [1, 2048].");
97 
98 auto& context = IoContext::current();
99 
100 kj::Own<api::Fetcher::OutgoingFactory> factory =
101 kj::heap<LocalActorOutgoingFactory>(channel, kj::mv(actorId));
102 auto outgoingFactory = context.addObject(kj::mv(factory));
103 
104 bool isInHouse = true;
105 return js.alloc<Fetcher>(
106 kj::mv(outgoingFactory), Fetcher::RequiresHostAndProtocol::YES, isInHouse);
107}
108 
109// =======================================================================================
110 
111kj::String DurableObjectId::toString() {
112 return id->toString();
113}
114 
115jsg::Ref<DurableObjectId> DurableObjectNamespace::newUniqueId(
116 jsg::Lock& js, jsg::Optional<NewUniqueIdOptions> options) {
117 return js.alloc<DurableObjectId>(
118 idFactory->newUniqueId(options.orDefault({}).jurisdiction.orDefault(kj::none)));
119}
120 
121jsg::Ref<DurableObjectId> DurableObjectNamespace::idFromName(jsg::Lock& js, kj::String name) {
122 return js.alloc<DurableObjectId>(idFactory->idFromName(kj::mv(name)));
123}
124 
125jsg::Ref<DurableObjectId> DurableObjectNamespace::idFromString(jsg::Lock& js, kj::String id) {
126 return js.alloc<DurableObjectId>(idFactory->idFromString(kj::mv(id)));
127}
128 
129jsg::Ref<DurableObject> DurableObjectNamespace::getByName(
130 jsg::Lock& js, kj::String name, jsg::Optional<GetDurableObjectOptions> options) {
131 auto id = js.alloc<DurableObjectId>(idFactory->idFromName(kj::mv(name)));
132 return getImpl(js, ActorGetMode::GET_OR_CREATE, kj::mv(id), kj::mv(options));
133}
134 
135jsg::Ref<DurableObject> DurableObjectNamespace::get(
136 jsg::Lock& js, jsg::Ref<DurableObjectId> id, jsg::Optional<GetDurableObjectOptions> options) {
137 return getImpl(js, ActorGetMode::GET_OR_CREATE, kj::mv(id), kj::mv(options));
138}
139 
140jsg::Ref<DurableObject> DurableObjectNamespace::getExisting(
141 jsg::Lock& js, jsg::Ref<DurableObjectId> id, jsg::Optional<GetDurableObjectOptions> options) {
142 return getImpl(js, ActorGetMode::GET_EXISTING, kj::mv(id), kj::mv(options));
143}
144 
145jsg::Ref<DurableObject> DurableObjectNamespace::getImpl(jsg::Lock& js,
146 ActorGetMode mode,
147 jsg::Ref<DurableObjectId> id,
148 jsg::Optional<GetDurableObjectOptions> options) {
149 JSG_REQUIRE(idFactory->matchesJurisdiction(id->getInner()), TypeError,
150 "get called on jurisdictional subnamespace with an ID from a different jurisdiction");
151 ActorRoutingMode routingMode = ActorRoutingMode::DEFAULT;
152 KJ_IF_SOME(o, options) {
153 KJ_IF_SOME(rm, o.routingMode) {
154 JSG_REQUIRE(rm == "primary-only", RangeError, "unknown routingMode: ", rm);
155 routingMode = ActorRoutingMode::PRIMARY_ONLY;
156 }
157 }
158 
159 auto& context = IoContext::current();
160 kj::Maybe<kj::String> locationHint;
161 kj::Maybe<ActorVersion> version;
162 KJ_IF_SOME(o, options) {
163 locationHint = kj::mv(o.locationHint);
164 if (FeatureFlags::get(js).getEnableVersionApi()) {
165 KJ_IF_SOME(v, o.version) {
166 version = ActorVersion{.cohort = kj::mv(v.cohort)};
167 }
168 }
169 }
170 
171 bool enableReplicaRouting = FeatureFlags::get(js).getReplicaRouting();
172 
173 kj::Own<Fetcher::OutgoingFactory> outgoingFactory;
174 KJ_SWITCH_ONEOF(channel) {
175 KJ_CASE_ONEOF(channelId, uint) {
176 outgoingFactory = kj::heap<GlobalActorOutgoingFactory>(channelId, id.addRef(),
177 kj::mv(locationHint), mode, enableReplicaRouting, routingMode, kj::mv(version));
178 }
179 KJ_CASE_ONEOF(channelFactory, IoOwn<ActorChannelFactory>) {
180 outgoingFactory =
181 kj::heap<GlobalActorOutgoingFactory>(kj::addRef(*channelFactory), id.addRef(),
182 kj::mv(locationHint), mode, enableReplicaRouting, routingMode, kj::mv(version));
183 }
184 }
185 
186 auto requiresHost = FeatureFlags::get(js).getDurableObjectFetchRequiresSchemeAuthority()
187 ? Fetcher::RequiresHostAndProtocol::YES
188 : Fetcher::RequiresHostAndProtocol::NO;
189 return js.alloc<DurableObject>(
190 kj::mv(id), context.addObject(kj::mv(outgoingFactory)), requiresHost);
191}
192 
193jsg::Ref<DurableObjectNamespace> DurableObjectNamespace::jurisdiction(
194 jsg::Lock& js, jsg::Optional<kj::Maybe<kj::String>> maybeJurisdiction) {
195 auto newIdFactory = idFactory->cloneWithJurisdiction(maybeJurisdiction.orDefault(kj::none));
196 
197 KJ_SWITCH_ONEOF(channel) {
198 KJ_CASE_ONEOF(channelId, uint) {
199 return js.alloc<api::DurableObjectNamespace>(channelId, kj::mv(newIdFactory));
200 }
201 KJ_CASE_ONEOF(channelFactory, IoOwn<ActorChannelFactory>) {
202 return js.alloc<api::DurableObjectNamespace>(
203 IoContext::current().addObject(kj::addRef(*channelFactory)), kj::mv(newIdFactory));
204 }
205 }
206 
207 KJ_UNREACHABLE;
208}
209 
210kj::Own<IoChannelFactory::ActorClassChannel> DurableObjectClass::getChannel(IoContext& ioctx) {
211 KJ_SWITCH_ONEOF(channel) {
212 KJ_CASE_ONEOF(number, uint) {
213 return ioctx.getIoChannelFactory().getActorClass(number);
214 }
215 KJ_CASE_ONEOF(object, IoOwn<IoChannelFactory::ActorClassChannel>) {
216 return kj::addRef(*object);
217 }
218 }
219 KJ_UNREACHABLE;
220}
221 
222void DurableObjectClass::serialize(jsg::Lock& js, jsg::Serializer& serializer) {
223 auto channel = getChannel(IoContext::current());
224 channel->requireAllowsTransfer();
225 
226 KJ_IF_SOME(handler, serializer.getExternalHandler()) {
227 KJ_IF_SOME(frankenvalueHandler, kj::tryDowncast<Frankenvalue::CapTableBuilder>(handler)) {
228 // Encoding a Frankenvalue (e.g. for dynamic loopback props or dynamic isolate env).
229 serializer.writeRawUint32(frankenvalueHandler.add(kj::mv(channel)));
230 return;
231 } else KJ_IF_SOME(rpcHandler, kj::tryDowncast<RpcSerializerExternalHandler>(handler)) {
232 JSG_REQUIRE(FeatureFlags::get(js).getWorkerdExperimental(), DOMDataCloneError,
233 "DurableObjectClass serialization requires the 'experimental' compat flag.");
234 
235 auto token = channel->getToken(IoChannelFactory::ChannelTokenUsage::RPC);
236 rpcHandler.write([token = kj::mv(token)](rpc::JsValue::External::Builder builder) {
237 builder.setActorClassChannelToken(token);
238 });
239 return;
240 }
241 // TODO(someday): structuredClone() should have special handling that just reproduces the same
242 // local object. At present we have no way to recognize structuredClone() here though.
243 }
244 
245 // The allow_irrevocable_stub_storage flag allows us to just embed the token inline. This format
246 // is temporary, anyone using this will lose their data later.
247 JSG_REQUIRE(FeatureFlags::get(js).getAllowIrrevocableStubStorage(), DOMDataCloneError,
248 "DurableObjectClass cannot be serialized in this context.");
249 serializer.writeLengthDelimited(channel->getToken(IoChannelFactory::ChannelTokenUsage::STORAGE));
250}
251 
252jsg::Ref<DurableObjectClass> DurableObjectClass::deserialize(
253 jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer) {
254 KJ_IF_SOME(handler, deserializer.getExternalHandler()) {
255 KJ_IF_SOME(frankenvalueHandler, kj::tryDowncast<Frankenvalue::CapTableReader>(handler)) {
256 // Decoding a Frankenvalue (e.g. for dynamic loopback props or dynamic isolate env).
257 auto& cap = KJ_REQUIRE_NONNULL(frankenvalueHandler.get(deserializer.readRawUint32()),
258 "serialized DurableObjectClass had invalid cap table index");
259 
260 KJ_IF_SOME(channel, kj::tryDowncast<IoChannelFactory::ActorClassChannel>(cap)) {
261 // Probably decoding dynamic ctx.props.
262 return js.alloc<DurableObjectClass>(IoContext::current().addObject(kj::addRef(channel)));
263 } else KJ_IF_SOME(channel, kj::tryDowncast<IoChannelCapTableEntry>(cap)) {
264 // Probably decoding dynamic isolate env.
265 return js.alloc<DurableObjectClass>(
266 channel.getChannelNumber(IoChannelCapTableEntry::Type::ACTOR_CLASS));
267 } else {
268 KJ_FAIL_REQUIRE(
269 "DurableObjectClass capability in Frankenvalue is not a ActorClassChannel?");
270 }
271 } else KJ_IF_SOME(rpcHandler, kj::tryDowncast<RpcDeserializerExternalHandler>(handler)) {
272 JSG_REQUIRE(FeatureFlags::get(js).getWorkerdExperimental(), DOMDataCloneError,
273 "DurableObjectClass serialization requires the 'experimental' compat flag.");
274 
275 auto external = rpcHandler.read();
276 KJ_REQUIRE(external.isActorClassChannelToken());
277 auto& ioctx = IoContext::current();
278 auto channel = ioctx.getIoChannelFactory().actorClassFromToken(
279 IoChannelFactory::ChannelTokenUsage::RPC, external.getActorClassChannelToken());
280 return js.alloc<DurableObjectClass>(ioctx.addObject(kj::mv(channel)));
281 }
282 }
283 
284 // The allow_irrevocable_stub_storage flag allows us to just embed the token inline. This format
285 // is temporary, anyone using this will lose their data later.
286 JSG_REQUIRE(FeatureFlags::get(js).getAllowIrrevocableStubStorage(), DOMDataCloneError,
287 "DOMDataCloneError cannot be deserialized in this context.");
288 auto& ioctx = IoContext::current();
289 auto channel = ioctx.getIoChannelFactory().actorClassFromToken(
290 IoChannelFactory::ChannelTokenUsage::STORAGE, deserializer.readLengthDelimitedBytes());
291 return js.alloc<DurableObjectClass>(ioctx.addObject(kj::mv(channel)));
292}
293 
294} // namespace workerd::api