src/corosio/src/local_connect_pair.cpp

96.2% Lines (125/130) 100.0% List of functions (6/6) 90.7% Branches (68/75)
local_connect_pair.cpp
f(x) Functions (6)
Line Branch 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 make_pair_fds(int type, int& a_fd, int& b_fd) noexcept
52 {
53 int fds[2];
54 if (::socketpair(AF_UNIX, type, 0, fds) != 0)
55 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 for (int i = 0; i < 2; ++i)
60 {
61 int flags = ::fcntl(fds[i], F_GETFL, 0);
62 if (flags < 0 || ::fcntl(fds[i], F_SETFL, flags | O_NONBLOCK) < 0)
63 {
64 auto ec = detail::make_err(errno);
65 ::close(fds[0]);
66 ::close(fds[1]);
67 return ec;
68 }
69 }
70
71 a_fd = fds[0];
72 b_fd = fds[1];
73 return {};
74 }
75
76 template<class Socket>
77 std::error_code
78 assign_pair(Socket& a, Socket& b, int a_fd, int b_fd) noexcept
79 {
80 if (auto ec = a.assign(a_fd))
81 {
82 ::close(a_fd);
83 ::close(b_fd);
84 return ec;
85 }
86
87 if (auto ec = b.assign(b_fd))
88 {
89 a.close();
90 ::close(b_fd);
91 return ec;
92 }
93
94 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 83x pick_pair_path(std::filesystem::path& dir_out)
103 {
104 namespace fs = std::filesystem;
105
106
5/5
✓ Branch 2 → 3 taken 32 times.
✓ Branch 2 → 8 taken 51 times.
✓ Branch 3 → 4 taken 32 times.
✓ Branch 4 → 5 taken 32 times.
✓ Branch 5 → 6 taken 32 times.
83x thread_local std::mt19937_64 gen{std::random_device{}()};
107
108
1/2
✓ Branch 36 → 9 taken 83 times.
✗ Branch 36 → 37 not taken.
83x for (int attempt = 0; attempt < 16; ++attempt)
109 {
110 auto candidate =
111
6/6
✓ Branch 9 → 10 taken 83 times.
✓ Branch 10 → 11 taken 83 times.
✓ Branch 11 → 12 taken 83 times.
✓ Branch 12 → 13 taken 83 times.
✓ Branch 13 → 14 taken 83 times.
✓ Branch 14 → 15 taken 83 times.
83x fs::temp_directory_path() / ("co_pair_" + std::to_string(gen()));
112 83x std::error_code ec;
113
1/2
✓ Branch 21 → 22 taken 83 times.
✗ Branch 21 → 30 not taken.
83x if (fs::create_directory(candidate, ec))
114 {
115
1/1
✓ Branch 22 → 23 taken 83 times.
83x dir_out = candidate;
116
3/3
✓ Branch 23 → 24 taken 83 times.
✓ Branch 24 → 25 taken 83 times.
✓ Branch 25 → 26 taken 83 times.
83x return (candidate / "s").string();
117 }
118 83x }
119 ✗ return {};
120 }
121
122 void
123 83x remove_pair_path(std::filesystem::path const& dir, std::string const& path)
124 {
125 83x std::error_code ec;
126
1/1
✓ Branch 3 → 4 taken 83 times.
83x std::filesystem::remove(std::filesystem::path(path), ec);
127 83x std::filesystem::remove(dir, ec);
128 83x }
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 83x make_pair_sockets(SOCKET& a_sock, SOCKET& b_sock) noexcept
137 {
138 namespace fs = std::filesystem;
139
140 83x a_sock = INVALID_SOCKET;
141 83x b_sock = INVALID_SOCKET;
142
143 83x fs::path dir;
144 83x std::string path = pick_pair_path(dir);
145
1/2
✗ Branch 5 → 6 not taken.
✓ Branch 5 → 7 taken 83 times.
83x if (path.empty())
146 ✗ return detail::make_err(ERROR_PATH_NOT_FOUND);
147
148 SOCKET listen_sock =
149 83x ::WSASocketW(AF_UNIX, SOCK_STREAM, 0, nullptr, 0, WSA_FLAG_OVERLAPPED);
150
2/2
✓ Branch 8 → 9 taken 1 time.
✓ Branch 8 → 13 taken 82 times.
83x if (listen_sock == INVALID_SOCKET)
151 {
152 1x auto ec = detail::make_err(::WSAGetLastError());
153 1x remove_pair_path(dir, path);
154 1x return ec;
155 }
156
157 82x detail::un_sa_t addr{};
158 82x addr.sun_family = AF_UNIX;
159 164x std::memcpy(
160 82x addr.sun_path, path.c_str(),
161 82x (std::min)(path.size(), sizeof(addr.sun_path) - 1));
162 int addr_len =
163 82x static_cast<int>(offsetof(detail::un_sa_t, sun_path) + path.size() + 1);
164
165
2/2
✓ Branch 18 → 19 taken 1 time.
✓ Branch 18 → 24 taken 81 times.
82x if (::bind(listen_sock, reinterpret_cast<sockaddr*>(&addr), addr_len) ==
166 SOCKET_ERROR)
167 {
168 1x auto ec = detail::make_err(::WSAGetLastError());
169 1x ::closesocket(listen_sock);
170 1x remove_pair_path(dir, path);
171 1x return ec;
172 }
173
174
2/2
✓ Branch 25 → 26 taken 1 time.
✓ Branch 25 → 31 taken 80 times.
81x if (::listen(listen_sock, 1) == SOCKET_ERROR)
175 {
176 1x auto ec = detail::make_err(::WSAGetLastError());
177 1x ::closesocket(listen_sock);
178 1x remove_pair_path(dir, path);
179 1x 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 80x u_long non_blocking = 1;
186
2/2
✓ Branch 32 → 33 taken 1 time.
✓ Branch 32 → 38 taken 79 times.
80x if (::ioctlsocket(listen_sock, FIONBIO, &non_blocking) == SOCKET_ERROR)
187 {
188 1x auto ec = detail::make_err(::WSAGetLastError());
189 1x ::closesocket(listen_sock);
190 1x remove_pair_path(dir, path);
191 1x return ec;
192 }
193
194 79x SOCKET worker_sock = INVALID_SOCKET;
195 79x std::error_code worker_ec;
196 79x 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 237x std::thread worker([&] {
201 79x worker_sock = ::WSASocketW(
202 AF_UNIX, SOCK_STREAM, 0, nullptr, 0, WSA_FLAG_OVERLAPPED);
203
2/2
✓ Branch 3 → 4 taken 9 times.
✓ Branch 3 → 6 taken 70 times.
79x if (worker_sock == INVALID_SOCKET)
204 {
205 9x worker_ec = detail::make_err(::WSAGetLastError());
206 }
207 else
208 {
209 70x detail::un_sa_t caddr{};
210 70x caddr.sun_family = AF_UNIX;
211 140x std::memcpy(
212 70x caddr.sun_path, path.c_str(),
213 70x (std::min)(path.size(), sizeof(caddr.sun_path) - 1));
214 int caddr_len = static_cast<int>(
215 70x offsetof(detail::un_sa_t, sun_path) + path.size() + 1);
216
217
1/1
✓ Branch 10 → 11 taken 70 times.
70x if (::connect(
218 worker_sock, reinterpret_cast<sockaddr*>(&caddr),
219
2/2
✓ Branch 11 → 12 taken 9 times.
✓ Branch 11 → 16 taken 61 times.
70x caddr_len) == SOCKET_ERROR)
220 {
221
1/1
✓ Branch 12 → 13 taken 9 times.
9x worker_ec = detail::make_err(::WSAGetLastError());
222
1/1
✓ Branch 14 → 15 taken 9 times.
9x ::closesocket(worker_sock);
223 9x worker_sock = INVALID_SOCKET;
224 }
225 }
226 // Released last so a reader that sees it also sees worker_ec.
227 79x worker_done.store(true, std::memory_order_release);
228 158x });
229
230 79x SOCKET accept_sock = INVALID_SOCKET;
231 79x 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
6/6
✓ Branch 42 → 43 taken 19 times.
✓ Branch 42 → 46 taken 79 times.
✓ Branch 44 → 45 taken 18 times.
✓ Branch 44 → 46 taken 1 time.
✓ Branch 47 → 48 taken 18 times.
✓ Branch 47 → 49 taken 80 times.
98x if (worker_done.load(std::memory_order_acquire) && worker_ec)
237 18x break;
238
239 80x WSAPOLLFD pfd{listen_sock, POLLRDNORM, 0};
240 80x int const n = ::WSAPoll(&pfd, 1, 100);
241
2/2
✓ Branch 50 → 51 taken 9 times.
✓ Branch 50 → 54 taken 71 times.
80x if (n == SOCKET_ERROR)
242 {
243 9x accept_ec = detail::make_err(::WSAGetLastError());
244 9x break;
245 }
246
2/2
✓ Branch 54 → 55 taken 18 times.
✓ Branch 54 → 56 taken 53 times.
71x if (n == 0)
247 18x 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
1/2
✗ Branch 56 → 57 not taken.
✓ Branch 56 → 59 taken 53 times.
53x if ((pfd.revents & POLLRDNORM) == 0)
254 {
255 ✗ accept_ec = std::make_error_code(std::errc::connection_aborted);
256 ✗ break;
257 }
258
259 53x accept_sock = ::accept(listen_sock, nullptr, nullptr);
260
2/2
✓ Branch 60 → 61 taken 51 times.
✓ Branch 60 → 62 taken 2 times.
53x if (accept_sock != INVALID_SOCKET)
261 51x break;
262 2x DWORD const err = ::WSAGetLastError();
263 // A readiness report with nothing left to accept: keep
264 // waiting for the worker's connection.
265
2/2
✓ Branch 63 → 64 taken 1 time.
✓ Branch 63 → 65 taken 1 time.
2x if (err == WSAEWOULDBLOCK)
266 1x continue;
267 1x accept_ec = detail::make_err(err);
268 1x break;
269 19x }
270
271 79x worker.join();
272
273 79x ::closesocket(listen_sock);
274 79x remove_pair_path(dir, path);
275
276
2/2
✓ Branch 71 → 73 taken 10 times.
✓ Branch 71 → 76 taken 69 times.
79x if (accept_ec)
277 {
278
1/2
✓ Branch 73 → 74 taken 10 times.
✗ Branch 73 → 75 not taken.
10x if (worker_sock != INVALID_SOCKET)
279 10x ::closesocket(worker_sock);
280 10x return accept_ec;
281 }
282
2/2
✓ Branch 77 → 78 taken 18 times.
✓ Branch 77 → 81 taken 51 times.
69x if (worker_ec)
283 {
284
1/2
✗ Branch 78 → 79 not taken.
✓ Branch 78 → 80 taken 18 times.
18x if (accept_sock != INVALID_SOCKET)
285 ✗ ::closesocket(accept_sock);
286 18x 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 51x non_blocking = 0;
293
2/2
✓ Branch 82 → 83 taken 9 times.
✓ Branch 82 → 88 taken 42 times.
51x if (::ioctlsocket(accept_sock, FIONBIO, &non_blocking) == SOCKET_ERROR)
294 {
295 9x auto ec = detail::make_err(::WSAGetLastError());
296 9x ::closesocket(accept_sock);
297 9x ::closesocket(worker_sock);
298 9x return ec;
299 }
300
301 42x a_sock = accept_sock;
302 42x b_sock = worker_sock;
303 42x return {};
304 83x }
305
306 std::error_code
307 42x assign_pair(
308 local_stream_socket& a,
309 local_stream_socket& b,
310 SOCKET a_sock,
311 SOCKET b_sock) noexcept
312 {
313
2/2
✓ Branch 4 → 5 taken 9 times.
✓ Branch 4 → 8 taken 33 times.
42x if (auto ec = a.assign(static_cast<native_handle_type>(a_sock)))
314 {
315 9x ::closesocket(a_sock);
316 9x ::closesocket(b_sock);
317 9x return ec;
318 }
319
320
2/2
✓ Branch 10 → 11 taken 9 times.
✓ Branch 10 → 14 taken 24 times.
33x if (auto ec = b.assign(static_cast<native_handle_type>(b_sock)))
321 {
322 9x a.close();
323 9x ::closesocket(b_sock);
324 9x return ec;
325 }
326
327 24x return {};
328 }
329
330 #endif
331
332 } // namespace
333
334 std::error_code
335 84x connect_pair(local_stream_socket& a, local_stream_socket& b) noexcept
336 {
337
5/6
✓ Branch 3 → 4 taken 83 times.
✓ Branch 3 → 6 taken 1 time.
✗ Branch 5 → 6 not taken.
✓ Branch 5 → 7 taken 83 times.
✓ Branch 8 → 9 taken 1 time.
✓ Branch 8 → 10 taken 83 times.
84x if (a.is_open() || b.is_open())
338 1x return make_error_code(error::already_open);
339
340 #if BOOST_COROSIO_POSIX
341 int a_fd = -1, b_fd = -1;
342 if (auto ec = make_pair_fds(SOCK_STREAM, a_fd, b_fd))
343 return ec;
344 return assign_pair(a, b, a_fd, b_fd);
345 #elif BOOST_COROSIO_HAS_IOCP
346 83x SOCKET a_sock = INVALID_SOCKET, b_sock = INVALID_SOCKET;
347
2/2
✓ Branch 12 → 13 taken 41 times.
✓ Branch 12 → 14 taken 42 times.
83x if (auto ec = make_pair_sockets(a_sock, b_sock))
348 41x return ec;
349 42x 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 connect_pair(local_datagram_socket& a, local_datagram_socket& b) noexcept
359 {
360 if (a.is_open() || b.is_open())
361 return make_error_code(error::already_open);
362
363 int a_fd = -1, b_fd = -1;
364 if (auto ec = make_pair_fds(SOCK_DGRAM, a_fd, b_fd))
365 return ec;
366 return assign_pair(a, b, a_fd, b_fd);
367 }
368
369 #endif
370
371 } // namespace boost::corosio
372