#include "global-scope.h" #include "sockets.h" #include #include #include #include #include #include namespace workerd::api { namespace { // Minimal WorkerInterface that tracks when connect() is called and exposes the pipe. class MockConnectWorkerInterface final: public WorkerInterface { public: MockConnectWorkerInterface( bool& connectCalled, kj::HttpHeaderTable& headerTable, kj::Maybe& pipeEnd) : connectCalled(connectCalled), headerTable(headerTable), pipeEnd(pipeEnd) {} kj::Promise connect(kj::StringPtr host, const kj::HttpHeaders& headers, kj::AsyncIoStream& connection, ConnectResponse& response, kj::HttpConnectSettings settings) override { connectCalled = true; pipeEnd = connection; kj::HttpHeaders responseHeaders(headerTable); response.accept(200, "OK"_kj, responseHeaders); return kj::NEVER_DONE; } kj::Promise request(kj::HttpMethod method, kj::StringPtr url, const kj::HttpHeaders& headers, kj::AsyncInputStream& requestBody, kj::HttpService::Response& response) override { KJ_UNIMPLEMENTED("not used in this test"); } kj::Promise prewarm(kj::StringPtr url) override { KJ_UNIMPLEMENTED("not used in this test"); } kj::Promise runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override { KJ_UNIMPLEMENTED("not used in this test"); } kj::Promise runAlarm(kj::Date scheduledTime, uint32_t retryCount) override { KJ_UNIMPLEMENTED("not used in this test"); } kj::Promise customEvent(kj::Own event) override { return event->notSupported(); } private: bool& connectCalled; kj::HttpHeaderTable& headerTable; kj::Maybe& pipeEnd; }; struct ConnectTestIoChannelFactory final: public TestFixture::DummyIoChannelFactory { ConnectTestIoChannelFactory(TimerChannel& timer, bool& connectCalled, kj::HttpHeaderTable& headerTable, kj::Maybe& pipeEnd) : DummyIoChannelFactory(timer), connectCalled(connectCalled), headerTable(headerTable), pipeEnd(pipeEnd) {} kj::Own startSubrequest(uint channel, SubrequestMetadata metadata) override { return kj::heap(connectCalled, headerTable, pipeEnd); } void abortIsolate(kj::StringPtr reason) override { JSG_FAIL_REQUIRE(Error, "abortIsolate() is not implemented for this runtime."); } bool& connectCalled; kj::HttpHeaderTable& headerTable; kj::Maybe& pipeEnd; }; // Turn-based timeout: resolves after n event loop turns, returning 0. kj::Promise turnTimeout(int n) { for (int i = 0; i < n; i++) { co_await kj::evalLater([]() {}); } co_return 0; } KJ_TEST("socket writes are blocked by output gate") { bool connectCalled = false; kj::HttpHeaderTable headerTable; kj::Maybe pipeEnd; Worker::Actor::Id actorId = kj::str("test-actor-write"); TestFixture fixture(TestFixture::SetupParams{ .actorId = kj::mv(actorId), .useRealTimers = false, .ioChannelFactory = kj::Function(TimerChannel&)>( [&](TimerChannel& timer) -> kj::Own { return kj::heap(timer, connectCalled, headerTable, pipeEnd); }), }); static constexpr kj::StringPtr errorsToIgnore[] = { "failed to invoke drain()"_kj, "no subrequests"_kj, }; fixture.runInIoContext(kj::Function(const TestFixture::Environment&)>( [&](const TestFixture::Environment& env) -> kj::Promise { auto& actor = env.context.getActorOrThrow(); // Step 1: Connect without gate lock so the pipe is established. auto socket = connectImplNoOutputLock(env.js, kj::none, kj::str("localhost:1234"), kj::none); env.js.runMicrotasks(); // Prepare write data and lock gate BEFORE any co_await (Worker lock still held). auto paf = kj::newPromiseAndFulfiller(); auto blocker = actor.getOutputGate().lockWhile(kj::mv(paf.promise), nullptr); auto writable = socket->getWritable(); auto data = kj::heapArray({'h', 'i'}); auto jsBuffer = env.js.bytes(kj::mv(data)).getHandle(env.js); writable->getController().write(env.js, jsBuffer).markAsHandled(env.js); // With autogate (@all-autogates), connect is deferred. Wait for it. // After co_await, Worker lock is released — no V8 calls allowed. for (int i = 0; i < 10 && pipeEnd == kj::none; i++) { co_await kj::evalLater([]() {}); } KJ_ASSERT(connectCalled); auto& pipe = KJ_ASSERT_NONNULL(pipeEnd); // Step 4: Race tryRead against a turn-based timeout. The output gate is locked, // so the write drain is stuck on outputLock — data cannot reach the pipe. auto buf = kj::heapArray(2); auto bytesRead = co_await pipe.tryRead(buf.begin(), 1, buf.size()).exclusiveJoin(turnTimeout(20)); KJ_EXPECT(bytesRead == 0, "read must time out while output gate is locked"); // Step 5: Release the gate. paf.fulfiller->fulfill(); // Step 6: Read again — data should arrive now. bytesRead = co_await pipe.tryRead(buf.begin(), 1, buf.size()); KJ_EXPECT(bytesRead == 2, "read must succeed after output gate releases"); KJ_EXPECT(buf[0] == 'h'); KJ_EXPECT(buf[1] == 'i'); }), errorsToIgnore); } // Connect deferral test runs last — its drain errors fire during process exit. KJ_TEST("connectImplNoOutputLock defers connect until output gate clears") { bool connectCalled = false; kj::HttpHeaderTable headerTable; kj::Maybe pipeEnd; Worker::Actor::Id actorId = kj::str("test-actor"); TestFixture fixture(TestFixture::SetupParams{ .actorId = kj::mv(actorId), .useRealTimers = false, .ioChannelFactory = kj::Function(TimerChannel&)>( [&](TimerChannel& timer) -> kj::Own { return kj::heap(timer, connectCalled, headerTable, pipeEnd); }), }); bool autogateOn = util::Autogate::isEnabled(util::AutogateKey::TCP_SOCKET_CONNECT_OUTPUT_GATE); static constexpr kj::StringPtr errorsToIgnore[] = { "failed to invoke drain()"_kj, "no subrequests"_kj, }; fixture.runInIoContext(kj::Function(const TestFixture::Environment&)>( [&](const TestFixture::Environment& env) -> kj::Promise { auto& actor = env.context.getActorOrThrow(); auto paf = kj::newPromiseAndFulfiller(); auto blocker = actor.getOutputGate().lockWhile(kj::mv(paf.promise), nullptr); auto socket = connectImplNoOutputLock(env.js, kj::none, kj::str("localhost:1234"), kj::none); if (autogateOn) { co_await kj::evalLater([]() {}); KJ_EXPECT(!connectCalled, "connect must not happen while output gate is locked"); paf.fulfiller->fulfill(); co_await kj::evalLater([]() {}); KJ_EXPECT(connectCalled, "connect must happen after output gate releases"); } else { KJ_EXPECT(connectCalled, "without autogate, connect must happen synchronously"); paf.fulfiller->fulfill(); co_await kj::evalLater([]() {}); } }), errorsToIgnore); } } // namespace } // namespace workerd::api