910 lines
38 KiB
C++
910 lines
38 KiB
C++
/** @file WebserverLive.test.cpp
|
|
*
|
|
* @author Roland Conybeare, Sep 2026
|
|
*
|
|
* A started Webserver on an OS-assigned port, and a real websocket client
|
|
* (WsTestClient). Covers what only runs with libwebsockets: session
|
|
* bookkeeping, the outbound queue, the service-thread wakeups. See
|
|
* .xo-backlog/xo-websock/issues/09.
|
|
*
|
|
* Kept out of utest.websock, which stays socket-free.
|
|
*
|
|
* Expectations are OBSERVED, never predicted.
|
|
**/
|
|
|
|
#include "WsTestClient.hpp"
|
|
#include "WebsockUtestAppcx.hpp"
|
|
#include "xo/websock/Webserver.hpp"
|
|
#include "xo/websock/WebsocketSink.hpp"
|
|
#include <xo/printjson/PrintJsonSingleton.hpp>
|
|
#include <xo/reflect/Reflect.hpp>
|
|
#include <catch2/catch.hpp>
|
|
#include <json/json.h>
|
|
#include <arpa/inet.h>
|
|
#include <netinet/in.h>
|
|
#include <sys/socket.h>
|
|
#include <unistd.h>
|
|
#include <chrono>
|
|
#include <cctype>
|
|
#include <condition_variable>
|
|
#include <cstdlib>
|
|
#include <fstream>
|
|
#include <functional>
|
|
#include <future>
|
|
#include <memory>
|
|
#include <mutex>
|
|
#include <sstream>
|
|
#include <string>
|
|
#include <thread>
|
|
#include <vector>
|
|
|
|
namespace xo {
|
|
using xo::web::Webserver;
|
|
using xo::web::WebserverConfig;
|
|
using xo::web::StreamEndpointDescr;
|
|
using xo::web::HttpEndpointDescr;
|
|
using xo::web::HttpRequest;
|
|
using xo::web::HttpResponse;
|
|
using xo::web::HttpStatus;
|
|
using xo::web::WebsocketSink;
|
|
using xo::web::StreamReceiver;
|
|
using xo::json::PrintJsonSingleton;
|
|
using xo::reflect::Reflect;
|
|
using xo::fn::CallbackId;
|
|
|
|
namespace ut {
|
|
namespace {
|
|
using namespace std::chrono_literals;
|
|
|
|
/* generous: these only bound a failure, never pace a pass */
|
|
constexpr auto c_timeout = 5s;
|
|
|
|
Json::Value parse(std::string const & text) {
|
|
Json::Value root;
|
|
JSONCPP_STRING err;
|
|
std::unique_ptr<Json::CharReader> rd(Json::CharReaderBuilder().newCharReader());
|
|
|
|
bool ok = rd->parse(text.data(), text.data() + text.size(), &root, &err);
|
|
|
|
INFO("text: " << text << " err: " << err);
|
|
REQUIRE(ok);
|
|
|
|
return root;
|
|
}
|
|
|
|
/** what a stream endpoint's subscribe function was handed, as
|
|
* seen from the test's thread (the function runs on the
|
|
* webserver's)
|
|
**/
|
|
struct SinkBox {
|
|
std::mutex mutex_;
|
|
std::condition_variable cv_;
|
|
std::vector<rp<WebsocketSink>> sink_v_;
|
|
/* callback ids the unsubscribe function was handed */
|
|
std::vector<std::uint32_t> unsub_v_;
|
|
/* messages the receiver was handed */
|
|
std::vector<Json::Value> msg_v_;
|
|
|
|
CallbackId subscribe(rp<WebsocketSink> const & sink) {
|
|
std::lock_guard<std::mutex> lock(mutex_);
|
|
|
|
sink_v_.push_back(sink);
|
|
cv_.notify_all();
|
|
|
|
return CallbackId(static_cast<uint32_t>(sink_v_.size()));
|
|
}
|
|
|
|
void unsubscribe(CallbackId id) {
|
|
std::lock_guard<std::mutex> lock(mutex_);
|
|
|
|
unsub_v_.push_back(id.id());
|
|
cv_.notify_all();
|
|
}
|
|
|
|
void receive(rp<WebsocketSink> const & sink, Json::Value const & msg) {
|
|
{
|
|
std::lock_guard<std::mutex> lock(mutex_);
|
|
|
|
msg_v_.push_back(msg);
|
|
cv_.notify_all();
|
|
}
|
|
|
|
/* reply through the sink handed in: reaches exactly the
|
|
* sender's session. Rendered synchronously, so a local
|
|
* is fine
|
|
*/
|
|
int reply = msg["n"].asInt() * 10;
|
|
sink->notify_ev_tp(Reflect::make_tp(&reply));
|
|
}
|
|
|
|
/* true once at least n unsubscribes have run */
|
|
bool wait_unsubscribed(std::size_t n) {
|
|
std::unique_lock<std::mutex> lock(mutex_);
|
|
|
|
return cv_.wait_for(lock, c_timeout, [this, n] { return unsub_v_.size() >= n; });
|
|
}
|
|
|
|
std::size_t n_unsubscribed() {
|
|
std::lock_guard<std::mutex> lock(mutex_);
|
|
|
|
return unsub_v_.size();
|
|
}
|
|
|
|
/* the n'th sink (0-based) once it exists; null on timeout */
|
|
rp<WebsocketSink> wait_sink(std::size_t n) {
|
|
std::unique_lock<std::mutex> lock(mutex_);
|
|
|
|
if (!cv_.wait_for(lock, c_timeout, [this, n] { return sink_v_.size() > n; }))
|
|
return nullptr;
|
|
|
|
return sink_v_[n];
|
|
}
|
|
};
|
|
|
|
/** hands each message to a SinkBox, which replies **/
|
|
class BoxReceiver : public StreamReceiver {
|
|
public:
|
|
explicit BoxReceiver(std::shared_ptr<SinkBox> box) : box_{std::move(box)} {}
|
|
|
|
/* a StreamReceiver is SelfTagging; not reflected in full */
|
|
xo::reflect::TaggedRcptr self_tp() override { return Reflect::make_rctp(this); }
|
|
|
|
void receive(rp<WebsocketSink> const & sink, Json::Value const & msg) override {
|
|
box_->receive(sink, msg);
|
|
}
|
|
|
|
private:
|
|
std::shared_ptr<SinkBox> box_;
|
|
};
|
|
|
|
/** stream endpoint on @p pattern, reporting to @p box **/
|
|
StreamEndpointDescr box_descr(std::string pattern, std::shared_ptr<SinkBox> const & box) {
|
|
return StreamEndpointDescr(std::move(pattern),
|
|
[box](rp<WebsocketSink> const & sink) { return box->subscribe(sink); },
|
|
[box](CallbackId id) { box->unsubscribe(id); },
|
|
new BoxReceiver(box));
|
|
}
|
|
|
|
/** true once @p pred holds; polls, since the server's session
|
|
* bookkeeping runs on its own thread with nothing to wait on
|
|
**/
|
|
template <typename Pred>
|
|
bool wait_until(Pred pred) {
|
|
auto deadline = std::chrono::steady_clock::now() + c_timeout;
|
|
|
|
while (!pred()) {
|
|
if (std::chrono::steady_clock::now() > deadline)
|
|
return false;
|
|
std::this_thread::sleep_for(1ms);
|
|
}
|
|
|
|
return true;
|
|
}
|
|
|
|
/** one http GET: status line code, Content-Type, body **/
|
|
struct HttpReply {
|
|
int status_ = 0;
|
|
std::string content_type_;
|
|
std::string body_;
|
|
};
|
|
|
|
/** GET @p path from localhost:@p port, over HTTP/1.0 (so the
|
|
* server closes when done). status_ 0 if no reply
|
|
**/
|
|
HttpReply http_get(std::int32_t port, std::string const & path) {
|
|
HttpReply reply;
|
|
|
|
int fd = ::socket(AF_INET, SOCK_STREAM, 0);
|
|
if (fd < 0)
|
|
return reply;
|
|
|
|
/* bound a hung server */
|
|
timeval tv{5, 0};
|
|
::setsockopt(fd, SOL_SOCKET, SO_RCVTIMEO, &tv, sizeof(tv));
|
|
|
|
sockaddr_in addr{};
|
|
addr.sin_family = AF_INET;
|
|
addr.sin_port = htons(static_cast<std::uint16_t>(port));
|
|
addr.sin_addr.s_addr = htonl(INADDR_LOOPBACK);
|
|
|
|
if (::connect(fd, reinterpret_cast<sockaddr *>(&addr), sizeof(addr)) != 0) {
|
|
::close(fd);
|
|
return reply;
|
|
}
|
|
|
|
std::string req = "GET " + path + " HTTP/1.0\r\nHost: localhost\r\n\r\n";
|
|
::send(fd, req.data(), req.size(), 0);
|
|
|
|
std::string text;
|
|
char buf[4096];
|
|
for (ssize_t n; (n = ::recv(fd, buf, sizeof(buf), 0)) > 0; )
|
|
text.append(buf, n);
|
|
::close(fd);
|
|
|
|
auto hdr_end = text.find("\r\n\r\n");
|
|
if (hdr_end == std::string::npos)
|
|
return reply;
|
|
|
|
std::string head = text.substr(0, hdr_end);
|
|
reply.body_ = text.substr(hdr_end + 4);
|
|
|
|
/* "HTTP/1.x NNN ..." */
|
|
auto sp = head.find(' ');
|
|
if (sp != std::string::npos)
|
|
reply.status_ = std::atoi(head.c_str() + sp + 1);
|
|
|
|
/* header names are case-insensitive: lws sends lowercase */
|
|
std::string lower = head;
|
|
for (char & c : lower)
|
|
c = static_cast<char>(std::tolower(static_cast<unsigned char>(c)));
|
|
auto ct = lower.find("\r\ncontent-type:");
|
|
if (ct != std::string::npos) {
|
|
auto v = ct + std::string("\r\ncontent-type:").size();
|
|
auto e = head.find("\r\n", v);
|
|
reply.content_type_ = head.substr(v, e - v);
|
|
while (!reply.content_type_.empty() && reply.content_type_.front() == ' ')
|
|
reply.content_type_.erase(0, 1);
|
|
}
|
|
|
|
return reply;
|
|
}
|
|
|
|
/** a started webserver on an OS-assigned port **/
|
|
struct LiveServer {
|
|
LiveServer() {
|
|
websrv_ = Webserver::make(WebsockUtestAppcx::appcx().cx<S_websock_tag>(), WebserverConfig());
|
|
}
|
|
|
|
~LiveServer() {
|
|
websrv_->stop_webserver();
|
|
websrv_->join_webserver();
|
|
}
|
|
|
|
/* start; the port once listening, 0 on timeout */
|
|
std::int32_t start() {
|
|
websrv_->start_webserver();
|
|
|
|
auto deadline = std::chrono::steady_clock::now() + c_timeout;
|
|
|
|
/* start_webserver() returns before the service thread has
|
|
* bound the port; nothing to wait on but the port itself
|
|
*/
|
|
while (websrv_->listen_port() == 0) {
|
|
if (std::chrono::steady_clock::now() > deadline)
|
|
return 0;
|
|
std::this_thread::sleep_for(1ms);
|
|
}
|
|
|
|
return websrv_->listen_port();
|
|
}
|
|
|
|
rp<Webserver> websrv_;
|
|
};
|
|
}
|
|
|
|
TEST_CASE("live-subscribe-then-frames", "[websock][live]")
|
|
{
|
|
auto box = std::make_shared<SinkBox>();
|
|
|
|
LiveServer srv;
|
|
srv.websrv_->register_stream_endpoint
|
|
(StreamEndpointDescr("/fw",
|
|
[box](rp<WebsocketSink> const & sink) { return box->subscribe(sink); },
|
|
[](CallbackId) {}));
|
|
|
|
std::int32_t port = srv.start();
|
|
REQUIRE(port > 0);
|
|
|
|
WsTestClient client(port);
|
|
REQUIRE(client.wait_connected(c_timeout));
|
|
|
|
client.send(R"({"cmd": "subscribe", "stream": "/fw"})");
|
|
|
|
REQUIRE(client.wait_received(1, c_timeout));
|
|
Json::Value subscribed = parse(client.received()[0]);
|
|
REQUIRE(subscribed["cmd"].asString() == "subscribed");
|
|
REQUIRE(subscribed["stream"].asString() == "/fw");
|
|
std::uint32_t sub_id = subscribed["sub_id"].asUInt();
|
|
|
|
/* the application's side: events pushed from its own thread, as
|
|
* a reactor source would
|
|
*/
|
|
rp<WebsocketSink> sink = box->wait_sink(0);
|
|
REQUIRE(sink);
|
|
|
|
int ev1 = 101;
|
|
int ev2 = 102;
|
|
sink->notify_ev_tp(Reflect::make_tp(&ev1));
|
|
sink->notify_ev_tp(Reflect::make_tp(&ev2));
|
|
|
|
REQUIRE(client.wait_received(3, c_timeout));
|
|
std::vector<std::string> msg_v = client.received();
|
|
|
|
for (std::size_t i = 1; i < 3; ++i) {
|
|
Json::Value env = parse(msg_v[i]);
|
|
|
|
INFO("frame: " << msg_v[i]);
|
|
REQUIRE(env["stream"].asString() == "/fw");
|
|
REQUIRE(env["sub_id"].asUInt() == sub_id);
|
|
REQUIRE(env["seq"].asInt() == static_cast<int>(i - 1));
|
|
REQUIRE(env["event"].asInt() == 100 + static_cast<int>(i));
|
|
}
|
|
|
|
REQUIRE(client.close(c_timeout));
|
|
}
|
|
|
|
TEST_CASE("live-http-status-and-content-type", "[websock][live][http]")
|
|
{
|
|
/* what a handler chooses is what the client gets: status line and
|
|
* Content-Type header; and the server's own answers -- no
|
|
* endpoint, no match, a handler's exception
|
|
*/
|
|
LiveServer srv;
|
|
srv.websrv_->register_http_endpoint
|
|
(HttpEndpointDescr("/j",
|
|
[](HttpRequest const &) { return HttpResponse::json("{\"a\": 1}"); }));
|
|
srv.websrv_->register_http_endpoint
|
|
(HttpEndpointDescr("/h/${a}",
|
|
[](HttpRequest const & req) {
|
|
return HttpResponse::html("<p>" + std::string(req.var("a")) + "</p>");
|
|
}));
|
|
srv.websrv_->register_http_endpoint
|
|
(HttpEndpointDescr("/boom",
|
|
[](HttpRequest const &) -> HttpResponse {
|
|
throw std::runtime_error("kaboom");
|
|
}));
|
|
|
|
std::int32_t port = srv.start();
|
|
REQUIRE(port > 0);
|
|
|
|
HttpReply j = http_get(port, "/dyn/j");
|
|
REQUIRE(j.status_ == 200);
|
|
REQUIRE(j.content_type_ == "application/json");
|
|
REQUIRE(j.body_ == "{\"a\": 1}");
|
|
|
|
HttpReply h = http_get(port, "/dyn/h/x-1.hpp");
|
|
REQUIRE(h.status_ == 200);
|
|
REQUIRE(h.content_type_ == "text/html; charset=utf-8");
|
|
REQUIRE(h.body_ == "<p>x-1.hpp</p>");
|
|
|
|
/* found by stem, pattern not matched */
|
|
HttpReply nm = http_get(port, "/dyn/h/x/y");
|
|
REQUIRE(nm.status_ == 404);
|
|
REQUIRE(nm.content_type_ == "text/html; charset=utf-8");
|
|
|
|
/* no endpoint at all */
|
|
HttpReply ne = http_get(port, "/dyn/nothing-here");
|
|
REQUIRE(ne.status_ == 404);
|
|
|
|
/* a handler's exception: 500, and the server lives on */
|
|
HttpReply b = http_get(port, "/dyn/boom");
|
|
REQUIRE(b.status_ == 500);
|
|
REQUIRE(b.body_.find("kaboom") != std::string::npos);
|
|
|
|
REQUIRE(http_get(port, "/dyn/j").status_ == 200);
|
|
}
|
|
|
|
TEST_CASE("live-a-server-that-cannot-start-still-joins", "[websock][live]")
|
|
{
|
|
/* lws_create_context fails when the port is taken; run() used to
|
|
* return without reporting stopped, so join_webserver() hung
|
|
*/
|
|
LiveServer first;
|
|
std::int32_t port = first.start();
|
|
REQUIRE(port > 0);
|
|
|
|
rp<Webserver> second = Webserver::make(WebsockUtestAppcx::appcx().cx<S_websock_tag>(),
|
|
WebserverConfig(port, false, false, false));
|
|
second->start_webserver();
|
|
|
|
/* join on a DETACHED thread, so a regression fails instead of
|
|
* hanging: a std::async future would block in its destructor.
|
|
* The thread holds its own rp, so on regression neither it nor
|
|
* ~WebserverImpl (which also joins) runs on the test's thread.
|
|
*/
|
|
std::promise<void> joined_promise;
|
|
std::future<void> joined = joined_promise.get_future();
|
|
|
|
std::thread([second, p = std::move(joined_promise)]() mutable
|
|
{
|
|
second->join_webserver();
|
|
p.set_value();
|
|
}).detach();
|
|
|
|
REQUIRE(joined.wait_for(c_timeout) == std::future_status::ready);
|
|
REQUIRE(second->listen_port() == 0);
|
|
REQUIRE(second->state() == xo::web::Runstate::stopped);
|
|
}
|
|
|
|
TEST_CASE("live-send-reaches-the-receiver-and-its-reply-comes-back", "[websock][live]")
|
|
{
|
|
auto box = std::make_shared<SinkBox>();
|
|
|
|
LiveServer srv;
|
|
srv.websrv_->register_stream_endpoint(box_descr("/fw", box));
|
|
|
|
std::int32_t port = srv.start();
|
|
REQUIRE(port > 0);
|
|
|
|
WsTestClient client(port);
|
|
REQUIRE(client.wait_connected(c_timeout));
|
|
|
|
client.send(R"({"cmd": "subscribe", "stream": "/fw"})");
|
|
REQUIRE(client.wait_received(1, c_timeout));
|
|
std::uint32_t sub_id = parse(client.received()[0])["sub_id"].asUInt();
|
|
|
|
client.send(std::string(R"({"cmd": "send", "sub_id": )")
|
|
+ std::to_string(sub_id)
|
|
+ R"(, "msg": {"op": "step", "n": 4}})");
|
|
|
|
/* the receiver's reply, as a frame of this subscription */
|
|
REQUIRE(client.wait_received(2, c_timeout));
|
|
Json::Value frame = parse(client.received()[1]);
|
|
|
|
REQUIRE(frame["sub_id"].asUInt() == sub_id);
|
|
REQUIRE(frame["seq"].asInt() == 0);
|
|
REQUIRE(frame["event"].asInt() == 40);
|
|
|
|
/* and the receiver saw the msg exactly as sent */
|
|
std::lock_guard<std::mutex> lock(box->mutex_);
|
|
REQUIRE(box->msg_v_.size() == 1);
|
|
REQUIRE(box->msg_v_[0]["op"].asString() == "step");
|
|
REQUIRE(box->msg_v_[0]["n"].asInt() == 4);
|
|
}
|
|
|
|
TEST_CASE("live-unregister-ends-a-live-subscription", "[websock][live]")
|
|
{
|
|
/* issue 07 over a socket: the unregister queue, the service
|
|
* thread's wakeup and drain, and the reply reaching the client
|
|
*/
|
|
auto box = std::make_shared<SinkBox>();
|
|
|
|
LiveServer srv;
|
|
srv.websrv_->register_stream_endpoint(box_descr("/fw/${id}", box));
|
|
|
|
std::int32_t port = srv.start();
|
|
REQUIRE(port > 0);
|
|
|
|
WsTestClient client(port);
|
|
REQUIRE(client.wait_connected(c_timeout));
|
|
|
|
client.send(R"({"cmd": "subscribe", "stream": "/fw/1"})");
|
|
client.send(R"({"cmd": "subscribe", "stream": "/fw/2"})");
|
|
REQUIRE(client.wait_received(2, c_timeout));
|
|
REQUIRE(box->wait_sink(1));
|
|
|
|
/* from the test's thread, as python would */
|
|
REQUIRE(srv.websrv_->unregister_stream_endpoint("/fw/${id}"));
|
|
|
|
REQUIRE(client.wait_received(4, c_timeout));
|
|
std::vector<std::string> msg_v = client.received();
|
|
|
|
for (std::size_t i = 2; i < 4; ++i) {
|
|
Json::Value r = parse(msg_v[i]);
|
|
|
|
INFO("reply: " << msg_v[i]);
|
|
REQUIRE(r["cmd"].asString() == "unsubscribed");
|
|
REQUIRE(r["sub_id"].asUInt() == i - 2);
|
|
REQUIRE(r["reason"].asString() == "endpoint removed");
|
|
}
|
|
|
|
/* the endpoint's unsubscribe ran once per subscription -- a real
|
|
* source would stop pushing here
|
|
*/
|
|
REQUIRE(box->wait_unsubscribed(2));
|
|
|
|
/* a new subscribe finds nothing */
|
|
client.send(R"({"cmd": "subscribe", "stream": "/fw/3"})");
|
|
REQUIRE(client.wait_received(5, c_timeout));
|
|
REQUIRE(parse(client.received()[4])["error"].asString() == "unknown stream");
|
|
|
|
/* and nothing ran twice */
|
|
REQUIRE(box->n_unsubscribed() == 2);
|
|
}
|
|
|
|
TEST_CASE("live-a-sink-kept-past-its-session-reaches-no-one", "[websock][live]")
|
|
{
|
|
/* issues 05 and 08: the kept sink's sender is closed with its
|
|
* session, and the session id is never reused, so a later client
|
|
* gets nothing from it
|
|
*/
|
|
auto box = std::make_shared<SinkBox>();
|
|
|
|
LiveServer srv;
|
|
srv.websrv_->register_stream_endpoint(box_descr("/fw", box));
|
|
|
|
std::int32_t port = srv.start();
|
|
REQUIRE(port > 0);
|
|
|
|
rp<WebsocketSink> kept;
|
|
|
|
{
|
|
WsTestClient first(port);
|
|
REQUIRE(first.wait_connected(c_timeout));
|
|
|
|
first.send(R"({"cmd": "subscribe", "stream": "/fw"})");
|
|
REQUIRE(first.wait_received(1, c_timeout));
|
|
|
|
kept = box->wait_sink(0);
|
|
REQUIRE(kept);
|
|
|
|
REQUIRE(first.close(c_timeout));
|
|
}
|
|
|
|
/* the server has handled the close once the session's
|
|
* subscriptions are unsubscribed (notify_ws_session_close)
|
|
*/
|
|
REQUIRE(box->wait_unsubscribed(1));
|
|
|
|
WsTestClient second(port);
|
|
REQUIRE(second.wait_connected(c_timeout));
|
|
|
|
/* the application pushes to the sink it kept */
|
|
int stale = 666;
|
|
kept->notify_ev_tp(Reflect::make_tp(&stale));
|
|
|
|
/* barrier: the second client's own traffic. Had the stale frame
|
|
* been delivered it would precede this reply -- a closed sender
|
|
* drops synchronously, before this subscribe is even sent
|
|
*/
|
|
second.send(R"({"cmd": "subscribe", "stream": "/fw"})");
|
|
REQUIRE(second.wait_received(1, c_timeout));
|
|
|
|
std::vector<std::string> msg_v = second.received();
|
|
|
|
INFO("first message: " << msg_v[0]);
|
|
REQUIRE(parse(msg_v[0])["cmd"].asString() == "subscribed");
|
|
for (auto const & m : msg_v)
|
|
REQUIRE(m.find("666") == std::string::npos);
|
|
}
|
|
|
|
TEST_CASE("live-sessions-lists-each-connection", "[websock][live]")
|
|
{
|
|
auto box = std::make_shared<SinkBox>();
|
|
|
|
LiveServer srv;
|
|
srv.websrv_->register_stream_endpoint(box_descr("/fw", box));
|
|
|
|
std::int32_t port = srv.start();
|
|
REQUIRE(port > 0);
|
|
|
|
/* the server as its json printer shows it */
|
|
auto server_json = [&srv] {
|
|
Webserver * server = srv.websrv_.get();
|
|
std::stringstream ss;
|
|
PrintJsonSingleton::instance()->print(server, &ss);
|
|
return parse(ss.str());
|
|
};
|
|
auto n_session = [&] { return server_json()["sessions"].size(); };
|
|
|
|
REQUIRE(n_session() == 0);
|
|
|
|
auto first = std::make_unique<WsTestClient>(port);
|
|
REQUIRE(first->wait_connected(c_timeout));
|
|
REQUIRE(wait_until([&] { return n_session() == 1; }));
|
|
|
|
WsTestClient second(port);
|
|
REQUIRE(second.wait_connected(c_timeout));
|
|
REQUIRE(wait_until([&] { return n_session() == 2; }));
|
|
|
|
/* the second subscribes; the reply precedes the subscribe
|
|
* function, so wait on the sink, not the reply
|
|
*/
|
|
second.send(R"({"cmd": "subscribe", "stream": "/fw"})");
|
|
REQUIRE(second.wait_received(1, c_timeout));
|
|
REQUIRE(box->wait_sink(0));
|
|
|
|
Json::Value const root = server_json();
|
|
Json::Value const & v = root["sessions"];
|
|
|
|
INFO("json: " << root.toStyledString());
|
|
|
|
/* by id, in connection order; distinct */
|
|
REQUIRE(v.size() == 2);
|
|
REQUIRE(v[0]["_name_"].asString() == "WsSession");
|
|
REQUIRE(v[0]["_canonical_type_"].asString() == "xo::web::WebsocketSessionRecd");
|
|
REQUIRE(v[0]["session_id"].asUInt64() < v[1]["session_id"].asUInt64());
|
|
REQUIRE(v[0]["_id_"].asInt() != v[1]["_id_"].asInt());
|
|
|
|
/* each session's sender, in full: open; its session's id; held
|
|
* by the session record and the router, plus one per sink
|
|
*/
|
|
for (Json::ArrayIndex k = 0; k < 2; ++k) {
|
|
Json::Value const & sender = v[k]["sender"];
|
|
|
|
REQUIRE(sender["_name_"].asString() == "WsSessionSender");
|
|
|
|
/* its chosen C++ members: target_ a ref to the server */
|
|
Json::Value const & smem = sender["_members_"];
|
|
/* reflected members first (session_id_, open_), then the rest */
|
|
REQUIRE(smem[0]["_name_"].asString() == "session_id_");
|
|
REQUIRE(smem[0]["_value_"].asUInt64() == v[k]["session_id"].asUInt64());
|
|
/* a std::atomic<bool>: its load()ed value (xo-reflect#04) */
|
|
REQUIRE(smem[1]["_name_"].asString() == "open_");
|
|
REQUIRE(smem[1]["_canonical_type_"].asString() == "std::atomic<bool>");
|
|
REQUIRE(smem[1]["_value_"].asBool());
|
|
REQUIRE(smem[2]["_name_"].asString() == "target_");
|
|
REQUIRE(smem[2]["_value_"]["_ref_"].asInt() == root["_id_"].asInt());
|
|
/* a template: its arguments follow */
|
|
REQUIRE(sender["_canonical_type_"].asString().starts_with("xo::web::WsSessionSender<"));
|
|
REQUIRE(sender["_short_type_"].asString().starts_with("WsSessionSender<"));
|
|
REQUIRE(sender["open"].asBool());
|
|
REQUIRE(sender["session_id"].asUInt64() == v[k]["session_id"].asUInt64());
|
|
}
|
|
REQUIRE(v[0]["sender"]["refcount"].asUInt() == 2);
|
|
REQUIRE(v[1]["sender"]["refcount"].asUInt() == 3);
|
|
|
|
REQUIRE(v[0]["subscriptions"].empty());
|
|
REQUIRE(v[1]["subscriptions"].size() == 1);
|
|
REQUIRE(v[1]["subscriptions"][0]["stream"].asString() == "/fw");
|
|
|
|
/* each session's chosen C++ members (.xo-backlog/xo-websock/issues/13):
|
|
* sender_ a ref to the sender printed in full above
|
|
*/
|
|
/* a "_members_" entry by name: robust to member order */
|
|
auto entry_of = [](Json::Value const & mem, std::string const & name) {
|
|
for (Json::Value const & m : mem)
|
|
if (m["_name_"].asString() == name)
|
|
return m;
|
|
return Json::Value();
|
|
};
|
|
|
|
for (Json::ArrayIndex k = 0; k < 2; ++k) {
|
|
Json::Value const & mem = v[k]["_members_"];
|
|
|
|
std::vector<std::string> names;
|
|
for (Json::Value const & m : mem)
|
|
names.push_back(m["_name_"].asString());
|
|
|
|
/* reflected members first -- router_, then output_buf_,
|
|
* last_msg_seq_ and outbound_q_ under the session's mutex
|
|
* (xo-websock#15) -- then sender_
|
|
*/
|
|
REQUIRE(names == std::vector<std::string>{"router_", "output_buf_", "last_msg_seq_",
|
|
"outbound_q_", "sender_"});
|
|
/* output_buf_ placed here (owning): its OutputBuffer, or null
|
|
* between writes -- not a ref
|
|
*/
|
|
{
|
|
Json::Value const ob = entry_of(mem, "output_buf_")["_value_"]; /* a copy: entry_of returns one */
|
|
REQUIRE((ob.isNull() || (ob.isObject() && ob.isMember("_id_"))));
|
|
}
|
|
REQUIRE(entry_of(mem, "sender_")["_metatype_"].asString() == "pointer");
|
|
REQUIRE(entry_of(mem, "sender_")["_value_"]["_ref_"].asInt() == v[k]["sender"]["_id_"].asInt());
|
|
REQUIRE(entry_of(mem, "router_")["_canonical_type_"].asString() == "xo::web::WsSessionRouter");
|
|
/* the queued messages themselves (xo-reflect#04): none, settled */
|
|
REQUIRE(entry_of(mem, "outbound_q_")["_metatype_"].asString() == "vector");
|
|
REQUIRE(entry_of(mem, "outbound_q_")["_value_"].isArray());
|
|
REQUIRE(entry_of(mem, "outbound_q_")["_value_"].size() == 0);
|
|
|
|
/* the router, nested: its own members. sender_ the same
|
|
* sender (by most-derived address); subscription_v_ refs to
|
|
* exactly the subscriptions printed under the session
|
|
*/
|
|
Json::Value const router = entry_of(mem, "router_"); /* a copy: entry_of returns one */
|
|
Json::Value const & rmem = router["_value_"]["_members_"];
|
|
|
|
std::vector<std::string> rnames;
|
|
for (Json::Value const & m : rmem)
|
|
rnames.push_back(m["_name_"].asString());
|
|
|
|
/* reflected members first (readjson_), then the rest */
|
|
REQUIRE(rnames == std::vector<std::string>{"readjson_", "url_router_", "sender_",
|
|
"pjson_", "subscription_v_"});
|
|
|
|
Json::Value const url_router = entry_of(rmem, "url_router_");
|
|
REQUIRE(url_router["_metatype_"].asString() == "pointer"); /* a reference */
|
|
/* ... to the server's url router, printed inside the server */
|
|
REQUIRE(url_router["_value_"]["_ref_"].asInt()
|
|
== entry_of(root["_members_"], "url_router_")["_value_"]["_id_"].asInt());
|
|
REQUIRE(entry_of(rmem, "sender_")["_value_"]["_ref_"].asInt() == v[k]["sender"]["_id_"].asInt());
|
|
|
|
/* a unique_ptr to jsoncpp's reader, reflected with no members:
|
|
* present, so an empty object (xo-reflect#04)
|
|
*/
|
|
Json::Value const readjson = entry_of(rmem, "readjson_");
|
|
REQUIRE(readjson["_metatype_"].asString() == "pointer");
|
|
REQUIRE(readjson["_value_"]["_name_"].asString() == "CharReader");
|
|
REQUIRE(readjson["_value_"]["_members_"].empty());
|
|
|
|
Json::Value const subscription_v = entry_of(rmem, "subscription_v_");
|
|
Json::Value const & slots = subscription_v["_value_"];
|
|
Json::Value const & subs = v[k]["subscriptions"];
|
|
|
|
REQUIRE(subscription_v["_metatype_"].asString() == "vector");
|
|
REQUIRE(slots.size() == subs.size());
|
|
for (Json::ArrayIndex i = 0; i < subs.size(); ++i)
|
|
REQUIRE(slots[i]["_ref_"].asInt() == subs[i]["_id_"].asInt());
|
|
}
|
|
|
|
/* the server's session table: session id -> a ref to that
|
|
* session, printed in full in "sessions"
|
|
*/
|
|
{
|
|
Json::Value const * st = nullptr;
|
|
for (Json::Value const & m : root["_members_"])
|
|
if (m["_name_"].asString() == "session_table_")
|
|
st = &m["_value_"];
|
|
REQUIRE(st);
|
|
|
|
Json::Value const & smap = (*st)["_members_"][1]["_value_"];
|
|
REQUIRE(smap.size() == v.size());
|
|
for (Json::ArrayIndex k = 0; k < v.size(); ++k) {
|
|
std::string const sid = std::to_string(v[k]["session_id"].asUInt64());
|
|
REQUIRE(smap[sid]["_ref_"].asInt() == v[k]["_id_"].asInt());
|
|
}
|
|
REQUIRE((*st)["_members_"][0]["_value_"].asUInt64()
|
|
> v[v.size() - 1]["session_id"].asUInt64());
|
|
}
|
|
|
|
/* nothing a printer opted in to is unprintable */
|
|
{
|
|
std::string const text = root.toStyledString();
|
|
INFO(text);
|
|
REQUIRE(text.find("\"_error_\"") == std::string::npos);
|
|
}
|
|
|
|
std::uint64_t second_id = v[1]["session_id"].asUInt64();
|
|
|
|
/* the /fw endpoint is held by the router's map and by the one
|
|
* subscription served: refcount 2
|
|
*/
|
|
Json::Value const & eps = root["endpoints"];
|
|
REQUIRE(eps.size() == 1);
|
|
REQUIRE(eps[0]["pattern"].asString() == "/fw");
|
|
REQUIRE(eps[0]["refcount"].asUInt() == 2);
|
|
|
|
/* the edges, joined by id: the subscription's endpoint is THAT
|
|
* endpoint; its sink's sender is ITS session's sender. The sink
|
|
* is held by the router's slot and by the endpoint's subscriber
|
|
* (the test's SinkBox)
|
|
*/
|
|
Json::Value const & sub = v[1]["subscriptions"][0];
|
|
|
|
REQUIRE(sub["endpoint"]["_ref_"].asInt() == eps[0]["_id_"].asInt());
|
|
REQUIRE(sub["sink"]["sender"]["_ref_"].asInt() == v[1]["sender"]["_id_"].asInt());
|
|
REQUIRE(sub["sink"]["refcount"].asUInt() == 2);
|
|
|
|
/* a closed session leaves the listing */
|
|
REQUIRE(first->close(c_timeout));
|
|
first.reset();
|
|
|
|
REQUIRE(wait_until([&] { return n_session() == 1; }));
|
|
REQUIRE(server_json()["sessions"][0]["session_id"].asUInt64() == second_id);
|
|
}
|
|
|
|
namespace {
|
|
/** values in a server snapshot that differ from run to run:
|
|
* replaced by "<redacted>" so the rest can be compared exactly.
|
|
* Top-level keys by name; _members_ entries by member name
|
|
**/
|
|
void redact(Json::Value & v) {
|
|
static const std::vector<std::string> c_keys = {
|
|
"listen_port",
|
|
"_address_", /* in "_unplaced_": an address */
|
|
};
|
|
static const std::vector<std::string> c_members = {
|
|
"listen_port_", /* the port the kernel chose */
|
|
"port_", /* the config's copy: 0, or the port chosen */
|
|
"mount_origin_", /* a path */
|
|
"output_buf_", /* its id, or null between writes */
|
|
};
|
|
|
|
if (v.isArray()) {
|
|
for (Json::Value & x : v)
|
|
redact(x);
|
|
} else if (v.isObject()) {
|
|
for (std::string const & k : c_keys)
|
|
if (v.isMember(k))
|
|
v[k] = "<redacted>";
|
|
if (v.isMember("_value_") && v.isMember("_name_")) {
|
|
std::string const name = v["_name_"].asString();
|
|
for (std::string const & m : c_members)
|
|
if (name == m)
|
|
v["_value_"] = "<redacted>";
|
|
}
|
|
for (std::string const & k : v.getMemberNames())
|
|
redact(v[k]);
|
|
}
|
|
}
|
|
|
|
/** @p text, styled for a line-per-field diff **/
|
|
std::string styled(Json::Value const & v) {
|
|
Json::StreamWriterBuilder wb;
|
|
wb["indentation"] = " ";
|
|
return Json::writeString(wb, v) + "\n";
|
|
}
|
|
|
|
/** the checked-in golden file @p name, under utest/golden/ **/
|
|
std::string golden_path(std::string const & name) {
|
|
return std::string(XO_WEBSOCK_UTEST_SOURCE_DIR) + "/golden/" + name;
|
|
}
|
|
|
|
std::string read_file(std::string const & path) {
|
|
std::ifstream in(path);
|
|
std::stringstream ss;
|
|
ss << in.rdbuf();
|
|
return ss.str();
|
|
}
|
|
}
|
|
|
|
TEST_CASE("live-server-snapshot-matches-golden", "[websock][live][json][golden]")
|
|
{
|
|
/* the whole server, as introspection prints it, against a
|
|
* checked-in expectation: a printer change shows as a diff to
|
|
* review (.xo-backlog/xo-websock/issues/14). Every websock
|
|
* printer appears: server, config, url router, endpoints (one
|
|
* with a receiver), session table, sessions, senders, routers,
|
|
* a subscription and its sink.
|
|
*
|
|
* To bless a deliberate change:
|
|
* XO_UPDATE_GOLDEN=1 utest.websock.live "[golden]"
|
|
* rewrites utest/golden/server-snapshot.json in the source tree
|
|
*/
|
|
auto box = std::make_shared<SinkBox>();
|
|
|
|
LiveServer srv;
|
|
srv.websrv_->register_stream_endpoint(box_descr("/fw", box));
|
|
srv.websrv_->register_http_endpoint
|
|
(HttpEndpointDescr("/status",
|
|
[](HttpRequest const &) { return HttpResponse::json("{}"); }));
|
|
|
|
std::int32_t port = srv.start();
|
|
REQUIRE(port > 0);
|
|
|
|
auto server_json = [&srv] {
|
|
Webserver * server = srv.websrv_.get();
|
|
std::stringstream ss;
|
|
PrintJsonSingleton::instance()->print(server, &ss);
|
|
return parse(ss.str());
|
|
};
|
|
|
|
/* two sessions, the second subscribed: wait on each step, so
|
|
* the snapshot is of a settled state
|
|
*/
|
|
WsTestClient first(port);
|
|
REQUIRE(first.wait_connected(c_timeout));
|
|
REQUIRE(wait_until([&] { return server_json()["sessions"].size() == 1; }));
|
|
|
|
WsTestClient second(port);
|
|
REQUIRE(second.wait_connected(c_timeout));
|
|
REQUIRE(wait_until([&] { return server_json()["sessions"].size() == 2; }));
|
|
|
|
second.send(R"({"cmd": "subscribe", "stream": "/fw"})");
|
|
REQUIRE(second.wait_received(1, c_timeout));
|
|
REQUIRE(box->wait_sink(0));
|
|
|
|
Json::Value snap = server_json();
|
|
|
|
/* every ref's target placed (.xo-backlog/xo-printjson/issues/08):
|
|
* no trailer
|
|
*/
|
|
REQUIRE(!snap.isMember("_unplaced_"));
|
|
|
|
redact(snap);
|
|
std::string const actual = styled(snap);
|
|
|
|
std::string const path = golden_path("server-snapshot.json");
|
|
|
|
if (char const * u = std::getenv("XO_UPDATE_GOLDEN"); u && *u && std::string(u) != "0") {
|
|
std::ofstream(path) << actual;
|
|
WARN("rewrote " << path);
|
|
return;
|
|
}
|
|
|
|
std::string const expected = read_file(path);
|
|
|
|
if (actual != expected) {
|
|
std::string const actual_path = "server-snapshot.actual.json";
|
|
std::ofstream(actual_path) << actual;
|
|
|
|
INFO("golden: " << path);
|
|
INFO("actual: " << actual_path << " (in the test's working directory)");
|
|
INFO("compare: diff -u <golden> <actual>; bless with XO_UPDATE_GOLDEN=1");
|
|
REQUIRE(actual == expected);
|
|
}
|
|
}
|
|
} /*namespace ut*/
|
|
} /*namespace xo*/
|
|
|
|
/* end WebserverLive.test.cpp */
|