LLVM 24.0.0git
raw_socket_stream.cpp
Go to the documentation of this file.
1//===-- llvm/Support/raw_socket_stream.cpp - Socket streams --*- C++ -*-===//
2//
3// Part of the LLVM Project, under the Apache License v2.0 with LLVM Exceptions.
4// See https://llvm.org/LICENSE.txt for license information.
5// SPDX-License-Identifier: Apache-2.0 WITH LLVM-exception
6//
7//===----------------------------------------------------------------------===//
8//
9// This file contains raw_ostream implementations for streams to communicate
10// via UNIX sockets
11//
12//===----------------------------------------------------------------------===//
13
15#include "llvm/Config/config.h"
16#include "llvm/Support/Error.h"
18
19#include <atomic>
20#include <fcntl.h>
21#include <functional>
22
23#ifndef _WIN32
24#include <poll.h>
25#include <sys/socket.h>
26#include <sys/un.h>
27#else
29// winsock2.h must be included before afunix.h. Briefly turn off clang-format to
30// avoid error.
31// clang-format off
32#include <winsock2.h>
33#include <afunix.h>
34// clang-format on
35#include <io.h>
36#endif // _WIN32
37
38#if defined(HAVE_UNISTD_H)
39#include <unistd.h>
40#endif
41
42using namespace llvm;
43
44#ifdef _WIN32
45WSABalancer::WSABalancer() {
46 WSADATA WsaData;
47 ::memset(&WsaData, 0, sizeof(WsaData));
48 if (WSAStartup(MAKEWORD(2, 2), &WsaData) != 0) {
49 llvm::report_fatal_error("WSAStartup failed");
50 }
51}
52
53WSABalancer::~WSABalancer() { WSACleanup(); }
54#endif // _WIN32
55
56static std::error_code getLastSocketErrorCode() {
57#ifdef _WIN32
58 return std::error_code(::WSAGetLastError(), std::system_category());
59#else
60 return errnoAsErrorCode();
61#endif
62}
63
64#ifdef _WIN32
65using NativeSocket = SOCKET;
66#else
67using NativeSocket = int;
68#define INVALID_SOCKET -1
69#endif
70
71static int closeSocket(NativeSocket Socket) {
72#ifdef _WIN32
73 return ::closesocket(Socket);
74#else
75 return ::close(Socket);
76#endif
77}
78
80 struct sockaddr_un Addr;
81 memset(&Addr, 0, sizeof(Addr));
82 Addr.sun_family = AF_UNIX;
83
84 if (sizeof(sockaddr_un::sun_path) <= SocketPath.size())
86 std::make_error_code(std::errc::filename_too_long),
87 "Socket path exceeds sockaddr_un::sun_path size limit");
88
89 strncpy(Addr.sun_path, SocketPath.str().c_str(), sizeof(Addr.sun_path) - 1);
90 return Addr;
91}
92
94 NativeSocket Socket = socket(AF_UNIX, SOCK_STREAM, 0);
95 if (Socket == INVALID_SOCKET) {
97 "Create socket failed");
98 }
99
100#ifdef __CYGWIN__
101 // On Cygwin, UNIX sockets involve a handshake between connect and accept
102 // to enable SO_PEERCRED/getpeereid handling. This necessitates accept being
103 // called before connect can return, but at least the tests in
104 // llvm/unittests/Support/raw_socket_stream_test do both on the same thread
105 // (first connect and then accept), resulting in a deadlock. This call turns
106 // off the handshake (and SO_PEERCRED/getpeereid support).
107 setsockopt(Socket, SOL_SOCKET, SO_PEERCRED, NULL, 0);
108#endif
110 if (!Addr) {
111 closeSocket(Socket);
112 return Addr.takeError();
113 }
114
115 if (::connect(Socket, (struct sockaddr *)&*Addr, sizeof(*Addr)) == -1) {
116 // Grab the error code before closing, which may overwrite it.
117 std::error_code EC = getLastSocketErrorCode();
118 closeSocket(Socket);
119 return llvm::make_error<StringError>(EC, "Connect socket failed");
120 }
121
122#ifdef _WIN32
123 return _open_osfhandle(Socket, 0);
124#else
125 return Socket;
126#endif // _WIN32
127}
128
129ListeningSocket::ListeningSocket(int SocketFD, StringRef SocketPath,
130 int PipeFD[2])
131 : FD(SocketFD), SocketPath(SocketPath), PipeFD{PipeFD[0], PipeFD[1]} {}
132
133ListeningSocket::ListeningSocket(ListeningSocket &&LS)
134 : FD(LS.FD.load()), SocketPath(LS.SocketPath),
135 PipeFD{LS.PipeFD[0], LS.PipeFD[1]} {
136
137 LS.FD = -1;
138 LS.SocketPath.clear();
139 LS.PipeFD[0] = -1;
140 LS.PipeFD[1] = -1;
141}
142
144 int MaxBacklog) {
145
146 // Handle instances where the target socket address already exists and
147 // differentiate between a preexisting file with and without a bound socket
148 //
149 // ::bind will return std::errc:address_in_use if a file at the socket address
150 // already exists (e.g., the file was not properly unlinked due to a crash)
151 // even if another socket has not yet binded to that address
152 if (llvm::sys::fs::exists(SocketPath)) {
153 Expected<int> MaybeFD = getSocketFD(SocketPath);
154 if (!MaybeFD) {
155
156 // Regardless of the error, notify the caller that a file already exists
157 // at the desired socket address and that there is no bound socket at that
158 // address. The file must be removed before ::bind can use the address
159 consumeError(MaybeFD.takeError());
161 std::make_error_code(std::errc::file_exists),
162 "Socket address unavailable");
163 }
164 ::close(std::move(*MaybeFD));
165
166 // Notify caller that the provided socket address already has a bound socket
168 std::make_error_code(std::errc::address_in_use),
169 "Socket address unavailable");
170 }
171
172#ifdef _WIN32
173 WSABalancer _;
174#endif
175 NativeSocket Socket = socket(AF_UNIX, SOCK_STREAM, 0);
176 if (Socket == INVALID_SOCKET)
178 "socket create failed");
179
180#ifdef __CYGWIN__
181 // On Cygwin, UNIX sockets involve a handshake between connect and accept
182 // to enable SO_PEERCRED/getpeereid handling. This necessitates accept being
183 // called before connect can return, but at least the tests in
184 // llvm/unittests/Support/raw_socket_stream_test do both on the same thread
185 // (first connect and then accept), resulting in a deadlock. This call turns
186 // off the handshake (and SO_PEERCRED/getpeereid support).
187 setsockopt(Socket, SOL_SOCKET, SO_PEERCRED, NULL, 0);
188#endif
190 if (!Addr) {
191 closeSocket(Socket);
192 return Addr.takeError();
193 }
194
195 if (::bind(Socket, (struct sockaddr *)&*Addr, sizeof(*Addr)) == -1) {
196 // Grab error code from call to ::bind before closing the socket
197 std::error_code EC = getLastSocketErrorCode();
198 closeSocket(Socket);
199 return llvm::make_error<StringError>(EC, "Bind error");
200 }
201
202 // Mark socket as passive so incoming connections can be accepted
203 if (::listen(Socket, MaxBacklog) == -1)
205 "Listen error");
206
207 int PipeFD[2];
208#ifdef _WIN32
209 // Reserve 1 byte for the pipe and use default textmode
210 if (::_pipe(PipeFD, 1, 0) == -1)
211#else
212 if (::pipe(PipeFD) == -1)
213#endif // _WIN32
215 "pipe failed");
216
217#ifdef _WIN32
218 return ListeningSocket{_open_osfhandle(Socket, 0), SocketPath, PipeFD};
219#else
220 return ListeningSocket{Socket, SocketPath, PipeFD};
221#endif // _WIN32
222}
223
224// If a file descriptor being monitored by ::poll is closed by another thread,
225// the result is unspecified. In the case ::poll does not unblock and return,
226// when ActiveFD is closed, you can provide another file descriptor via CancelFD
227// that when written to will cause poll to return. Typically CancelFD is the
228// read end of a unidirectional pipe.
229//
230// Timeout should be -1 to block indefinitly
231//
232// getActiveFD is a callback to handle ActiveFD's of std::atomic<int> and int
233static std::error_code
234manageTimeout(const std::chrono::milliseconds &Timeout,
235 const std::function<int()> &getActiveFD,
236 const std::optional<int> &CancelFD = std::nullopt) {
237 struct pollfd FD[2];
238 FD[0].events = POLLIN;
239#ifdef _WIN32
240 SOCKET WinServerSock = _get_osfhandle(getActiveFD());
241 FD[0].fd = WinServerSock;
242#else
243 FD[0].fd = getActiveFD();
244#endif
245 uint8_t FDCount = 1;
246 if (CancelFD.has_value()) {
247 FD[1].events = POLLIN;
248 FD[1].fd = CancelFD.value();
249 FDCount++;
250 }
251
252 // Keep track of how much time has passed in case ::poll or WSAPoll are
253 // interupted by a signal and need to be recalled
254 auto Start = std::chrono::steady_clock::now();
255 auto RemainingTimeout = Timeout;
256 int PollStatus = 0;
257 do {
258 // If Timeout is -1 then poll should block and RemainingTimeout does not
259 // need to be recalculated
260 if (PollStatus != 0 && Timeout != std::chrono::milliseconds(-1)) {
261 auto TotalElapsedTime =
262 std::chrono::duration_cast<std::chrono::milliseconds>(
263 std::chrono::steady_clock::now() - Start);
264
265 if (TotalElapsedTime >= Timeout)
266 return std::make_error_code(std::errc::operation_would_block);
267
268 RemainingTimeout = Timeout - TotalElapsedTime;
269 }
270#ifdef _WIN32
271 PollStatus = WSAPoll(FD, FDCount, RemainingTimeout.count());
272 } while (PollStatus == SOCKET_ERROR &&
273 getLastSocketErrorCode() == std::errc::interrupted);
274#else
275 PollStatus = ::poll(FD, FDCount, RemainingTimeout.count());
276 } while (PollStatus == -1 &&
277 getLastSocketErrorCode() == std::errc::interrupted);
278#endif
279
280 // If ActiveFD equals -1 or CancelFD has data to be read then the operation
281 // has been canceled by another thread
282 if (getActiveFD() == -1 || (CancelFD.has_value() && FD[1].revents & POLLIN))
283 return std::make_error_code(std::errc::operation_canceled);
284#ifdef _WIN32
285 if (PollStatus == SOCKET_ERROR)
286#else
287 if (PollStatus == -1)
288#endif
289 return getLastSocketErrorCode();
290 if (PollStatus == 0)
291 return std::make_error_code(std::errc::timed_out);
292 if (FD[0].revents & POLLNVAL)
293 return std::make_error_code(std::errc::bad_file_descriptor);
294 return std::error_code();
295}
296
298ListeningSocket::accept(const std::chrono::milliseconds &Timeout) {
299 auto getActiveFD = [this]() -> int { return FD; };
300 std::error_code TimeoutErr = manageTimeout(Timeout, getActiveFD, PipeFD[0]);
301 if (TimeoutErr)
302 return llvm::make_error<StringError>(TimeoutErr, "Timeout error");
303
304 int AcceptFD;
305#ifdef _WIN32
306 SOCKET WinAcceptSock = ::accept(_get_osfhandle(FD), NULL, NULL);
307 AcceptFD = _open_osfhandle(WinAcceptSock, 0);
308#else
309 AcceptFD = ::accept(FD, NULL, NULL);
310#endif
311
312 if (AcceptFD == -1)
314 "Socket accept failed");
315 return std::make_unique<raw_socket_stream>(AcceptFD);
316}
317
319 int ObservedFD = FD.load();
320
321 if (ObservedFD == -1)
322 return;
323
324 // If FD equals ObservedFD set FD to -1; If FD doesn't equal ObservedFD then
325 // another thread is responsible for shutdown so return
326 if (!FD.compare_exchange_strong(ObservedFD, -1))
327 return;
328
329 ::close(ObservedFD);
330 ::unlink(SocketPath.c_str());
331
332 // Ensure ::poll returns if shutdown is called by a separate thread
333 char Byte = 'A';
334 ssize_t written = ::write(PipeFD[1], &Byte, 1);
335
336 // Ignore any write() error
337 (void)written;
338}
339
341 shutdown();
342
343 // Close the pipe's FDs in the destructor instead of within
344 // ListeningSocket::shutdown to avoid unnecessary synchronization issues that
345 // would occur as PipeFD's values would have to be changed to -1
346 //
347 // The move constructor sets PipeFD to -1
348 if (PipeFD[0] != -1)
349 ::close(PipeFD[0]);
350 if (PipeFD[1] != -1)
351 ::close(PipeFD[1]);
352}
353
354//===----------------------------------------------------------------------===//
355// raw_socket_stream
356//===----------------------------------------------------------------------===//
357
360
362
365#ifdef _WIN32
366 WSABalancer _;
367#endif // _WIN32
368 Expected<int> FD = getSocketFD(SocketPath);
369 if (!FD)
370 return FD.takeError();
371 return std::make_unique<raw_socket_stream>(*FD);
372}
373
374ssize_t raw_socket_stream::read(char *Ptr, size_t Size,
375 const std::chrono::milliseconds &Timeout) {
376 auto getActiveFD = [this]() -> int { return this->get_fd(); };
377 std::error_code Err = manageTimeout(Timeout, getActiveFD);
378 // Mimic raw_fd_stream::read error handling behavior
379 if (Err) {
381 return -1;
382 }
383 return raw_fd_stream::read(Ptr, Size);
384}
AMDGPU Mark last scratch load
#define _
Tagged union holding either a T or a Error.
Definition Error.h:485
Error takeError()
Take ownership of the stored error.
Definition Error.h:612
static LLVM_ABI Expected< ListeningSocket > createUnix(StringRef SocketPath, int MaxBacklog=llvm::hardware_concurrency().compute_thread_count())
Creates a listening socket bound to the specified file system path.
LLVM_ABI void shutdown()
Closes the FD, unlinks the socket file, and writes to PipeFD.
LLVM_ABI Expected< std::unique_ptr< raw_socket_stream > > accept(const std::chrono::milliseconds &Timeout=std::chrono::milliseconds(-1))
Accepts an incoming connection on the listening socket.
Represent a constant reference to a string, i.e.
Definition StringRef.h:56
std::string str() const
Get the contents as an std::string.
Definition StringRef.h:222
constexpr size_t size() const
Get the string size.
Definition StringRef.h:144
int get_fd() const
Return the file descriptor.
void error_detected(std::error_code EC)
Set the flag indicating that an output error has been encountered.
LLVM_ABI raw_fd_stream(StringRef Filename, std::error_code &EC)
Open the specified file for reading/writing/seeking.
LLVM_ABI ssize_t read(char *Ptr, size_t Size)
This reads the Size bytes into a buffer pointed by Ptr.
static Expected< std::unique_ptr< raw_socket_stream > > createConnectedUnix(StringRef SocketPath)
Create a raw_socket_stream connected to the UNIX domain socket at SocketPath.
~raw_socket_stream() override
ssize_t read(char *Ptr, size_t Size, const std::chrono::milliseconds &Timeout=std::chrono::milliseconds(-1))
Attempt to read from the raw_socket_stream's file descriptor.
LLVM_ABI bool exists(const basic_file_status &status)
Does file exist?
Definition Path.cpp:1107
This is an optimization pass for GlobalISel generic memory operations.
LLVM_ABI void report_fatal_error(Error Err, bool gen_crash_diag=true)
Definition Error.cpp:163
@ Timeout
Reached timeout while waiting for the owner to release the lock.
Error make_error(ArgTs &&... Args)
Make a Error instance representing failure using the given error info type.
Definition Error.h:340
std::error_code errnoAsErrorCode()
Helper to get errno as an std::error_code.
Definition Error.h:1256
void consumeError(Error Err)
Consume a Error without doing anything.
Definition Error.h:1106
LLVM_ABI Error write(DWPWriter &Out, ArrayRef< std::string > Inputs, OnCuIndexOverflow OverflowOptValue, Dwarf64StrOffsetsPromotion StrOffsetsOptValue, raw_pwrite_stream *OS=nullptr)
Definition DWP.cpp:746
static Expected< int > getSocketFD(StringRef SocketPath)
int NativeSocket
#define INVALID_SOCKET
static Expected< sockaddr_un > setSocketAddr(StringRef SocketPath)
static int closeSocket(NativeSocket Socket)
static std::error_code getLastSocketErrorCode()
static std::error_code manageTimeout(const std::chrono::milliseconds &Timeout, const std::function< int()> &getActiveFD, const std::optional< int > &CancelFD=std::nullopt)