File
Blob: src/workerd/api/actor.c++
| 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 | |
| 16 | namespace workerd::api { |
| 17 | |
| 18 | kj::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 | |
| 42 | kj::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 | |
| 75 | kj::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 | |
| 94 | jsg::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 | |
| 111 | kj::String DurableObjectId::toString() { |
| 112 | return id->toString(); |
| 113 | } |
| 114 | |
| 115 | jsg::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 | |
| 121 | jsg::Ref<DurableObjectId> DurableObjectNamespace::idFromName(jsg::Lock& js, kj::String name) { |
| 122 | return js.alloc<DurableObjectId>(idFactory->idFromName(kj::mv(name))); |
| 123 | } |
| 124 | |
| 125 | jsg::Ref<DurableObjectId> DurableObjectNamespace::idFromString(jsg::Lock& js, kj::String id) { |
| 126 | return js.alloc<DurableObjectId>(idFactory->idFromString(kj::mv(id))); |
| 127 | } |
| 128 | |
| 129 | jsg::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 | |
| 135 | jsg::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 | |
| 140 | jsg::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 | |
| 145 | jsg::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 | |
| 193 | jsg::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 | |
| 210 | kj::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 | |
| 222 | void 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 | |
| 252 | jsg::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 |