#include "screencast/network/transport.h" #include #include #include #include #include #include #include #include #include #include #include #include namespace sc { namespace { constexpr int kInvalidSocket = -1; // A UDP datagram cannot exceed 64 KiB on IPv4; one buffer fits any RTP // packet the receiver will ever see. constexpr std::size_t kReceiveBufferSize = 65536; std::optional resolve_ipv4(const Endpoint& endpoint) { addrinfo hints{}; hints.ai_family = AF_INET; hints.ai_socktype = SOCK_DGRAM; addrinfo* result = nullptr; const std::string port = std::to_string(endpoint.port); const char* node = endpoint.address.empty() ? nullptr : endpoint.address.c_str(); if (getaddrinfo(node, port.c_str(), &hints, &result) != 0) { return std::nullopt; } std::optional address; if (result != nullptr && result->ai_family == AF_INET && result->ai_addrlen >= sizeof(sockaddr_in)) { sockaddr_in resolved{}; std::memcpy(&resolved, result->ai_addr, sizeof(sockaddr_in)); address = resolved; } freeaddrinfo(result); return address; } } // namespace class UdpRtpTransport final : public RtpTransport { public: UdpRtpTransport() = default; ~UdpRtpTransport() override { stop(); } UdpRtpTransport(const UdpRtpTransport&) = delete; UdpRtpTransport& operator=(const UdpRtpTransport&) = delete; bool start(const Endpoint& local_endpoint, ReceiveCallback on_receive) override { if (running_.load()) { return false; } socket_ = ::socket(AF_INET, SOCK_DGRAM, 0); if (socket_ < 0) { return false; } int reuse = 1; (void)::setsockopt(socket_, SOL_SOCKET, SO_REUSEADDR, &reuse, sizeof(reuse)); // A zero port skips binding: the OS picks the source port on send. if (local_endpoint.port != 0) { const std::optional address = resolve_ipv4(local_endpoint); if (!address.has_value()) { (void)::close(socket_); socket_ = kInvalidSocket; return false; } if (::bind(socket_, reinterpret_cast(&*address), sizeof(*address)) < 0) { (void)::close(socket_); socket_ = kInvalidSocket; return false; } } running_.store(true); receive_thread_ = std::jthread([this, callback = std::move(on_receive)]() mutable { receive_loop(callback); }); return true; } bool send(const RtpPacket& packet) override { if (!running_.load() || !has_peer_.load()) { return false; } const std::vector bytes = packet.serialize(); if (bytes.empty()) { return false; } sockaddr_in peer{}; { // Snapshot the peer so set_peer() can be called concurrently. std::lock_guard lock(peer_mutex_); peer = peer_; } const ssize_t sent = ::sendto(socket_, bytes.data(), bytes.size(), 0, reinterpret_cast(&peer), sizeof(peer)); return sent == static_cast(bytes.size()); } void set_peer(const Endpoint& peer) override { const std::optional address = resolve_ipv4(peer); if (!address.has_value()) { return; } std::lock_guard lock(peer_mutex_); peer_ = *address; has_peer_.store(true); } void stop() override { if (!running_.exchange(false)) { return; } if (socket_ >= 0) { (void)::shutdown(socket_, SHUT_RDWR); (void)::close(socket_); socket_ = kInvalidSocket; } // Closing the socket unblocks recvfrom; joining happens implicitly // when the jthread assignment destroys the previous thread. receive_thread_ = std::jthread{}; } private: void receive_loop(ReceiveCallback& callback) { std::array buffer{}; while (running_.load()) { const ssize_t received = ::recvfrom(socket_, buffer.data(), buffer.size(), 0, nullptr, nullptr); if (received <= 0) { continue; // closed socket while running_ still true, or error } const std::optional packet = RtpPacket::parse(std::span{buffer.data(), static_cast(received)}); if (packet.has_value()) { callback(*packet); } } } int socket_ = kInvalidSocket; std::atomic running_{false}; std::atomic has_peer_{false}; std::mutex peer_mutex_; sockaddr_in peer_{}; std::jthread receive_thread_; }; std::unique_ptr RtpTransportFactory::create() { return std::make_unique(); } } // namespace sc