src/corosio/src/local_connect_pair.cpp

100.0% Lines (37 / 37) 100.0% Functions (5 / 5)
local_connect_pair.cpp
f(x) Functions (5)
Line TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Steve Gerbino
3 // Copyright (c) 2026 Michael Vandeberg
4 //
5 // Distributed under the Boost Software License, Version 1.0. (See accompanying
6 // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7 //
8 // Official repository: https://github.com/cppalliance/corosio
9 //
10
11 #include <boost/corosio/local_connect_pair.hpp>
12 #include <boost/corosio/error.hpp>
13 #include <boost/corosio/detail/platform.hpp>
14 #include <boost/corosio/native/detail/make_err.hpp>
15
16 #include <system_error>
17
18 #if BOOST_COROSIO_POSIX
19 #include <fcntl.h>
20 #include <sys/socket.h>
21 #include <sys/un.h>
22 #include <unistd.h>
23 #elif BOOST_COROSIO_HAS_IOCP
24 #include <boost/corosio/native/detail/endpoint_convert.hpp>
25
26 #include <algorithm>
27 #include <atomic>
28 #include <cstring>
29 #include <filesystem>
30 #include <random>
31 #include <string>
32 #include <thread>
33
34 #ifndef WIN32_LEAN_AND_MEAN
35 #define WIN32_LEAN_AND_MEAN
36 #endif
37 #include <winsock2.h>
38
39 #ifndef AF_UNIX
40 #define AF_UNIX 1
41 #endif
42 #endif
43
44 namespace boost::corosio {
45
46 namespace {
47
48 #if BOOST_COROSIO_POSIX
49
50 std::error_code
51 260x make_pair_fds(int type, int& a_fd, int& b_fd) noexcept
52 {
53 int fds[2];
54 260x if (::socketpair(AF_UNIX, type, 0, fds) != 0)
55 18x return detail::make_err(errno);
56
57 // assign() is documented "adopt-only" and will not mutate the fd;
58 // set O_NONBLOCK before transferring ownership.
59 690x for (int i = 0; i < 2; ++i)
60 {
61 466x int flags = ::fcntl(fds[i], F_GETFL, 0);
62 466x if (flags < 0 || ::fcntl(fds[i], F_SETFL, flags | O_NONBLOCK) < 0)
63 {
64 18x auto ec = detail::make_err(errno);
65 18x ::close(fds[0]);
66 18x ::close(fds[1]);
67 18x return ec;
68 }
69 }
70
71 224x a_fd = fds[0];
72 224x b_fd = fds[1];
73 224x return {};
74 }
75
76 template<class Socket>
77 std::error_code
78 224x assign_pair(Socket& a, Socket& b, int a_fd, int b_fd) noexcept
79 {
80 224x if (auto ec = a.assign(a_fd))
81 {
82 9x ::close(a_fd);
83 9x ::close(b_fd);
84 9x return ec;
85 }
86
87 215x if (auto ec = b.assign(b_fd))
88 {
89 9x a.close();
90 9x ::close(b_fd);
91 9x return ec;
92 }
93
94 206x return {};
95 }
96
97 #elif BOOST_COROSIO_HAS_IOCP
98
99 // Build a unique sub-directory under temp and return the full socket
100 // path inside it. Empty string on failure.
101 std::string
102 pick_pair_path(std::filesystem::path& dir_out)
103 {
104 namespace fs = std::filesystem;
105
106 thread_local std::mt19937_64 gen{std::random_device{}()};
107
108 for (int attempt = 0; attempt < 16; ++attempt)
109 {
110 auto candidate =
111 fs::temp_directory_path() / ("co_pair_" + std::to_string(gen()));
112 std::error_code ec;
113 if (fs::create_directory(candidate, ec))
114 {
115 dir_out = candidate;
116 return (candidate / "s").string();
117 }
118 }
119 return {};
120 }
121
122 void
123 remove_pair_path(std::filesystem::path const& dir, std::string const& path)
124 {
125 std::error_code ec;
126 std::filesystem::remove(std::filesystem::path(path), ec);
127 std::filesystem::remove(dir, ec);
128 }
129
130 // Synchronously rendezvous two AF_UNIX SOCK_STREAM sockets. The
131 // listener and accept happen on the caller's thread; the connect
132 // runs on a short-lived worker to avoid a deadlock. The returned
133 // sockets are created with WSA_FLAG_OVERLAPPED so they can be
134 // registered with IOCP by assign_socket().
135 std::error_code
136 make_pair_sockets(SOCKET& a_sock, SOCKET& b_sock) noexcept
137 {
138 namespace fs = std::filesystem;
139
140 a_sock = INVALID_SOCKET;
141 b_sock = INVALID_SOCKET;
142
143 fs::path dir;
144 std::string path = pick_pair_path(dir);
145 if (path.empty())
146 return detail::make_err(ERROR_PATH_NOT_FOUND);
147
148 SOCKET listen_sock =
149 ::WSASocketW(AF_UNIX, SOCK_STREAM, 0, nullptr, 0, WSA_FLAG_OVERLAPPED);
150 if (listen_sock == INVALID_SOCKET)
151 {
152 auto ec = detail::make_err(::WSAGetLastError());
153 remove_pair_path(dir, path);
154 return ec;
155 }
156
157 detail::un_sa_t addr{};
158 addr.sun_family = AF_UNIX;
159 std::memcpy(
160 addr.sun_path, path.c_str(),
161 (std::min)(path.size(), sizeof(addr.sun_path) - 1));
162 int addr_len =
163 static_cast<int>(offsetof(detail::un_sa_t, sun_path) + path.size() + 1);
164
165 if (::bind(listen_sock, reinterpret_cast<sockaddr*>(&addr), addr_len) ==
166 SOCKET_ERROR)
167 {
168 auto ec = detail::make_err(::WSAGetLastError());
169 ::closesocket(listen_sock);
170 remove_pair_path(dir, path);
171 return ec;
172 }
173
174 if (::listen(listen_sock, 1) == SOCKET_ERROR)
175 {
176 auto ec = detail::make_err(::WSAGetLastError());
177 ::closesocket(listen_sock);
178 remove_pair_path(dir, path);
179 return ec;
180 }
181
182 // A worker that fails before connecting produces no connection at
183 // all, so the accept below must be able to give up. Poll the
184 // listener instead of blocking in accept() forever.
185 u_long non_blocking = 1;
186 if (::ioctlsocket(listen_sock, FIONBIO, &non_blocking) == SOCKET_ERROR)
187 {
188 auto ec = detail::make_err(::WSAGetLastError());
189 ::closesocket(listen_sock);
190 remove_pair_path(dir, path);
191 return ec;
192 }
193
194 SOCKET worker_sock = INVALID_SOCKET;
195 std::error_code worker_ec;
196 std::atomic<bool> worker_done{false};
197
198 // One exit, so worker_done is published on every path: the accept
199 // below waits on it, and a path that skipped it would hang.
200 std::thread worker([&] {
201 worker_sock = ::WSASocketW(
202 AF_UNIX, SOCK_STREAM, 0, nullptr, 0, WSA_FLAG_OVERLAPPED);
203 if (worker_sock == INVALID_SOCKET)
204 {
205 worker_ec = detail::make_err(::WSAGetLastError());
206 }
207 else
208 {
209 detail::un_sa_t caddr{};
210 caddr.sun_family = AF_UNIX;
211 std::memcpy(
212 caddr.sun_path, path.c_str(),
213 (std::min)(path.size(), sizeof(caddr.sun_path) - 1));
214 int caddr_len = static_cast<int>(
215 offsetof(detail::un_sa_t, sun_path) + path.size() + 1);
216
217 if (::connect(
218 worker_sock, reinterpret_cast<sockaddr*>(&caddr),
219 caddr_len) == SOCKET_ERROR)
220 {
221 worker_ec = detail::make_err(::WSAGetLastError());
222 ::closesocket(worker_sock);
223 worker_sock = INVALID_SOCKET;
224 }
225 }
226 // Released last so a reader that sees it also sees worker_ec.
227 worker_done.store(true, std::memory_order_release);
228 });
229
230 SOCKET accept_sock = INVALID_SOCKET;
231 std::error_code accept_ec;
232 for (;;)
233 {
234 // A worker that succeeded has left a connection in the
235 // backlog, so only a failed one means nothing is coming.
236 if (worker_done.load(std::memory_order_acquire) && worker_ec)
237 break;
238
239 WSAPOLLFD pfd{listen_sock, POLLRDNORM, 0};
240 int const n = ::WSAPoll(&pfd, 1, 100);
241 if (n == SOCKET_ERROR)
242 {
243 accept_ec = detail::make_err(::WSAGetLastError());
244 break;
245 }
246 if (n == 0)
247 continue;
248
249 // Readiness that is not "a connection is waiting" is an error
250 // condition on the listener; accepting on it would spin. The
251 // condition carries no retrievable code, so this one is
252 // corosio's own and has to compare equal on every toolchain.
253 if ((pfd.revents & POLLRDNORM) == 0)
254 {
255 accept_ec = std::make_error_code(std::errc::connection_aborted);
256 break;
257 }
258
259 accept_sock = ::accept(listen_sock, nullptr, nullptr);
260 if (accept_sock != INVALID_SOCKET)
261 break;
262 DWORD const err = ::WSAGetLastError();
263 // A readiness report with nothing left to accept: keep
264 // waiting for the worker's connection.
265 if (err == WSAEWOULDBLOCK)
266 continue;
267 accept_ec = detail::make_err(err);
268 break;
269 }
270
271 worker.join();
272
273 ::closesocket(listen_sock);
274 remove_pair_path(dir, path);
275
276 if (accept_ec)
277 {
278 if (worker_sock != INVALID_SOCKET)
279 ::closesocket(worker_sock);
280 return accept_ec;
281 }
282 if (worker_ec)
283 {
284 if (accept_sock != INVALID_SOCKET)
285 ::closesocket(accept_sock);
286 return worker_ec;
287 }
288
289 // accept() inherits the listener's non-blocking mode; the rest of
290 // the IOCP backend hands out blocking sockets and drives them
291 // through overlapped I/O.
292 non_blocking = 0;
293 if (::ioctlsocket(accept_sock, FIONBIO, &non_blocking) == SOCKET_ERROR)
294 {
295 auto ec = detail::make_err(::WSAGetLastError());
296 ::closesocket(accept_sock);
297 ::closesocket(worker_sock);
298 return ec;
299 }
300
301 a_sock = accept_sock;
302 b_sock = worker_sock;
303 return {};
304 }
305
306 std::error_code
307 assign_pair(
308 local_stream_socket& a,
309 local_stream_socket& b,
310 SOCKET a_sock,
311 SOCKET b_sock) noexcept
312 {
313 if (auto ec = a.assign(static_cast<native_handle_type>(a_sock)))
314 {
315 ::closesocket(a_sock);
316 ::closesocket(b_sock);
317 return ec;
318 }
319
320 if (auto ec = b.assign(static_cast<native_handle_type>(b_sock)))
321 {
322 a.close();
323 ::closesocket(b_sock);
324 return ec;
325 }
326
327 return {};
328 }
329
330 #endif
331
332 } // namespace
333
334 std::error_code
335 155x connect_pair(local_stream_socket& a, local_stream_socket& b) noexcept
336 {
337 155x if (a.is_open() || b.is_open())
338 3x return make_error_code(error::already_open);
339
340 #if BOOST_COROSIO_POSIX
341 152x int a_fd = -1, b_fd = -1;
342 152x if (auto ec = make_pair_fds(SOCK_STREAM, a_fd, b_fd))
343 27x return ec;
344 125x return assign_pair(a, b, a_fd, b_fd);
345 #elif BOOST_COROSIO_HAS_IOCP
346 SOCKET a_sock = INVALID_SOCKET, b_sock = INVALID_SOCKET;
347 if (auto ec = make_pair_sockets(a_sock, b_sock))
348 return ec;
349 return assign_pair(a, b, a_sock, b_sock);
350 #else
351 return detail::make_err(ENOSYS);
352 #endif
353 }
354
355 #if BOOST_COROSIO_POSIX
356
357 std::error_code
358 111x connect_pair(local_datagram_socket& a, local_datagram_socket& b) noexcept
359 {
360 111x if (a.is_open() || b.is_open())
361 3x return make_error_code(error::already_open);
362
363 108x int a_fd = -1, b_fd = -1;
364 108x if (auto ec = make_pair_fds(SOCK_DGRAM, a_fd, b_fd))
365 9x return ec;
366 99x return assign_pair(a, b, a_fd, b_fd);
367 }
368
369 #endif
370
371 } // namespace boost::corosio
372