//SPDX-License-Identifier: LGPL-3.0-only //SPDX-FileCopyrightText: Copyright (C) 2026 Catcrafts® module; #include #include #include #include module Crafter.Network:Stream_impl; import :Stream; import std; using namespace Crafter; void Crafter::SetNonBlocking(int descriptor) { const int flags = fcntl(descriptor, F_GETFL, 0); if (flags == -1 || fcntl(descriptor, F_SETFL, flags | O_NONBLOCK) == -1) { throw std::runtime_error(std::string("could not make the socket non-blocking: ") + std::strerror(errno)); } } bool Crafter::PollDescriptor(int descriptor, short events, std::chrono::steady_clock::time_point deadline) { for (;;) { const auto left = std::chrono::duration_cast( deadline - std::chrono::steady_clock::now()); // A deadline already in the past still gets one non-blocking look, so // a zero timeout means "is it ready right now" rather than "give up". const int wait = left.count() > 0 ? static_cast(left.count()) : 0; pollfd descriptors{ .fd = descriptor, .events = events, .revents = 0 }; const int ready = poll(&descriptors, 1, wait); if (ready < 0) { if (errno == EINTR) continue; throw std::runtime_error(std::string("poll failed: ") + std::strerror(errno)); } if (ready == 0) return false; return true; } } PlainStream::PlainStream(int descriptor) : descriptor(descriptor) { SetNonBlocking(descriptor); } StreamStatus PlainStream::ReadSome(char* buffer, std::size_t size, std::chrono::milliseconds timeout, std::size_t& read) { read = 0; const auto deadline = std::chrono::steady_clock::now() + timeout; for (;;) { const auto got = recv(descriptor, buffer, size, 0); if (got > 0) { read = static_cast(got); return StreamStatus::Data; } if (got == 0) return StreamStatus::Closed; if (errno == EINTR) continue; if (errno == EAGAIN || errno == EWOULDBLOCK) { if (!PollDescriptor(descriptor, POLLIN, deadline)) return StreamStatus::TimedOut; continue; } // A reset is how a peer that stopped caring shows up; it is an end of // connection rather than something worth a diagnostic. if (errno == ECONNRESET) return StreamStatus::Closed; throw std::runtime_error(std::string("recv failed: ") + std::strerror(errno)); } } void PlainStream::Write(const void* buffer, std::size_t size, std::chrono::milliseconds timeout) { const auto deadline = std::chrono::steady_clock::now() + timeout; const char* data = reinterpret_cast(buffer); std::size_t sent = 0; while (sent < size) { // MSG_NOSIGNAL: a peer that closed early must surface as EPIPE here, // not as a SIGPIPE that takes the process down. const auto wrote = send(descriptor, data + sent, size - sent, MSG_NOSIGNAL); if (wrote > 0) { sent += static_cast(wrote); continue; } if (wrote == 0) throw std::runtime_error("the peer closed the connection"); if (errno == EINTR) continue; if (errno == EAGAIN || errno == EWOULDBLOCK) { if (!PollDescriptor(descriptor, POLLOUT, deadline)) { throw std::runtime_error("timed out writing to the peer"); } continue; } throw std::runtime_error(std::string("send failed: ") + std::strerror(errno)); } } void PlainStream::Shutdown() noexcept { shutdown(descriptor, SHUT_WR); }