From 04da3d7f151fe401c86a96ef3c23be55515e71c4 Mon Sep 17 00:00:00 2001 From: toxicphreAK Date: Tue, 18 Aug 2026 15:07:57 +0200 Subject: [PATCH] Reassemble socket API messages across reads readSocket() read at most 1023 bytes per call and left message reassembly to a FIXME. The socket is a SOCK_STREAM and carries no message boundaries, so any reply longer than the buffer -- or merely split across two reads by the kernel -- was parsed as two independent messages. Both halves then failed to parse, the original reply was lost, and the waiting open() sat out its full backoff. Worse, the effect is self-sustaining: once the stream desynchronises every subsequent message on the connection is misparsed, so a single long reply broke hydration until restart. A V2/HYDRATE_FILE_RESULT carrying a path plus a free-form arguments.error string clears 1 KB without trying. processSocketInput() now drains the socket into a persistent buffer, dispatches only complete newline terminated messages, and keeps the remainder for the next read. handleReceivedMsg() correspondingly handles a single message and no longer splits on newlines itself, which also retires its trailing-empty-fragment case. The buffer is bounded so a peer that never sends a newline cannot grow it without limit. --- src/openvfsfuse/socketthread.cpp | 168 ++++++++++++++++++------------- src/openvfsfuse/socketthread.h | 9 +- 2 files changed, 107 insertions(+), 70 deletions(-) diff --git a/src/openvfsfuse/socketthread.cpp b/src/openvfsfuse/socketthread.cpp index c253417..7639989 100644 --- a/src/openvfsfuse/socketthread.cpp +++ b/src/openvfsfuse/socketthread.cpp @@ -24,6 +24,7 @@ #include "sharedmap.h" #include "strtools.h" +#include #include #include #include @@ -42,6 +43,13 @@ using namespace std; #define MSG_POST_USER_DATA 2 #define MSG_TIMER 3 +namespace { +/// Upper bound for the receive buffer. A message from the socket API is a +/// single JSON line and stays far below this; anything larger means the peer +/// is not speaking the protocol. +constexpr size_t MaxRxBufferSize = 1024 * 1024; +} + using json = nlohmann::json; struct ThreadMsg @@ -159,86 +167,112 @@ bool SocketThread::socketSendMsg(std::shared_ptr msgData) // openvfsfuse_log(socket_path.c_str(), "socket send", value, "Message: %s", msg.c_str()); } -std::string SocketThread::readSocket() +void SocketThread::processSocketInput() { - // read answer FIXME: Split messages by \n and keep the rest - char buf[1024]; - ssize_t n = read(_socket, buf, sizeof(buf) - 1); - if (n <= 0) - return std::string(); - return std::string(buf, n); + // The socket is a SOCK_STREAM and carries no message boundaries: a single + // read may return a fragment of a message, several messages at once, or + // both. Accumulate into _rxBuffer and only dispatch complete lines. + char buf[4096]; + while (true) { + const ssize_t n = read(_socket, buf, sizeof(buf)); + if (n > 0) { + _rxBuffer.append(buf, static_cast(n)); + continue; + } + if (n == 0) { + // Peer closed the connection. Whatever is left in the buffer can + // never be completed, so drop it rather than misparsing it later. + if (!_rxBuffer.empty()) { + std::cerr << "Socket closed with " << _rxBuffer.size() << " bytes of incomplete message, discarding" << std::endl; + _rxBuffer.clear(); + } + return; + } + if (errno == EINTR) { + continue; + } + // EAGAIN on the non-blocking socket simply means there is nothing more + // to read right now. (EWOULDBLOCK is an alias for it on Linux and macOS.) + if (errno != EAGAIN) { + perror("read"); + return; + } + break; + } + + size_t pos; + while ((pos = _rxBuffer.find('\n')) != std::string::npos) { + handleReceivedMsg(_rxBuffer.substr(0, pos)); + _rxBuffer.erase(0, pos + 1); + } + + // A peer that never sends a newline must not be able to grow our buffer + // without bound. + if (_rxBuffer.size() > MaxRxBufferSize) { + std::cerr << "Discarding " << _rxBuffer.size() << " bytes of unterminated message from the socket API" << std::endl; + _rxBuffer.clear(); + } } -void SocketThread::handleReceivedMsg(const std::string &rawmsg) +void SocketThread::handleReceivedMsg(const std::string &msg) { - if (rawmsg.empty()) { - cout << "Received Message empty" << endl; + string msgType, msgAttr; + if (msg.empty()) { return; } - auto copies = StrTools::split(rawmsg, 0x000A); + cout << "Handle single message " << msg << endl; - for (const string &msg : copies) { - string msgType, msgAttr; - if (msg.empty()) { - continue; - } + size_t found = msg.find(':'); + if (found != string::npos) { + msgType = msg.substr(0, found); + msgAttr = msg.substr(found + 1, string::npos); + } else { + std::cerr << "Invalid message format: " << msg << std::endl; + return; + } - cout << "Handle single message " << msg << endl; + if (msgType == "V2/HYDRATE_FILE_RESULT") { + int id = -1; + std::string status; - size_t found = msg.find(':'); - if (found != string::npos) { - msgType = msg.substr(0, found); - msgAttr = msg.substr(found + 1, string::npos); - } else { - std::cerr << "Invalid message format: " << msg << std::endl; - continue; + try { + const auto j = json::parse(msgAttr); + id = std::stoi(j["id"].get()); + const auto arguments = j["arguments"].get(); + if (arguments.contains("error")) { + std::cerr << "Error from socket API for Id " << id << ": " << arguments["error"].get() << std::endl; + } else { + status = arguments["status"].get(); + } + } catch (json::exception &e) { + std::cerr << "Invalid JSON message: " << msgAttr << e.what() << std::endl; + return; } - // FIXME: Think if splitting by newline makes sense - - if (msgType == "V2/HYDRATE_FILE_RESULT") { - int id = -1; - std::string status; - - try { - const auto j = json::parse(msgAttr); - id = std::stoi(j["id"].get()); - const auto arguments = j["arguments"].get(); - if (arguments.contains("error")) { - std::cerr << "Error from socket API for Id " << id << ": " << arguments["error"].get() << std::endl; - } else { - status = arguments["status"].get(); - } - } catch (json::exception &e) { - std::cerr << "Invalid JSON message: " << msgAttr << e.what() << std::endl; - continue; + if (id > 0) { + int res{-1}; // Default set to fail + if (status == "OK") { + res = 0; // good! + } else { + cout << "ERROR from socket API for Id" << id << endl; } - if (id > 0) { - int res{-1}; // Default set to fail - if (status == "OK") { - res = 0; // good! - } else { - cout << "ERROR from socket API for Id" << id << endl; - } - - const HydJob hj{.state = res}; - bool ok = _sharedMap.set(id, hj); - if (!ok) { - // the id could not be set. That means, the job was not inserted. - cout << "Job not found:" << id << endl; - } else { - cout << "Setting Job ID " << id << " to result " << res << endl; - } - } - } else if (msgType == "VERSION") { - vector attribs = StrTools::split(msgAttr, ':'); - if (attribs.size() == 3) { - cout << "Got PID of the Desktop Client: " << attribs.at(2) << endl; - _sharedMap.setDesktopClientPid(std::stol(attribs.at(2))); + const HydJob hj{.state = res}; + bool ok = _sharedMap.set(id, hj); + if (!ok) { + // the id could not be set. That means, the job was not inserted. + cout << "Job not found:" << id << endl; + } else { + cout << "Setting Job ID " << id << " to result " << res << endl; } } + } else if (msgType == "VERSION") { + vector attribs = StrTools::split(msgAttr, ':'); + if (attribs.size() == 3) { + cout << "Got PID of the Desktop Client: " << attribs.at(2) << endl; + _sharedMap.setDesktopClientPid(std::stol(attribs.at(2))); + } } } @@ -402,11 +436,7 @@ void SocketThread::Process() case MSG_TIMER: { // cout << "Timer expired on " << THREAD_NAME << endl; - const std::string msg = readSocket(); - if (!msg.empty()) { - cout << "Message received: " << msg << endl; - handleReceivedMsg(msg); - } + processSocketInput(); break; } diff --git a/src/openvfsfuse/socketthread.h b/src/openvfsfuse/socketthread.h index ed4734c..93b81ed 100644 --- a/src/openvfsfuse/socketthread.h +++ b/src/openvfsfuse/socketthread.h @@ -89,7 +89,10 @@ class SocketThread int initSocket(const std::string& socketPath); bool socketSendMsg(std::shared_ptr); - std::string readSocket(); + + /// Drain everything readable from the socket into _rxBuffer and dispatch + /// every complete (newline terminated) message it contains. + void processSocketInput(); void handleReceivedMsg(const std::string &msg); /// Entry point for the worker thread @@ -114,6 +117,10 @@ class SocketThread std::atomic _socket; + /// Receive buffer holding the bytes read from the socket that do not form + /// a complete message yet. Only touched by the worker thread. + std::string _rxBuffer; + SharedMap &_sharedMap; };