Repository navigation
Expand file tree
/
Copy pathsimple_socket.cpp
More file actions
340 lines (286 loc) · 6.74 KB
/
Copy pathsimple_socket.cpp
File metadata and controls
340 lines (286 loc) · 6.74 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
#include "simple_socket.hpp"
#include <stdio.h>
#include <stdint.h>
#include <string.h>
#ifdef _WIN32
#include <ws2tcpip.h>
#include <winsock2.h>
#undef SHUT_RD
#define SHUT_RD SD_RECEIVE
#else
#include <sys/types.h>
#include <sys/socket.h>
#include <netdb.h>
#include <unistd.h>
#define closesocket(x) ::close(x)
#endif
namespace PyroFling
{
Socket::~Socket()
{
if (thr.joinable())
{
// Unblock the thread in case it's waiting for us to read data.
{
std::lock_guard<std::mutex> holder{lock};
ring.read_count = ring.write_count;
ring.dead = true;
cond.notify_one();
}
#ifdef _WIN32
// Dirty hack since shutdown doesn't work and cba to use more complicated APIs.
if (fd >= 0)
{
closesocket(fd);
fd = -1;
}
#else
// If thread is blocking on a read, it should unblock now.
if (fd >= 0)
shutdown(fd, SHUT_RD);
#endif
thr.join();
}
if (fd >= 0)
closesocket(fd);
}
bool Socket::connect(Proto proto, const char *addr, const char *port)
{
#ifdef _WIN32
WSADATA wsaData;
if (WSAStartup(MAKEWORD(2, 2), &wsaData) != 0)
{
fprintf(stderr, "Failed to initialize WSA.\n");
return false;
}
#endif
addrinfo hints = {};
addrinfo *servinfo;
hints.ai_family = AF_UNSPEC;
if (proto == Proto::TCP)
{
hints.ai_socktype = SOCK_STREAM;
#ifdef __ANDROID__
hints.ai_protocol = 0;
#else
hints.ai_protocol = IPPROTO_TCP;
#endif
}
else
{
hints.ai_socktype = SOCK_DGRAM;
#ifdef __ANDROID__
hints.ai_protocol = 0;
#else
hints.ai_protocol = IPPROTO_UDP;
#endif
}
int res = getaddrinfo(addr, port, &hints, &servinfo);
if (res < 0)
return false;
addrinfo *walk;
for (walk = servinfo; walk; walk = walk->ai_next)
{
int new_fd = socket(walk->ai_family, walk->ai_socktype, walk->ai_protocol);
if (new_fd < 0)
return false;
if (::connect(new_fd, walk->ai_addr, walk->ai_addrlen) < 0)
{
closesocket(new_fd);
continue;
}
fd = new_fd;
break;
}
freeaddrinfo(servinfo);
if (proto == Proto::UDP)
{
// Keep the rcvbuf healthy so we don't drop packets too easily.
int size = 4 * 1024 * 1024;
if (setsockopt(fd, SOL_SOCKET, SO_RCVBUF, reinterpret_cast<const char *>(&size), sizeof(size)) < 0)
return false;
int actual_size = 0;
socklen_t sizelen = sizeof(size);
getsockopt(fd, SOL_SOCKET, SO_RCVBUF, reinterpret_cast<char *>(&actual_size), &sizelen);
fprintf(stderr, "Actual UDP rcvbuf size: %d bytes\n", actual_size);
}
if (!walk)
return false;
return true;
}
bool Socket::init_recv_thread(size_t max_packet_size, size_t num_packets)
{
if (thr.joinable())
return false;
if (num_packets & (num_packets - 1))
{
fprintf(stderr, "num_packets must be POT.\n");
return false;
}
ring.packets.clear();
ring.packets.reserve(num_packets);
for (size_t i = 0; i < num_packets; i++)
ring.packets.push_back({ std::unique_ptr<char []>{new char[max_packet_size]}, 0 });
ring.max_packet_size = max_packet_size;
try
{
thr = std::thread(&Socket::recv_thread, this);
}
catch (const std::exception &e)
{
fprintf(stderr, "Failed to create thread.\n");
return false;
}
return true;
}
void Socket::recv_thread()
{
uint32_t mask = ring.packets.size() - 1;
for (;;)
{
{
std::unique_lock<std::mutex> holder{lock};
cond.wait(holder, [this]() {
uint32_t queued = ring.write_count - ring.read_count;
return queued < ring.packets.size() || ring.dead;
});
if (ring.dead)
break;
}
auto &packet = ring.packets[ring.write_count & mask];
int ret = int(::recv(fd, packet.data.get(), ring.max_packet_size, 0));
if (ret <= 0)
break;
packet.size = ret;
std::lock_guard<std::mutex> holder{lock};
ring.write_count++;
cond.notify_one();
}
std::lock_guard<std::mutex> holder{lock};
ring.dead = true;
cond.notify_one();
}
size_t Socket::read_thread_packet(void *data, size_t size)
{
bool has_packet = false;
{
std::unique_lock<std::mutex> holder{lock};
// This functions more like a flush input queue.
if (ring.write_count == ring.read_count && !data)
return 0;
auto deadline = std::chrono::steady_clock::now() + std::chrono::seconds(5);
if (!cond.wait_until(holder, deadline, [this]() {
return ring.dead || ring.write_count != ring.read_count; }))
{
return 0;
}
has_packet = ring.write_count != ring.read_count;
}
if (!has_packet)
return 0;
auto &packet = ring.packets[ring.read_count & (ring.packets.size() - 1)];
size = std::min<size_t>(size, packet.size);
if (data)
memcpy(data, packet.data.get(), size);
std::lock_guard<std::mutex> holder{lock};
ring.read_count++;
cond.notify_one();
return size;
}
bool Socket::read(void *data_, size_t size, const Socket *sentinel)
{
auto *data = static_cast<uint8_t *>(data_);
while (size)
{
// Unified to be compat with Windows as well.
timeval tv = {};
fd_set fds;
FD_ZERO(&fds);
FD_SET(fd, &fds);
int nfds = fd + 1;
if (sentinel)
{
FD_SET(sentinel->fd, &fds);
if (sentinel->fd > fd)
nfds = sentinel->fd + 1;
}
tv.tv_sec = 5;
if (select(nfds, &fds, nullptr, nullptr, &tv) > 0 && FD_ISSET(fd, &fds))
{
int ret;
if ((ret = int(::recv(fd, reinterpret_cast<char *>(data), size, 0))) <= 0)
return false;
size -= ret;
data += ret;
}
else
return false;
}
return true;
}
size_t Socket::read_partial(void *data, size_t size, const Socket *sentinel)
{
// Unified to be compat with Windows as well.
timeval tv = {};
fd_set fds;
FD_ZERO(&fds);
FD_SET(fd, &fds);
int nfds = fd + 1;
if (sentinel)
{
FD_SET(sentinel->fd, &fds);
if (sentinel->fd > fd)
nfds = sentinel->fd + 1;
}
tv.tv_sec = 5;
if (select(nfds, &fds, nullptr, nullptr, &tv) > 0 && FD_ISSET(fd, &fds))
{
int ret;
if ((ret = int(::recv(fd, reinterpret_cast<char *>(data), size, 0))) <= 0)
return false;
else
return size_t(ret);
}
else
return false;
}
#ifdef __linux__
static constexpr int MSG_FLAG = MSG_NOSIGNAL;
#else
static constexpr int MSG_FLAG = 0;
#endif
bool Socket::write(const void *data_, size_t size)
{
auto *data = static_cast<const uint8_t *>(data_);
while (size)
{
int ret;
if ((ret = int(::send(fd, reinterpret_cast<const char *>(data), size, MSG_FLAG))) <= 0)
return false;
size -= ret;
data += ret;
}
return true;
}
bool Socket::write_message(const void *header, size_t header_size, const void *data, size_t size)
{
#ifdef _WIN32
uint8_t buffer[64 * 1024];
if (header_size + size > sizeof(buffer))
return false;
memcpy(buffer, header, header_size);
memcpy(buffer + header_size, data, size);
return write(buffer, header_size + size);
#else
struct msghdr msg = {};
struct iovec iv[2] = {};
iv[0].iov_base = const_cast<void *>(header);
iv[0].iov_len = header_size;
iv[1].iov_base = const_cast<void *>(data);
iv[1].iov_len = size;
msg.msg_iovlen = 2;
msg.msg_iov = iv;
return ::sendmsg(fd, &msg, MSG_FLAG) == ssize_t(header_size + size);
#endif
}
}