File
Blob: src/workerd/api/sockets-test.c++
| 1 | #include "global-scope.h" |
| 2 | #include "sockets.h" |
| 3 | |
| 4 | #include <workerd/io/io-context.h> |
| 5 | #include <workerd/io/worker-interface.h> |
| 6 | #include <workerd/io/worker.h> |
| 7 | #include <workerd/tests/test-fixture.h> |
| 8 | #include <workerd/util/autogate.h> |
| 9 | |
| 10 | #include <kj/test.h> |
| 11 | |
| 12 | namespace workerd::api { |
| 13 | namespace { |
| 14 | |
| 15 | // Minimal WorkerInterface that tracks when connect() is called and exposes the pipe. |
| 16 | class MockConnectWorkerInterface final: public WorkerInterface { |
| 17 | public: |
| 18 | MockConnectWorkerInterface( |
| 19 | bool& connectCalled, kj::HttpHeaderTable& headerTable, kj::Maybe<kj::AsyncIoStream&>& pipeEnd) |
| 20 | : connectCalled(connectCalled), |
| 21 | headerTable(headerTable), |
| 22 | pipeEnd(pipeEnd) {} |
| 23 | |
| 24 | kj::Promise<void> connect(kj::StringPtr host, |
| 25 | const kj::HttpHeaders& headers, |
| 26 | kj::AsyncIoStream& connection, |
| 27 | ConnectResponse& response, |
| 28 | kj::HttpConnectSettings settings) override { |
| 29 | connectCalled = true; |
| 30 | pipeEnd = connection; |
| 31 | kj::HttpHeaders responseHeaders(headerTable); |
| 32 | response.accept(200, "OK"_kj, responseHeaders); |
| 33 | return kj::NEVER_DONE; |
| 34 | } |
| 35 | |
| 36 | kj::Promise<void> request(kj::HttpMethod method, |
| 37 | kj::StringPtr url, |
| 38 | const kj::HttpHeaders& headers, |
| 39 | kj::AsyncInputStream& requestBody, |
| 40 | kj::HttpService::Response& response) override { |
| 41 | KJ_UNIMPLEMENTED("not used in this test"); |
| 42 | } |
| 43 | kj::Promise<void> prewarm(kj::StringPtr url) override { |
| 44 | KJ_UNIMPLEMENTED("not used in this test"); |
| 45 | } |
| 46 | kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override { |
| 47 | KJ_UNIMPLEMENTED("not used in this test"); |
| 48 | } |
| 49 | kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override { |
| 50 | KJ_UNIMPLEMENTED("not used in this test"); |
| 51 | } |
| 52 | kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override { |
| 53 | return event->notSupported(); |
| 54 | } |
| 55 | |
| 56 | private: |
| 57 | bool& connectCalled; |
| 58 | kj::HttpHeaderTable& headerTable; |
| 59 | kj::Maybe<kj::AsyncIoStream&>& pipeEnd; |
| 60 | }; |
| 61 | |
| 62 | struct ConnectTestIoChannelFactory final: public TestFixture::DummyIoChannelFactory { |
| 63 | ConnectTestIoChannelFactory(TimerChannel& timer, |
| 64 | bool& connectCalled, |
| 65 | kj::HttpHeaderTable& headerTable, |
| 66 | kj::Maybe<kj::AsyncIoStream&>& pipeEnd) |
| 67 | : DummyIoChannelFactory(timer), |
| 68 | connectCalled(connectCalled), |
| 69 | headerTable(headerTable), |
| 70 | pipeEnd(pipeEnd) {} |
| 71 | |
| 72 | kj::Own<WorkerInterface> startSubrequest(uint channel, SubrequestMetadata metadata) override { |
| 73 | return kj::heap<MockConnectWorkerInterface>(connectCalled, headerTable, pipeEnd); |
| 74 | } |
| 75 | |
| 76 | void abortIsolate(kj::StringPtr reason) override { |
| 77 | JSG_FAIL_REQUIRE(Error, "abortIsolate() is not implemented for this runtime."); |
| 78 | } |
| 79 | |
| 80 | bool& connectCalled; |
| 81 | kj::HttpHeaderTable& headerTable; |
| 82 | kj::Maybe<kj::AsyncIoStream&>& pipeEnd; |
| 83 | }; |
| 84 | |
| 85 | // Turn-based timeout: resolves after n event loop turns, returning 0. |
| 86 | kj::Promise<size_t> turnTimeout(int n) { |
| 87 | for (int i = 0; i < n; i++) { |
| 88 | co_await kj::evalLater([]() {}); |
| 89 | } |
| 90 | co_return 0; |
| 91 | } |
| 92 | |
| 93 | KJ_TEST("socket writes are blocked by output gate") { |
| 94 | bool connectCalled = false; |
| 95 | kj::HttpHeaderTable headerTable; |
| 96 | kj::Maybe<kj::AsyncIoStream&> pipeEnd; |
| 97 | |
| 98 | Worker::Actor::Id actorId = kj::str("test-actor-write"); |
| 99 | TestFixture fixture(TestFixture::SetupParams{ |
| 100 | .actorId = kj::mv(actorId), |
| 101 | .useRealTimers = false, |
| 102 | .ioChannelFactory = kj::Function<kj::Own<IoChannelFactory>(TimerChannel&)>( |
| 103 | [&](TimerChannel& timer) -> kj::Own<IoChannelFactory> { |
| 104 | return kj::heap<ConnectTestIoChannelFactory>(timer, connectCalled, headerTable, pipeEnd); |
| 105 | }), |
| 106 | }); |
| 107 | |
| 108 | static constexpr kj::StringPtr errorsToIgnore[] = { |
| 109 | "failed to invoke drain()"_kj, |
| 110 | "no subrequests"_kj, |
| 111 | }; |
| 112 | |
| 113 | fixture.runInIoContext(kj::Function<kj::Promise<void>(const TestFixture::Environment&)>( |
| 114 | [&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 115 | auto& actor = env.context.getActorOrThrow(); |
| 116 | |
| 117 | // Step 1: Connect without gate lock so the pipe is established. |
| 118 | auto socket = connectImplNoOutputLock(env.js, kj::none, kj::str("localhost:1234"), kj::none); |
| 119 | env.js.runMicrotasks(); |
| 120 | |
| 121 | // Prepare write data and lock gate BEFORE any co_await (Worker lock still held). |
| 122 | auto paf = kj::newPromiseAndFulfiller<void>(); |
| 123 | auto blocker = actor.getOutputGate().lockWhile(kj::mv(paf.promise), nullptr); |
| 124 | auto writable = socket->getWritable(); |
| 125 | auto data = kj::heapArray<kj::byte>({'h', 'i'}); |
| 126 | auto jsBuffer = env.js.bytes(kj::mv(data)).getHandle(env.js); |
| 127 | writable->getController().write(env.js, jsBuffer).markAsHandled(env.js); |
| 128 | |
| 129 | // With autogate (@all-autogates), connect is deferred. Wait for it. |
| 130 | // After co_await, Worker lock is released — no V8 calls allowed. |
| 131 | for (int i = 0; i < 10 && pipeEnd == kj::none; i++) { |
| 132 | co_await kj::evalLater([]() {}); |
| 133 | } |
| 134 | KJ_ASSERT(connectCalled); |
| 135 | auto& pipe = KJ_ASSERT_NONNULL(pipeEnd); |
| 136 | |
| 137 | // Step 4: Race tryRead against a turn-based timeout. The output gate is locked, |
| 138 | // so the write drain is stuck on outputLock — data cannot reach the pipe. |
| 139 | auto buf = kj::heapArray<kj::byte>(2); |
| 140 | auto bytesRead = |
| 141 | co_await pipe.tryRead(buf.begin(), 1, buf.size()).exclusiveJoin(turnTimeout(20)); |
| 142 | KJ_EXPECT(bytesRead == 0, "read must time out while output gate is locked"); |
| 143 | |
| 144 | // Step 5: Release the gate. |
| 145 | paf.fulfiller->fulfill(); |
| 146 | |
| 147 | // Step 6: Read again — data should arrive now. |
| 148 | bytesRead = co_await pipe.tryRead(buf.begin(), 1, buf.size()); |
| 149 | KJ_EXPECT(bytesRead == 2, "read must succeed after output gate releases"); |
| 150 | KJ_EXPECT(buf[0] == 'h'); |
| 151 | KJ_EXPECT(buf[1] == 'i'); |
| 152 | }), |
| 153 | errorsToIgnore); |
| 154 | } |
| 155 | |
| 156 | // Connect deferral test runs last — its drain errors fire during process exit. |
| 157 | KJ_TEST("connectImplNoOutputLock defers connect until output gate clears") { |
| 158 | bool connectCalled = false; |
| 159 | kj::HttpHeaderTable headerTable; |
| 160 | kj::Maybe<kj::AsyncIoStream&> pipeEnd; |
| 161 | |
| 162 | Worker::Actor::Id actorId = kj::str("test-actor"); |
| 163 | TestFixture fixture(TestFixture::SetupParams{ |
| 164 | .actorId = kj::mv(actorId), |
| 165 | .useRealTimers = false, |
| 166 | .ioChannelFactory = kj::Function<kj::Own<IoChannelFactory>(TimerChannel&)>( |
| 167 | [&](TimerChannel& timer) -> kj::Own<IoChannelFactory> { |
| 168 | return kj::heap<ConnectTestIoChannelFactory>(timer, connectCalled, headerTable, pipeEnd); |
| 169 | }), |
| 170 | }); |
| 171 | |
| 172 | bool autogateOn = util::Autogate::isEnabled(util::AutogateKey::TCP_SOCKET_CONNECT_OUTPUT_GATE); |
| 173 | |
| 174 | static constexpr kj::StringPtr errorsToIgnore[] = { |
| 175 | "failed to invoke drain()"_kj, |
| 176 | "no subrequests"_kj, |
| 177 | }; |
| 178 | |
| 179 | fixture.runInIoContext(kj::Function<kj::Promise<void>(const TestFixture::Environment&)>( |
| 180 | [&](const TestFixture::Environment& env) -> kj::Promise<void> { |
| 181 | auto& actor = env.context.getActorOrThrow(); |
| 182 | auto paf = kj::newPromiseAndFulfiller<void>(); |
| 183 | auto blocker = actor.getOutputGate().lockWhile(kj::mv(paf.promise), nullptr); |
| 184 | |
| 185 | auto socket = connectImplNoOutputLock(env.js, kj::none, kj::str("localhost:1234"), kj::none); |
| 186 | |
| 187 | if (autogateOn) { |
| 188 | co_await kj::evalLater([]() {}); |
| 189 | KJ_EXPECT(!connectCalled, "connect must not happen while output gate is locked"); |
| 190 | paf.fulfiller->fulfill(); |
| 191 | co_await kj::evalLater([]() {}); |
| 192 | KJ_EXPECT(connectCalled, "connect must happen after output gate releases"); |
| 193 | } else { |
| 194 | KJ_EXPECT(connectCalled, "without autogate, connect must happen synchronously"); |
| 195 | paf.fulfiller->fulfill(); |
| 196 | co_await kj::evalLater([]() {}); |
| 197 | } |
| 198 | }), |
| 199 | errorsToIgnore); |
| 200 | } |
| 201 | |
| 202 | } // namespace |
| 203 | } // namespace workerd::api |