diff --git a/src/workerd/api/BUILD.bazel b/src/workerd/api/BUILD.bazel index bf98c2cd66d..ddb70ebdf8f 100644 --- a/src/workerd/api/BUILD.bazel +++ b/src/workerd/api/BUILD.bazel @@ -849,6 +849,15 @@ kj_test( ], ) +kj_test( + src = "jsrpc-call-retry-test.c++", + deps = [ + "//src/workerd/io", + "//src/workerd/io:worker-interface", + "//src/workerd/tests:test-fixture", + ], +) + kj_test( src = "fetch-retry-claim-hook-test.c++", deps = [ diff --git a/src/workerd/api/actor-call-retry-test.c++ b/src/workerd/api/actor-call-retry-test.c++ index 2f4f8d75130..2a0d9f0f5ca 100644 --- a/src/workerd/api/actor-call-retry-test.c++ +++ b/src/workerd/api/actor-call-retry-test.c++ @@ -172,6 +172,52 @@ KJ_TEST("actor retries return the original disconnect after claim rejection") { KJ_EXPECT(observer->outcomes[0] == ActorRetryOutcome::CLAIM_REJECTED); } +KJ_TEST("actor retries preserve ambiguity through later not-delivered and rejected attempts") { + TestTimerChannel timer; + auto observer = kj::refcounted(); + auto state = newRetryState(timer, *observer); + + auto first = startAttempt(*state); + auto nonce = KJ_ASSERT_NONNULL(first.getMetadata()).nonce; + KJ_EXPECT( + handleFailure(*state, makeDisconnect("original ambiguous disconnect"_kj)).is()); + + auto second = startAttempt(*state); + KJ_EXPECT(KJ_ASSERT_NONNULL(second.getMetadata()).nonce == nonce); + KJ_EXPECT(KJ_ASSERT_NONNULL(second.getMetadata()).isRetry == IsActorRetry::YES); + auto notDelivered = makeDisconnect("later not-delivered disconnect"_kj); + jsg::markActorRequestNotDelivered(notDelivered); + KJ_EXPECT(handleFailure(*state, kj::mv(notDelivered)).is()); + + auto third = startAttempt(*state); + KJ_EXPECT(KJ_ASSERT_NONNULL(third.getMetadata()).nonce == nonce); + KJ_EXPECT(KJ_ASSERT_NONNULL(third.getMetadata()).isRetry == IsActorRetry::YES); + auto rejected = KJ_EXCEPTION(FAILED, "claim rejected"); + rejected.setDetail(jsg::ACTOR_RETRY_CLAIM_REJECTED_DETAIL_ID, kj::heapArray(0)); + auto result = handleFailure(*state, kj::mv(rejected)); + auto& failure = KJ_ASSERT_NONNULL(result.tryGet()); + KJ_EXPECT(failure.getDescription().contains("original ambiguous disconnect")); + KJ_ASSERT(observer->outcomes.size() == 1); + KJ_EXPECT(observer->outcomes[0] == ActorRetryOutcome::CLAIM_REJECTED); +} + +KJ_TEST("a committed retry preserves the original disconnect after claim rejection") { + TestTimerChannel timer; + auto observer = kj::refcounted(); + auto state = newRetryState(timer, *observer); + + startAttempt(*state); + KJ_EXPECT(handleFailure(*state, makeDisconnect("original disconnect"_kj)).is()); + startAttempt(*state); + auto rejected = KJ_EXCEPTION(FAILED, "claim rejected"); + rejected.setDetail(jsg::ACTOR_RETRY_CLAIM_REJECTED_DETAIL_ID, kj::heapArray(0)); + auto failure = state->handleCommittedAttemptFailure(kj::mv(rejected)); + + KJ_EXPECT(failure.getDescription().contains("original disconnect")); + KJ_ASSERT(observer->outcomes.size() == 1); + KJ_EXPECT(observer->outcomes[0] == ActorRetryOutcome::CLAIM_REJECTED); +} + KJ_TEST("actor retries stop after five total attempts") { TestTimerChannel timer; auto observer = kj::refcounted(); @@ -245,5 +291,40 @@ KJ_TEST("actor calls do not retry when retry requests are disabled") { KJ_EXPECT(observer->outcomes.size() == 0); } +KJ_TEST("disabled actor retries return failures unchanged") { + auto expectDisabled = [](ActorCallRetryState::Config config) { + TestTimerChannel timer; + auto observer = kj::refcounted(); + auto state = kj::rc(timer, *observer, config); + + startAttempt(*state); + auto result = handleFailure(*state, makeDisconnect("original disconnect"_kj)); + auto& failure = KJ_ASSERT_NONNULL(result.tryGet()); + KJ_EXPECT(failure.getType() == kj::Exception::Type::DISCONNECTED); + KJ_EXPECT(failure.getDescription().contains("original disconnect")); + KJ_EXPECT(observer->retryCallTypes.empty()); + KJ_EXPECT(observer->outcomes.empty()); + }; + + expectDisabled({ + .callType = ActorRetryCallType::JSRPC, + .observationEnabled = ActorRetryGateEnabled::YES, + .enforcementEnabled = ActorRetryGateEnabled::NO, + .payloadReplayable = ActorCallPayloadReplayable::YES, + }); + expectDisabled({ + .callType = ActorRetryCallType::JSRPC, + .observationEnabled = ActorRetryGateEnabled::NO, + .enforcementEnabled = ActorRetryGateEnabled::YES, + .payloadReplayable = ActorCallPayloadReplayable::YES, + }); + expectDisabled({ + .callType = ActorRetryCallType::JSRPC, + .observationEnabled = ActorRetryGateEnabled::YES, + .enforcementEnabled = ActorRetryGateEnabled::YES, + .payloadReplayable = ActorCallPayloadReplayable::NO, + }); +} + } // namespace } // namespace workerd::api diff --git a/src/workerd/api/actor-call-retry.c++ b/src/workerd/api/actor-call-retry.c++ index 5bf320159fb..e723b9d21ed 100644 --- a/src/workerd/api/actor-call-retry.c++ +++ b/src/workerd/api/actor-call-retry.c++ @@ -51,6 +51,8 @@ kj::OneOf ActorCallRetryState::star kj::OneOf ActorCallRetryState::handleAttemptFailure( kj::Exception exception) { + if (!retriesEnabled) return kj::mv(exception); + maybeStartRetryLatencyTimer(exception); KJ_IF_SOME(claimRejection, handleClaimRejection(exception)) { return kj::mv(claimRejection); @@ -59,6 +61,17 @@ kj::OneOf ActorCallRetryState::handleAttemptFailure return checkCanRetry(kj::mv(exception)); } +kj::Exception ActorCallRetryState::handleCommittedAttemptFailure(kj::Exception exception) { + if (!retriesEnabled) return kj::mv(exception); + + maybeStartRetryLatencyTimer(exception); + KJ_IF_SOME(claimRejection, handleClaimRejection(exception)) { + return kj::mv(claimRejection); + } + recordOutcome(ActorRetryOutcome::UNABLE_TO_RETRY); + return kj::mv(exception); +} + void ActorCallRetryState::maybeStartRetryLatencyTimer(const kj::Exception& exception) { if (exception.getType() == kj::Exception::Type::DISCONNECTED && retryStartTime == kj::none) { retryStartTime = kj::systemPreciseMonotonicClock().now(); diff --git a/src/workerd/api/actor-call-retry.h b/src/workerd/api/actor-call-retry.h index ae233a3675f..c9f9e0eaa51 100644 --- a/src/workerd/api/actor-call-retry.h +++ b/src/workerd/api/actor-call-retry.h @@ -63,6 +63,11 @@ class ActorCallRetryState final: public kj::Refcounted { kj::OneOf startAttempt(); kj::OneOf handleAttemptFailure(kj::Exception exception); + kj::Exception handleCommittedAttemptFailure(kj::Exception exception); + + kj::Exception getOriginalDisconnect() const { + return KJ_ASSERT_NONNULL(originalDisconnect).clone(); + } bool isRetryEnabled() const { return retriesEnabled; diff --git a/src/workerd/api/http.h b/src/workerd/api/http.h index fe229b1968d..279a6e16fcf 100644 --- a/src/workerd/api/http.h +++ b/src/workerd/api/http.h @@ -349,7 +349,7 @@ class Fetcher: public JsRpcClientProvider { MakeUserSpanParent makeUserSpanParent); kj::Maybe getActorTargetRetryability() override; - void onActorCallRetry(); + void onActorCallRetry() override; // Get a SubrequestChannel representing this Fetcher. kj::Own getSubrequestChannel(IoContext& ioContext); diff --git a/src/workerd/api/jsrpc-call-observation-test.c++ b/src/workerd/api/jsrpc-call-observation-test.c++ index 60ec5780d66..54f7d6b760d 100644 --- a/src/workerd/api/jsrpc-call-observation-test.c++ +++ b/src/workerd/api/jsrpc-call-observation-test.c++ @@ -37,6 +37,7 @@ struct Observation { ActorCallPayloadReplayable payloadReplayable; ActorCallTargetRetryable targetRetryable; Settlement settlement = Settlement::PENDING; + bool pipelineCommitted = false; }; struct ObservationState { @@ -55,6 +56,9 @@ class RecordingCallObserver final: public OutgoingActorCallObserver { void recordSuccess() override { settle(Settlement::SUCCESS); } + void markPipelineCommitted() override { + observation.pipelineCommitted = true; + } void recordFailure(kj::Exception&) override { settle(Settlement::FAILURE); } @@ -215,7 +219,7 @@ class Harness { auto& targetHandler = KJ_REQUIRE_NONNULL(env.js.tryGetTypeHandler>()); auto target = KJ_REQUIRE_NONNULL(jsg::JsValue(targetHandler.wrap(env.js, env.js.alloc())) - .tryCast()); + .tryCast()); return JsRpcStub::constructor(env.js, target); } @@ -345,6 +349,8 @@ KJ_TEST("using an RPC result pipeline releases projected replay memory") { return kj::mv(pipelined).attach(kj::mv(child), kj::mv(fetcher)); }); + KJ_ASSERT(harness.state.observations.size() == 2); + KJ_EXPECT(harness.state.observations[0]->pipelineCommitted); KJ_EXPECT(harness.state.replayMemoryBytes == 0); } @@ -366,7 +372,7 @@ KJ_TEST("disposing an RPC promise does not release projected replay memory early KJ_EXPECT(harness.state.replayMemoryBytes == 0); } -KJ_TEST("an RPC property get is observed with a non-replayable payload") { +KJ_TEST("an RPC property get is observed with a replayable payload") { Harness harness; harness.sender->runInIoContext([&](const TestFixture::Environment& env) { auto fetcher = harness.makeFetcher(env).fetcher; @@ -375,7 +381,7 @@ KJ_TEST("an RPC property get is observed with a non-replayable payload") { auto& observations = harness.state.observations; KJ_ASSERT(observations.size() == 1); - KJ_EXPECT(observations[0]->payloadReplayable == ActorCallPayloadReplayable::NO); + KJ_EXPECT(observations[0]->payloadReplayable == ActorCallPayloadReplayable::YES); KJ_EXPECT(observations[0]->targetRetryable == ActorCallTargetRetryable::YES); KJ_EXPECT(observations[0]->settlement == Settlement::SUCCESS); } diff --git a/src/workerd/api/jsrpc-call-plan-test.c++ b/src/workerd/api/jsrpc-call-plan-test.c++ index 0bd477c2792..826d92e4e62 100644 --- a/src/workerd/api/jsrpc-call-plan-test.c++ +++ b/src/workerd/api/jsrpc-call-plan-test.c++ @@ -128,7 +128,7 @@ KJ_TEST("JS RPC call plan copies calls and property accesses") { builder.setMethodName("value"); builder.getOperation().setGetProperty(); }); - KJ_EXPECT(!property.getReplayable()); + KJ_EXPECT(property.getReplayable()); capnp::MallocMessageBuilder propertyAttempt; property.copyTo(propertyAttempt.initRoot()); auto propertyParams = propertyAttempt.getRoot(); @@ -158,6 +158,23 @@ KJ_TEST("JS RPC call plan copies calls and property accesses") { } } +KJ_TEST("JS RPC call plan replay memory includes variable-sized metadata") { + auto smallPlan = makePlan( + [](rpc::JsRpcTarget::CallParams::Builder builder) { builder.setMethodName("method"); }); + KJ_EXPECT(smallPlan.getReplayMemoryEstimate() >= JsRpcCallPlan::REPLAY_MEMORY_OVERHEAD + + JsRpcCallPlan::METADATA_SEGMENT_WORDS * sizeof(capnp::word)); + + auto methodName = kj::heapString(4096); + for (auto& character: methodName) { + character = 'x'; + } + auto plan = makePlan( + [&](rpc::JsRpcTarget::CallParams::Builder builder) { builder.setMethodName(methodName); }); + + KJ_EXPECT( + plan.getReplayMemoryEstimate() >= JsRpcCallPlan::REPLAY_MEMORY_OVERHEAD + methodName.size()); +} + KJ_TEST("JS RPC call plan copies usable external capabilities but rejects replay") { kj::EventLoop loop; kj::WaitScope waitScope(loop); @@ -644,6 +661,25 @@ KJ_TEST("replayable actor RPC calls carry observe-only retry metadata") { KJ_EXPECT(recordedMetadata.retryGateEnabled == ActorRetryGateEnabled::NO); } +KJ_TEST("JSRPC enforcement remains disabled without fetch enforcement") { + auto dispatch = makeReplayableActorCall(kj::arr("durable-object-retries-fetch"_kj, + "durable-object-retries-jsrpc"_kj, "durable-object-retries-jsrpc-retry-requests"_kj)); + + KJ_EXPECT(dispatch.singleUseCount == 0); + KJ_EXPECT(dispatch.actorAttemptCount == 1); + KJ_EXPECT(KJ_ASSERT_NONNULL(dispatch.metadata).retryGateEnabled == ActorRetryGateEnabled::NO); +} + +KJ_TEST("JSRPC enforcement remains disabled when replay memory is not reserved") { + auto dispatch = makeReplayableActorCall( + kj::arr("durable-object-retries-fetch"_kj, "durable-object-retries-fetch-retry-requests"_kj, + "durable-object-retries-jsrpc"_kj, "durable-object-retries-jsrpc-retry-requests"_kj)); + + KJ_EXPECT(dispatch.singleUseCount == 0); + KJ_EXPECT(dispatch.actorAttemptCount == 1); + KJ_EXPECT(KJ_ASSERT_NONNULL(dispatch.metadata).retryGateEnabled == ActorRetryGateEnabled::NO); +} + KJ_TEST("replayable actor RPC calls carry no retry metadata without the JSRPC gate") { auto dispatch = makeReplayableActorCall( kj::arr("durable-object-retries-fetch"_kj, "durable-object-retries-fetch-retry-requests"_kj)); @@ -661,13 +697,13 @@ KJ_TEST("replayable actor RPC calls carry no retry metadata without the fetch ga KJ_EXPECT(dispatch.metadata == kj::none); } -KJ_TEST("actor RPC property reads carry no retry metadata") { +KJ_TEST("actor RPC property reads carry observe-only retry metadata") { auto dispatch = makeActorPropertyRead( kj::arr("durable-object-retries-fetch"_kj, "durable-object-retries-jsrpc"_kj)); - KJ_EXPECT(dispatch.singleUseCount == 1); - KJ_EXPECT(dispatch.actorAttemptCount == 0); - KJ_EXPECT(dispatch.metadata == kj::none); + KJ_EXPECT(dispatch.singleUseCount == 0); + KJ_EXPECT(dispatch.actorAttemptCount == 1); + KJ_EXPECT(KJ_ASSERT_NONNULL(dispatch.metadata).retryGateEnabled == ActorRetryGateEnabled::NO); } // A Durable Object whose methods fail in the ways the receiver must classify as delivered. diff --git a/src/workerd/api/jsrpc-call-retry-test.c++ b/src/workerd/api/jsrpc-call-retry-test.c++ new file mode 100644 index 00000000000..07974355c6d --- /dev/null +++ b/src/workerd/api/jsrpc-call-retry-test.c++ @@ -0,0 +1,536 @@ +// Copyright (c) 2026 Cloudflare, Inc. +// Licensed under the Apache 2.0 license found in the LICENSE file or at: +// https://opensource.org/licenses/Apache-2.0 + +#include "http.h" +#include "worker-rpc.h" + +#include + +#include +#include + +namespace workerd::api { +namespace { + +enum class FirstFailure { + AMBIGUOUS, + NOT_DELIVERED, +}; + +// What the attempt after the failing ones connects to. +enum class Replacement { + RECEIVER, + HANGING, +}; + +// How a KJ DISCONNECTED exception surfaces to JS; its description is not exposed. +constexpr kj::StringPtr DISCONNECT_JS_MESSAGE = "Error: Network connection lost."_kj; + +constexpr kj::StringPtr RECEIVER_SOURCE = R"JS( + import { WorkerEntrypoint } from "cloudflare:workers"; + + export default class extends WorkerEntrypoint { + echo(value) { return value; } + get value() { return 42; } + } +)JS"_kj; + +struct RetryTestState { + kj::Vector metadata; + kj::Vector countSubrequests; + kj::Vector outcomes; + uint acceptedRetries = 0; + uint observedRetries = 0; + uint observedAttempts = 0; + uint committedAttempts = 0; + size_t replayMemoryBytes = 0; + // `replayMemoryBytes` when the retry outcome was recorded, the first step after reentry. + kj::Maybe replayMemoryBytesAtOutcome; + kj::Maybe>> replacementStarted; +}; + +class ImmediateTimerChannel final: public TimerChannel { + public: + void syncTime() override {} + + kj::Date now(kj::Maybe) override { + return kj::UNIX_EPOCH; + } + + kj::Promise atTime(kj::Date) override { + return kj::READY_NOW; + } + + kj::Promise afterLimitTimeout(kj::Duration) override { + return kj::READY_NOW; + } + + kj::TimePoint nowForLimitTimeout() override { + return kj::origin(); + } +}; + +class PausingTimerChannel final: public TimerChannel { + public: + explicit PausingTimerChannel(uint pauseAtBackoff): pauseAtBackoff(pauseAtBackoff) { + auto started = kj::newPromiseAndFulfiller(); + backoffStarted = kj::mv(started.promise); + backoffStartedFulfiller = kj::mv(started.fulfiller); + } + + void syncTime() override {} + + kj::Date now(kj::Maybe) override { + return kj::UNIX_EPOCH; + } + + kj::Promise atTime(kj::Date) override { + return kj::READY_NOW; + } + + kj::Promise afterLimitTimeout(kj::Duration) override { + if (++backoffCount == pauseAtBackoff) { + KJ_REQUIRE_NONNULL(backoffStartedFulfiller)->fulfill(); + backoffStartedFulfiller = kj::none; + return kj::NEVER_DONE; + } + return kj::READY_NOW; + } + + kj::TimePoint nowForLimitTimeout() override { + return kj::origin(); + } + + kj::Promise onBackoffStarted() { + return kj::mv(backoffStarted); + } + + private: + kj::Promise backoffStarted = nullptr; + kj::Maybe>> backoffStartedFulfiller; + uint pauseAtBackoff; + uint backoffCount = 0; +}; + +class RetryCallObserver final: public OutgoingActorCallObserver { + public: + explicit RetryCallObserver(RetryTestState& state): state(state) {} + + void markPipelineCommitted() override { + if (!pipelineCommitted) { + pipelineCommitted = true; + ++state.committedAttempts; + } + } + + private: + RetryTestState& state; + bool pipelineCommitted = false; +}; + +class RetryObserver final: public RequestObserver { + public: + explicit RetryObserver(RetryTestState& state): state(state) {} + + kj::Maybe> observeOutgoingActorRpcCall( + ActorCallPayloadReplayable payloadReplayable, ActorCallTargetRetryable) override { + KJ_EXPECT(payloadReplayable == ActorCallPayloadReplayable::YES); + ++state.observedAttempts; + return kj::heap(state); + } + + void recordActorRetry(ActorRetryCallType callType) override { + KJ_EXPECT(callType == ActorRetryCallType::JSRPC); + ++state.observedRetries; + } + + void recordActorRetryOutcome( + ActorRetryCallType callType, ActorRetryOutcome outcome, kj::Duration) override { + KJ_EXPECT(callType == ActorRetryCallType::JSRPC); + state.outcomes.add(outcome); + state.replayMemoryBytesAtOutcome = state.replayMemoryBytes; + } + + kj::Own trackActorCallReplayMemory(size_t bytes) override { + return trackMemory(bytes); + } + + kj::Maybe> tryReserveActorCallReplayMemory(size_t bytes) override { + return trackMemory(bytes); + } + + private: + kj::Own trackMemory(size_t bytes) { + state.replayMemoryBytes += bytes; + return kj::heap(kj::defer([&state = state, bytes]() { state.replayMemoryBytes -= bytes; })); + } + + RetryTestState& state; +}; + +// A session that fails before delivering its event with the given disconnect. +kj::Own newFailingSession(FirstFailure firstFailure) { + auto exception = KJ_EXCEPTION(DISCONNECTED, "ambiguous JSRPC session failure"); + if (firstFailure == FirstFailure::NOT_DELIVERED) { + jsg::markActorRequestNotDelivered(exception); + } + return newPromisedWorkerInterface(kj::mv(exception)); +} + +// A session that never delivers its event. +kj::Own newHangingSession() { + return newPromisedWorkerInterface(kj::NEVER_DONE); +} + +class RetryOutgoingFactory final: public Fetcher::OutgoingFactory { + public: + RetryOutgoingFactory(TestFixture& receiver, + RetryTestState& state, + FirstFailure firstFailure, + uint failingAttempts, + Replacement replacement) + : receiver(receiver), + state(state), + firstFailure(firstFailure), + failingAttempts(failingAttempts), + replacement(replacement) {} + + Result newSingleUseClient(kj::Maybe, MakeUserSpanParent) override { + KJ_FAIL_REQUIRE("retryable JSRPC call bypassed actor attempt plumbing"); + } + + kj::Maybe getActorTargetRetryability() const override { + return ActorCallTargetRetryable::YES; + } + + void onActorCallRetry() override { + ++state.acceptedRetries; + } + + Result newActorCallAttempt( + kj::Maybe, ActorCallRetryState::Attempt attempt, MakeUserSpanParent) override { + state.countSubrequests.add(attempt.getCountSubrequest()); + state.metadata.add(KJ_REQUIRE_NONNULL(attempt.takeMetadata())); + if (state.metadata.size() <= failingAttempts) { + return {.client = newFailingSession(firstFailure), .spanParents = kj::none}; + } + KJ_IF_SOME(fulfiller, state.replacementStarted) { + fulfiller->fulfill(); + state.replacementStarted = kj::none; + } + if (replacement == Replacement::HANGING) { + return {.client = newHangingSession(), .spanParents = kj::none}; + } + return {.client = receiver.makeWorkerEntrypoint(), .spanParents = kj::none}; + } + + private: + TestFixture& receiver; + RetryTestState& state; + FirstFailure firstFailure; + uint failingAttempts; + Replacement replacement; +}; + +CompatibilityFlags::Reader makeRetryFlags(capnp::MallocMessageBuilder& message) { + auto flags = message.initRoot(); + flags.setFetcherRpc(true); + return flags.asReader(); +} + +TestFixture::SetupParams makeReceiverParams(kj::WaitScope& waitScope) { + return { + .waitScope = waitScope, + .mainModuleSource = RECEIVER_SOURCE, + .useRealTimers = false, + }; +} + +TestFixture::SetupParams makeSenderParams(kj::WaitScope& waitScope, + CompatibilityFlags::Reader flags, + TimerChannel& timer, + RetryTestState& state) { + return { + .waitScope = waitScope, + .featureFlags = flags, + .useRealTimers = false, + .autogates = + kj::arr("durable-object-retries-fetch"_kj, "durable-object-retries-fetch-retry-requests"_kj, + "durable-object-retries-jsrpc"_kj, "durable-object-retries-jsrpc-retry-requests"_kj), + .ioChannelFactory = kj::Function(TimerChannel&)>( + [&timer](TimerChannel&) -> kj::Rc { + return kj::rc(timer); + }), + .requestObserverFactory = kj::Function()>( + [&state]() -> kj::Own { return kj::refcounted(state); }), + }; +} + +jsg::Ref makeRetryFetcher(const TestFixture::Environment& env, + TestFixture& receiver, + RetryTestState& state, + FirstFailure firstFailure, + uint failingAttempts, + Replacement replacement = Replacement::RECEIVER) { + return env.js.alloc( + env.context.addObject(kj::heap( + receiver, state, firstFailure, failingAttempts, replacement)), + Fetcher::RequiresHostAndProtocol::YES); +} + +// Pipelines a call to `child()` on `parent`, returning the child's result promise. +jsg::JsValue callChild(jsg::Lock& js, JsRpcPromise& parent) { + auto childMethod = KJ_REQUIRE_NONNULL(parent.getProperty(js, kj::str("child"))); + auto& handler = KJ_REQUIRE_NONNULL(js.tryGetTypeHandler>()); + auto childFunction = KJ_REQUIRE_NONNULL( + jsg::JsValue(handler.wrap(js, kj::mv(childMethod))).tryCast()); + return childFunction.call(js, js.undefined()); +} + +// Resolves once `value` rejects with a disconnect, as the stored first-attempt failure is. +jsg::Promise expectDisconnect(jsg::Lock& js, jsg::JsValue value) { + return js.toPromise(value).then(js, [](jsg::Lock& js, jsg::Value) { + KJ_FAIL_ASSERT("committed RPC attempt unexpectedly succeeded"); + }, [](jsg::Lock& js, jsg::Value error) { + auto message = jsg::JsValue(error.getHandle(js)).toString(js); + KJ_EXPECT(message == DISCONNECT_JS_MESSAGE, message); + }); +} + +jsg::JsFunction getRpcFunction(jsg::Lock& js, Fetcher& fetcher, kj::StringPtr name) { + auto method = KJ_REQUIRE_NONNULL(fetcher.getRpcMethodForTestOnly(js, kj::str(name))); + auto& handler = KJ_REQUIRE_NONNULL(js.tryGetTypeHandler>()); + return KJ_REQUIRE_NONNULL( + jsg::JsValue(handler.wrap(js, kj::mv(method))).tryCast()); +} + +KJ_TEST("replayable actor RPC retries an ambiguous session failure") { + auto io = kj::setupAsyncIo(); + capnp::MallocMessageBuilder flagsMessage; + RetryTestState state; + ImmediateTimerChannel timer; + TestFixture receiver(makeReceiverParams(io.waitScope)); + TestFixture sender(makeSenderParams(io.waitScope, makeRetryFlags(flagsMessage), timer, state)); + + sender.runInIoContext([&](const TestFixture::Environment& env) { + auto fetcher = makeRetryFetcher(env, receiver, state, FirstFailure::AMBIGUOUS, 1); + auto function = getRpcFunction(env.js, *fetcher, "echo"_kj); + auto result = function.call(env.js, env.js.undefined(), env.js.num(42)); + auto checked = env.js.toPromise(result).then(env.js, [](jsg::Lock& js, jsg::Value value) { + KJ_EXPECT(jsg::JsValue(value.getHandle(js)).strictEquals(js.num(42))); + }); + return env.context.awaitJs(env.js, kj::mv(checked)).attach(kj::mv(fetcher)); + }); + + KJ_ASSERT(state.metadata.size() == 2); + KJ_EXPECT(state.metadata[0].nonce == state.metadata[1].nonce); + KJ_EXPECT(state.metadata[0].isRetry == IsActorRetry::NO); + KJ_EXPECT(state.metadata[1].isRetry == IsActorRetry::YES); + KJ_ASSERT(state.countSubrequests.size() == 2); + KJ_EXPECT(state.countSubrequests[0] == CountSubrequest::YES); + KJ_EXPECT(state.countSubrequests[1] == CountSubrequest::NO); + KJ_EXPECT(state.acceptedRetries == 1); + KJ_EXPECT(state.observedRetries == 1); + KJ_EXPECT(state.observedAttempts == 2); + KJ_ASSERT(state.outcomes.size() == 1); + KJ_EXPECT(state.outcomes[0] == ActorRetryOutcome::RECOVERED); + // The payload is released at native success, before the call reenters the isolate. + KJ_EXPECT(KJ_ASSERT_NONNULL(state.replayMemoryBytesAtOutcome) == 0); + KJ_EXPECT(state.replayMemoryBytes == 0); +} + +KJ_TEST("actor RPC retries a not-delivered failure with a fresh token") { + auto io = kj::setupAsyncIo(); + capnp::MallocMessageBuilder flagsMessage; + RetryTestState state; + ImmediateTimerChannel timer; + TestFixture receiver(makeReceiverParams(io.waitScope)); + TestFixture sender(makeSenderParams(io.waitScope, makeRetryFlags(flagsMessage), timer, state)); + + sender.runInIoContext([&](const TestFixture::Environment& env) { + auto fetcher = makeRetryFetcher(env, receiver, state, FirstFailure::NOT_DELIVERED, 1); + auto function = getRpcFunction(env.js, *fetcher, "echo"_kj); + auto result = function.call(env.js, env.js.undefined(), env.js.num(42)); + auto checked = env.js.toPromise(result).then(env.js, [](jsg::Lock& js, jsg::Value value) { + KJ_EXPECT(jsg::JsValue(value.getHandle(js)).strictEquals(js.num(42))); + }); + return env.context.awaitJs(env.js, kj::mv(checked)).attach(kj::mv(fetcher)); + }); + + KJ_ASSERT(state.metadata.size() == 2); + KJ_EXPECT(state.metadata[0].nonce != state.metadata[1].nonce); + KJ_EXPECT(state.metadata[0].isRetry == IsActorRetry::NO); + KJ_EXPECT(state.metadata[1].isRetry == IsActorRetry::NO); + KJ_EXPECT(state.acceptedRetries == 1); + KJ_EXPECT(state.observedRetries == 1); + KJ_EXPECT(state.observedAttempts == 2); + KJ_ASSERT(state.outcomes.size() == 1); + KJ_EXPECT(state.outcomes[0] == ActorRetryOutcome::RECOVERED); + KJ_EXPECT(state.replayMemoryBytes == 0); +} + +KJ_TEST("pipelining before failure commits the actor RPC to its first attempt") { + auto io = kj::setupAsyncIo(); + capnp::MallocMessageBuilder flagsMessage; + RetryTestState state; + ImmediateTimerChannel timer; + TestFixture receiver(makeReceiverParams(io.waitScope)); + TestFixture sender(makeSenderParams(io.waitScope, makeRetryFlags(flagsMessage), timer, state)); + + sender.runInIoContext([&](const TestFixture::Environment& env) { + auto fetcher = makeRetryFetcher(env, receiver, state, FirstFailure::AMBIGUOUS, 1); + auto function = getRpcFunction(env.js, *fetcher, "echo"_kj); + auto parentValue = function.call(env.js, env.js.undefined(), env.js.num(42)); + auto parentObject = KJ_REQUIRE_NONNULL(parentValue.tryCast()); + auto parent = KJ_REQUIRE_NONNULL(parentObject.tryUnwrapAs(env.js)); + + auto childValue = callChild(env.js, *parent); + KJ_EXPECT(state.replayMemoryBytes == 0); + + auto expectRejected = [&](jsg::JsValue value) { + return env.context.awaitJs(env.js, expectDisconnect(env.js, value)); + }; + return kj::joinPromises(kj::arr(expectRejected(parentValue), expectRejected(childValue))) + .attach(kj::mv(parent), kj::mv(fetcher)); + }); + + KJ_EXPECT(state.metadata.size() == 1); + KJ_EXPECT(state.acceptedRetries == 0); + KJ_EXPECT(state.observedRetries == 0); + KJ_EXPECT(state.committedAttempts == 1); + KJ_EXPECT(state.outcomes.empty()); + KJ_EXPECT(state.replayMemoryBytes == 0); +} + +KJ_TEST("pipelining during backoff prevents a replacement actor RPC attempt") { + auto io = kj::setupAsyncIo(); + capnp::MallocMessageBuilder flagsMessage; + RetryTestState state; + PausingTimerChannel timer(1); + TestFixture receiver(makeReceiverParams(io.waitScope)); + TestFixture sender(makeSenderParams(io.waitScope, makeRetryFlags(flagsMessage), timer, state)); + + sender.runInIoContext([&](const TestFixture::Environment& env) { + auto fetcher = makeRetryFetcher(env, receiver, state, FirstFailure::AMBIGUOUS, 1); + auto function = getRpcFunction(env.js, *fetcher, "echo"_kj); + auto parentValue = function.call(env.js, env.js.undefined(), env.js.num(42)); + auto parentObject = KJ_REQUIRE_NONNULL(parentValue.tryCast()); + auto parent = KJ_REQUIRE_NONNULL(parentObject.tryUnwrapAs(env.js)); + + auto parentRejected = env.context.awaitJs(env.js, expectDisconnect(env.js, parentValue)); + auto commitDuringBackoff = env.context.awaitIo(env.js, timer.onBackoffStarted(), + [parent = parent.addRef(), &state](jsg::Lock& js) mutable { + auto childValue = callChild(js, *parent); + KJ_EXPECT(state.replayMemoryBytes == 0); + return expectDisconnect(js, childValue); + }); + auto childRejected = env.context.awaitJs(env.js, kj::mv(commitDuringBackoff)); + return kj::joinPromises(kj::arr(kj::mv(parentRejected), kj::mv(childRejected))) + .attach(kj::mv(parent), kj::mv(fetcher)); + }); + + KJ_EXPECT(state.metadata.size() == 1); + KJ_EXPECT(state.acceptedRetries == 1); + KJ_EXPECT(state.observedRetries == 0); + KJ_EXPECT(state.outcomes.empty()); + KJ_EXPECT(state.replayMemoryBytes == 0); +} + +KJ_TEST("pipelining during a later backoff records a terminal actor RPC retry outcome") { + auto io = kj::setupAsyncIo(); + capnp::MallocMessageBuilder flagsMessage; + RetryTestState state; + PausingTimerChannel timer(2); + TestFixture receiver(makeReceiverParams(io.waitScope)); + TestFixture sender(makeSenderParams(io.waitScope, makeRetryFlags(flagsMessage), timer, state)); + + sender.runInIoContext([&](const TestFixture::Environment& env) { + auto fetcher = makeRetryFetcher(env, receiver, state, FirstFailure::AMBIGUOUS, 2); + auto function = getRpcFunction(env.js, *fetcher, "echo"_kj); + auto parentValue = function.call(env.js, env.js.undefined(), env.js.num(42)); + auto parentObject = KJ_REQUIRE_NONNULL(parentValue.tryCast()); + auto parent = KJ_REQUIRE_NONNULL(parentObject.tryUnwrapAs(env.js)); + + auto parentRejected = env.context.awaitJs(env.js, expectDisconnect(env.js, parentValue)); + auto commitDuringBackoff = env.context.awaitIo(env.js, timer.onBackoffStarted(), + [parent = parent.addRef()]( + jsg::Lock& js) mutable { return expectDisconnect(js, callChild(js, *parent)); }); + auto childRejected = env.context.awaitJs(env.js, kj::mv(commitDuringBackoff)); + return kj::joinPromises(kj::arr(kj::mv(parentRejected), kj::mv(childRejected))) + .attach(kj::mv(parent), kj::mv(fetcher)); + }); + + KJ_EXPECT(state.metadata.size() == 2); + KJ_EXPECT(state.acceptedRetries == 2); + KJ_EXPECT(state.observedRetries == 1); + KJ_ASSERT(state.outcomes.size() == 1); + KJ_EXPECT(state.outcomes[0] == ActorRetryOutcome::UNABLE_TO_RETRY); + KJ_EXPECT(state.replayMemoryBytes == 0); +} + +KJ_TEST("dropping a committed replacement actor RPC attempt records cancellation") { + auto io = kj::setupAsyncIo(); + capnp::MallocMessageBuilder flagsMessage; + RetryTestState state; + ImmediateTimerChannel timer; + TestFixture receiver(makeReceiverParams(io.waitScope)); + TestFixture sender(makeSenderParams(io.waitScope, makeRetryFlags(flagsMessage), timer, state)); + + auto replacementStarted = kj::newPromiseAndFulfiller(); + state.replacementStarted = kj::mv(replacementStarted.fulfiller); + + sender.runInIoContext([&](const TestFixture::Environment& env) { + auto fetcher = + makeRetryFetcher(env, receiver, state, FirstFailure::AMBIGUOUS, 1, Replacement::HANGING); + auto function = getRpcFunction(env.js, *fetcher, "echo"_kj); + auto parentValue = function.call(env.js, env.js.undefined(), env.js.num(42)); + auto parentObject = KJ_REQUIRE_NONNULL(parentValue.tryCast()); + auto parent = KJ_REQUIRE_NONNULL(parentObject.tryUnwrapAs(env.js)); + + // Commit to the hanging replacement attempt, then let the context end with the call in flight. + auto commitToReplacement = env.context.awaitIo(env.js, kj::mv(replacementStarted.promise), + [parent = parent.addRef()](jsg::Lock& js) mutable { callChild(js, *parent); }); + return env.context.awaitJs(env.js, kj::mv(commitToReplacement)) + .attach(kj::mv(parent), kj::mv(fetcher)); + }); + + KJ_EXPECT(state.metadata.size() == 2); + KJ_EXPECT(state.acceptedRetries == 1); + KJ_EXPECT(state.observedRetries == 1); + KJ_ASSERT(state.outcomes.size() == 1); + KJ_EXPECT(state.outcomes[0] == ActorRetryOutcome::CANCELED); + KJ_EXPECT(state.replayMemoryBytes == 0); +} + +KJ_TEST("actor RPC property reads retry on a fresh session") { + auto io = kj::setupAsyncIo(); + capnp::MallocMessageBuilder flagsMessage; + RetryTestState state; + ImmediateTimerChannel timer; + TestFixture receiver(makeReceiverParams(io.waitScope)); + TestFixture sender(makeSenderParams(io.waitScope, makeRetryFlags(flagsMessage), timer, state)); + + sender.runInIoContext([&](const TestFixture::Environment& env) { + auto fetcher = makeRetryFetcher(env, receiver, state, FirstFailure::AMBIGUOUS, 1); + auto property = KJ_REQUIRE_NONNULL(fetcher->getRpcMethodForTestOnly(env.js, kj::str("value"))); + auto& handler = KJ_REQUIRE_NONNULL(env.js.tryGetTypeHandler>()); + auto value = jsg::JsValue(handler.wrap(env.js, kj::mv(property))); + auto checked = env.js.toPromise(value).then(env.js, [](jsg::Lock& js, jsg::Value value) { + KJ_EXPECT(jsg::JsValue(value.getHandle(js)).strictEquals(js.num(42))); + }); + return env.context.awaitJs(env.js, kj::mv(checked)).attach(kj::mv(fetcher)); + }); + + KJ_ASSERT(state.metadata.size() == 2); + KJ_EXPECT(state.metadata[0].nonce == state.metadata[1].nonce); + KJ_EXPECT(state.metadata[1].isRetry == IsActorRetry::YES); + KJ_EXPECT(state.acceptedRetries == 1); + KJ_EXPECT(state.observedRetries == 1); + KJ_ASSERT(state.outcomes.size() == 1); + KJ_EXPECT(state.outcomes[0] == ActorRetryOutcome::RECOVERED); + KJ_EXPECT(state.replayMemoryBytes == 0); +} + +} // namespace +} // namespace workerd::api diff --git a/src/workerd/api/worker-rpc.c++ b/src/workerd/api/worker-rpc.c++ index 25ac419865f..27eaa1c5b8b 100644 --- a/src/workerd/api/worker-rpc.c++ +++ b/src/workerd/api/worker-rpc.c++ @@ -102,7 +102,12 @@ JsRpcCallPlan::JsRpcCallPlan(kj::Own message, using Replayability = RpcSerializerExternalHandler::Replayability; auto operation = getParams().getOperation(); - if (!operation.isCallWithArgs()) return; + if (operation.isGetProperty()) { + KJ_REQUIRE(this->serializedData.size() == 0, + "RPC property call plan unexpectedly contains serialized arguments"); + replayable = true; + return; + } KJ_REQUIRE(operation.hasCallWithArgs() == (this->serializedData.size() > 0), "RPC call plan has inconsistent serialized arguments"); @@ -135,6 +140,20 @@ void JsRpcCallPlan::copyTo(rpc::JsRpcTarget::CallParams::Builder builder) { } } +size_t JsRpcCallPlan::getReplayMemoryEstimate() { + size_t result = serializedData.size() + REPLAY_MEMORY_OVERHEAD; + for (auto segment: message->getSegmentsForOutput()) { + result += kj::max(segment.size(), size_t(METADATA_SEGMENT_WORDS)) * sizeof(capnp::word); + } + return result; +} + +void JsRpcCallPlan::makeSerializedDataExactSizedForReplay() { + if (serializedData.size() > 0) { + serializedData = kj::heapArray(serializedData.asPtr()); + } +} + RpcDeserializerExternalHandler::~RpcDeserializerExternalHandler() noexcept(false) { if (!unwindDetector.isUnwinding()) { KJ_ASSERT(i == externals.size(), "deserialization did not consume all of the externals"); @@ -408,17 +427,178 @@ kj::Maybe> ownOriginatingCall( } // namespace +enum class JsRpcOperation { + CALL, + GET_PROPERTY, +}; + +class JsRpcCallAttemptObserver final: public kj::Refcounted { + public: + void markPipelineCommitted() { + pipelineCommitted = true; + } + + void applyPipelineCommitment(OutgoingActorCallObserver& observer) { + if (pipelineCommitted) observer.markPipelineCommitted(); + } + + private: + bool pipelineCommitted = false; +}; + +class JsRpcCallRetryState final: public kj::Refcounted { + public: + struct StartedAttempt { + kj::Promise> promise; + TraceContext callSpan; + }; + + JsRpcCallRetryState(jsg::Ref parent, + kj::Maybe name, + JsRpcOperation operation, + JsRpcCallPlan callPlan, + kj::Rc retryState, + kj::Own replayMemoryTracker, + rpc::JsRpcTarget::CallResults::Pipeline pipeline, + kj::Maybe> attemptObserver, + kj::Maybe callSpanParents) + : parent(kj::mv(parent)), + name(kj::mv(name)), + operation(operation), + callPlan(kj::mv(callPlan)), + retryState(kj::mv(retryState)), + replayMemoryTracker(kj::mv(replayMemoryTracker)), + pipeline(kj::mv(pipeline)), + attemptObserver(kj::mv(attemptObserver)), + callSpanParents(kj::mv(callSpanParents)) {} + + // Commitment only prevents replay; the logical call is still in flight until it settles, so + // dropping a committed attempt is still a cancellation. + ~JsRpcCallRetryState() noexcept(false) { + if (!settled) { + retryState->recordCanceled(); + } + } + + JsRpcClientProvider::ClientForOneCall commitToCurrentAttempt() { + KJ_IF_SOME(observer, attemptObserver) { + observer->markPipelineCommitted(); + } + if (status == Status::BACKOFF) { + setBrokenPipeline(retryState->getOriginalDisconnect()); + KJ_ASSERT_NONNULL(backoffCommitFulfiller)->fulfill(); + backoffCommitFulfiller = kj::none; + } else { + KJ_REQUIRE( + status == Status::ACTIVE || status == Status::COMMITTED || status == Status::TERMINAL, + "RPC call has no attempt pipeline"); + } + status = Status::COMMITTED; + releaseReplayState(); + return { + .client = pipeline.getCallPipeline(), + .callSpanParents = + callSpanParents.map([](TraceContextParent& value) { return value.addRef(); }), + }; + } + + bool isCommitted() const { + return status == Status::COMMITTED; + } + + kj::Promise beginBackoff() { + KJ_REQUIRE(status == Status::ACTIVE); + auto paf = kj::newPromiseAndFulfiller(); + backoffCommitFulfiller = kj::mv(paf.fulfiller); + status = Status::BACKOFF; + return kj::mv(paf.promise); + } + + void finishBackoffWait() { + backoffCommitFulfiller = kj::none; + } + + void finishSuccess() { + settled = true; + if (status != Status::COMMITTED) { + status = Status::TERMINAL; + } + releaseReplayState(); + } + + void finishFailure(const kj::Exception& exception) { + settled = true; + if (status != Status::COMMITTED) { + status = Status::TERMINAL; + setBrokenPipeline(exception.clone()); + } + releaseReplayState(); + } + + ActorCallRetryState& getRetryState() { + return *retryState; + } + + // Frees the retained call plan and its memory reservation. Holds no JS heap references, so it + // may run at native settlement, before the isolate lock is reacquired. + void releaseReplayPayload() { + name = kj::none; + callPlan = kj::none; + replayMemoryTracker->release(); + } + + void onRetryAccepted() { + KJ_ASSERT_NONNULL(parent)->onActorCallRetry(); + } + + StartedAttempt startAttempt(jsg::Lock& js, ActorCallRetryState::Attempt attempt); + + private: + enum class Status { + ACTIVE, + BACKOFF, + COMMITTED, + TERMINAL, + }; + + void releaseReplayState() { + parent = kj::none; + attemptObserver = kj::none; + releaseReplayPayload(); + } + + void setBrokenPipeline(kj::Exception exception) { + auto broken = capnp::newBrokenPipeline(kj::mv(exception)); + pipeline = rpc::JsRpcTarget::CallResults::Pipeline(capnp::AnyPointer::Pipeline(kj::mv(broken))); + } + + kj::Maybe> parent; + kj::Maybe name; + JsRpcOperation operation; + kj::Maybe callPlan; + kj::Rc retryState; + kj::Own replayMemoryTracker; + rpc::JsRpcTarget::CallResults::Pipeline pipeline; + kj::Maybe> attemptObserver; + kj::Maybe callSpanParents; + kj::Maybe>> backoffCommitFulfiller; + Status status = Status::ACTIVE; + bool settled = false; +}; + JsRpcPromise::JsRpcPromise(jsg::JsRef inner, kj::Own weakRefParam, - IoOwn pipeline, + PendingPipeline pipeline, kj::Maybe originatingCall, kj::Maybe actorTargetRetryability, - kj::Maybe> replayMemoryTracker) + kj::Maybe> replayMemoryTracker, + kj::Maybe> attemptObserver) : inner(kj::mv(inner)), weakRef(kj::mv(weakRefParam)), originatingCall(ownOriginatingCall(kj::mv(originatingCall))), actorTargetRetryability(actorTargetRetryability), replayMemoryTracker(kj::mv(replayMemoryTracker)), + attemptObserver(kj::mv(attemptObserver)), state(Pending{kj::mv(pipeline)}) { KJ_REQUIRE(weakRef->ref == kj::none); weakRef->ref = *this; @@ -439,6 +619,10 @@ void JsRpcPromise::resolve(jsg::Lock& js, jsg::JsValue result) { } } +void JsRpcPromise::setOriginatingCall(kj::Maybe value) { + originatingCall = ownOriginatingCall(kj::mv(value)); +} + void JsRpcPromise::dispose(jsg::Lock& js) { KJ_IF_SOME(resolved, state.tryGet()) { // Disposing the promise implies disposing the final result. @@ -464,10 +648,21 @@ JsRpcClientProvider::ClientForOneCall JsRpcPromise::getClientForOneCall( KJ_IF_SOME(tracker, replayMemoryTracker) { tracker->release(); } - return { - .client = pending.pipeline->getCallPipeline(), - .callSpanParents = kj::mv(callSpanParents), - }; + KJ_SWITCH_ONEOF(pending.pipeline) { + KJ_CASE_ONEOF(pipeline, IoOwn) { + KJ_IF_SOME(observer, attemptObserver) { + observer->markPipelineCommitted(); + } + return { + .client = pipeline->getCallPipeline(), + .callSpanParents = kj::mv(callSpanParents), + }; + } + KJ_CASE_ONEOF(retryState, IoOwn) { + return retryState->commitToCurrentAttempt(); + } + } + KJ_UNREACHABLE; } KJ_CASE_ONEOF(resolved, Resolved) { // Dereference `ctxCheck` just to verify we're running in the correct context. (If not, @@ -540,26 +735,36 @@ namespace { struct JsRpcPromiseAndPipeline { jsg::JsPromise promise; kj::Own weakRef; - rpc::JsRpcTarget::CallResults::Pipeline pipeline; + kj::OneOf> pipeline; // The jsRpcCall of the call that produced this promise, so calls pipelined on the promise // nest under it. Absent when untraced, and on the error paths where no call span was opened. kj::Maybe originatingCall; kj::Maybe actorTargetRetryability; kj::Maybe> replayMemoryTracker; + kj::Maybe> attemptObserver; jsg::Ref asJsRpcPromise(jsg::Lock& js) && { - return js.alloc(jsg::JsRef(js, promise), kj::mv(weakRef), - IoContext::current().addObject(kj::heap(kj::mv(pipeline))), kj::mv(originatingCall), - actorTargetRetryability, kj::mv(replayMemoryTracker)); + auto makePromise = [&](JsRpcPromise::PendingPipeline pendingPipeline) { + return js.alloc(jsg::JsRef(js, promise), kj::mv(weakRef), + kj::mv(pendingPipeline), kj::mv(originatingCall), actorTargetRetryability, + kj::mv(replayMemoryTracker), + attemptObserver.map([](kj::Rc& observer) { + return IoContext::current().addObject(observer.addRef()); + })); + }; + KJ_SWITCH_ONEOF(pipeline) { + KJ_CASE_ONEOF(pipeline, rpc::JsRpcTarget::CallResults::Pipeline) { + return makePromise(IoContext::current().addObject(kj::heap(kj::mv(pipeline)))); + } + KJ_CASE_ONEOF(retryState, kj::Rc) { + return makePromise(IoContext::current().addObject(kj::mv(retryState))); + } + } + KJ_UNREACHABLE; } }; -enum class JsRpcOperation { - CALL, - GET_PROPERTY, -}; - static void setJsRpcCallSpanTags(TraceContext& span, JsRpcClientProvider& parent, kj::Maybe name, @@ -607,17 +812,277 @@ static TraceContext makeJsRpcCallSpan(IoContext& ioContext, return span; } +TraceContext prepareJsRpcCallAttempt(IoContext& ioContext, + JsRpcClientProvider& parent, + kj::Maybe name, + kj::ArrayPtr path, + JsRpcOperation operation, + JsRpcClientProvider::ClientForOneCall& oneCall) { + TraceContext callSpan; + if (util::Autogate::isEnabled(util::AutogateKey::JSRPC_TRACING)) { + KJ_IF_SOME(span, oneCall.callSpan) { + callSpan = kj::mv(span); + setJsRpcCallSpanTags(callSpan, parent, name, path, operation); + } else { + callSpan = makeJsRpcCallSpan( + ioContext, parent, name, path, kj::mv(oneCall.callSpanParents), operation); + } + } + KJ_IF_SOME(lock, ioContext.waitForOutputLocksIfNecessary()) { + oneCall.client = + lock.then([client = kj::mv(oneCall.client)]() mutable { return kj::mv(client); }); + } + return callSpan; +} + +kj::Promise> observeActorRpcCallAttempt( + kj::Own observer, + kj::Rc observerState, + kj::Promise> promise); + +struct JsRpcCallDispatch { + kj::Promise> promise; + rpc::JsRpcTarget::CallResults::Pipeline pipeline; + kj::Maybe> attemptObserver; +}; + +JsRpcCallDispatch dispatchJsRpcCall(IoContext& ioContext, + rpc::JsRpcTarget::Client client, + JsRpcCallPlan& callPlan, + TraceContext& callSpan, + kj::Maybe targetRetryability) { + auto builder = client.callRequest(); + callPlan.copyTo(builder); + KJ_IF_SOME(callerSpanContext, callSpan.getUserSpanParent().toSpanContext()) { + callerSpanContext.toCapnp(builder.initCallerSpanContext()); + } + builder.getResultsStreamHandler().setExternalPusher(ioContext.getExternalPusher()); + + auto callResult = builder.send(); + kj::Promise> promise = kj::mv(callResult); + KJ_IF_SOME(targetRetryable, targetRetryability) { + KJ_IF_SOME(observer, + ioContext.getMetrics().observeOutgoingActorRpcCall( + ActorCallPayloadReplayable(callPlan.getReplayable()), targetRetryable)) { + auto attemptObserver = kj::rc(); + promise = observeActorRpcCallAttempt( + kj::mv(observer), attemptObserver.addRef(), kj::mv(promise)); + return {.promise = kj::mv(promise), + .pipeline = kj::mv(callResult), + .attemptObserver = kj::mv(attemptObserver)}; + } + } + return {.promise = kj::mv(promise), .pipeline = kj::mv(callResult)}; +} + +struct JsRpcRetrySetup { + kj::Maybe> replayMemoryTracker; + kj::Maybe> state; + kj::Maybe attempt; +}; + +JsRpcRetrySetup setupJsRpcRetries(IoContext& ioContext, + bool destinationSupportsRetries, + JsRpcCallPlan& callPlan, + kj::Maybe name) { + JsRpcRetrySetup result; + if (!destinationSupportsRetries || !callPlan.getReplayable() || + !util::Autogate::isEnabled(util::AutogateKey::DURABLE_OBJECT_RETRIES_FETCH) || + !util::Autogate::isEnabled(util::AutogateKey::DURABLE_OBJECT_RETRIES_JSRPC)) { + return result; + } + + auto enforcementRequested = + util::Autogate::isEnabled(util::AutogateKey::DURABLE_OBJECT_RETRIES_FETCH_RETRY_REQUESTS) && + util::Autogate::isEnabled(util::AutogateKey::DURABLE_OBJECT_RETRIES_JSRPC_RETRY_REQUESTS); + auto replayMemoryEstimate = callPlan.getReplayMemoryEstimate(); + KJ_IF_SOME(value, name) { + replayMemoryEstimate += value.size() + 1; + } + auto reservation = enforcementRequested + ? ioContext.getMetrics().tryReserveActorCallReplayMemory(replayMemoryEstimate) + : kj::Maybe>(kj::none); + auto enforcementEnabled = ActorRetryGateEnabled(reservation != kj::none); + KJ_IF_SOME(reservedMemory, reservation) { + callPlan.makeSerializedDataExactSizedForReplay(); + result.replayMemoryTracker = kj::refcounted(kj::mv(reservedMemory)); + } else { + result.replayMemoryTracker = kj::refcounted( + ioContext.getMetrics().trackActorCallReplayMemory(replayMemoryEstimate)); + } + result.state = kj::rc(ioContext.getIoChannelFactory().getTimer(), + ioContext.getMetrics(), + ActorCallRetryState::Config{ + .callType = ActorRetryCallType::JSRPC, + .observationEnabled = ActorRetryGateEnabled::YES, + .enforcementEnabled = enforcementEnabled, + .payloadReplayable = ActorCallPayloadReplayable::YES, + }); + auto attemptOrException = KJ_ASSERT_NONNULL(result.state)->startAttempt(); + result.attempt = + kj::mv(KJ_ASSERT_NONNULL(attemptOrException.tryGet())); + return result; +} + +} // namespace + +JsRpcCallRetryState::StartedAttempt JsRpcCallRetryState::startAttempt( + jsg::Lock& js, ActorCallRetryState::Attempt attempt) { + auto& ioContext = IoContext::current(); + auto& parent = KJ_ASSERT_NONNULL(this->parent); + auto oneCall = parent->getClientForOneCall(js, kj::mv(attempt)); + kj::Vector pathViews; + if (util::Autogate::isEnabled(util::AutogateKey::JSRPC_TRACING)) { + parent->appendPath(pathViews); + } + auto nameRef = name.map([](kj::String& value) -> const kj::String& { return value; }); + auto callSpan = + prepareJsRpcCallAttempt(ioContext, *parent, nameRef, pathViews.asPtr(), operation, oneCall); + + auto dispatch = dispatchJsRpcCall(ioContext, kj::mv(oneCall.client), KJ_ASSERT_NONNULL(callPlan), + callSpan, parent->getActorTargetRetryability()); + callSpanParents = callSpan.getSpanParentsIfObserved(); + pipeline = kj::mv(dispatch.pipeline); + attemptObserver = kj::mv(dispatch.attemptObserver); + status = Status::ACTIVE; + return { + .promise = kj::mv(dispatch.promise), + .callSpan = kj::mv(callSpan), + }; +} + +namespace { + +// Both members are KJ I/O state (the span observer may own an RPC client), so this struct only +// crosses onto the JS heap wrapped in an `IoOwn`. +struct JsRpcCallAttemptSuccess { + capnp::Response response; + TraceContext callSpan; +}; + +struct JsRpcCallAttemptFailure { + kj::Exception exception; +}; + +using JsRpcCallAttemptResult = kj::OneOf; +using JsRpcCallAttemptPromise = jsg::Promise>; + +kj::Promise captureJsRpcCallAttempt( + kj::Promise> promise, + TraceContext callSpan, + kj::Rc state) { + try { + auto response = co_await promise; + // The call can no longer be replayed, so drop its payload now rather than after reentry, + // which another critical section may delay. + state->releaseReplayPayload(); + co_return JsRpcCallAttemptSuccess{kj::mv(response), kj::mv(callSpan)}; + } catch (...) { + auto exception = kj::getCaughtExceptionAsKj(); + state->getRetryState().maybeStartRetryLatencyTimer(exception); + co_return JsRpcCallAttemptFailure{kj::mv(exception)}; + } +} + +JsRpcCallAttemptPromise awaitJsRpcCallAttempt(jsg::Lock& js, + kj::Rc state, + kj::Promise> promise, + TraceContext callSpan); + +JsRpcCallAttemptPromise rejectJsRpcCall(jsg::Lock& js, kj::Exception exception) { + return js.rejectedPromise>( + js.exceptionToJsValue(kj::mv(exception))); +} + +JsRpcCallAttemptPromise handleFailedJsRpcCallAttempt( + jsg::Lock& js, kj::Rc state, kj::Exception failure) { + if (state->isCommitted()) { + auto exception = state->getRetryState().handleCommittedAttemptFailure(kj::mv(failure)); + state->finishFailure(exception); + return rejectJsRpcCall(js, kj::mv(exception)); + } + + auto delayOrException = state->getRetryState().handleAttemptFailure(kj::mv(failure)); + KJ_IF_SOME(exception, delayOrException.tryGet()) { + state->finishFailure(exception); + return rejectJsRpcCall(js, kj::mv(exception)); + } + + auto backoffCommitted = state->beginBackoff(); + state->onRetryAccepted(); + auto delay = KJ_ASSERT_NONNULL(delayOrException.tryGet()); + auto backoff = + IoContext::current().afterLimitTimeout(delay).exclusiveJoin(kj::mv(backoffCommitted)); + return IoContext::current().awaitIo(js, kj::mv(backoff), + [state = kj::mv(state)](jsg::Lock& js) mutable -> JsRpcCallAttemptPromise { + state->finishBackoffWait(); + if (state->isCommitted()) { + auto exception = state->getRetryState().handleCommittedAttemptFailure( + state->getRetryState().getOriginalDisconnect()); + state->finishFailure(exception); + return rejectJsRpcCall(js, kj::mv(exception)); + } + + auto attemptOrException = state->getRetryState().startAttempt(); + KJ_IF_SOME(exception, attemptOrException.tryGet()) { + state->finishFailure(exception); + return rejectJsRpcCall(js, kj::mv(exception)); + } + auto attempt = + kj::mv(KJ_ASSERT_NONNULL(attemptOrException.tryGet())); + // `awaitJsRpcCallAttempt()` takes its own reference so `state` stays valid for the handler. + // Termination exceptions propagate unchanged; the retry state is torn down with the context. + try { + auto started = state->startAttempt(js, kj::mv(attempt)); + return awaitJsRpcCallAttempt( + js, state.addRef(), kj::mv(started.promise), kj::mv(started.callSpan)); + } catch (jsg::JsExceptionThrown&) { + throw; + } catch (...) { + auto exception = kj::getCaughtExceptionAsKj(); + state->finishFailure(exception); + kj::throwFatalException(kj::mv(exception)); + } + }); +} + +JsRpcCallAttemptPromise awaitJsRpcCallAttempt(jsg::Lock& js, + kj::Rc state, + kj::Promise> promise, + TraceContext callSpan) { + // Built before the callback below moves `state`; argument evaluation order is unspecified. + auto attempt = captureJsRpcCallAttempt(kj::mv(promise), kj::mv(callSpan), state.addRef()); + return IoContext::current().awaitIo(js, kj::mv(attempt), + [state = kj::mv(state)]( + jsg::Lock& js, JsRpcCallAttemptResult result) mutable -> JsRpcCallAttemptPromise { + KJ_SWITCH_ONEOF(result) { + KJ_CASE_ONEOF(success, JsRpcCallAttemptSuccess) { + state->getRetryState().recordRecovered(); + state->finishSuccess(); + return js.resolvedPromise(IoContext::current().addObject(kj::heap(kj::mv(success)))); + } + KJ_CASE_ONEOF(failure, JsRpcCallAttemptFailure) { + return handleFailedJsRpcCallAttempt(js, kj::mv(state), kj::mv(failure.exception)); + } + } + KJ_UNREACHABLE; + }); +} + // Reports the settlement of one `JsRpcTarget.call()` attempt to `observer`. Dropping the returned // promise drops the observer without a result, which it reports as cancellation. kj::Promise> observeActorRpcCallAttempt( kj::Own observer, + kj::Rc observerState, kj::Promise> promise) { try { auto response = co_await promise; + observerState->applyPipelineCommitment(*observer); observer->recordSuccess(); co_return kj::mv(response); } catch (...) { auto exception = kj::getCaughtExceptionAsKj(); + observerState->applyPipelineCommitment(*observer); observer->recordFailure(exception); kj::throwFatalException(kj::mv(exception)); } @@ -625,7 +1090,7 @@ kj::Promise> observeActorRpcCallA // Core implementation of making an RPC call, reusable for many cases below. JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, - JsRpcClientProvider& parent, + jsg::Ref parent, kj::Maybe name, // If `maybeArgs` is provided, this is a call, otherwise it is a property access. kj::Maybe&> maybeArgs) { @@ -652,36 +1117,22 @@ JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, // `path` is filled in with the chain of property names leading to the method. kj::Vector path; - parent.appendPath(path); + parent->appendPath(path); auto operation = maybeArgs != kj::none ? JsRpcOperation::CALL : JsRpcOperation::GET_PROPERTY; - auto actorTargetRetryability = parent.getActorTargetRetryability(); + auto actorTargetRetryability = parent->getActorTargetRetryability(); auto destinationSupportsRetries = actorTargetRetryability.orDefault(ActorCallTargetRetryable::NO).toBool(); TraceContext jsRpcCallSpan; kj::Maybe oneCall; - kj::Maybe actorCallAttempt; + JsRpcRetrySetup retrySetup; auto resolveOneCall = [&]() -> JsRpcClientProvider::ClientForOneCall& { if (oneCall == kj::none) { - oneCall = parent.getClientForOneCall(js, kj::mv(actorCallAttempt)); + oneCall = parent->getClientForOneCall(js, kj::mv(retrySetup.attempt)); auto& result = KJ_ASSERT_NONNULL(oneCall); - if (util::Autogate::isEnabled(util::AutogateKey::JSRPC_TRACING)) { - // Per-call dispatch span, captured into the awaitIo callback below so it stays open - // until the response settles. - KJ_IF_SOME(span, result.callSpan) { - jsRpcCallSpan = kj::mv(span); - setJsRpcCallSpanTags(jsRpcCallSpan, parent, name, path.asPtr(), operation); - } else { - jsRpcCallSpan = makeJsRpcCallSpan( - ioContext, parent, name, path.asPtr(), kj::mv(result.callSpanParents), operation); - } - } - KJ_IF_SOME(lock, ioContext.waitForOutputLocksIfNecessary()) { - // A promise client keeps calls made while serializing externals behind the output gate. - result.client = - lock.then([client = kj::mv(result.client)]() mutable { return kj::mv(client); }); - } + jsRpcCallSpan = + prepareJsRpcCallAttempt(ioContext, *parent, name, path.asPtr(), operation, result); } return KJ_ASSERT_NONNULL(oneCall); }; @@ -750,54 +1201,11 @@ JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, // JSRPC retries build on the fetch retry machinery, so the fetch gate remains a shared // prerequisite while the JSRPC gate controls this event type's separate rollout. - kj::Maybe> replayMemoryTracker; - if (destinationSupportsRetries && callPlan.getReplayable() && - util::Autogate::isEnabled(util::AutogateKey::DURABLE_OBJECT_RETRIES_FETCH) && - util::Autogate::isEnabled(util::AutogateKey::DURABLE_OBJECT_RETRIES_JSRPC)) { - replayMemoryTracker = kj::refcounted( - ioContext.getMetrics().trackActorCallReplayMemory(callPlan.getReplayMemoryEstimate())); - actorCallAttempt.emplace( - generateActorRetryRequestMetadata( - kj::systemCoarseCalendarClock().now(), ActorRetryGateEnabled::NO), - IsFirstActorCallAttempt::YES); - } - - auto client = kj::mv(resolveOneCall().client); - auto builder = client.callRequest(); - callPlan.copyTo(builder); + retrySetup = setupJsRpcRetries(ioContext, destinationSupportsRetries, callPlan, name); - // Tell the callee which caller span corresponds to this dispatch. A session carries many - // calls (e.g. calls pipelined on a returned stub), so the context propagated when the session - // opened identifies only the first call. Yields kj::none when untraced. - KJ_IF_SOME(callerSpanContext, jsRpcCallSpan.getUserSpanParent().toSpanContext()) { - callerSpanContext.toCapnp(builder.initCallerSpanContext()); - } - - // Unfortunately, we always have to send the ExternalPusher since we don't know whether the - // call will return any streams (or other pushed externals). Luckily, it's a - // one-per-IoContext object, not a big deal. (It'll take a slot on the capnp export table - // though.) - builder.getResultsStreamHandler().setExternalPusher(ioContext.getExternalPusher()); - - auto callResult = builder.send(); - - // RemotePromise lets us consume its pipeline and promise portions independently; we consume - // the promise here and the pipeline below, both via kj::mv(). - kj::Promise> resultPromise = - kj::mv(callResult); - KJ_IF_SOME(tracker, replayMemoryTracker) { - resultPromise = resultPromise.attach( - kj::defer([tracker = kj::addRef(*tracker)]() mutable { tracker->release(); })); - } - KJ_IF_SOME(targetRetryable, actorTargetRetryability) { - // Observe the individual call rather than its session, whose lifetime ends with capability - // teardown rather than with this result. - KJ_IF_SOME(observer, - ioContext.getMetrics().observeOutgoingActorRpcCall( - ActorCallPayloadReplayable(callPlan.getReplayable()), targetRetryable)) { - resultPromise = observeActorRpcCallAttempt(kj::mv(observer), kj::mv(resultPromise)); - } - } + auto dispatch = dispatchJsRpcCall(ioContext, kj::mv(resolveOneCall().client), callPlan, + jsRpcCallSpan, actorTargetRetryability); + auto resultPromise = kj::mv(dispatch.promise); // A follow-up call through this result's pending pipeline still targets the actor, but cannot // create a fresh attempt without replaying the call that produced the pipeline. @@ -817,19 +1225,41 @@ JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, // retained when traced, so the untraced path holds no span state. auto originatingCall = jsRpcCallSpan.getSpanParentsIfObserved(); + kj::Maybe> callRetryState; + KJ_IF_SOME(retryState, retrySetup.state) { + if (retryState->isRetryEnabled()) { + auto ownedName = name.map([](const kj::String& value) { return kj::str(value); }); + callRetryState = kj::rc(parent.addRef(), kj::mv(ownedName), + operation, kj::mv(callPlan), retryState.addRef(), + kj::mv(KJ_ASSERT_NONNULL(retrySetup.replayMemoryTracker)), kj::mv(dispatch.pipeline), + dispatch.attemptObserver.map( + [](kj::Rc& observer) { return observer.addRef(); }), + originatingCall.map([](TraceContextParent& value) { return value.addRef(); })); + } + } + + if (callRetryState == kj::none) { + KJ_IF_SOME(tracker, retrySetup.replayMemoryTracker) { + resultPromise = resultPromise.attach( + kj::defer([tracker = kj::addRef(*tracker)]() mutable { tracker->release(); })); + } + } + kj::Maybe> promiseReplayMemoryTracker; - KJ_IF_SOME(tracker, replayMemoryTracker) { - promiseReplayMemoryTracker = ioContext.addObject(kj::mv(tracker)); + if (callRetryState == kj::none) { + KJ_IF_SOME(tracker, retrySetup.replayMemoryTracker) { + promiseReplayMemoryTracker = ioContext.addObject(kj::mv(tracker)); + } } - auto jsPromise = ioContext.awaitIo(js, kj::mv(resultPromise), - [weakRef = kj::atomicAddRef(*weakRef), jsRpcCallSpan = kj::mv(jsRpcCallSpan)]( - jsg::Lock& js, - capnp::Response response) mutable -> jsg::Value { + auto resolveResult = [weakRef = kj::atomicAddRef(*weakRef)](jsg::Lock& js, + rpc::JsRpcTarget::CallResults::Reader response, + TraceContext& jsRpcCallSpan) mutable -> jsg::Value { // Stubs in the response record this call as their originating call so that // follow-up calls on those stubs nest under it (only when traced). - auto jsResult = - deserializeRpcReturnValue(js, response, jsRpcCallSpan.getSpanParentsIfObserved()); + auto originatingCall = jsRpcCallSpan.getSpanParentsIfObserved(); + auto jsResult = deserializeRpcReturnValue(js, response, + originatingCall.map([](TraceContextParent& value) { return value.addRef(); })); if (weakRef->disposed) { // The promise was explicitly disposed before it even resolved. This means we must dispose @@ -837,20 +1267,53 @@ JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, tryCallDisposeMethod(js, jsResult); } else { KJ_IF_SOME(r, weakRef->ref) { + r.setOriginatingCall(kj::mv(originatingCall)); r.resolve(js, jsResult); } } return jsg::Value(js.v8Isolate, jsResult); - }); + }; + + auto jsPromise = [&]() -> jsg::Promise { + KJ_IF_SOME(retryState, callRetryState) { + return awaitJsRpcCallAttempt( + js, retryState.addRef(), kj::mv(resultPromise), kj::mv(jsRpcCallSpan)) + .then(js, + [resolveResult = kj::mv(resolveResult)]( + jsg::Lock& js, IoOwn success) mutable { + // `success` is released when this continuation returns, closing the call span. + return resolveResult(js, success->response, success->callSpan); + }); + } + return ioContext.awaitIo(js, kj::mv(resultPromise), + [resolveResult = kj::mv(resolveResult), jsRpcCallSpan = kj::mv(jsRpcCallSpan)]( + jsg::Lock& js, + capnp::Response response) mutable -> jsg::Value { + return resolveResult(js, response, jsRpcCallSpan); + }); + }(); + + auto pendingPipeline = + [&]() -> kj::OneOf> { + KJ_IF_SOME(retryState, callRetryState) { + return retryState.addRef(); + } + return kj::mv(dispatch.pipeline); + }(); + kj::Maybe> pendingAttemptObserver; + if (callRetryState == kj::none) { + pendingAttemptObserver = kj::mv(dispatch.attemptObserver); + } return { .promise = jsg::JsPromise(js.wrapSimplePromise(kj::mv(jsPromise))), .weakRef = kj::mv(weakRef), - .pipeline = kj::mv(callResult), + .pipeline = kj::mv(pendingPipeline), .originatingCall = kj::mv(originatingCall), .actorTargetRetryability = pipelineActorTargetRetryability, .replayMemoryTracker = kj::mv(promiseReplayMemoryTracker), + .attemptObserver = kj::mv(pendingAttemptObserver), }; }, [&](jsg::Value error) -> JsRpcPromiseAndPipeline { // Probably a serialization error. Need to convert to an async error since we never throw @@ -883,19 +1346,19 @@ JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, jsg::Ref JsRpcProperty::call(const v8::FunctionCallbackInfo& args) { jsg::Lock& js = jsg::Lock::from(args.GetIsolate()); - return callImpl(js, *parent, name, args).asJsRpcPromise(js); + return callImpl(js, parent.addRef(), name, args).asJsRpcPromise(js); } jsg::Ref JsRpcStub::call(const v8::FunctionCallbackInfo& args) { jsg::Lock& js = jsg::Lock::from(args.GetIsolate()); - return callImpl(js, *this, kj::none, args).asJsRpcPromise(js); + return callImpl(js, JSG_THIS, kj::none, args).asJsRpcPromise(js); } jsg::Ref JsRpcPromise::call(const v8::FunctionCallbackInfo& args) { jsg::Lock& js = jsg::Lock::from(args.GetIsolate()); - return callImpl(js, *this, kj::none, args).asJsRpcPromise(js); + return callImpl(js, JSG_THIS, kj::none, args).asJsRpcPromise(js); } namespace { @@ -935,19 +1398,19 @@ jsg::JsValue finallyImpl( jsg::JsValue JsRpcProperty::then(jsg::Lock& js, v8::Local handler, jsg::Optional> errorHandler) { - auto promise = callImpl(js, *parent, name, kj::none).promise; + auto promise = callImpl(js, parent.addRef(), name, kj::none).promise; return thenImpl(js, promise, handler, errorHandler); } jsg::JsValue JsRpcProperty::catch_(jsg::Lock& js, v8::Local errorHandler) { - auto promise = callImpl(js, *parent, name, kj::none).promise; + auto promise = callImpl(js, parent.addRef(), name, kj::none).promise; return catchImpl(js, promise, errorHandler); } jsg::JsValue JsRpcProperty::finally(jsg::Lock& js, v8::Local onFinally) { - auto promise = callImpl(js, *parent, name, kj::none).promise; + auto promise = callImpl(js, parent.addRef(), name, kj::none).promise; return finallyImpl(js, promise, onFinally); } diff --git a/src/workerd/api/worker-rpc.h b/src/workerd/api/worker-rpc.h index 8b3d121e748..041f2996ed2 100644 --- a/src/workerd/api/worker-rpc.h +++ b/src/workerd/api/worker-rpc.h @@ -167,15 +167,14 @@ class JsRpcCallPlan { kj::Array serializedData, RpcSerializerExternalHandler::Replayability serializerReplayability); - // True only for method calls whose arguments are absent or contain no externals and no - // serializer-handled ineligible values. Property reads are not replayable. + // True only for property reads and method calls whose arguments are absent or contain no + // externals and no serializer-handled ineligible values. bool getReplayable() const { return replayable; } - size_t getReplayMemoryEstimate() const { - return serializedData.size() + REPLAY_MEMORY_OVERHEAD; - } + size_t getReplayMemoryEstimate(); + void makeSerializedDataExactSizedForReplay(); void copyTo(rpc::JsRpcTarget::CallParams::Builder builder); @@ -319,6 +318,10 @@ class JsRpcClientProvider: public jsg::Object { return getActorTargetRetryability().orDefault(ActorCallTargetRetryable::NO).toBool(); } + virtual void onActorCallRetry() { + KJ_FAIL_REQUIRE("actor call retry requested from an unsupported RPC target"); + } + // Get a capnp client that can be used to dispatch one call. virtual ClientForOneCall getClientForOneCall( jsg::Lock& js, kj::Maybe actorCallAttempt) = 0; @@ -329,6 +332,8 @@ class JsRpcClientProvider: public jsg::Object { }; class JsRpcProperty; +class JsRpcCallRetryState; +class JsRpcCallAttemptObserver; class JsRpcReplayMemoryTracker final: public kj::Refcounted { public: @@ -347,6 +352,9 @@ class JsRpcReplayMemoryTracker final: public kj::Refcounted { // but rather our own custom thenable, so that we can support pipelining on it. class JsRpcPromise: public JsRpcClientProvider { public: + using PendingPipeline = + kj::OneOf, IoOwn>; + // A weak reference to this JsRpcPromise. Unlike the usual WeakRef pattern, though, this ref is // allocated before the promise itself is actually created, and filled in later. This is needed // to solve cyclic initialization challenges in `callImpl()`. @@ -364,13 +372,15 @@ class JsRpcPromise: public JsRpcClientProvider { JsRpcPromise(jsg::JsRef inner, kj::Own weakRef, - IoOwn pipeline, + PendingPipeline pipeline, kj::Maybe originatingCall, kj::Maybe actorTargetRetryability, - kj::Maybe> replayMemoryTracker); + kj::Maybe> replayMemoryTracker, + kj::Maybe> attemptObserver); ~JsRpcPromise() noexcept(false); void resolve(jsg::Lock& js, jsg::JsValue result); + void setOriginatingCall(kj::Maybe value); void dispose(jsg::Lock& js); ClientForOneCall getClientForOneCall( @@ -424,9 +434,10 @@ class JsRpcPromise: public JsRpcClientProvider { kj::Maybe> originatingCall; kj::Maybe actorTargetRetryability; kj::Maybe> replayMemoryTracker; + kj::Maybe> attemptObserver; struct Pending { - IoOwn pipeline; + PendingPipeline pipeline; }; struct Resolved { jsg::Value result; @@ -473,6 +484,9 @@ class JsRpcProperty: public JsRpcClientProvider { kj::Maybe getActorTargetRetryability() override { return parent->getActorTargetRetryability(); } + void onActorCallRetry() override { + parent->onActorCallRetry(); + } ClientForOneCall getClientForOneCall( jsg::Lock& js, kj::Maybe actorCallAttempt) override; diff --git a/src/workerd/io/observer.h b/src/workerd/io/observer.h index 3215588ffa8..261d849aba2 100644 --- a/src/workerd/io/observer.h +++ b/src/workerd/io/observer.h @@ -82,6 +82,7 @@ class ByteStreamObserver { class OutgoingActorCallObserver { public: virtual ~OutgoingActorCallObserver() noexcept(false) = default; + virtual void markPipelineCommitted() {} virtual void recordSuccess() {} virtual void recordFailure(kj::Exception& e) {} }; @@ -185,6 +186,13 @@ class RequestObserver: public kj::Refcounted { return kj::Own(); } + // Attempts to reserve platform memory for retained actor-call replay state. Returning none keeps + // the call observe-only. Production observers must enforce an aggregate bound before returning a + // reservation handle. + virtual kj::Maybe> tryReserveActorCallReplayMemory(size_t bytes) { + return kj::none; + } + // Records an additional outgoing actor call started by a runtime retry loop. virtual void recordActorRetry(ActorRetryCallType callType) {} diff --git a/src/workerd/util/autogate.h b/src/workerd/util/autogate.h index 96e4ca11b5a..4f78ea8e38c 100644 --- a/src/workerd/util/autogate.h +++ b/src/workerd/util/autogate.h @@ -94,10 +94,11 @@ namespace workerd::util { prerequisite. */ \ V(DURABLE_OBJECT_RETRIES_FETCH_RETRY_REQUESTS) \ /* Extends observe-only retry-token claiming to Durable Object JSRPC calls: senders attach \ - tokens and receivers claim them. Requires DURABLE_OBJECT_RETRIES_FETCH. A JSRPC retry \ - request gate, the counterpart of DURABLE_OBJECT_RETRIES_FETCH_RETRY_REQUESTS, is added with \ - sender replay. */ \ + tokens and receivers claim them. Requires DURABLE_OBJECT_RETRIES_FETCH. */ \ V(DURABLE_OBJECT_RETRIES_JSRPC) \ + /* Enables Durable Object JSRPC retry requests. Requires both fetch retry gates and the JSRPC \ + observe gate. */ \ + V(DURABLE_OBJECT_RETRIES_JSRPC_RETRY_REQUESTS) \ /* When enabled, the native `node-internal:url` module is provided by the Rust \ implementation (api::node UrlUtil ported to src/rust/api) instead of the \ C++ implementation. The C++ implementation is retained for rollback.*/ \