Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 9 additions & 0 deletions src/workerd/api/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -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 = [
Expand Down
81 changes: 81 additions & 0 deletions src/workerd/api/actor-call-retry-test.c++
Original file line number Diff line number Diff line change
Expand Up @@ -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<RecordingObserver>();
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<kj::Duration>());

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<kj::Duration>());

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<kj::byte>(0));
auto result = handleFailure(*state, kj::mv(rejected));
auto& failure = KJ_ASSERT_NONNULL(result.tryGet<kj::Exception>());
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<RecordingObserver>();
auto state = newRetryState(timer, *observer);

startAttempt(*state);
KJ_EXPECT(handleFailure(*state, makeDisconnect("original disconnect"_kj)).is<kj::Duration>());
startAttempt(*state);
auto rejected = KJ_EXCEPTION(FAILED, "claim rejected");
rejected.setDetail(jsg::ACTOR_RETRY_CLAIM_REJECTED_DETAIL_ID, kj::heapArray<kj::byte>(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<RecordingObserver>();
Expand Down Expand Up @@ -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<RecordingObserver>();
auto state = kj::rc<ActorCallRetryState>(timer, *observer, config);

startAttempt(*state);
auto result = handleFailure(*state, makeDisconnect("original disconnect"_kj));
auto& failure = KJ_ASSERT_NONNULL(result.tryGet<kj::Exception>());
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
13 changes: 13 additions & 0 deletions src/workerd/api/actor-call-retry.c++
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ kj::OneOf<ActorCallRetryState::Attempt, kj::Exception> ActorCallRetryState::star

kj::OneOf<kj::Duration, kj::Exception> ActorCallRetryState::handleAttemptFailure(
kj::Exception exception) {
if (!retriesEnabled) return kj::mv(exception);

maybeStartRetryLatencyTimer(exception);
KJ_IF_SOME(claimRejection, handleClaimRejection(exception)) {
return kj::mv(claimRejection);
Expand All @@ -59,6 +61,17 @@ kj::OneOf<kj::Duration, kj::Exception> 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();
Expand Down
5 changes: 5 additions & 0 deletions src/workerd/api/actor-call-retry.h
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,11 @@ class ActorCallRetryState final: public kj::Refcounted {

kj::OneOf<Attempt, kj::Exception> startAttempt();
kj::OneOf<kj::Duration, kj::Exception> 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;
Expand Down
2 changes: 1 addition & 1 deletion src/workerd/api/http.h
Original file line number Diff line number Diff line change
Expand Up @@ -349,7 +349,7 @@ class Fetcher: public JsRpcClientProvider {
MakeUserSpanParent makeUserSpanParent);

kj::Maybe<ActorCallTargetRetryable> getActorTargetRetryability() override;
void onActorCallRetry();
void onActorCallRetry() override;

// Get a SubrequestChannel representing this Fetcher.
kj::Own<IoChannelFactory::SubrequestChannel> getSubrequestChannel(IoContext& ioContext);
Expand Down
12 changes: 9 additions & 3 deletions src/workerd/api/jsrpc-call-observation-test.c++
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ struct Observation {
ActorCallPayloadReplayable payloadReplayable;
ActorCallTargetRetryable targetRetryable;
Settlement settlement = Settlement::PENDING;
bool pipelineCommitted = false;
};

struct ObservationState {
Expand All @@ -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);
}
Expand Down Expand Up @@ -215,7 +219,7 @@ class Harness {
auto& targetHandler = KJ_REQUIRE_NONNULL(env.js.tryGetTypeHandler<jsg::Ref<JsRpcTarget>>());
auto target =
KJ_REQUIRE_NONNULL(jsg::JsValue(targetHandler.wrap(env.js, env.js.alloc<JsRpcTarget>()))
.tryCast<jsg::JsObject>());
.tryCast<jsg::JsObject>());
return JsRpcStub::constructor(env.js, target);
}

Expand Down Expand Up @@ -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);
}

Expand All @@ -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;
Expand All @@ -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);
}
Expand Down
46 changes: 41 additions & 5 deletions src/workerd/api/jsrpc-call-plan-test.c++
Original file line number Diff line number Diff line change
Expand Up @@ -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<rpc::JsRpcTarget::CallParams>());
auto propertyParams = propertyAttempt.getRoot<rpc::JsRpcTarget::CallParams>();
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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));
Expand All @@ -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.
Expand Down
Loading
Loading