/* streamer-tools OBS Camera Plugin - minimal loopback HTTP server for tests Copyright (C) 2026 CyberCoveLLC Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0 */ #pragma once // A single-threaded, one-request-at-a-time HTTP/1.1 server on 127.0.0.1, used // to exercise the *real* platform HTTP backend (libcurl on Linux/macOS, // WinHTTP on Windows) rather than only a fake. The handler returns raw bytes, // so tests can serve deliberately malformed responses and half-closed // connections -- the cases that must not hang or crash OBS. // // Plain HTTP only: a TLS listener would need a certificate and would test // libcurl/SChannel rather than this plugin. #include #include #include #include #include #include #include #ifdef _WIN32 #include #include using st_socket_t = SOCKET; #define ST_INVALID_SOCKET INVALID_SOCKET #define ST_CLOSE_SOCKET closesocket #else #include #include #include #include #include using st_socket_t = int; #define ST_INVALID_SOCKET (-1) #define ST_CLOSE_SOCKET ::close #endif namespace sttest { /// Returns raw response bytes for a received raw request. Returning an empty /// string means "close the connection without replying". using LoopbackHandler = std::function; class LoopbackServer { public: explicit LoopbackServer(LoopbackHandler handler) : handler_(std::move(handler)) { #ifdef _WIN32 WSADATA wsa; WSAStartup(MAKEWORD(2, 2), &wsa); #endif listen_ = ::socket(AF_INET, SOCK_STREAM, 0); if (listen_ == ST_INVALID_SOCKET) return; int reuse = 1; ::setsockopt(listen_, SOL_SOCKET, SO_REUSEADDR, reinterpret_cast(&reuse), sizeof(reuse)); sockaddr_in addr{}; addr.sin_family = AF_INET; addr.sin_addr.s_addr = htonl(INADDR_LOOPBACK); addr.sin_port = 0; // let the OS pick a free port if (::bind(listen_, reinterpret_cast(&addr), sizeof(addr)) != 0) { ST_CLOSE_SOCKET(listen_); listen_ = ST_INVALID_SOCKET; return; } if (::listen(listen_, 4) != 0) { ST_CLOSE_SOCKET(listen_); listen_ = ST_INVALID_SOCKET; return; } sockaddr_in bound{}; #ifdef _WIN32 int len = sizeof(bound); #else socklen_t len = sizeof(bound); #endif if (::getsockname(listen_, reinterpret_cast(&bound), &len) != 0) { ST_CLOSE_SOCKET(listen_); listen_ = ST_INVALID_SOCKET; return; } port_ = ntohs(bound.sin_port); thread_ = std::thread([this] { run(); }); } ~LoopbackServer() { stop_.store(true); if (thread_.joinable()) thread_.join(); if (listen_ != ST_INVALID_SOCKET) ST_CLOSE_SOCKET(listen_); #ifdef _WIN32 WSACleanup(); #endif } LoopbackServer(const LoopbackServer &) = delete; LoopbackServer &operator=(const LoopbackServer &) = delete; bool valid() const { return listen_ != ST_INVALID_SOCKET; } int port() const { return port_; } std::string baseUrl() const { return "http://127.0.0.1:" + std::to_string(port_); } int requestCount() const { return requests_.load(); } /// The most recent raw request, for asserting on method/path/body. std::string lastRequest() const { std::lock_guard guard(mutex_); return last_request_; } private: void run() { while (!stop_.load()) { // select() with a short timeout rather than a blocking accept(), // so the destructor's stop flag is honoured promptly on every // platform (closing a socket another thread is blocked in // accept() on is not portable). fd_set readable; FD_ZERO(&readable); FD_SET(listen_, &readable); timeval tv{}; tv.tv_sec = 0; tv.tv_usec = 50000; // 50ms const int ready = ::select(static_cast(listen_) + 1, &readable, nullptr, nullptr, &tv); if (ready <= 0) continue; st_socket_t client = ::accept(listen_, nullptr, nullptr); if (client == ST_INVALID_SOCKET) continue; const std::string request = readRequest(client); { std::lock_guard guard(mutex_); last_request_ = request; } requests_.fetch_add(1); const std::string response = handler_ ? handler_(request) : std::string(); if (!response.empty()) sendAll(client, response); ST_CLOSE_SOCKET(client); } } static std::string readRequest(st_socket_t client) { std::string data; char buffer[4096]; std::size_t header_end = std::string::npos; long content_length = 0; for (;;) { #ifdef _WIN32 const int n = ::recv(client, buffer, static_cast(sizeof(buffer)), 0); #else const ssize_t n = ::recv(client, buffer, sizeof(buffer), 0); #endif if (n <= 0) break; data.append(buffer, static_cast(n)); if (header_end == std::string::npos) { header_end = data.find("\r\n\r\n"); if (header_end != std::string::npos) content_length = parseContentLength(data.substr(0, header_end)); } if (header_end != std::string::npos && data.size() >= header_end + 4 + static_cast(content_length)) break; } return data; } static long parseContentLength(const std::string &headers) { std::string lower; lower.reserve(headers.size()); for (char c : headers) lower.push_back(static_cast(c >= 'A' && c <= 'Z' ? c + 32 : c)); const std::size_t at = lower.find("content-length:"); if (at == std::string::npos) return 0; return std::strtol(headers.c_str() + at + 15, nullptr, 10); } static void sendAll(st_socket_t client, const std::string &data) { std::size_t sent = 0; while (sent < data.size()) { #ifdef _WIN32 const int n = ::send(client, data.data() + sent, static_cast(data.size() - sent), 0); #else const ssize_t n = ::send(client, data.data() + sent, data.size() - sent, 0); #endif if (n <= 0) return; sent += static_cast(n); } } LoopbackHandler handler_; st_socket_t listen_ = ST_INVALID_SOCKET; int port_ = 0; std::thread thread_; std::atomic stop_{false}; std::atomic requests_{0}; mutable std::mutex mutex_; std::string last_request_; }; /// Build a well-formed HTTP/1.1 response with an explicit Content-Length and /// Connection: close, so the client never waits for keep-alive reuse. inline std::string httpResponse(int status, const std::string &reason, const std::string &body, const std::string &content_type = "application/json") { return "HTTP/1.1 " + std::to_string(status) + " " + reason + "\r\n" + "Content-Type: " + content_type + "\r\n" + "Content-Length: " + std::to_string(body.size()) + "\r\n" + "Connection: close\r\n\r\n" + body; } } // namespace sttest