diff --git a/CHANGELOG.md b/CHANGELOG.md index fb78a13b..27db488d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -20,7 +20,11 @@ policy in [docs/versioning.md](docs/versioning.md). instantiates with a small driver. Both transports carried the identical terminal-ordering bug precisely because their suites mirrored each other by hand; anything an implementation must do regardless of its wire now - lives in one body that runs against all of them. + lives in one body that runs against all of them. Published as the + `websocket_contract_test_support` target, so an out-of-tree `WebSocket` + can be held to the same contract; the consumer module implements one and + does exactly that. Four implementations run it today — both transports, + the `JsonRpcStreamSocket` decorator, and the consumer's. ### Fixed diff --git a/examples/bazel-consumer/BUILD.bazel b/examples/bazel-consumer/BUILD.bazel index e879f097..e2a97929 100644 --- a/examples/bazel-consumer/BUILD.bazel +++ b/examples/bazel-consumer/BUILD.bazel @@ -126,6 +126,26 @@ cc_test( ], ) +# Out-of-tree proof of the ADR-0019 implementor contract (#173): a +# WebSocket written in consumer code, held to the same shared contract +# suite the in-repo transports are. Depends on the runtime's testonly +# contract-suite target, which is the point — that seam is published for +# implementors, so a consumer must be able to reach it. No wire, no +# network: this one runs everywhere. +cc_test( + name = "websocket_contract_consumer_test", + size = "small", + srcs = ["websocket_contract_consumer_test.cc"], + copts = SMITHY_COPTS, + deps = [ + "@googletest//:gtest_main", + "@smithy_cpp//runtime:core", + "@smithy_cpp//runtime:eventstream", + "@smithy_cpp//runtime:http", + "@smithy_cpp//runtime:websocket_contract_test_support", + ], +) + # Streaming-transport acceptance (ADR-0015, the e2e ADR-0014's amendment # requires): a consumer dials the upgraded server and drains real frames # through the module boundary. Beast-backed like todo_beast_acceptance_test; diff --git a/examples/bazel-consumer/websocket_contract_consumer_test.cc b/examples/bazel-consumer/websocket_contract_consumer_test.cc new file mode 100644 index 00000000..4587e78a --- /dev/null +++ b/examples/bazel-consumer/websocket_contract_consumer_test.cc @@ -0,0 +1,158 @@ +// Out-of-tree proof for the ADR-0019 implementor contract (#173): a +// WebSocket written entirely in consumer code adopts the async primitives, +// runs its terminal transition through WebSocket::TerminalWaiters, and is +// held to the same shared contract suite the in-repo transports are. +// +// This is what makes the seam's claim real rather than aspirational. The +// runtime says the async methods are public virtuals a third party may +// override, and that overriding them accepts the send-before-receive rule; +// this test is the only place that a third party actually does so — across +// the module boundary, against the published targets alone, with no +// in-repo transport in the loop. + +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "smithy/eventstream/frame.h" +#include "smithy/http/websocket.h" +#include "smithy/testing/websocket_contract_test.h" + +namespace { + +using smithy::Outcome; +using smithy::Unit; +using smithy::eventstream::Message; +using smithy::http::WebSocket; + +// A third-party session: a bounded outbound wire nobody drains, one parked +// receive, one parked send. Deliberately minimal — the point is not the +// wire but that the terminal transition is expressed with TerminalWaiters, +// so this implementation inherits the ordering rule without its author +// having to rediscover why the rule exists. +class ConsumerSocket final : public WebSocket { + public: + // Small on purpose: the contract suite wedges the wire by sending, and a + // shallow queue gets there in a few messages. + static constexpr std::size_t kDepth = 4; + + Outcome> Receive() override { + std::unique_lock lock(mutex_); + changed_.wait(lock, [this] { return closed_; }); + return std::optional(); // this socket's peer only ever ends it + } + + Outcome> Receive(std::chrono::milliseconds timeout) override { + std::unique_lock lock(mutex_); + if (!changed_.wait_for(lock, timeout, [this] { return closed_; })) { + return smithy::Error::Timeout("consumer socket: no message within the deadline"); + } + return std::optional(); + } + + Outcome Send(const Message& message) override { + std::unique_lock lock(mutex_); + changed_.wait(lock, [this] { return queued_ < kDepth || closed_; }); + if (closed_) return smithy::Error::Transport("consumer socket: session is closed"); + ++queued_; + (void)message; + return Unit{}; + } + + void Close() override { EndSession(); } + + bool SupportsAsync() const override { return true; } + + // Both async twins: park under the lock, or complete inline once the + // lock is released — the seam's documented shapes, nothing more. + void ReceiveAsync(ReceiveCallback callback) override { + Outcome> immediate = std::optional(); // the clean end + { + const std::lock_guard lock(mutex_); + if (!closed_ && !pending_receive_) { + pending_receive_ = std::move(callback); + return; // EndSession completes it + } + if (!closed_) { + immediate = smithy::Error::Validation("consumer socket: a receive is already outstanding"); + } + } + callback(std::move(immediate)); + } + + void SendAsync(const Message& message, SendCallback callback) override { + (void)message; + Outcome immediate = Unit{}; + { + const std::lock_guard lock(mutex_); + if (closed_) { + immediate = smithy::Error::Transport("consumer socket: session is closed"); + } else if (pending_send_) { + immediate = smithy::Error::Validation("consumer socket: a send is already in flight"); + } else if (queued_ >= kDepth) { + pending_send_ = std::move(callback); // parked on the full wire + return; + } else { + ++queued_; + } + } + callback(std::move(immediate)); + } + + // The far side ending the session — a peer close, a reset, whatever this + // implementation's wire calls it. Takes both parked completions under the + // lock, releases it, and fires through TerminalWaiters, which is what + // puts the send ahead of the receive. + void EndSession() { + WebSocket::TerminalWaiters waiters; + { + const std::lock_guard lock(mutex_); + if (closed_) return; + closed_ = true; + waiters = WebSocket::TerminalWaiters(std::exchange(pending_receive_, nullptr), + std::exchange(pending_send_, nullptr)); + changed_.notify_all(); + } + std::move(waiters).Fire( + smithy::Error::Transport("consumer socket: session is closed"), std::optional(), + [](const char*, const auto& callback, auto outcome) { callback(std::move(outcome)); }); + } + + private: + std::mutex mutex_; + std::condition_variable changed_; + ReceiveCallback pending_receive_; + SendCallback pending_send_; + std::size_t queued_ = 0; + bool closed_ = false; +}; + +struct ConsumerContractDriver { + static constexpr int kWedgeAttempts = static_cast(ConsumerSocket::kDepth) + 4; + + std::shared_ptr Socket() { return socket_; } + + Message BulkMessage(int n) { + return Message{.headers = {{":event-type", "bulk"}}, + .payload = smithy::Blob::FromString(std::to_string(n))}; + } + + void EndSessionFromPeer() { socket_->EndSession(); } + + std::shared_ptr socket_ = std::make_shared(); +}; + +} // namespace + +// gtest builds the registration symbols from the bare suite name, so the +// instantiation lives in the namespace the suite was registered in. +namespace smithy::testing { +INSTANTIATE_TYPED_TEST_SUITE_P(ConsumerSocket, WebSocketContractTest, ConsumerContractDriver); +} // namespace smithy::testing diff --git a/runtime/BUILD.bazel b/runtime/BUILD.bazel index 9054474d..ae01b678 100644 --- a/runtime/BUILD.bazel +++ b/runtime/BUILD.bazel @@ -161,6 +161,7 @@ cc_test( ":eventstream_jsonrpc", ":http", ":json", + ":websocket_contract_test_support", "@googletest//:gtest_main", ], ) diff --git a/runtime/testing/include/smithy/testing/websocket_contract_test.h b/runtime/testing/include/smithy/testing/websocket_contract_test.h index 56fd50e4..04ba78c3 100644 --- a/runtime/testing/include/smithy/testing/websocket_contract_test.h +++ b/runtime/testing/include/smithy/testing/websocket_contract_test.h @@ -137,6 +137,52 @@ std::shared_ptr>> WedgeThenPark(Driver& driver, Is return nullptr; } +// Waits for a session loop to unwind. Aborts rather than returning false: +// the loop holds pointers into the test's frame, so a test that gives up +// and returns would unwind that frame underneath a live coroutine. +inline void AwaitLoopEnd(const std::atomic& ended, const char* what) { + for (int i = 0; i < 3000 && !ended.load(); ++i) { + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + } + if (!ended.load()) { + ADD_FAILURE() << what; + std::abort(); + } +} + +// Waits until a counter stops advancing — how a loop parked in its own +// `co_await Send` announces itself. Stability rather than a fixed count, +// because how many writes a wire accepts before wedging is the transport's +// business, not this suite's. +// +// Telling "parked" from "slow" is the whole difficulty, and it is the +// driver's job to make it easy: a BulkMessage larger than the wire's +// buffers wedges decisively, so the counter stops dead rather than +// crawling. Three consecutive equal readings a second apart is the margin +// for that; a driver whose messages are too small to wedge fails here +// instead, which is what "check the driver" means. +// +// Zero is a perfectly good stable value — with a message big enough, the +// FIRST send parks and the counter never leaves zero. (Requiring progress +// before believing the count made this unable to see the most decisively +// parked case there is.) A Detached coroutine starts eagerly and issues +// its first send before the launch expression returns, so by the time +// anyone polls, "stable" cannot mean "not started yet"; it means parked — +// or finished, which the caller rules out separately. +inline bool WaitUntilStable(const std::atomic& counter) { + constexpr auto kSettle = std::chrono::seconds(1); + int last = -1; + int stable = 0; + for (int i = 0; i < 20; ++i) { + std::this_thread::sleep_for(kSettle); + const int now = counter.load(); + stable = (now == last) ? stable + 1 : 0; + if (stable >= 3) return true; + last = now; + } + return false; +} + template class WebSocketContractTest : public ::testing::Test {}; @@ -188,10 +234,59 @@ TYPED_TEST_P(WebSocketContractTest, ATerminalTransitionFiresTheParkedSendBeforeT auto dead = parked->Wait(); ASSERT_FALSE(dead.ok()) << "a send on an ended session must fail"; EXPECT_EQ(dead.error().kind(), ErrorKind::kTransport); - for (int i = 0; i < 1000 && !loop_ended.load(); ++i) { - std::this_thread::sleep_for(std::chrono::milliseconds(10)); - } - EXPECT_TRUE(loop_ended.load()) << "the session loop never unwound"; + AwaitLoopEnd(loop_ended, "the session loop never unwound"); +} + +// A terminal transition whose ONLY parked waiter is the loop's own +// `co_await Send`: it must complete that send, and the loop must unwind +// through its whole teardown without wedging the thread that fired it. +// +// What makes this worth its own test is where the teardown runs. Firing +// the send resumes the loop inline, so the loop's exit — ~AsyncEventStream +// revoking its shared view, closing the session, draining pins — all +// happens underneath Fire, on the completing thread, before Fire returns. +// The session's Close() reenters the very transition that is running. So +// the stream Share()s: without a shared view End() is a no-op and this +// exercises nothing but a resume. +// +// It is NOT the ordering tripwire, despite being the send-first case — +// with no receive parked, Fire's order is unobservable here. Ordering is +// pinned by ATerminalTransitionFiresTheParkedSendBeforeTheParkedReceive, +// which is the test to look at if the send-before-receive rule regresses. +TYPED_TEST_P(WebSocketContractTest, ATerminalTransitionCompletesALoneParkedCoroutineSend) { + TypeParam driver; + ContractMailbox torn_down; + std::atomic sent{0}; + std::atomic loop_ended{false}; + + [](std::shared_ptr socket, TypeParam* driver, std::atomic* sent, + std::atomic* ended) -> eventstream::Detached { + ContractStream stream(std::move(socket), IdentityCodec, IdentityCodec); + // A live shared view, so the stream's destructor really revokes and + // drains instead of returning early (ADR-0017) — the hub shape, and + // the whole point of ending the session from inside Fire. + (void)stream.Share(); + for (int i = 0; i <= TypeParam::kWedgeAttempts; ++i) { + auto ok = co_await stream.Send(driver->BulkMessage(i)); + if (!ok.ok()) break; // the session ended under the parked send + ++*sent; + } + *ended = true; + }(driver.Socket(), &driver, &sent, &loop_ended); + + ASSERT_TRUE(WaitUntilStable(sent)) << "the loop never parked in Send; check the driver"; + // Stable-and-finished is the other way a counter stops moving, and it + // would make this test pass with nothing parked at all. + ASSERT_FALSE(loop_ended.load()) << "the loop ran to completion instead of parking in Send"; + + std::thread ender([&] { + driver.EndSessionFromPeer(); + torn_down.Post(Unit{}); + }); + torn_down.Wait(); + ender.join(); + + AwaitLoopEnd(loop_ended, "the loop never unwound after its parked send completed"); } // The same transition with only a receive parked: the send-first ordering @@ -214,10 +309,7 @@ TYPED_TEST_P(WebSocketContractTest, ATerminalTransitionCompletesALoneParkedRecei torn_down.Wait(); ender.join(); - for (int i = 0; i < 1000 && !loop_ended.load(); ++i) { - std::this_thread::sleep_for(std::chrono::milliseconds(10)); - } - EXPECT_TRUE(loop_ended.load()) << "a lone parked receive was left hanging"; + AwaitLoopEnd(loop_ended, "a lone parked receive was left hanging"); } // One outstanding send-class operation per session: the second refuses @@ -268,6 +360,7 @@ TYPED_TEST_P(WebSocketContractTest, ASecondReceiveClassOperationRefusesWhileOneI REGISTER_TYPED_TEST_SUITE_P(WebSocketContractTest, ATerminalTransitionFiresTheParkedSendBeforeTheParkedReceive, + ATerminalTransitionCompletesALoneParkedCoroutineSend, ATerminalTransitionCompletesALoneParkedReceive, ASecondSendClassOperationRefusesWhileOneIsParked, ASecondReceiveClassOperationRefusesWhileOneIsParked); diff --git a/runtime/tests/eventstream/jsonrpc_stream_socket_test.cc b/runtime/tests/eventstream/jsonrpc_stream_socket_test.cc index 413b44fd..5cb86152 100644 --- a/runtime/tests/eventstream/jsonrpc_stream_socket_test.cc +++ b/runtime/tests/eventstream/jsonrpc_stream_socket_test.cc @@ -22,6 +22,7 @@ #include "smithy/eventstream/envelope.h" #include "smithy/eventstream/frame.h" #include "smithy/http/websocket_pair.h" +#include "smithy/testing/websocket_contract_test.h" namespace smithy::eventstream { namespace { @@ -330,5 +331,45 @@ TEST(JsonRpcStreamSocketTest, CloseAndThePeersCleanCloseDelegate) { EXPECT_FALSE(received->has_value()); // the peer's clean close, untranslated } +// --------------------------------------------------------------------------- +// The shared WebSocket contract (websocket_contract_test.h), decorator half. +// --------------------------------------------------------------------------- + +// This wrapper parks nothing of its own — both async twins forward to the +// socket underneath — which is why the #173 terminal-ordering fix had +// nothing to change here. That was a claim about this file's code, checked +// by reading; this instantiation is what keeps it true. If the wrapper ever +// grows slots of its own, the ordering rule starts applying to it, and the +// suite says so instead of the next reader having to notice. +struct JsonRpcContractDriver { + static constexpr int kWedgeAttempts = + static_cast(http::InMemoryWebSocketPair::kQueueDepth) + 4; + + JsonRpcContractDriver() { + auto [peer, wrapped] = http::InMemoryWebSocketPair::Create(); + peer_ = peer; // never receives, so the wire behind the wrapper wedges + socket_ = Wrap(wrapped, JsonRpcStreamSocket::Role::kServer); + } + + std::shared_ptr Socket() { return socket_; } + + Message BulkMessage(int n) { + return MakeEventMessage("message", "application/json", + Blob::FromString(R"({"n":)" + std::to_string(n) + "}")); + } + + void EndSessionFromPeer() { peer_->Close(); } + + std::shared_ptr peer_; + std::shared_ptr socket_; +}; + } // namespace } // namespace smithy::eventstream + +// gtest builds the registration symbols from the bare suite name, so the +// instantiation lives in the namespace the suite was registered in. +namespace smithy::testing { +INSTANTIATE_TYPED_TEST_SUITE_P(JsonRpcStreamSocket, WebSocketContractTest, + eventstream::JsonRpcContractDriver); +} // namespace smithy::testing diff --git a/runtime/tests/http/beast_websocket_test.cc b/runtime/tests/http/beast_websocket_test.cc index c6caa58e..31a70df0 100644 --- a/runtime/tests/http/beast_websocket_test.cc +++ b/runtime/tests/http/beast_websocket_test.cc @@ -1951,10 +1951,14 @@ struct BeastContractDriver { return ready.get(); } - // A megabyte a go: the wire wedges once the socket buffers fill, well - // inside kWedgeAttempts on any loopback. + // Eight mebibytes a go — one message larger than any loopback's socket + // buffers, so the wedge is decisive rather than a race against how fast + // the kernel drains. A megabyte was not: macOS absorbed 1 MiB writes for + // 30s straight without ever parking, and the 2s-per-attempt detector in + // WedgeThenPark read those slow-but-progressing writes as parked ones. + // Sized against the 16 MiB frame limit, with room for the envelope. eventstream::Message BulkMessage(int /*n*/) { - return Text("bulk", std::string(1024 * 1024, 'x')); + return Text("bulk", std::string(8 * 1024 * 1024, 'x')); } void EndSessionFromPeer() { peer_->SendText("boom"); } diff --git a/runtime/tests/server/session_registry_test.cc b/runtime/tests/server/session_registry_test.cc index b557f13a..fc93888c 100644 --- a/runtime/tests/server/session_registry_test.cc +++ b/runtime/tests/server/session_registry_test.cc @@ -30,6 +30,7 @@ #include #include +#include "smithy/eventstream/async_event_stream.h" #include "smithy/eventstream/event_stream.h" #include "smithy/eventstream/frame.h" #include "smithy/http/websocket.h" @@ -573,6 +574,70 @@ TEST(SessionRegistryAsyncTest, TeardownWithAParkedChainNeverHangs) { SUCCEED(); } +// The whole composition a hub actually runs (ADR-0019 + ADR-0017): a +// Detached loop owning the session, its Share() handle in an +// async_delivery registry, and a peer that closes while a fan-out delivery +// is still parked on the wire. This is the shape #173 wedged, and the one +// no test had: the transport suites park waiters without a registry above +// them, and the registry suites park chains without a coroutine loop +// below. The bug needed both halves at once. +using AsyncServerStream = eventstream::AsyncEventStream; + +Outcome DecodeAnything(const Message& message) { return message; } + +TEST(SessionRegistryAsyncTest, APeerCloseWithAParkedFanOutEndsTheLoopAndTheDelivery) { + // The loop parks in Receive; the chain parks inside SendAsync on the full + // wire, holding its revocation pin. The peer's close then has to complete + // the delivery, resume the loop, and let ~AsyncEventStream drain that pin + // — all on the closing thread, in that order. This test hanging is the + // failure mode. + auto [client_end, server_end] = http::InMemoryWebSocketPair::Create(); + Registry registry = AsyncRegistry(); + std::atomic loop_ended{false}; + std::promise> minted; + auto handle_ready = minted.get_future(); + + [](std::shared_ptr socket, std::promise>* minted, + std::atomic* ended) -> eventstream::Detached { + AsyncServerStream stream(std::move(socket), EncodeNote, DecodeAnything); + minted->set_value(stream.Share()); + (void)co_await stream.Receive(); // parked until the peer closes + *ended = true; + }(server_end, &minted, &loop_ended); + + ASSERT_EQ(handle_ready.wait_for(std::chrono::seconds(5)), std::future_status::ready); + ASSERT_TRUE(registry.Add("ada", handle_ready.get())); + + // Past the wire bound on purpose: the wire takes what it can and the + // chain parks on the next delivery, which is where the pin lives. + for (std::size_t i = 0; i < 2 * kWireDepth; ++i) { + ASSERT_TRUE(registry.SendTo("ada", Note{"pile-up"})); + } + + // The pair runs the entire teardown on whoever calls Close, so keep it + // off the test thread — a wedge there cannot even be joined. + std::promise closed; + auto closed_future = closed.get_future(); + std::thread closer([&] { + client_end->Close(); + closed.set_value(); + }); + if (closed_future.wait_for(std::chrono::seconds(30)) != std::future_status::ready) { + ADD_FAILURE() << "the close wedged: a parked fan-out delivery was never completed"; + std::abort(); // the closer cannot be joined, and the loop still holds this frame + } + closer.join(); + + for (int i = 0; i < 1000 && !loop_ended.load(); ++i) { + std::this_thread::sleep_for(std::chrono::milliseconds(10)); + } + if (!loop_ended.load()) { + ADD_FAILURE() << "the session loop never unwound"; + std::abort(); + } + EXPECT_TRUE(registry.Remove("ada")); // the entry outlived the session, harmlessly +} + TEST(SessionRegistryAsyncTest, ChurnUnderBroadcastStaysSafe) { Registry registry = AsyncRegistry(/*queue_capacity=*/2); std::atomic stop{false};