Skip to content
Merged
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
6 changes: 5 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
20 changes: 20 additions & 0 deletions examples/bazel-consumer/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
158 changes: 158 additions & 0 deletions examples/bazel-consumer/websocket_contract_consumer_test.cc
Original file line number Diff line number Diff line change
@@ -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 <gtest/gtest.h>

#include <chrono>
#include <condition_variable>
#include <cstddef>
#include <memory>
#include <mutex>
#include <optional>
#include <string>
#include <utility>

#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<std::optional<Message>> Receive() override {
std::unique_lock<std::mutex> lock(mutex_);
changed_.wait(lock, [this] { return closed_; });
return std::optional<Message>(); // this socket's peer only ever ends it
}

Outcome<std::optional<Message>> Receive(std::chrono::milliseconds timeout) override {
std::unique_lock<std::mutex> 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<Message>();
}

Outcome<Unit> Send(const Message& message) override {
std::unique_lock<std::mutex> 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<std::optional<Message>> immediate = std::optional<Message>(); // the clean end
{
const std::lock_guard<std::mutex> 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<Unit> immediate = Unit{};
{
const std::lock_guard<std::mutex> 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<std::mutex> 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<Message>(),
[](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<int>(ConsumerSocket::kDepth) + 4;

std::shared_ptr<WebSocket> 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<ConsumerSocket> socket_ = std::make_shared<ConsumerSocket>();
};

} // 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
1 change: 1 addition & 0 deletions runtime/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,7 @@ cc_test(
":eventstream_jsonrpc",
":http",
":json",
":websocket_contract_test_support",
"@googletest//:gtest_main",
],
)
Expand Down
109 changes: 101 additions & 8 deletions runtime/testing/include/smithy/testing/websocket_contract_test.h
Original file line number Diff line number Diff line change
Expand Up @@ -137,6 +137,52 @@ std::shared_ptr<ContractMailbox<Outcome<Unit>>> 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<bool>& 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<int>& 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 <typename Driver>
class WebSocketContractTest : public ::testing::Test {};

Expand Down Expand Up @@ -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<Unit> torn_down;
std::atomic<int> sent{0};
std::atomic<bool> loop_ended{false};

[](std::shared_ptr<http::WebSocket> socket, TypeParam* driver, std::atomic<int>* sent,
std::atomic<bool>* 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
Expand All @@ -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
Expand Down Expand Up @@ -268,6 +360,7 @@ TYPED_TEST_P(WebSocketContractTest, ASecondReceiveClassOperationRefusesWhileOneI

REGISTER_TYPED_TEST_SUITE_P(WebSocketContractTest,
ATerminalTransitionFiresTheParkedSendBeforeTheParkedReceive,
ATerminalTransitionCompletesALoneParkedCoroutineSend,
ATerminalTransitionCompletesALoneParkedReceive,
ASecondSendClassOperationRefusesWhileOneIsParked,
ASecondReceiveClassOperationRefusesWhileOneIsParked);
Expand Down
41 changes: 41 additions & 0 deletions runtime/tests/eventstream/jsonrpc_stream_socket_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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<int>(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<http::WebSocket> 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<http::WebSocket> peer_;
std::shared_ptr<JsonRpcStreamSocket> 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
Loading