[llvm] [orc-rt] Add a SimpleRemote transport over POSIX sockets (PR #225393)
Lang Hames via llvm-commits
llvm-commits at lists.llvm.org
Tue Sep 22 06:13:55 PDT 2026
https://github.com/lhames created https://github.com/llvm/llvm-project/pull/225393
A ControllerAccess speaking SimpleRemote over a connected stream socket. The wire format matches LLVM's SimpleRemoteEPC. POSIX only -- Windows needs a counterpart alongside the other sys/posix sources.
It owns a reactor thread built on poll(2) over the connection and a wake socket. IO is non-blocking, so reads and writes resume wherever a partial transfer stopped. The reactor is the only thread that sends; callers on other threads append to a queue and wake it.
Reached through createSimpleRemoteCAOverSocket; the implementation is file-local to the POSIX translation unit, so a Windows version needs no header change.
The unit test drives it over a real socket pair, playing the controller on the far end: framing, partial transfers, teardown ordering and the malformed-message paths. The fixture owns the pair and the Session, so a test body starts at the interesting part. It installs no on-disconnect handler: Session reports the disconnect reason through reportError when none is set, so a test that ends abnormally without asking for the error fails rather than passing quietly.
>From 1dfa1e362b1c3a264c5dce8eb3ed0853b7e064e7 Mon Sep 17 00:00:00 2001
From: Lang Hames <lhames at gmail.com>
Date: Mon, 14 Sep 2026 15:10:16 +1000
Subject: [PATCH] [orc-rt] Add a SimpleRemote transport over POSIX sockets
A ControllerAccess speaking SimpleRemote over a connected stream socket.
The wire format matches LLVM's SimpleRemoteEPC. POSIX only -- Windows
needs a counterpart alongside the other sys/posix sources.
It owns a reactor thread built on poll(2) over the connection and a wake
socket. IO is non-blocking, so reads and writes resume wherever a
partial transfer stopped. The reactor is the only thread that sends;
callers on other threads append to a queue and wake it.
Reached through createSimpleRemoteCAOverSocket; the implementation is
file-local to the POSIX translation unit, so a Windows version needs no
header change.
The unit test drives it over a real socket pair, playing the controller
on the far end: framing, partial transfers, teardown ordering and the
malformed-message paths. The fixture owns the pair and the Session, so a
test body starts at the interesting part. It installs no on-disconnect
handler: Session reports the disconnect reason through reportError when
none is set, so a test that ends abnormally without asking for the error
fails rather than passing quietly.
---
.../bedrock/sps/SimpleRemoteCAOverSocket.h | 41 ++
orc-rt/lib/bedrock/CMakeLists.txt | 1 +
.../posix/sps/SimpleRemoteCAOverSocket.cpp | 547 +++++++++++++++
orc-rt/test/unit/CMakeLists.txt | 1 +
.../sps/SimpleRemoteCAOverSocketTest.cpp | 651 ++++++++++++++++++
5 files changed, 1241 insertions(+)
create mode 100644 orc-rt/include/orc-rt/bedrock/sps/SimpleRemoteCAOverSocket.h
create mode 100644 orc-rt/lib/bedrock/sys/posix/sps/SimpleRemoteCAOverSocket.cpp
create mode 100644 orc-rt/test/unit/bedrock/sps/SimpleRemoteCAOverSocketTest.cpp
diff --git a/orc-rt/include/orc-rt/bedrock/sps/SimpleRemoteCAOverSocket.h b/orc-rt/include/orc-rt/bedrock/sps/SimpleRemoteCAOverSocket.h
new file mode 100644
index 0000000000000..e523933e5cd48
--- /dev/null
+++ b/orc-rt/include/orc-rt/bedrock/sps/SimpleRemoteCAOverSocket.h
@@ -0,0 +1,41 @@
+//===- SimpleRemoteCAOverSocket.h - SimpleRemote CA over socket -*- C++ -*-===//
+//
+// Part of the LLVM Project, under the Apache License v2.0 with LLVM Exceptions.
+// See https://llvm.org/LICENSE.txt for license information.
+// SPDX-License-Identifier: Apache-2.0 WITH LLVM-exception
+//
+//===----------------------------------------------------------------------===//
+//
+// A ControllerAccess speaking the SimpleRemote protocol over a connected
+// socket.
+//
+//===----------------------------------------------------------------------===//
+
+#ifndef ORC_RT_BEDROCK_SPS_SIMPLEREMOTECAOVERSOCKET_H
+#define ORC_RT_BEDROCK_SPS_SIMPLEREMOTECAOVERSOCKET_H
+
+#include "orc-rt/bedrock/Session.h"
+#include "orc-rt/bedrock/SocketHandle.h"
+#include "orc-rt/support/Error.h"
+
+#include <memory>
+
+namespace orc_rt {
+
+/// Creates a ControllerAccess that carries SimpleRemote messages over Sock,
+/// taking ownership of it. Sock must be a connected stream socket.
+///
+/// The result is ready to hand to Session::attach, which is what starts the
+/// conversation; nothing is sent before then.
+///
+/// A factory rather than a class, because how the messages are pumped is the
+/// platform's business and not the caller's: the POSIX implementation owns a
+/// reactor thread built on poll(2) and a wake socket, and a Windows one will
+/// need something else entirely. The wire format is the same either way, and
+/// matches LLVM's SimpleRemoteEPC.
+Expected<std::shared_ptr<Session::ControllerAccess>>
+createSimpleRemoteCAOverSocket(Session &S, SocketHandle Sock);
+
+} // namespace orc_rt
+
+#endif // ORC_RT_BEDROCK_SPS_SIMPLEREMOTECAOVERSOCKET_H
diff --git a/orc-rt/lib/bedrock/CMakeLists.txt b/orc-rt/lib/bedrock/CMakeLists.txt
index 99016c5866d41..48f77c664b526 100644
--- a/orc-rt/lib/bedrock/CMakeLists.txt
+++ b/orc-rt/lib/bedrock/CMakeLists.txt
@@ -41,6 +41,7 @@ set(ORC_RT_BEDROCK_POSIX_SOURCES
sys/posix/Memory.cpp
sys/posix/PageSize.cpp
sys/posix/SocketHandle.cpp
+ sys/posix/sps/SimpleRemoteCAOverSocket.cpp
)
set(ORC_RT_BEDROCK_DARWIN_SOURCES
diff --git a/orc-rt/lib/bedrock/sys/posix/sps/SimpleRemoteCAOverSocket.cpp b/orc-rt/lib/bedrock/sys/posix/sps/SimpleRemoteCAOverSocket.cpp
new file mode 100644
index 0000000000000..50bcde3c6536b
--- /dev/null
+++ b/orc-rt/lib/bedrock/sys/posix/sps/SimpleRemoteCAOverSocket.cpp
@@ -0,0 +1,547 @@
+//===- SimpleRemoteCAOverSocket.cpp ---------------------------------------===//
+//
+// Part of the LLVM Project, under the Apache License v2.0 with LLVM Exceptions.
+// See https://llvm.org/LICENSE.txt for license information.
+// SPDX-License-Identifier: Apache-2.0 WITH LLVM-exception
+//
+//===----------------------------------------------------------------------===//
+//
+// SimpleRemote protocol over a connected socket, on POSIX.
+//
+//===----------------------------------------------------------------------===//
+
+#include "orc-rt/bedrock/sps/SimpleRemoteCAOverSocket.h"
+
+#include "orc-rt-c/support/Logging.h"
+#include "orc-rt-internal/support/sys/Errno.h"
+#include "orc-rt/bedrock/sps/SimpleRemoteCA.h"
+#include "orc-rt/support/Compiler.h"
+#include "orc-rt/support/span.h"
+
+#include <algorithm>
+#include <cassert>
+#include <cerrno>
+#include <chrono>
+#include <fcntl.h>
+#include <optional>
+#include <poll.h>
+#include <string>
+#include <sys/socket.h>
+#include <thread>
+#include <unistd.h>
+#include <utility>
+
+namespace orc_rt {
+
+namespace {
+
+class SocketSimpleRemoteCA : public SimpleRemoteCA {
+public:
+ static Expected<std::shared_ptr<SocketSimpleRemoteCA>>
+ Create(Session &S, SocketHandle Sock);
+
+private:
+ enum class State {
+ NotConnected, ///< Before connect. Nothing may be registered or queued.
+ Running, ///< Calls may be registered, messages queued.
+ Draining, ///< Hang-up queued and latched: nothing may follow it.
+ Closed, ///< Reactor stopped and descriptors released.
+ };
+
+ /// A framed message waiting to go out. Owns the payload. Supports interrupted
+ /// sends via pending/advance/complete.
+ class OutgoingMessage {
+ public:
+ OutgoingMessage(Opcode Op, uint64_t SeqNo, uint64_t Tag,
+ WrapperFunctionBuffer Payload);
+
+ /// The next bytes to send. Empty once the message has gone.
+ span<const char> pending() const;
+
+ void advance(size_t N) { Sent += N; }
+ bool complete() const { return Sent == MsgHeader::Size + Payload.size(); }
+
+ private:
+ char Header[MsgHeader::Size];
+ WrapperFunctionBuffer Payload;
+ size_t Sent = 0;
+ };
+
+ /// Incoming message. Supports interrupted reads via pending/advance/complete.
+ /// The header is filled first then decoded to determine the allocation size
+ /// for the payload buffer.
+ class IncomingMessage {
+ public:
+ /// Where the next bytes read should land. Empty once the current buffer is
+ /// full: the header is then ready to decode, or the message is complete.
+ span<char> pending();
+
+ void advance(size_t N) { Filled += N; }
+
+ /// True once the payload buffer is full and the message can be dispatched.
+ bool complete() const {
+ return Decoded && Filled == MsgHeader::Size + Payload.size();
+ }
+
+ /// Decodes the completed header and allocates space for the payload.
+ /// Fails on a header that describes an impossible message.
+ Error decodeHeader();
+
+ const MsgHeader::Fields &fields() const { return F; }
+
+ /// Take the payload. Resets this value for the next message.
+ WrapperFunctionBuffer take();
+
+ private:
+ char Header[MsgHeader::Size];
+ WrapperFunctionBuffer Payload;
+ size_t Filled = 0;
+ bool Decoded = false;
+ MsgHeader::Fields F;
+ };
+
+ SocketSimpleRemoteCA(Session &S, SocketHandle Sock, SocketHandle WakeRead,
+ SocketHandle WakeWrite)
+ : SimpleRemoteCA(S), Sock(std::move(Sock)), WakeRead(std::move(WakeRead)),
+ WakeWrite(std::move(WakeWrite)) {}
+
+ // Session::ControllerAccess.
+ void connect(BootstrapInfo BI) override;
+ void disconnect() override;
+ void callController(OnControllerCallReturn OnComplete,
+ orc_rt_ControllerHandlerTag T,
+ WrapperFunctionBuffer ArgBytes) override;
+ void sendWrapperResult(WrapperFunctionBuffer ResultBytes,
+ uint64_t CallId) override;
+
+ /// Wake the reactor thread.
+ ///
+ /// Makes a poll in progress return, or the next one return at once.
+ ///
+ /// Caller must hold M and must not be in the Closed state: the
+ /// reactor releases the descriptors under the same lock, so holding
+ /// it is what keeps this from writing to a closed one.
+ void wakeReactorLocked();
+
+ /// Reactor thread entry point. Runs the reactor loop, then performs cleanup:
+ /// releasing descriptors, failing outstanding calls, and notifying the
+ /// session.
+ void runReactor();
+
+ /// Message IO loop: polls, reads, dispatches and sends until the connection
+ /// ends. Returns the reason it stopped: success if either side hung up.
+ Error reactorLoop();
+
+ /// Sends as much of the queue as the socket will take. Reactor thread only.
+ Error drainSends();
+
+ /// Receives what is available and hands each complete message to
+ /// handleMessage. Incoming carries any part-assembled message across calls.
+ /// Reactor thread only.
+ ///
+ /// Only a hang-up ends the session cleanly. Unexpected end of stream is an
+ /// error.
+ Expected<Action> readAndDispatch(IncomingMessage &Incoming);
+
+ /// Takes the handler for SeqNo under M.
+ OnControllerCallReturn takePendingCall(uint64_t SeqNo) override;
+
+ SocketHandle Sock;
+ SocketHandle WakeRead, WakeWrite;
+
+ std::mutex M;
+ State CurState = State::NotConnected;
+
+ /// Framed messages awaiting send. Only the reactor pops, so it may hold a
+ /// reference to the front across a send without the lock.
+ ///
+ /// TODO: Currently unbounded. We may want to add a bound on this.
+ std::deque<OutgoingMessage> Queue;
+};
+
+Error makeError(const char *Op, int ErrNum) {
+ return make_error<StringError>(std::string(Op) +
+ " failed: " + sys::strError(ErrNum));
+}
+
+bool isWouldBlock(int ErrNum) {
+ return ErrNum == EAGAIN || ErrNum == EWOULDBLOCK;
+}
+
+/// O_NONBLOCK rather than MSG_DONTWAIT on each send/recv: Darwin defines
+/// MSG_DONTWAIT but does not honour it on sends.
+Error setNonBlocking(int FD) {
+ int Flags = fcntl(FD, F_GETFL, 0);
+ if (Flags == -1)
+ return makeError("fcntl(F_GETFL)", errno);
+ if ((Flags & O_NONBLOCK) == 0 && fcntl(FD, F_SETFL, Flags | O_NONBLOCK) == -1)
+ return makeError("fcntl(F_SETFL)", errno);
+ return Error::success();
+}
+
+template <typename OpT> ssize_t retryOnEINTR(OpT Op) {
+ for (;;) {
+ ssize_t N = Op();
+ if (N >= 0 || errno != EINTR)
+ return N;
+ }
+}
+
+} // namespace
+
+SocketSimpleRemoteCA::OutgoingMessage::OutgoingMessage(
+ Opcode Op, uint64_t SeqNo, uint64_t Tag, WrapperFunctionBuffer Payload)
+ : Payload(std::move(Payload)) {
+ assert(!this->Payload.getOutOfBandError() &&
+ "Out-of-band errors have no byte representation: encode as a result "
+ "kind before framing");
+ MsgHeader::encode(Header, Op, SeqNo, Tag, this->Payload.size());
+}
+
+span<const char> SocketSimpleRemoteCA::OutgoingMessage::pending() const {
+ if (Sent < MsgHeader::Size)
+ return {Header + Sent, MsgHeader::Size - Sent};
+ size_t InPayload = Sent - MsgHeader::Size;
+ return {Payload.data() + InPayload, Payload.size() - InPayload};
+}
+
+span<char> SocketSimpleRemoteCA::IncomingMessage::pending() {
+ if (!Decoded)
+ return {Header + Filled, MsgHeader::Size - Filled};
+ size_t InPayload = Filled - MsgHeader::Size;
+ return {Payload.data() + InPayload, Payload.size() - InPayload};
+}
+
+Error SocketSimpleRemoteCA::IncomingMessage::decodeHeader() {
+ F = MsgHeader::decode(Header);
+ if (F.MsgSize < MsgHeader::Size)
+ return make_error<StringError>("Message size smaller than its header");
+
+ Payload = WrapperFunctionBuffer::allocate(F.MsgSize - MsgHeader::Size);
+ Decoded = true;
+ return Error::success();
+}
+
+WrapperFunctionBuffer SocketSimpleRemoteCA::IncomingMessage::take() {
+ auto P = std::move(Payload);
+ Payload = WrapperFunctionBuffer();
+ Filled = 0;
+ Decoded = false;
+ return P;
+}
+
+Expected<std::shared_ptr<SocketSimpleRemoteCA>>
+SocketSimpleRemoteCA::Create(Session &S, SocketHandle Sock) {
+ // Sock is owned here, so every early return below closes it.
+ if (auto Err = setNonBlocking(Sock.get()))
+ return std::move(Err);
+
+ int Pair[2];
+ if (::socketpair(AF_UNIX, SOCK_STREAM, 0, Pair) != 0)
+ return makeError("socketpair", errno);
+ SocketHandle WakeRead(Pair[0]), WakeWrite(Pair[1]);
+
+ // Both ends non-blocking. A blocking write end would stall a sender inside M,
+ // which the reactor needs before it can drain -- a deadlock; a blocking read
+ // end would leave the drain loop waiting for a byte after emptying the queue.
+ for (int W : {WakeRead.get(), WakeWrite.get()})
+ if (auto Err = setNonBlocking(W))
+ return std::move(Err);
+
+ // Not make_shared: the constructor is private.
+ return std::shared_ptr<SocketSimpleRemoteCA>(new SocketSimpleRemoteCA(
+ S, std::move(Sock), std::move(WakeRead), std::move(WakeWrite)));
+}
+
+void SocketSimpleRemoteCA::wakeReactorLocked() {
+ assert(CurState != State::Closed && "wake on a released descriptor");
+ char C = 0;
+ ssize_t N = retryOnEINTR(
+ [&] { return ::send(WakeWrite.get(), &C, 1, MSG_NOSIGNAL); });
+ int ErrNum = errno;
+ if (N < 0 && !isWouldBlock(ErrNum)) {
+ // send to wait socket failed. Log the reason in case this jams up the
+ // reactor.
+ ORC_RT_LOG(Info, ControllerAccess,
+ "SimpleRemoteCA/socket wake-send error: " ORC_RT_LOG_PUB_S,
+ sys::strError(ErrNum).c_str());
+ }
+}
+
+void SocketSimpleRemoteCA::connect(BootstrapInfo BI) {
+ {
+ std::scoped_lock<std::mutex> Lock(M);
+ assert(CurState == State::NotConnected && "connect called twice");
+ CurState = State::Running;
+ // Queued before the reactor exists, and the reactor is the only sender, so
+ // setup is the first message on the wire.
+ Queue.emplace_back(Opcode::Setup, 0, 0, encodeSetup(BI));
+ }
+
+ // The Session holds this alive until notifyDisconnected, which the reactor
+ // reaches only as it exits, so the thread needs no reference of its own.
+ //
+ // FIXME: Report a spawn failure, and offer a pumped mode that borrows the
+ // caller's thread. std::thread's constructor aborts rather than reporting
+ // under -fno-exceptions, so connect cannot surface one today.
+ std::thread([this] { runReactor(); }).detach();
+}
+
+void SocketSimpleRemoteCA::disconnect() {
+ std::scoped_lock<std::mutex> Lock(M);
+ // Anything but Running means teardown is under way or done, or the connection
+ // never opened. The Session tolerates a disconnect racing a remote one.
+ if (CurState != State::Running)
+ return;
+
+ // Queued and latched in one lock hold. Draining is what keeps anything from
+ // landing behind the hang-up, so it must take effect with the same atomicity
+ // as the queueing.
+ Queue.emplace_back(Opcode::Hangup, 0, 0, encodeHangup(Error::success()));
+ CurState = State::Draining;
+ wakeReactorLocked();
+}
+
+void SocketSimpleRemoteCA::callController(OnControllerCallReturn OnComplete,
+ orc_rt_ControllerHandlerTag T,
+ WrapperFunctionBuffer ArgBytes) {
+ {
+ std::scoped_lock<std::mutex> Lock(M);
+ if (CurState == State::Running) {
+ // Registered and queued in one lock hold, so a call is never left pending
+ // with nothing to answer it, nor sent with no handler to complete.
+ uint64_t SeqNo = registerCall(std::move(OnComplete));
+ Queue.emplace_back(Opcode::Call, SeqNo,
+ ExecutorAddr::fromPtr(T).getValue(),
+ std::move(ArgBytes));
+ wakeReactorLocked();
+ return;
+ }
+ }
+
+ // The connection is gone, so no result can arrive. The caller is still on the
+ // stack, so fail the handler there.
+ failControllerCallInline(std::move(OnComplete));
+}
+
+void SocketSimpleRemoteCA::sendWrapperResult(WrapperFunctionBuffer ResultBytes,
+ uint64_t CallId) {
+ // Encoded before the lock: an out-of-band error has no byte representation,
+ // so it travels as a distinct result kind with its message as the payload.
+ auto [Kind, Payload] = encodeResult(std::move(ResultBytes));
+
+ // No state check of its own: a result has no pending call on this side, so a
+ // departed connection drops it with nothing left unsettled.
+ std::scoped_lock<std::mutex> Lock(M);
+ if (CurState != State::Running)
+ return;
+ Queue.emplace_back(Opcode::Result, CallId, static_cast<uint64_t>(Kind),
+ std::move(Payload));
+ wakeReactorLocked();
+}
+
+void SocketSimpleRemoteCA::runReactor() {
+ Error Err = reactorLoop();
+
+ {
+ std::scoped_lock<std::mutex> Lock(M);
+ // Closed refuses every further send and wake.
+ CurState = State::Closed;
+ Sock.reset();
+ WakeRead.reset();
+ WakeWrite.reset();
+ }
+
+ // Before the notification, while the managed-code group is still open, or the
+ // handlers are dropped rather than dispatched.
+ PendingCallsMap Failed;
+ {
+ std::scoped_lock<std::mutex> Lock(M);
+ Failed = takeAllCalls();
+ }
+ for (auto &[SeqNo, OnComplete] : Failed)
+ failPendingControllerCall(std::move(OnComplete));
+
+ // Last thing to touch this: it may run the destructor.
+ notifyDisconnected(std::move(Err));
+}
+
+Error SocketSimpleRemoteCA::reactorLoop() {
+ // Part-assembled incoming message.
+ IncomingMessage Incoming;
+
+ // Bounds the drain after a local disconnect so that a controller that has
+ // stopped reading doesn't prevent the reactor from exiting.
+ // The deadline is established on the first call. Returns the milliseconds
+ // left to wait.
+ auto DrainTimeRemaining =
+ [DrainDeadline =
+ std::optional<std::chrono::steady_clock::time_point>()]() mutable
+ -> int {
+ constexpr auto DrainTimeout = std::chrono::seconds(5);
+ auto Now = std::chrono::steady_clock::now();
+ if (!DrainDeadline)
+ DrainDeadline = Now + DrainTimeout;
+ auto Remaining = std::chrono::duration_cast<std::chrono::milliseconds>(
+ *DrainDeadline - Now);
+ return std::max(static_cast<int>(Remaining.count()), 0);
+ };
+
+ for (;;) {
+ pollfd PollFDs[2];
+ PollFDs[0].fd = Sock.get();
+ PollFDs[0].events = POLLIN;
+ PollFDs[0].revents = 0;
+ PollFDs[1].fd = WakeRead.get();
+ PollFDs[1].events = POLLIN;
+ PollFDs[1].revents = 0;
+
+ // Check for outgoing messages / draining-state.
+ bool Draining = false;
+ {
+ std::scoped_lock<std::mutex> Lock(M);
+ if (!Queue.empty())
+ PollFDs[0].events |= POLLOUT;
+ else if (CurState == State::Draining)
+ return Error::success(); // The hang-up has gone out.
+ Draining = CurState == State::Draining;
+ }
+
+ // Poll timeout.
+ int Timeout = -1;
+
+ // If draining, bound the poll by the timeout window or exit if the window
+ // has closed.
+ if (Draining) {
+ Timeout = DrainTimeRemaining();
+ if (Timeout == 0) {
+ ORC_RT_LOG(Info, ControllerAccess,
+ "SimpleRemoteCA/socket drain timed out with messages still "
+ "queued; dropping them");
+ return Error::success();
+ }
+ }
+
+ // Check readiness (with timout if draining).
+ while (::poll(PollFDs, 2, Timeout) < 0) {
+ if (errno == EINTR)
+ continue;
+ return makeError("poll", errno);
+ }
+
+ // Drain the wake notification so that next poll blocks.
+ if (PollFDs[1].revents & POLLIN) {
+ char Buf[64];
+ while (::recv(WakeRead.get(), Buf, sizeof(Buf), 0) > 0)
+ ;
+ }
+
+ // Read before sending, so that a peer's parting hang-up reaches us rather
+ // than the EPIPE from a send to a peer that has already gone.
+ //
+ // POLLHUP and POLLERR read too. A closed peer may have left buffered
+ // messages, and recv returning zero is what separates an orderly close from
+ // a truncated one; POLLERR says only that something is pending, so recv is
+ // what reports what.
+ auto Next = Action::Continue;
+ if (PollFDs[0].revents & (POLLIN | POLLHUP | POLLERR | POLLNVAL)) {
+ auto A = readAndDispatch(Incoming);
+ if (!A)
+ return A.takeError();
+ Next = *A;
+ }
+
+ // If readAndDispatch got a hang-up then just do a best-effort send of the
+ // remaining queue items, then return.
+ if (Next == Action::End) {
+ if (auto Err = drainSends()) {
+ [[maybe_unused]] std::string Msg = toString(std::move(Err));
+ ORC_RT_LOG(Info, ControllerAccess,
+ "SimpleRemoteCA/socket final-flush: " ORC_RT_LOG_PUB_S,
+ Msg.c_str());
+ }
+ return Error::success();
+ }
+
+ // Otherwise send any remaining queue items.
+ if (PollFDs[0].revents & POLLOUT)
+ if (auto Err = drainSends())
+ return Err;
+ }
+}
+
+Error SocketSimpleRemoteCA::drainSends() {
+ for (;;) {
+ OutgoingMessage *Msg = nullptr;
+ {
+ std::scoped_lock<std::mutex> Lock(M);
+ if (Queue.empty())
+ return Error::success();
+ Msg = &Queue.front();
+ }
+
+ // Sent without the lock: only the reactor pops, and appending to a deque
+ // never moves an existing element.
+ auto Bytes = Msg->pending();
+ ssize_t N = retryOnEINTR([&] {
+ return ::send(Sock.get(), Bytes.data(), Bytes.size(), MSG_NOSIGNAL);
+ });
+ if (N < 0)
+ return isWouldBlock(errno) ? Error::success() : makeError("send", errno);
+ Msg->advance(N);
+
+ if (Msg->complete()) {
+ std::scoped_lock<std::mutex> Lock(M);
+ Queue.pop_front();
+ }
+ }
+}
+
+Expected<SocketSimpleRemoteCA::Action>
+SocketSimpleRemoteCA::readAndDispatch(IncomingMessage &Incoming) {
+ for (;;) {
+ // A non-blocking recv can stop anywhere, including mid-header, so this
+ // resumes wherever the last left off.
+ if (auto Bytes = Incoming.pending(); !Bytes.empty()) {
+ ssize_t N = retryOnEINTR(
+ [&] { return ::recv(Sock.get(), Bytes.data(), Bytes.size(), 0); });
+ if (N < 0) {
+ if (isWouldBlock(errno))
+ return Action::Continue;
+ return makeError("recv", errno);
+ }
+ if (N == 0)
+ return make_error<StringError>(
+ "Connection closed without a hang-up message");
+ Incoming.advance(N);
+ continue;
+ }
+
+ // Header full: decode it, then read the payload it describes. An empty
+ // payload leaves nothing pending, so the next pass dispatches.
+ if (!Incoming.complete()) {
+ if (auto Err = Incoming.decodeHeader())
+ return Err;
+ continue;
+ }
+
+ // Taken before dispatching: a handler may send, re-entering this object.
+ auto F = Incoming.fields();
+ auto A = handleMessage(F.OpC, F.SeqNo, F.Tag, Incoming.take());
+ if (!A || *A == Action::End)
+ return A;
+ }
+}
+
+SocketSimpleRemoteCA::OnControllerCallReturn
+SocketSimpleRemoteCA::takePendingCall(uint64_t SeqNo) {
+ std::scoped_lock<std::mutex> Lock(M);
+ return takeCall(SeqNo);
+}
+
+Expected<std::shared_ptr<Session::ControllerAccess>>
+createSimpleRemoteCAOverSocket(Session &S, SocketHandle Sock) {
+ return SocketSimpleRemoteCA::Create(S, std::move(Sock));
+}
+
+} // namespace orc_rt
diff --git a/orc-rt/test/unit/CMakeLists.txt b/orc-rt/test/unit/CMakeLists.txt
index 208be6313036d..80434ea6bdc22 100644
--- a/orc-rt/test/unit/CMakeLists.txt
+++ b/orc-rt/test/unit/CMakeLists.txt
@@ -123,6 +123,7 @@ add_orc_rt_unittest(BedrockTests
bedrock/sps/NativeDylibManagerSPSCITest.cpp
bedrock/sps/SimpleNativeMemoryMapSPSCITest.cpp
bedrock/sps/SimpleRemoteCATest.cpp
+ bedrock/sps/SimpleRemoteCAOverSocketTest.cpp
bedrock/sys/CPUFeaturesTest.cpp
bedrock/sys/TargetTripleTest.cpp
diff --git a/orc-rt/test/unit/bedrock/sps/SimpleRemoteCAOverSocketTest.cpp b/orc-rt/test/unit/bedrock/sps/SimpleRemoteCAOverSocketTest.cpp
new file mode 100644
index 0000000000000..c14d344111683
--- /dev/null
+++ b/orc-rt/test/unit/bedrock/sps/SimpleRemoteCAOverSocketTest.cpp
@@ -0,0 +1,651 @@
+//===- SimpleRemoteCAOverSocketTest.cpp -----------------------------------===//
+//
+// Part of the LLVM Project, under the Apache License v2.0 with LLVM Exceptions.
+// See https://llvm.org/LICENSE.txt for license information.
+// SPDX-License-Identifier: Apache-2.0 WITH LLVM-exception
+//
+//===----------------------------------------------------------------------===//
+//
+// Tests for createSimpleRemoteCAOverSocket.
+//
+// Each test attaches a CA to one end of a socket pair and plays the controller
+// on the other, so the reactor, the connection state machine and teardown are
+// covered without a second process.
+//
+// The controller side frames messages with the CA's own encoder, reached
+// through the friend fixture. That makes these tests about behaviour rather
+// than byte layout; conformance to LLVM's SimpleRemoteEPC is established by the
+// cross-process regression tests.
+//
+//===----------------------------------------------------------------------===//
+
+#include "orc-rt/bedrock/sps/SimpleRemoteCAOverSocket.h"
+
+#include "orc-rt/bedrock/sps/SimpleRemoteCA.h"
+
+#include "gtest/gtest.h"
+
+#include "BedrockTestUtils.h"
+#include "CommonTestUtils.h"
+
+#include "orc-rt-internal/support/Endian.h"
+
+#include <cassert>
+#include <cerrno>
+#include <chrono>
+#include <cstring>
+#include <future>
+#include <string>
+#include <string_view>
+#include <sys/socket.h>
+#include <thread>
+#include <unistd.h>
+#include <utility>
+#include <vector>
+
+using namespace orc_rt;
+
+namespace orc_rt {
+
+/// Plays the controller against a SimpleRemote CA running over a socket.
+///
+/// Frames messages with the protocol's own encoders rather than reimplementing
+/// them. Those are protected on SimpleRemoteCA, so a derived type re-exports
+/// the pieces the controller side needs; the transport itself is not a type a
+/// test can reach, which is the point of it being a factory.
+class SimpleRemoteCAOverSocketTest : public ::testing::Test {
+ struct CA : SimpleRemoteCA {
+ using SimpleRemoteCA::decodeResult;
+ using SimpleRemoteCA::encodeHangup;
+ using SimpleRemoteCA::encodeResult;
+ using MsgHeader = SimpleRemoteCA::MsgHeader;
+ using Opcode = SimpleRemoteCA::Opcode;
+ using ResultKind = SimpleRemoteCA::ResultKind;
+ };
+
+public:
+ using MsgHeader = CA::MsgHeader;
+ using Opcode = CA::Opcode;
+
+ /// A connected socket pair, established before each test: Near is the end the
+ /// CA adopts, Far the end the test drives as the controller.
+ ///
+ /// Declared before S so that the Session is torn down first, while Far is
+ /// still open -- that is the order a test's own locals had.
+ SocketHandle Near, Far;
+
+ /// The Session under test. noErrors is not boilerplate: Session routes the
+ /// disconnect reason to reportError when no on-disconnect handler was
+ /// installed, so any test that ends abnormally without asking for the error
+ /// aborts here rather than passing quietly. Tests that want the error install
+ /// their own handler before attaching.
+ Session S{mockExecutorProcessInfo(), inlineDispatch, noErrors};
+
+ void SetUp() override {
+ auto P = makePair();
+ ASSERT_TRUE(!!P) << toString(P.takeError());
+ Near = SocketHandle(P->first);
+ Far = SocketHandle(P->second);
+ }
+
+ /// Creates a CA over Near and attaches it, the way a connector would.
+ Error attachOverSocket() {
+ auto CA = createSimpleRemoteCAOverSocket(S, std::move(Near));
+ if (!CA)
+ return CA.takeError();
+ S.attach(std::move(*CA), BootstrapInfo(S));
+ return Error::success();
+ }
+
+ static constexpr size_t HeaderSize = MsgHeader::Size;
+
+ /// A message as it appears on the wire, with the size field consumed.
+ struct Frame {
+ MsgHeader::Fields Fields;
+ std::vector<char> Payload;
+
+ Opcode opcode() const { return static_cast<Opcode>(Fields.OpC); }
+ uint64_t seqNo() const { return Fields.SeqNo; }
+ uint64_t tag() const { return Fields.Tag; }
+
+ /// The payload as a view, for comparison against test data.
+ std::string_view payload() const {
+ return {Payload.data(), Payload.size()};
+ }
+ };
+
+ /// A connected pair of blocking stream sockets: one end for the CA to adopt,
+ /// one for the test to drive.
+ static Expected<std::pair<int, int>> makePair() {
+ int FDs[2];
+ if (::socketpair(AF_UNIX, SOCK_STREAM, 0, FDs) != 0)
+ return make_error<StringError>(std::string("socketpair: ") +
+ strerror(errno));
+ return std::make_pair(FDs[0], FDs[1]);
+ }
+
+ /// The test's end stays blocking, so these just loop until done.
+ static Error sendAll(int FDNum, const char *Buf, size_t Size) {
+ while (Size) {
+ ssize_t N = ::send(FDNum, Buf, Size, MSG_NOSIGNAL);
+ if (N < 0) {
+ if (errno == EINTR)
+ continue;
+ return make_error<StringError>(std::string("send: ") + strerror(errno));
+ }
+ Buf += N;
+ Size -= N;
+ }
+ return Error::success();
+ }
+
+ /// Returns short only if the peer closed first.
+ static Expected<size_t> recvAll(int FDNum, char *Buf, size_t Size) {
+ size_t Got = 0;
+ while (Got < Size) {
+ ssize_t N = ::recv(FDNum, Buf + Got, Size - Got, 0);
+ if (N == 0)
+ return Got;
+ if (N < 0) {
+ if (errno == EINTR)
+ continue;
+ return make_error<StringError>(std::string("recv: ") + strerror(errno));
+ }
+ Got += N;
+ }
+ return Got;
+ }
+
+ /// A framed message, ready to write to the socket.
+ static std::vector<char> frame(Opcode Op, uint64_t SeqNo, uint64_t Tag,
+ std::string_view Payload) {
+ std::vector<char> Buf(HeaderSize + Payload.size());
+ MsgHeader::encode(Buf.data(), Op, SeqNo, Tag, Payload.size());
+ if (!Payload.empty())
+ memcpy(Buf.data() + HeaderSize, Payload.data(), Payload.size());
+ return Buf;
+ }
+
+ static Error writeFrame(int FDNum, Opcode Op, uint64_t SeqNo,
+ uint64_t Tag = 0, std::string_view Payload = {}) {
+ auto Buf = frame(Op, SeqNo, Tag, Payload);
+ return sendAll(FDNum, Buf.data(), Buf.size());
+ }
+
+ static Expected<Frame> readFrame(int FDNum) {
+ char H[HeaderSize];
+ auto N = recvAll(FDNum, H, HeaderSize);
+ if (!N)
+ return N.takeError();
+ if (*N != HeaderSize)
+ return make_error<StringError>(
+ "peer closed before a full header arrived");
+
+ Frame F;
+ F.Fields = MsgHeader::decode(H);
+ if (F.Fields.MsgSize < HeaderSize)
+ return make_error<StringError>("framed message size is below the header");
+
+ if (size_t PayloadSize = F.Fields.MsgSize - HeaderSize) {
+ F.Payload.resize(PayloadSize);
+ auto M = recvAll(FDNum, F.Payload.data(), PayloadSize);
+ if (!M)
+ return M.takeError();
+ if (*M != PayloadSize)
+ return make_error<StringError>("peer closed mid-payload");
+ }
+ return F;
+ }
+
+ /// The payload a controller sends to hang up, via the CA's own encoder.
+ static std::vector<char> hangupPayload(Error Err) {
+ auto P = CA::encodeHangup(std::move(Err));
+ return std::vector<char>(P.data(), P.data() + P.size());
+ }
+
+ using ResultKind = CA::ResultKind;
+
+ /// The wire tag naming a wrapper function in this process.
+ static uint64_t wrapperTag(void *Fn) {
+ return ExecutorAddr::fromPtr(Fn).getValue();
+ }
+
+ /// The tag value that marks a result message as carrying a result of Kind.
+ static uint64_t resultTag(ResultKind Kind) {
+ return static_cast<uint64_t>(Kind);
+ }
+
+ /// The tag and payload a controller sends to return ResultBytes, via the CA's
+ /// own encoder.
+ static std::pair<uint64_t, std::vector<char>>
+ resultMessage(WrapperFunctionBuffer ResultBytes) {
+ auto [Kind, P] = CA::encodeResult(std::move(ResultBytes));
+ return {resultTag(Kind), std::vector<char>(P.data(), P.data() + P.size())};
+ }
+
+ /// Decodes a result message the CA sent, as the controller would.
+ static WrapperFunctionBuffer decodeResult(ResultKind Kind,
+ std::string_view Payload) {
+ return CA::decodeResult(
+ Kind, WrapperFunctionBuffer::copyFrom(Payload.data(), Payload.size()));
+ }
+};
+
+} // namespace orc_rt
+
+namespace {
+
+std::string_view view(const std::vector<char> &V) {
+ return {V.data(), V.size()};
+}
+
+// Wrapper that echoes its arguments back as the result.
+void echoWrapper(orc_rt_SessionRef S, orc_rt_WrapperFunctionBuffer ArgBytes,
+ orc_rt_WrapperFunctionReturn Return, uint64_t CallId) {
+ Return(S, ArgBytes, CallId);
+}
+
+// A wrapper that defers its result: it stashes everything needed to return, so
+// a test can complete the call later. A plain function pointer has nowhere to
+// put context, so the caller passes the address of its own DeferredCall as the
+// call's payload.
+//
+// The wrapper runs on the reactor thread and signals completion once its fields
+// are populated. The test can wait on this by calling waitForCall.
+//
+// A caller's DeferredCall has to outlive any chance of the wrapper running,
+// which for a stack one means waiting on waitForCall before leaving the scope.
+struct DeferredCall {
+ orc_rt_SessionRef S = nullptr;
+ orc_rt_WrapperFunctionBuffer ArgBytes{};
+ orc_rt_WrapperFunctionReturn Return = nullptr;
+ uint64_t CallId = 0;
+
+ /// The payload that names this object to deferringWrapper.
+ std::string_view asCallPayload() const {
+ return {reinterpret_cast<const char *>(&Self), sizeof(Self)};
+ }
+
+ void publish(orc_rt_SessionRef S, orc_rt_WrapperFunctionBuffer ArgBytes,
+ orc_rt_WrapperFunctionReturn Return, uint64_t CallId) {
+ this->S = S;
+ this->ArgBytes = ArgBytes;
+ this->Return = Return;
+ this->CallId = CallId;
+ Published.set_value();
+ }
+
+ /// Blocks until the wrapper has run, then returns its return function. The
+ /// other fields are safe to read once this has returned; going through here
+ /// is what makes them so.
+ orc_rt_WrapperFunctionReturn waitForCall() {
+ PublishedF.get();
+ return Return;
+ }
+
+ static void wrapper(orc_rt_SessionRef S,
+ orc_rt_WrapperFunctionBuffer ArgBytes,
+ orc_rt_WrapperFunctionReturn Return, uint64_t CallId) {
+ WrapperFunctionBuffer Args(ArgBytes);
+ assert(Args.size() == sizeof(DeferredCall *) &&
+ "deferringWrapper expects a DeferredCall address as its payload");
+ DeferredCall *D = nullptr;
+ memcpy(&D, Args.data(), sizeof(D));
+
+ // ArgBytes is handed on rather than disposed: it goes back out as the
+ // result.
+ D->publish(S, Args.release(), Return, CallId);
+ }
+
+private:
+ /// Storage for the address, so asCallPayload has something stable to point
+ /// at: a view over a temporary would dangle.
+ DeferredCall *const Self = this;
+
+ // Declared in this order: PublishedF's initializer reads Published.
+ std::promise<void> Published;
+ std::future<void> PublishedF = Published.get_future();
+};
+
+// A wrapper that fails the way the wrapper machinery itself fails: an
+// out-of-band error says the call could not be made sense of, rather than
+// carrying a result the wrapper chose to return.
+void outOfBandErrorWrapper(orc_rt_SessionRef S,
+ orc_rt_WrapperFunctionBuffer ArgBytes,
+ orc_rt_WrapperFunctionReturn Return,
+ uint64_t CallId) {
+ WrapperFunctionBuffer Args(ArgBytes); // Disposed on the way out.
+ Return(S,
+ WrapperFunctionBuffer::createOutOfBandError(
+ "Could not deserialize wrapper function arg data")
+ .release(),
+ CallId);
+}
+
+} // namespace
+
+TEST_F(SimpleRemoteCAOverSocketTest, SetupIsSentOnConnect) {
+ ASSERT_FALSE(!!attachOverSocket());
+
+ auto F = readFrame(Far.get());
+ ASSERT_TRUE(!!F) << toString(F.takeError());
+ EXPECT_EQ(F->opcode(), Opcode::Setup);
+ EXPECT_EQ(F->seqNo(), 0u) << "setup carries no sequence number";
+ EXPECT_EQ(F->tag(), 0u) << "setup carries no handler tag";
+ EXPECT_FALSE(F->Payload.empty()) << "setup carries the bootstrap payload";
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest,
+ ControllerCallIsFramedAndResultCompletesIt) {
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ std::future<std::string> Result;
+ S.callController(
+ [SetResult = waitFor(Result)](WrapperFunctionBuffer R) mutable {
+ SetResult(std::string(R.data(), R.size()));
+ },
+ nullptr, WrapperFunctionBuffer::copyFrom("args", 4));
+
+ auto Call = readFrame(Far.get());
+ ASSERT_TRUE(!!Call) << toString(Call.takeError());
+ EXPECT_EQ(Call->opcode(), Opcode::Call);
+ EXPECT_NE(Call->seqNo(), 0u)
+ << "a call awaiting a result needs a sequence no.";
+ EXPECT_EQ(Call->payload(), "args");
+
+ // Reply under the same sequence number; the handler completes.
+ ASSERT_FALSE(
+ !!writeFrame(Far.get(), Opcode::Result, Call->seqNo(), 0, "reply"));
+ EXPECT_EQ(Result.get(), "reply");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, ControllerInitiatedCallReturnsAResult) {
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ // A Call names the wrapper by tag, and its sequence number is the call id.
+ auto Tag = wrapperTag(reinterpret_cast<void *>(echoWrapper));
+ ASSERT_FALSE(
+ !!writeFrame(Far.get(), Opcode::Call, /*SeqNo=*/42, Tag, "world"));
+
+ auto R = readFrame(Far.get());
+ ASSERT_TRUE(!!R) << toString(R.takeError());
+ EXPECT_EQ(R->opcode(), Opcode::Result);
+ EXPECT_EQ(R->seqNo(), 42u) << "the result must carry the call id back";
+ EXPECT_EQ(R->payload(), "world");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, OutOfBandErrorResultIsFramedAsItsOwnKind) {
+ // An out-of-band error is a pointer to a message, not a serialized value, so
+ // it has no bytes to put in a payload. Before it had a result kind of its own
+ // there was no way to send one at all: framing it aborted the reactor thread.
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ auto Tag = wrapperTag(reinterpret_cast<void *>(outOfBandErrorWrapper));
+ ASSERT_FALSE(!!writeFrame(Far.get(), Opcode::Call, /*SeqNo=*/7, Tag, "junk"));
+
+ auto R = readFrame(Far.get());
+ ASSERT_TRUE(!!R) << toString(R.takeError());
+ EXPECT_EQ(R->opcode(), Opcode::Result);
+ EXPECT_EQ(R->seqNo(), 7u) << "the result must carry the call id back";
+ EXPECT_EQ(R->tag(), resultTag(ResultKind::OutOfBandError));
+
+ // The controller recovers the message the wrapper produced, rather than the
+ // empty result it would otherwise report as a deserialization failure.
+ auto Decoded = decodeResult(ResultKind::OutOfBandError, R->payload());
+ ASSERT_NE(Decoded.getOutOfBandError(), nullptr);
+ EXPECT_STREQ(Decoded.getOutOfBandError(),
+ "Could not deserialize wrapper function arg data");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest,
+ OutOfBandErrorFromControllerCompletesTheCall) {
+ // The other direction: a controller-side wrapper fails the same way, and the
+ // handler waiting on the result must see an out-of-band error rather than an
+ // empty result.
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ std::future<std::string> Result;
+ S.callController(
+ [SetResult = waitFor(Result)](WrapperFunctionBuffer R) mutable {
+ const char *Msg = R.getOutOfBandError();
+ SetResult(Msg ? std::string(Msg) : std::string("<not out-of-band>"));
+ },
+ nullptr, WrapperFunctionBuffer::copyFrom("args", 4));
+
+ auto Call = readFrame(Far.get());
+ ASSERT_TRUE(!!Call) << toString(Call.takeError());
+
+ auto [Tag, Payload] = resultMessage(
+ WrapperFunctionBuffer::createOutOfBandError("controller lost its mind"));
+ ASSERT_FALSE(!!writeFrame(Far.get(), Opcode::Result, Call->seqNo(), Tag,
+ view(Payload)));
+
+ EXPECT_EQ(Result.get(), "controller lost its mind");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, UnknownResultKindEndsTheSession) {
+ // An unrecognized kind means the controller is speaking a dialect this build
+ // does not know, so the payload cannot be interpreted: terminal, as an
+ // unrecognized opcode is.
+ std::future<Error> Disconnected;
+ S.setOnDisconnect(waitFor(Disconnected));
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ std::future<std::string> Result;
+ S.callController(
+ [SetResult = waitFor(Result)](WrapperFunctionBuffer R) mutable {
+ const char *Msg = R.getOutOfBandError();
+ SetResult(Msg ? std::string(Msg) : std::string("<not out-of-band>"));
+ },
+ nullptr, WrapperFunctionBuffer::copyFrom("args", 4));
+
+ auto Call = readFrame(Far.get());
+ ASSERT_TRUE(!!Call) << toString(Call.takeError());
+
+ uint64_t UnknownKind = static_cast<uint64_t>(ResultKind::LastResultKind) + 1;
+ ASSERT_FALSE(
+ !!writeFrame(Far.get(), Opcode::Result, Call->seqNo(), UnknownKind, ""));
+
+ auto Err = Disconnected.get();
+ ASSERT_TRUE(!!Err) << "an uninterpretable result must not end cleanly";
+ EXPECT_EQ(toString(std::move(Err)),
+ "Malformed result message: invalid kind 2");
+
+ // The call it could not answer is still settled on the way out, rather than
+ // left waiting forever.
+ EXPECT_NE(Result.get(), "<not out-of-band>");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, EmptyPayloadsRoundTrip) {
+ // A message whose size is exactly the header. Nothing is pending the moment
+ // the header is decoded, which is the case that stops "header full" and
+ // "header decoded" being the same state, in both directions: the echoed
+ // result is empty too.
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ auto Tag = wrapperTag(reinterpret_cast<void *>(echoWrapper));
+ ASSERT_FALSE(!!writeFrame(Far.get(), Opcode::Call, /*SeqNo=*/8, Tag));
+
+ auto R = readFrame(Far.get());
+ ASSERT_TRUE(!!R) << toString(R.takeError());
+ EXPECT_EQ(R->opcode(), Opcode::Result);
+ EXPECT_EQ(R->seqNo(), 8u);
+ EXPECT_TRUE(R->Payload.empty());
+
+ // Still framing correctly afterwards: an empty message must not desynchronise
+ // the stream.
+ ASSERT_FALSE(
+ !!writeFrame(Far.get(), Opcode::Call, /*SeqNo=*/9, Tag, "after"));
+ auto R2 = readFrame(Far.get());
+ ASSERT_TRUE(!!R2) << toString(R2.takeError());
+ EXPECT_EQ(R2->seqNo(), 9u);
+ EXPECT_EQ(R2->payload(), "after");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, PartialWritesAreReassembled) {
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ // One byte at a time, so the reactor sees the message split in every possible
+ // place -- including mid-header.
+ auto Tag = wrapperTag(reinterpret_cast<void *>(echoWrapper));
+ auto Buf = frame(Opcode::Call, /*SeqNo=*/7, Tag, "drip");
+ for (char C : Buf)
+ ASSERT_FALSE(!!sendAll(Far.get(), &C, 1));
+
+ auto R = readFrame(Far.get());
+ ASSERT_TRUE(!!R) << toString(R.takeError());
+ EXPECT_EQ(R->opcode(), Opcode::Result);
+ EXPECT_EQ(R->seqNo(), 7u);
+ EXPECT_EQ(R->payload(), "drip");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, NothingIsQueuedBehindTheHangup) {
+ // The hang-up must be the last message on the wire, and sendWrapperResult
+ // does no state check -- the base leaves that to the transport. So a result
+ // that arrives once teardown has begun must be dropped rather than appended.
+ //
+ // The window is between beginTeardown queueing the hang-up and the reactor
+ // finishing, which is narrow. It is held open here by never reading the far
+ // end until the very end: the reactor parks in would-block with a part-sent
+ // message, so nothing drains while the late result is submitted.
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ // Echo back more than the socket buffer can hold, so the reactor stalls
+ // part-way through sending the result.
+ const std::string Big(1 << 20, 'x');
+ auto EchoTag = wrapperTag(reinterpret_cast<void *>(echoWrapper));
+ ASSERT_FALSE(
+ !!writeFrame(Far.get(), Opcode::Call, /*SeqNo=*/1, EchoTag, Big));
+
+ // A second call that will not have returned when teardown starts. The
+ // wrapper is told where to park its state by being handed this object's
+ // address as the call payload.
+ DeferredCall Deferred;
+ auto DeferTag = wrapperTag(reinterpret_cast<void *>(DeferredCall::wrapper));
+ ASSERT_FALSE(!!writeFrame(Far.get(), Opcode::Call, /*SeqNo=*/2, DeferTag,
+ Deferred.asCallPayload()));
+
+ // Wait for the wrapper to run without draining the socket.
+ auto DeferredReturn = Deferred.waitForCall();
+
+ // Queues the hang-up and latches the queue. Returns without waiting for the
+ // reactor, which is still stalled.
+ S.detach([] {});
+
+ // Too late: this must not reach the wire.
+ DeferredReturn(Deferred.S, Deferred.ArgBytes, Deferred.CallId);
+
+ // Now drain. The stalled result completes, then the hang-up, then EOF.
+ auto R = readFrame(Far.get());
+ ASSERT_TRUE(!!R) << toString(R.takeError());
+ EXPECT_EQ(R->opcode(), Opcode::Result);
+ EXPECT_EQ(R->Payload.size(), Big.size());
+
+ auto H = readFrame(Far.get());
+ ASSERT_TRUE(!!H) << toString(H.takeError());
+ EXPECT_EQ(H->opcode(), Opcode::Hangup) << "the hang-up did not come last";
+
+ char Byte = 0;
+ auto N = recvAll(Far.get(), &Byte, 1);
+ ASSERT_TRUE(!!N) << toString(N.takeError());
+ EXPECT_EQ(*N, 0u) << "a message was queued behind the hang-up";
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, HangupFromControllerEndsTheSession) {
+ std::future<Error> Disconnected;
+ S.setOnDisconnect(waitFor(Disconnected));
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ // An orderly hang-up carries a success Error as its reason.
+ ASSERT_FALSE(!!writeFrame(Far.get(), Opcode::Hangup, 0, 0,
+ view(hangupPayload(Error::success()))));
+
+ EXPECT_FALSE(!!Disconnected.get()) << "an orderly hang-up is not an error";
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, PeerReasonSurvivesAStalledSendQueue) {
+ // A hang-up reason from the controller must be what ends the session, even
+ // when we have messages queued that can no longer be delivered. Sending first
+ // would fail with EPIPE and report that instead, losing the reason.
+ std::future<Error> Disconnected;
+ S.setOnDisconnect(waitFor(Disconnected));
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ // Echo back more than the socket will hold, and never read it, so the reactor
+ // is left with a part-sent message queued.
+ const std::string Big(1 << 20, 'x');
+ auto Tag = wrapperTag(reinterpret_cast<void *>(echoWrapper));
+ ASSERT_FALSE(!!writeFrame(Far.get(), Opcode::Call, /*SeqNo=*/1, Tag, Big));
+
+ // Hang up with a reason, then vanish. The reason is buffered on our side of
+ // the socket and survives the close.
+ ASSERT_FALSE(!!writeFrame(
+ Far.get(), Opcode::Hangup, 0, 0,
+ view(hangupPayload(make_error<StringError>("controller ran out of x")))));
+ ::close(Far.release());
+
+ auto Err = Disconnected.get();
+ ASSERT_TRUE(!!Err) << "a hang-up carrying a reason ends with that reason";
+ EXPECT_EQ(toString(std::move(Err)), "controller ran out of x");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, TruncatedMessageIsReportedAsAnError) {
+ std::future<Error> Disconnected;
+ S.setOnDisconnect(waitFor(Disconnected));
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ // Half a header, then gone: distinguishable from a close at a boundary.
+ char Half[HeaderSize / 2] = {};
+ ASSERT_FALSE(!!sendAll(Far.get(), Half, sizeof(Half)));
+ ::close(Far.release());
+
+ auto Err = Disconnected.get();
+ EXPECT_TRUE(!!Err) << "a truncated message must not look like a clean end";
+ EXPECT_EQ(toString(std::move(Err)),
+ "Connection closed without a hang-up message");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, PeerCloseWithoutAHangupIsAnError) {
+ // A peer that means to end the session says so with a hang-up. Vanishing
+ // instead means it crashed or was killed, which must not be reported as a
+ // clean shutdown -- a controller that dies has to fail the session.
+ std::future<Error> Disconnected;
+ S.setOnDisconnect(waitFor(Disconnected));
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ // Closed between messages, so nothing is truncated -- it is simply gone.
+ ::close(Far.release());
+
+ auto Err = Disconnected.get();
+ ASSERT_TRUE(!!Err) << "a silent close is not an orderly end";
+ EXPECT_EQ(toString(std::move(Err)),
+ "Connection closed without a hang-up message");
+}
+
+TEST_F(SimpleRemoteCAOverSocketTest, MessageSizeBelowHeaderIsRejected) {
+ std::future<Error> Disconnected;
+ S.setOnDisconnect(waitFor(Disconnected));
+ ASSERT_FALSE(!!attachOverSocket());
+ ASSERT_TRUE(!!readFrame(Far.get())) << "expected setup first";
+
+ // A size that excludes its own header would make the payload length negative.
+ char H[HeaderSize] = {};
+ endian_write<uint64_t>(H, HeaderSize - 1, endian::little);
+ ASSERT_FALSE(!!sendAll(Far.get(), H, sizeof(H)));
+
+ auto Err = Disconnected.get();
+ EXPECT_TRUE(!!Err);
+ EXPECT_EQ(toString(std::move(Err)), "Message size smaller than its header");
+}
More information about the llvm-commits
mailing list