Skip to content

Commit b831f89

Browse files
committed
Move retry claim before actor construction
Claim retry tokens in WorkerEntrypoint before delivered() constructs the actor. Keep the hook synchronous so capability pipelining is unchanged. Update the test to pin the claim-before-delivery ordering.
1 parent b73ce3a commit b831f89

6 files changed

Lines changed: 93 additions & 77 deletions

File tree

src/workerd/api/BUILD.bazel

Lines changed: 0 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -812,14 +812,6 @@ kj_test(
812812
],
813813
)
814814

815-
kj_test(
816-
src = "fetch-retry-claim-hook-test.c++",
817-
deps = [
818-
"//src/workerd/io",
819-
"//src/workerd/tests:test-fixture",
820-
],
821-
)
822-
823815
wd_test(
824816
src = "streams/streams-test.wd-test",
825817
args = ["--experimental"],

src/workerd/api/fetch-retry-claim-hook-test.c++

Lines changed: 0 additions & 56 deletions
This file was deleted.

src/workerd/api/global-scope.c++

Lines changed: 0 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -378,11 +378,6 @@ kj::Promise<DeferredProxy<void>> ServiceWorkerGlobalScope::request(kj::HttpMetho
378378
bool useDefaultHandling;
379379
KJ_IF_SOME(h, exportedHandler) {
380380
KJ_IF_SOME(f, h.fetch) {
381-
// Immediately before dispatching into user code, give the observer a chance to claim the
382-
// request's retry-token nonce. No-op unless the request is to an actor and the observer
383-
// overrides the hook. Only the exported-handler path is hooked: Durable Objects are always
384-
// class-based, so actor fetches never take the service-worker dispatchEventImpl() path below.
385-
ioContext.getMetrics().claimRetryTokenBeforeUserCode();
386381
auto promise = f(lock, event->getRequest(), h.env.addRef(js), h.getCtx());
387382
event->respondWith(lock, kj::mv(promise));
388383
useDefaultHandling = false;

src/workerd/io/observer.h

Lines changed: 3 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -143,14 +143,9 @@ class RequestObserver: public kj::Refcounted {
143143
// the base observer; edgeworker overrides it to feed retry classification.
144144
virtual void setNextSubrequestBodyRewindable(SubrequestBodyRewindable bodyRewindable) {}
145145

146-
// Fired immediately before an actor fetch dispatches into user code, so an observer can claim the
147-
// request's retry-token nonce against the actor's claim store. No-op in the base observer;
148-
// edgeworker overrides it. Gating and fetch-only scoping are the override's concern.
149-
//
150-
// This is only fired on the exported-handler (ES modules) dispatch path, not the service-worker
151-
// addEventListener('fetch') path. That is deliberate, not an omission: Durable Objects must be
152-
// class-based and so are always invoked via an exported handler, meaning an actor fetch never
153-
// reaches the service-worker path -- where this hook would be a guaranteed no-op anyway.
146+
// Fired before a fetch request is delivered, so an observer can claim an actor request's
147+
// retry-token nonce before actor construction or user code. This also fires for non-actor and
148+
// service-worker fetches; observers are responsible for treating those as no-ops.
154149
virtual void claimRetryTokenBeforeUserCode() {}
155150

156151
// Used to record when a worker has used a dynamic dispatch binding.

src/workerd/io/worker-entrypoint-test.c++

Lines changed: 87 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,12 +80,34 @@ class TestResponse final: public kj::HttpService::Response {
8080
kj::StringPtr statusText,
8181
const kj::HttpHeaders& headers,
8282
kj::Maybe<uint64_t> expectedBodySize) override {
83+
this->statusCode = statusCode;
8384
return kj::heap<kj::NullStream>();
8485
}
8586

8687
kj::Own<kj::WebSocket> acceptWebSocket(const kj::HttpHeaders& headers) override {
8788
KJ_FAIL_ASSERT("request unexpectedly returned a WebSocket");
8889
}
90+
91+
uint statusCode = 0;
92+
};
93+
94+
class RetryClaimObserver final: public RequestObserver {
95+
public:
96+
void claimRetryTokenBeforeUserCode() override {
97+
KJ_EXPECT(stage == 0);
98+
stage = 1;
99+
if (rejectClaim) {
100+
kj::throwFatalException(KJ_EXCEPTION(FAILED, "claim rejected"));
101+
}
102+
}
103+
104+
void delivered() override {
105+
KJ_EXPECT(stage == 1);
106+
stage = 2;
107+
}
108+
109+
uint stage = 0;
110+
bool rejectClaim = false;
89111
};
90112

91113
class ThrowingResponse final: public kj::HttpService::Response {
@@ -195,6 +217,71 @@ class RecordingObserver final: public RequestObserver, public WorkerInterface {
195217
kj::Maybe<WorkerInterface&> inner;
196218
};
197219

220+
KJ_TEST("retry claim fires synchronously before fetch delivery") {
221+
auto observer = kj::refcounted<RetryClaimObserver>();
222+
TestFixture fixture(TestFixture::SetupParams{
223+
.mainModuleSource = R"SCRIPT(
224+
export default {
225+
async fetch() {
226+
return new Response("OK");
227+
},
228+
};
229+
)SCRIPT"_kj,
230+
.requestObserverFactory = kj::Function<kj::Own<RequestObserver>()>(
231+
[&observer]() -> kj::Own<RequestObserver> { return kj::addRef(*observer); }),
232+
});
233+
234+
auto entrypoint = fixture.makeWorkerEntrypoint();
235+
kj::HttpHeaderTable headerTable;
236+
kj::HttpHeaders headers(headerTable);
237+
kj::NullStream requestBody;
238+
TestResponse response;
239+
240+
auto request = entrypoint->request(
241+
kj::HttpMethod::GET, "https://example.com", headers, requestBody, response);
242+
KJ_EXPECT(observer->stage == 2, "claim and delivery hooks did not fire synchronously");
243+
request.wait(fixture.getWaitScope());
244+
245+
KJ_EXPECT(response.statusCode == 200);
246+
KJ_EXPECT(observer->stage == 2, "claim hook fired more than once");
247+
}
248+
249+
KJ_TEST("rejected retry claim prevents actor construction and delivery") {
250+
auto observer = kj::refcounted<RetryClaimObserver>();
251+
observer->rejectClaim = true;
252+
TestFixture fixture(TestFixture::SetupParams{
253+
.mainModuleSource = R"SCRIPT(
254+
export default class {
255+
constructor() {
256+
throw new Error("actor was constructed");
257+
}
258+
259+
async fetch() {
260+
return new Response("unexpected");
261+
}
262+
}
263+
)SCRIPT"_kj,
264+
.actorId = Worker::Actor::Id(kj::str("retry-claim-test")),
265+
.requestObserverFactory = kj::Function<kj::Own<RequestObserver>()>(
266+
[&observer]() -> kj::Own<RequestObserver> { return kj::addRef(*observer); }),
267+
});
268+
269+
auto entrypoint = fixture.makeWorkerEntrypoint();
270+
kj::HttpHeaderTable headerTable;
271+
kj::HttpHeaders headers(headerTable);
272+
kj::NullStream requestBody;
273+
TestResponse response;
274+
275+
auto exception = kj::runCatchingExceptions([&]() {
276+
entrypoint->request(kj::HttpMethod::GET, "https://example.com", headers, requestBody, response)
277+
.wait(fixture.getWaitScope());
278+
});
279+
280+
KJ_EXPECT(KJ_ASSERT_NONNULL(exception).getDescription().contains("claim rejected"));
281+
KJ_EXPECT(observer->stage == 1, "request was delivered after its retry claim was rejected");
282+
KJ_EXPECT(response.statusCode == 0);
283+
}
284+
198285
KJ_TEST("connect pass-through tags failures after delivery") {
199286
capnp::MallocMessageBuilder flagsMessage;
200287
auto flags = flagsMessage.initRoot<CompatibilityFlags>();

src/workerd/io/worker-entrypoint.c++

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -399,6 +399,9 @@ kj::Promise<void> WorkerEntrypoint::requestImpl(kj::HttpMethod method,
399399
workerTracer = t;
400400
}
401401

402+
// Claim before delivered() constructs an actor. This introduces no asynchronous boundary, so
403+
// capability pipelining remains unchanged.
404+
incomingRequest->getMetrics().claimRetryTokenBeforeUserCode();
402405
incomingRequest->delivered();
403406

404407
auto metricsForCatch = kj::addRef(incomingRequest->getMetrics());

0 commit comments

Comments
 (0)