From 93280aa0d7c24b1fc2b3d22272bc00daf032506f Mon Sep 17 00:00:00 2001 From: Ashley Peacock Date: Thu, 17 Sep 2026 13:47:59 +0100 Subject: [PATCH] STOR-5489: Measure Durable Object RPC replay memory Track projected memory for replayable JSRPC calls without retaining another payload copy. Release accounting when the call settles, its caller context ends, or pipeline use commits it to the current attempt. --- .../api/jsrpc-call-observation-test.c++ | 92 ++++++++++++++++++- src/workerd/api/worker-rpc.c++ | 23 ++++- src/workerd/api/worker-rpc.h | 21 ++++- src/workerd/io/observer.h | 6 ++ 4 files changed, 137 insertions(+), 5 deletions(-) diff --git a/src/workerd/api/jsrpc-call-observation-test.c++ b/src/workerd/api/jsrpc-call-observation-test.c++ index 39793689244..9e28225ee3f 100644 --- a/src/workerd/api/jsrpc-call-observation-test.c++ +++ b/src/workerd/api/jsrpc-call-observation-test.c++ @@ -48,6 +48,8 @@ struct Observation { struct ObservationState { kj::Vector> observations; + size_t replayMemoryBytes = 0; + size_t peakReplayMemoryBytes = 0; }; class RecordingCallObserver final: public OutgoingActorCallObserver { @@ -85,6 +87,12 @@ class RecordingObserver final: public RequestObserver { return kj::heap(observation); } + kj::Own trackActorCallReplayMemory(size_t bytes) override { + state.replayMemoryBytes += bytes; + state.peakReplayMemoryBytes = kj::max(state.peakReplayMemoryBytes, state.replayMemoryBytes); + return kj::heap(kj::defer([&state = state, bytes]() { state.replayMemoryBytes -= bytes; })); + } + private: ObservationState& state; }; @@ -148,6 +156,7 @@ class Harness { sender = kj::heap(TestFixture::SetupParams{ .waitScope = io.waitScope, .featureFlags = flags.asReader(), + .autogates = kj::arr("durable-object-retries-fetch"_kj, "durable-object-retries-jsrpc"_kj), .mainModuleSource = SENDER_SOURCE, .useRealTimers = false, .requestObserverFactory = @@ -210,7 +219,9 @@ KJ_TEST("concurrent RPC calls settling out of order keep their own observations" auto [fetcher, factory] = harness.makeFetcher(env); auto release = factory.holdNextSession(); auto held = Harness::awaitScript(env, fetcher.addRef(), "fetcher.echo(1)"_kj); + auto heldBytes = harness.state.replayMemoryBytes; auto immediate = Harness::awaitScript(env, kj::mv(fetcher), "fetcher.echo(makeStub())"_kj); + KJ_EXPECT(harness.state.replayMemoryBytes == heldBytes); return kj::mv(immediate) .then([&harness, release = kj::mv(release)]() mutable { auto& observations = harness.state.observations; @@ -229,13 +240,18 @@ KJ_TEST("concurrent RPC calls settling out of order keep their own observations" KJ_EXPECT(observations[1]->payloadReplayable == ActorCallPayloadReplayable::NO); KJ_EXPECT(observations[1]->targetRetryable == ActorCallTargetRetryable::YES); KJ_EXPECT(observations[1]->settlement == Settlement::SUCCESS); + KJ_EXPECT(harness.state.replayMemoryBytes == 0); } KJ_TEST("an RPC call to an actor that cannot retry is observed as such") { Harness harness; harness.sender->runInIoContext([&](const TestFixture::Environment& env) { - auto fetcher = harness.makeFetcher(env, ActorCallTargetRetryable::NO).fetcher; - return Harness::awaitScript(env, kj::mv(fetcher), "fetcher.echo(1)"_kj); + auto [fetcher, factory] = harness.makeFetcher(env, ActorCallTargetRetryable::NO); + auto release = factory.holdNextSession(); + auto call = Harness::awaitScript(env, kj::mv(fetcher), "fetcher.echo(1)"_kj); + KJ_EXPECT(harness.state.replayMemoryBytes == 0); + release->fulfill(); + return kj::mv(call); }); auto& observations = harness.state.observations; @@ -245,6 +261,33 @@ KJ_TEST("an RPC call to an actor that cannot retry is observed as such") { KJ_EXPECT(observations[0]->settlement == Settlement::SUCCESS); } +KJ_TEST("concurrent replayable RPC calls track projected replay memory independently") { + Harness harness; + harness.sender->runInIoContext([&](const TestFixture::Environment& env) { + auto [fetcher, factory] = harness.makeFetcher(env); + auto releaseFirst = factory.holdNextSession(); + auto first = Harness::awaitScript(env, fetcher.addRef(), "fetcher.echo('first')"_kj); + auto firstBytes = harness.state.replayMemoryBytes; + auto releaseSecond = factory.holdNextSession(); + auto second = Harness::awaitScript(env, kj::mv(fetcher), + "fetcher.echo('a payload large enough to produce a different serialized size')"_kj); + auto secondBytes = harness.state.replayMemoryBytes - firstBytes; + + KJ_EXPECT(firstBytes > 0); + KJ_EXPECT(secondBytes > firstBytes); + KJ_EXPECT(harness.state.peakReplayMemoryBytes == firstBytes + secondBytes); + + releaseFirst->fulfill(); + return kj::mv(first) + .then([&harness, secondBytes, releaseSecond = kj::mv(releaseSecond)]() mutable { + KJ_EXPECT(harness.state.replayMemoryBytes == secondBytes); + releaseSecond->fulfill(); + }).then([second = kj::mv(second)]() mutable { return kj::mv(second); }); + }); + + KJ_EXPECT(harness.state.replayMemoryBytes == 0); +} + KJ_TEST("an RPC call through an actor result pipeline is observed as non-retryable") { Harness harness; harness.sender->runInIoContext([&](const TestFixture::Environment& env) { @@ -288,6 +331,50 @@ KJ_TEST("actor RPC calls remain observed after their parent promise resolves") { } } +KJ_TEST("using an RPC result pipeline releases projected replay memory") { + Harness harness; + harness.sender->runInIoContext([&](const TestFixture::Environment& env) { + auto [fetcher, factory] = harness.makeFetcher(env); + auto release = factory.holdNextSession(); + Harness::runScript(env, kj::mv(fetcher), "globalThis.child = fetcher.makeChild('payload')"_kj); + KJ_EXPECT(harness.state.replayMemoryBytes > 0); + + auto pipelined = Harness::awaitScript(env, R"JS( + (() => { + const result = child.echo(1); + delete globalThis.child; + return result; + })() + )JS"_kj); + KJ_EXPECT(harness.state.replayMemoryBytes == 0); + + release->fulfill(); + return kj::mv(pipelined); + }); + + KJ_EXPECT(harness.state.replayMemoryBytes == 0); +} + +KJ_TEST("disposing an RPC promise does not release projected replay memory early") { + Harness harness; + harness.sender->runInIoContext([&](const TestFixture::Environment& env) { + auto [fetcher, factory] = harness.makeFetcher(env); + auto release = factory.holdNextSession(); + auto result = Harness::runScript(env, kj::mv(fetcher), "fetcher.echo('payload')"_kj); + auto object = KJ_REQUIRE_NONNULL(result.tryCast()); + auto call = KJ_REQUIRE_NONNULL(object.tryUnwrapAs(env.js)); + auto trackedBytes = harness.state.replayMemoryBytes; + KJ_EXPECT(trackedBytes > 0); + + call->dispose(env.js); + KJ_EXPECT(harness.state.replayMemoryBytes == trackedBytes); + + return kj::Promise(kj::READY_NOW).attach(kj::mv(call), kj::mv(release)); + }); + + KJ_EXPECT(harness.state.replayMemoryBytes == 0); +} + KJ_TEST("an RPC property get is observed with a non-replayable payload") { Harness harness; harness.sender->runInIoContext([&](const TestFixture::Environment& env) { @@ -360,6 +447,7 @@ KJ_TEST("tearing down the caller records cancellation rather than a disconnect") KJ_ASSERT(harness.state.observations.size() == 1); KJ_EXPECT(harness.state.observations[0]->settlement == Settlement::CANCELED); + KJ_EXPECT(harness.state.replayMemoryBytes == 0); } } // namespace diff --git a/src/workerd/api/worker-rpc.c++ b/src/workerd/api/worker-rpc.c++ index 80e05090c64..42ead0d97cc 100644 --- a/src/workerd/api/worker-rpc.c++ +++ b/src/workerd/api/worker-rpc.c++ @@ -412,11 +412,13 @@ JsRpcPromise::JsRpcPromise(jsg::JsRef inner, kj::Own weakRefParam, IoOwn pipeline, kj::Maybe originatingCall, - kj::Maybe actorTargetRetryability) + kj::Maybe actorTargetRetryability, + kj::Maybe> replayMemoryTracker) : inner(kj::mv(inner)), weakRef(kj::mv(weakRefParam)), originatingCall(ownOriginatingCall(kj::mv(originatingCall))), actorTargetRetryability(actorTargetRetryability), + replayMemoryTracker(kj::mv(replayMemoryTracker)), state(Pending{kj::mv(pipeline)}) { KJ_REQUIRE(weakRef->ref == kj::none); weakRef->ref = *this; @@ -469,6 +471,9 @@ JsRpcClientProvider::ClientForOneCall JsRpcPromise::getClientForOneCall( originatingCall.map([](IoOwn& p) { return p->addRef(); }); KJ_SWITCH_ONEOF(state) { KJ_CASE_ONEOF(pending, Pending) { + KJ_IF_SOME(tracker, replayMemoryTracker) { + tracker->release(); + } return { .client = pending.pipeline->getCallPipeline(), .callSpanParents = kj::mv(callSpanParents), @@ -548,11 +553,12 @@ struct JsRpcPromiseAndPipeline { // 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; 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); + actorTargetRetryability, kj::mv(replayMemoryTracker)); } }; @@ -751,9 +757,12 @@ 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.getReplayMemoryBytes())); actorCallAttempt.emplace( generateActorRetryRequestMetadata( kj::systemCoarseCalendarClock().now(), ActorRetryGateEnabled::NO), @@ -783,6 +792,10 @@ JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, // 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. @@ -811,6 +824,11 @@ JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, // retained when traced, so the untraced path holds no span state. auto originatingCall = jsRpcCallSpan.getSpanParentsIfObserved(); + kj::Maybe> promiseReplayMemoryTracker; + KJ_IF_SOME(tracker, 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, @@ -839,6 +857,7 @@ JsRpcPromiseAndPipeline callImpl(jsg::Lock& js, .pipeline = kj::mv(callResult), .originatingCall = kj::mv(originatingCall), .actorTargetRetryability = pipelineActorTargetRetryability, + .replayMemoryTracker = kj::mv(promiseReplayMemoryTracker), }; }, [&](jsg::Value error) -> JsRpcPromiseAndPipeline { // Probably a serialization error. Need to convert to an async error since we never throw diff --git a/src/workerd/api/worker-rpc.h b/src/workerd/api/worker-rpc.h index 6555d2dee47..ab2e4b1fb5c 100644 --- a/src/workerd/api/worker-rpc.h +++ b/src/workerd/api/worker-rpc.h @@ -172,6 +172,10 @@ class JsRpcCallPlan { return replayable; } + size_t getReplayMemoryBytes() const { + return serializedData.size(); + } + void copyTo(rpc::JsRpcTarget::CallParams::Builder builder); private: @@ -325,6 +329,19 @@ class JsRpcClientProvider: public jsg::Object { class JsRpcProperty; +class JsRpcReplayMemoryTracker final: public kj::Refcounted { + public: + explicit JsRpcReplayMemoryTracker(kj::Own trackedMemory) + : trackedMemory(kj::mv(trackedMemory)) {} + + void release() { + trackedMemory = kj::Own(); + } + + private: + kj::Own trackedMemory; +}; + // Represents the promise returned by calling an RPC method. We don't use a regular Promise object, // but rather our own custom thenable, so that we can support pipelining on it. class JsRpcPromise: public JsRpcClientProvider { @@ -348,7 +365,8 @@ class JsRpcPromise: public JsRpcClientProvider { kj::Own weakRef, IoOwn pipeline, kj::Maybe originatingCall, - kj::Maybe actorTargetRetryability); + kj::Maybe actorTargetRetryability, + kj::Maybe> replayMemoryTracker); ~JsRpcPromise() noexcept(false); void resolve(jsg::Lock& js, jsg::JsValue result); @@ -404,6 +422,7 @@ class JsRpcPromise: public JsRpcClientProvider { // pipelined on the promise under it (mirrors JsRpcStub::originatingCall). Only set when traced. kj::Maybe> originatingCall; kj::Maybe actorTargetRetryability; + kj::Maybe> replayMemoryTracker; struct Pending { IoOwn pipeline; diff --git a/src/workerd/io/observer.h b/src/workerd/io/observer.h index 7081aabbfa7..b8e97a11533 100644 --- a/src/workerd/io/observer.h +++ b/src/workerd/io/observer.h @@ -180,6 +180,12 @@ class RequestObserver: public kj::Refcounted { return kj::none; } + // Tracks serialized argument bytes retained while a replayable actor call's retry state is live. + // The returned handle releases the tracked bytes when destroyed. + virtual kj::Own trackActorCallReplayMemory(size_t bytes) { + return kj::Own(); + } + // Records an additional outgoing actor call started by a runtime retry loop. virtual void recordActorRetry(ActorRetryCallType callType) {}