/src/h2o/lib/common/socket/evloop.c.h
Line | Count | Source |
1 | | /* |
2 | | * Copyright (c) 2014-2016 DeNA Co., Ltd., Kazuho Oku, Fastly, Inc. |
3 | | * |
4 | | * Permission is hereby granted, free of charge, to any person obtaining a copy |
5 | | * of this software and associated documentation files (the "Software"), to |
6 | | * deal in the Software without restriction, including without limitation the |
7 | | * rights to use, copy, modify, merge, publish, distribute, sublicense, and/or |
8 | | * sell copies of the Software, and to permit persons to whom the Software is |
9 | | * furnished to do so, subject to the following conditions: |
10 | | * |
11 | | * The above copyright notice and this permission notice shall be included in |
12 | | * all copies or substantial portions of the Software. |
13 | | * |
14 | | * THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR |
15 | | * IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, |
16 | | * FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE |
17 | | * AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER |
18 | | * LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING |
19 | | * FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS |
20 | | * IN THE SOFTWARE. |
21 | | */ |
22 | | #include <netinet/in.h> |
23 | | #include <netinet/tcp.h> |
24 | | #include <stdlib.h> |
25 | | #include <sys/time.h> |
26 | | #include <sys/uio.h> |
27 | | #include <unistd.h> |
28 | | #if H2O_USE_KTLS |
29 | | #include <linux/tls.h> |
30 | | #endif |
31 | | #include "cloexec.h" |
32 | | #include "h2o/linklist.h" |
33 | | |
34 | | #if !defined(H2O_USE_ACCEPT4) |
35 | | #ifdef __linux__ |
36 | | #if defined(__ANDROID__) && __ANDROID_API__ < 21 |
37 | | #define H2O_USE_ACCEPT4 0 |
38 | | #else |
39 | | #define H2O_USE_ACCEPT4 1 |
40 | | #endif |
41 | | #elif __FreeBSD__ >= 10 |
42 | | #define H2O_USE_ACCEPT4 1 |
43 | | #else |
44 | | #define H2O_USE_ACCEPT4 0 |
45 | | #endif |
46 | | #endif |
47 | | |
48 | | struct st_h2o_evloop_socket_t { |
49 | | h2o_socket_t super; |
50 | | int fd; |
51 | | int _flags; |
52 | | h2o_evloop_t *loop; |
53 | | size_t max_read_size; |
54 | | struct st_h2o_evloop_socket_t *_next_pending; |
55 | | struct st_h2o_evloop_socket_t *_next_statechanged; |
56 | | struct { |
57 | | uint64_t prev_loop; |
58 | | uint64_t cur_loop; |
59 | | uint64_t cur_run_count; |
60 | | } bytes_written; |
61 | | /** |
62 | | * vector to be sent (or vec.callbacks is NULL when not used) |
63 | | */ |
64 | | h2o_sendvec_t sendvec; |
65 | | }; |
66 | | |
67 | | static void link_to_pending(struct st_h2o_evloop_socket_t *sock); |
68 | | static void link_to_statechanged(struct st_h2o_evloop_socket_t *sock); |
69 | | static void write_pending(struct st_h2o_evloop_socket_t *sock); |
70 | | static h2o_evloop_t *create_evloop(size_t sz); |
71 | | static void update_now(h2o_evloop_t *loop); |
72 | | static int32_t adjust_max_wait(h2o_evloop_t *loop, int32_t max_wait); |
73 | | |
74 | | static void notify_write_progress(struct st_h2o_evloop_socket_t *sock) |
75 | 30.9k | { |
76 | 30.9k | if ((sock->super._cb.write_flags & H2O_SOCKET_SENDVEC_FLAG_REPORT_PROGRESS) != 0) |
77 | 11.2k | sock->super._cb.write(&sock->super, h2o_socket_error_write_progress); |
78 | 30.9k | } |
79 | | |
80 | | /* functions to be defined in the backends */ |
81 | | static int evloop_do_proceed(h2o_evloop_t *loop, int32_t max_wait); |
82 | | static void evloop_do_dispose(h2o_evloop_t *loop); |
83 | | static void evloop_do_on_socket_create(struct st_h2o_evloop_socket_t *sock); |
84 | | static int evloop_do_on_socket_close(struct st_h2o_evloop_socket_t *sock); |
85 | | static void evloop_do_on_socket_export(struct st_h2o_evloop_socket_t *sock); |
86 | | |
87 | | #if H2O_USE_POLL || H2O_USE_EPOLL || H2O_USE_KQUEUE |
88 | | /* explicitly specified */ |
89 | | #else |
90 | | #if defined(__APPLE__) || defined(__FreeBSD__) || defined(__NetBSD__) || defined(__OpenBSD__) |
91 | | #define H2O_USE_KQUEUE 1 |
92 | | #elif defined(__linux) |
93 | | #define H2O_USE_EPOLL 1 |
94 | | #if defined(SO_ZEROCOPY) && defined(SO_EE_ORIGIN_ZEROCOPY) && defined(MSG_ZEROCOPY) |
95 | | #define H2O_USE_MSG_ZEROCOPY 1 |
96 | | #endif |
97 | | #else |
98 | | #define H2O_USE_POLL 1 |
99 | | #endif |
100 | | #endif |
101 | | #if !defined(H2O_USE_MSG_ZEROCOPY) |
102 | | #define H2O_USE_MSG_ZEROCOPY 0 |
103 | | #endif |
104 | | |
105 | | #if H2O_USE_POLL |
106 | | #include "evloop/poll.c.h" |
107 | | #elif H2O_USE_EPOLL |
108 | | #include "evloop/epoll.c.h" |
109 | | #elif H2O_USE_KQUEUE |
110 | | #include "evloop/kqueue.c.h" |
111 | | #else |
112 | | #error "poller not specified" |
113 | | #endif |
114 | | |
115 | | size_t h2o_evloop_socket_max_read_size = 1024 * 1024; /* by default, we read up to 1MB at once */ |
116 | | size_t h2o_evloop_socket_max_write_size = 1024 * 1024; /* by default, we write up to 1MB at once */ |
117 | | |
118 | | void link_to_pending(struct st_h2o_evloop_socket_t *sock) |
119 | 75.8k | { |
120 | 75.8k | if (sock->_next_pending == sock) { |
121 | 75.2k | struct st_h2o_evloop_socket_t **slot = (sock->_flags & H2O_SOCKET_FLAG_IS_ACCEPTED_CONNECTION) != 0 |
122 | 75.2k | ? &sock->loop->_pending_as_server |
123 | 75.2k | : &sock->loop->_pending_as_client; |
124 | 75.2k | sock->_next_pending = *slot; |
125 | 75.2k | *slot = sock; |
126 | 75.2k | } |
127 | 75.8k | } |
128 | | |
129 | | void link_to_statechanged(struct st_h2o_evloop_socket_t *sock) |
130 | 78.2k | { |
131 | 78.2k | if (sock->_next_statechanged == sock) { |
132 | 56.7k | sock->_next_statechanged = NULL; |
133 | 56.7k | *sock->loop->_statechanged.tail_ref = sock; |
134 | 56.7k | sock->loop->_statechanged.tail_ref = &sock->_next_statechanged; |
135 | 56.7k | } |
136 | 78.2k | } |
137 | | |
138 | | static const char *on_read_core(int fd, h2o_buffer_t **input, size_t max_bytes) |
139 | 39.9k | { |
140 | 39.9k | ssize_t read_so_far = 0; |
141 | | |
142 | 40.0k | while (1) { |
143 | 40.0k | ssize_t rret; |
144 | 40.0k | h2o_iovec_t buf = h2o_buffer_try_reserve(input, max_bytes < 4096 ? max_bytes : 4096); |
145 | 40.0k | if (buf.base == NULL) { |
146 | | /* memory allocation failed */ |
147 | 0 | return h2o_socket_error_out_of_memory; |
148 | 0 | } |
149 | 40.0k | size_t read_size = buf.len <= INT_MAX / 2 ? buf.len : INT_MAX / 2 + 1; |
150 | 40.0k | if (read_size > max_bytes) |
151 | 0 | read_size = max_bytes; |
152 | 40.0k | while ((rret = read(fd, buf.base, read_size)) == -1 && errno == EINTR) |
153 | 0 | ; |
154 | 40.0k | if (rret == -1) { |
155 | 34 | if (errno == EAGAIN) |
156 | 6 | break; |
157 | 28 | else |
158 | 28 | return h2o_socket_error_io; |
159 | 40.0k | } else if (rret == 0) { |
160 | 10.9k | if (read_so_far == 0) |
161 | 10.9k | return h2o_socket_error_closed; /* TODO notify close */ |
162 | 5 | break; |
163 | 10.9k | } |
164 | 29.0k | (*input)->size += rret; |
165 | 29.0k | if (buf.len != rret) |
166 | 28.8k | break; |
167 | 168 | read_so_far += rret; |
168 | 168 | if (read_so_far >= max_bytes) |
169 | 0 | break; |
170 | 168 | } |
171 | 28.9k | return NULL; |
172 | 39.9k | } |
173 | | |
174 | | static size_t write_vecs(struct st_h2o_evloop_socket_t *sock, h2o_iovec_t **bufs, size_t *bufcnt, int sendmsg_flags) |
175 | 31.8k | { |
176 | 31.8k | ssize_t wret; |
177 | | |
178 | 31.8k | while (*bufcnt != 0) { |
179 | | /* write */ |
180 | 31.5k | int iovcnt = *bufcnt < IOV_MAX ? (int)*bufcnt : IOV_MAX; |
181 | 31.5k | struct msghdr msg; |
182 | 31.5k | do { |
183 | 31.5k | msg = (struct msghdr){.msg_iov = (struct iovec *)*bufs, .msg_iovlen = iovcnt}; |
184 | 31.5k | } while ((wret = sendmsg(sock->fd, &msg, sendmsg_flags)) == -1 && errno == EINTR); |
185 | 31.5k | SOCKET_PROBE(WRITEV, &sock->super, wret); |
186 | 31.5k | H2O_LOG_SOCK(writev, &sock->super, { PTLS_LOG_ELEMENT_SIGNED(ret, wret); }); |
187 | | |
188 | 31.5k | if (wret == -1) |
189 | 596 | return errno == EAGAIN ? 0 : SIZE_MAX; |
190 | 30.9k | if (wret != 0) |
191 | 30.9k | notify_write_progress(sock); |
192 | | |
193 | | /* adjust the buffer, doing the write once again only if all IOV_MAX buffers being supplied were fully written */ |
194 | 46.5k | while ((*bufs)->len <= wret) { |
195 | 46.5k | wret -= (*bufs)->len; |
196 | 46.5k | ++*bufs; |
197 | 46.5k | --*bufcnt; |
198 | 46.5k | if (*bufcnt == 0) { |
199 | 30.9k | assert(wret == 0); |
200 | 30.9k | return 0; |
201 | 30.9k | } |
202 | 46.5k | } |
203 | 0 | if (wret != 0) { |
204 | 0 | return wret; |
205 | 0 | } else if (iovcnt < IOV_MAX) { |
206 | 0 | return 0; |
207 | 0 | } |
208 | 0 | } |
209 | | |
210 | 340 | return 0; |
211 | 31.8k | } |
212 | | |
213 | | static size_t write_core(struct st_h2o_evloop_socket_t *sock, h2o_iovec_t **bufs, size_t *bufcnt) |
214 | 31.8k | { |
215 | 31.8k | if (sock->super.ssl == NULL || sock->super.ssl->offload == H2O_SOCKET_SSL_OFFLOAD_ON) { |
216 | 31.8k | if (sock->super.ssl != NULL) |
217 | 31.8k | assert(!has_pending_ssl_bytes(sock->super.ssl)); |
218 | 31.8k | return write_vecs(sock, bufs, bufcnt, 0); |
219 | 31.8k | } |
220 | | |
221 | | /* SSL: flatten given vector if that has not been done yet; `*bufs` is guaranteed to have one slot available at the end; see |
222 | | * `do_write_with_sendvec`, `init_write_buf`. */ |
223 | 0 | if (sock->sendvec.callbacks != NULL) { |
224 | 0 | size_t veclen = flatten_sendvec(&sock->super, &sock->sendvec); |
225 | 0 | if (veclen == SIZE_MAX) |
226 | 0 | return SIZE_MAX; |
227 | 0 | sock->sendvec.callbacks = NULL; |
228 | 0 | (*bufs)[(*bufcnt)++] = h2o_iovec_init(sock->super._write_buf.flattened, veclen); |
229 | 0 | } |
230 | | |
231 | | /* continue encrypting and writing, until we run out of data */ |
232 | 0 | size_t first_buf_written = 0; |
233 | 0 | while (1) { |
234 | | /* write bytes already encrypted, if any */ |
235 | 0 | if (has_pending_ssl_bytes(sock->super.ssl)) { |
236 | 0 | h2o_iovec_t encbuf = h2o_iovec_init(sock->super.ssl->output.buf.base + sock->super.ssl->output.pending_off, |
237 | 0 | sock->super.ssl->output.buf.off - sock->super.ssl->output.pending_off); |
238 | 0 | h2o_iovec_t *encbufs = &encbuf; |
239 | 0 | size_t encbufcnt = 1, enc_written; |
240 | 0 | int sendmsg_flags = 0; |
241 | 0 | #if H2O_USE_MSG_ZEROCOPY |
242 | | /* Use zero copy if amount of data to be written is no less than 4KB, and if the memory can be returned to |
243 | | * `h2o_socket_zerocopy_buffer_allocator`. Latter is a short-cut. It is only under exceptional conditions (e.g., TLS |
244 | | * stack adding a post-handshake message) that we'd see the buffer grow to a size that cannot be returned to the |
245 | | * recycling allocator. |
246 | | * Even though https://www.kernel.org/doc/html/v5.17/networking/msg_zerocopy.html recommends 10KB, 4KB has been chosen |
247 | | * as the threshold, because we are likely to be using the non-temporal aesgcm engine and tx-nocache-copy, in which case |
248 | | * copying sendmsg is going to be more costly than what the kernel documentation assumes. In a synthetic benchmark, |
249 | | * changing from 16KB to 4KB increased the throughput by ~10%. */ |
250 | 0 | if (sock->super.ssl->output.allocated_for_zerocopy && encbuf.len >= 4096 && |
251 | 0 | sock->super.ssl->output.buf.capacity == h2o_socket_zerocopy_buffer_allocator.conf->memsize) |
252 | 0 | sendmsg_flags = MSG_ZEROCOPY; |
253 | 0 | #endif |
254 | 0 | if ((enc_written = write_vecs(sock, &encbufs, &encbufcnt, sendmsg_flags)) == SIZE_MAX) { |
255 | 0 | dispose_ssl_output_buffer(sock->super.ssl); |
256 | 0 | return SIZE_MAX; |
257 | 0 | } |
258 | 0 | if (sendmsg_flags != 0 && (encbufcnt == 0 || enc_written > 0)) { |
259 | 0 | zerocopy_buffers_push(sock->super._zerocopy, sock->super.ssl->output.buf.base); |
260 | 0 | if (!sock->super.ssl->output.zerocopy_owned) { |
261 | 0 | sock->super.ssl->output.zerocopy_owned = 1; |
262 | 0 | ++h2o_socket_num_zerocopy_buffers_inflight; |
263 | 0 | } |
264 | 0 | } |
265 | | /* if write is incomplete, record the advance and bail out */ |
266 | 0 | if (encbufcnt != 0) { |
267 | 0 | sock->super.ssl->output.pending_off += enc_written; |
268 | 0 | break; |
269 | 0 | } |
270 | | /* succeeded in writing all the encrypted data; free the buffer */ |
271 | 0 | dispose_ssl_output_buffer(sock->super.ssl); |
272 | 0 | } |
273 | | /* bail out if complete */ |
274 | 0 | if (*bufcnt == 0 && sock->sendvec.callbacks == NULL) |
275 | 0 | break; |
276 | | /* convert more cleartext to TLS records if possible, or bail out on fatal error */ |
277 | 0 | if ((first_buf_written = generate_tls_records(&sock->super, bufs, bufcnt, first_buf_written)) == SIZE_MAX) |
278 | 0 | break; |
279 | | /* as an optimization, if we have a flattened vector, release memory as soon as they have been encrypted */ |
280 | 0 | if (*bufcnt == 0 && sock->super._write_buf.flattened != NULL) { |
281 | 0 | h2o_mem_free_recycle(&h2o_socket_ssl_buffer_allocator, sock->super._write_buf.flattened); |
282 | 0 | sock->super._write_buf.flattened = NULL; |
283 | 0 | } |
284 | 0 | } |
285 | | |
286 | 0 | return first_buf_written; |
287 | 0 | } |
288 | | |
289 | | /** |
290 | | * Sends contents of sendvec, and returns if operation has been successful, either completely or partially. Upon completion, |
291 | | * `sendvec.vec.callbacks` is reset to NULL. |
292 | | */ |
293 | | static int sendvec_core(struct st_h2o_evloop_socket_t *sock) |
294 | 0 | { |
295 | 0 | size_t bytes_sent; |
296 | |
|
297 | 0 | assert(sock->sendvec.len != 0); |
298 | | |
299 | | /* send, and return an error if failed */ |
300 | 0 | if ((bytes_sent = sock->sendvec.callbacks->send_(&sock->sendvec, sock->fd, sock->sendvec.len)) == SIZE_MAX) |
301 | 0 | return 0; |
302 | 0 | if (bytes_sent != 0) |
303 | 0 | notify_write_progress(sock); |
304 | | |
305 | | /* update offset, and return if we are not done yet */ |
306 | 0 | if (sock->sendvec.len != 0) |
307 | 0 | return 1; |
308 | | |
309 | | /* operation complete; mark as such */ |
310 | 0 | sock->sendvec.callbacks = NULL; |
311 | 0 | return 1; |
312 | 0 | } |
313 | | |
314 | | void write_pending(struct st_h2o_evloop_socket_t *sock) |
315 | 119 | { |
316 | 119 | assert(sock->super._cb.write != NULL); |
317 | | |
318 | | /* write from buffer, if we have anything */ |
319 | 119 | int ssl_needs_flatten = sock->sendvec.callbacks != NULL && sock->super.ssl != NULL |
320 | | #if H2O_USE_KTLS |
321 | | && sock->super.ssl->offload != H2O_SOCKET_SSL_OFFLOAD_ON |
322 | | #endif |
323 | 119 | ; |
324 | 119 | if (sock->super._write_buf.cnt != 0 || has_pending_ssl_bytes(sock->super.ssl) || ssl_needs_flatten) { |
325 | 0 | size_t first_buf_written; |
326 | 0 | if ((first_buf_written = write_core(sock, &sock->super._write_buf.bufs, &sock->super._write_buf.cnt)) != SIZE_MAX) { |
327 | | /* return if there's still pending data, adjusting buf[0] if necessary */ |
328 | 0 | if (sock->super._write_buf.cnt != 0) { |
329 | 0 | sock->super._write_buf.bufs[0].base += first_buf_written; |
330 | 0 | sock->super._write_buf.bufs[0].len -= first_buf_written; |
331 | 0 | return; |
332 | 0 | } else if (has_pending_ssl_bytes(sock->super.ssl)) { |
333 | 0 | return; |
334 | 0 | } |
335 | 0 | } |
336 | 0 | } |
337 | | |
338 | | /* either completed or failed */ |
339 | 119 | dispose_write_buf(&sock->super); |
340 | | |
341 | | /* send the vector, if we have one and if all buffered writes are complete */ |
342 | 119 | if (sock->sendvec.callbacks != NULL && sock->super._write_buf.cnt == 0 && !has_pending_ssl_bytes(sock->super.ssl)) { |
343 | | /* send, and upon partial send, return without changing state for another round */ |
344 | 0 | if (sendvec_core(sock) && sock->sendvec.callbacks != NULL) |
345 | 0 | return; |
346 | 0 | } |
347 | | |
348 | | /* operation completed or failed, schedule notification */ |
349 | 119 | SOCKET_PROBE(WRITE_COMPLETE, &sock->super, sock->super._write_buf.cnt == 0 && !has_pending_ssl_bytes(sock->super.ssl)); |
350 | 119 | H2O_LOG_SOCK(write_complete, &sock->super, |
351 | 119 | { PTLS_LOG_ELEMENT_BOOL(success, sock->super._write_buf.cnt == 0 && !has_pending_ssl_bytes(sock->super.ssl)); }); |
352 | 119 | sock->bytes_written.cur_loop = sock->super.bytes_written; |
353 | 119 | sock->_flags |= H2O_SOCKET_FLAG_IS_WRITE_NOTIFY; |
354 | 119 | link_to_pending(sock); |
355 | 119 | link_to_statechanged(sock); /* might need to disable the write polling */ |
356 | 119 | } |
357 | | |
358 | | static void read_on_ready(struct st_h2o_evloop_socket_t *sock) |
359 | 39.9k | { |
360 | 39.9k | const char *err = 0; |
361 | 39.9k | size_t prev_size = sock->super.input->size; |
362 | | |
363 | 39.9k | if ((sock->_flags & H2O_SOCKET_FLAG_DONT_READ) != 0) |
364 | 0 | goto Notify; |
365 | | |
366 | 39.9k | if ((err = on_read_core(sock->fd, sock->super.ssl == NULL ? &sock->super.input : &sock->super.ssl->input.encrypted, |
367 | 39.9k | sock->max_read_size)) != NULL) |
368 | 11.0k | goto Notify; |
369 | | |
370 | 28.9k | if (sock->super.ssl != NULL && sock->super.ssl->handshake.cb == NULL) |
371 | 0 | err = decode_ssl_input(&sock->super); |
372 | | |
373 | 39.9k | Notify: |
374 | | /* the application may get notified even if no new data is avaiable. The |
375 | | * behavior is intentional; it is designed as such so that the applications |
376 | | * can update their timeout counters when a partial SSL record arrives. |
377 | | */ |
378 | 39.9k | sock->super.bytes_read += sock->super.input->size - prev_size; |
379 | 39.9k | sock->super._cb.read(&sock->super, err); |
380 | 39.9k | } |
381 | | |
382 | | void do_dispose_socket(h2o_socket_t *_sock) |
383 | 22.1k | { |
384 | 22.1k | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
385 | | |
386 | 22.1k | dispose_write_buf(&sock->super); |
387 | | |
388 | 22.1k | sock->_flags = H2O_SOCKET_FLAG_IS_DISPOSED | (sock->_flags & H2O_SOCKET_FLAG__EPOLL_IS_REGISTERED); |
389 | | |
390 | | /* Give backends chance to do the necessary cleanup, as well as giving them chance to switch to their own disposal method; e.g., |
391 | | * shutdown(SHUT_RDWR) with delays to reclaim all zero copy buffers. */ |
392 | 22.1k | if (evloop_do_on_socket_close(sock)) |
393 | 0 | return; |
394 | | |
395 | | /* immediate close */ |
396 | 22.1k | if (sock->fd != -1) { |
397 | 22.1k | close(sock->fd); |
398 | 22.1k | sock->fd = -1; |
399 | 22.1k | } |
400 | 22.1k | link_to_statechanged(sock); |
401 | 22.1k | } |
402 | | |
403 | | void report_early_write_error(h2o_socket_t *_sock) |
404 | 596 | { |
405 | 596 | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
406 | | |
407 | | /* fill in _wreq.bufs with fake data to indicate error */ |
408 | 596 | sock->super._write_buf.bufs = sock->super._write_buf.smallbufs; |
409 | 596 | sock->super._write_buf.cnt = 1; |
410 | 596 | *sock->super._write_buf.bufs = h2o_iovec_init(H2O_STRLIT("deadbeef")); |
411 | 596 | sock->_flags |= H2O_SOCKET_FLAG_IS_WRITE_NOTIFY; |
412 | 596 | link_to_pending(sock); |
413 | 596 | } |
414 | | |
415 | | void do_write(h2o_socket_t *_sock, h2o_iovec_t *bufs, size_t bufcnt) |
416 | 31.8k | { |
417 | 31.8k | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
418 | 31.8k | size_t first_buf_written; |
419 | | |
420 | | /* Don't write too much; if more than 1MB have been already written in the current invocation of `h2o_evloop_run`, wait until |
421 | | * the event loop notifies us that the socket is writable. */ |
422 | 31.8k | if (sock->bytes_written.cur_run_count != sock->loop->run_count) { |
423 | 29.4k | sock->bytes_written.prev_loop = sock->bytes_written.cur_loop; |
424 | 29.4k | sock->bytes_written.cur_run_count = sock->loop->run_count; |
425 | 29.4k | } else if (sock->bytes_written.cur_loop - sock->bytes_written.prev_loop >= h2o_evloop_socket_max_write_size) { |
426 | 0 | init_write_buf(&sock->super, bufs, bufcnt, 0); |
427 | 0 | goto Schedule_Write; |
428 | 0 | } |
429 | | |
430 | | /* try to write now */ |
431 | 31.8k | if ((first_buf_written = write_core(sock, &bufs, &bufcnt)) == SIZE_MAX) { |
432 | 596 | report_early_write_error(&sock->super); |
433 | 596 | return; |
434 | 596 | } |
435 | 31.2k | if (bufcnt == 0 && !has_pending_ssl_bytes(sock->super.ssl)) { |
436 | | /* write complete, schedule the callback */ |
437 | 31.2k | if (sock->super._write_buf.flattened != NULL) { |
438 | 0 | h2o_mem_free_recycle(&h2o_socket_ssl_buffer_allocator, sock->super._write_buf.flattened); |
439 | 0 | sock->super._write_buf.flattened = NULL; |
440 | 0 | } |
441 | 31.2k | if (sock->sendvec.callbacks != NULL) { |
442 | 0 | if (!sendvec_core(sock)) { |
443 | 0 | report_early_write_error(&sock->super); |
444 | 0 | return; |
445 | 0 | } |
446 | 0 | if (sock->sendvec.callbacks != NULL) |
447 | 0 | goto Schedule_Write; |
448 | 0 | } |
449 | 31.2k | sock->bytes_written.cur_loop = sock->super.bytes_written; |
450 | 31.2k | sock->_flags |= H2O_SOCKET_FLAG_IS_WRITE_NOTIFY; |
451 | 31.2k | link_to_pending(sock); |
452 | 31.2k | return; |
453 | 31.2k | } |
454 | | |
455 | | /* setup the buffer to send pending data */ |
456 | 0 | init_write_buf(&sock->super, bufs, bufcnt, first_buf_written); |
457 | |
|
458 | 0 | Schedule_Write: |
459 | 0 | link_to_statechanged(sock); |
460 | 0 | } |
461 | | |
462 | | static int can_tls_offload(h2o_socket_t *sock) |
463 | 0 | { |
464 | | #if H2O_USE_KTLS |
465 | | if (sock->ssl->offload != H2O_SOCKET_SSL_OFFLOAD_NONE && sock->ssl->ptls != NULL) { |
466 | | ptls_cipher_suite_t *cipher = ptls_get_cipher(sock->ssl->ptls); |
467 | | switch (cipher->id) { |
468 | | case PTLS_CIPHER_SUITE_AES_128_GCM_SHA256: |
469 | | case PTLS_CIPHER_SUITE_AES_256_GCM_SHA384: |
470 | | return 1; |
471 | | default: |
472 | | break; |
473 | | } |
474 | | } |
475 | | #endif |
476 | |
|
477 | 0 | return 0; |
478 | 0 | } |
479 | | |
480 | | #if H2O_USE_KTLS |
481 | | static void switch_to_ktls(struct st_h2o_evloop_socket_t *sock) |
482 | | { |
483 | | assert(sock->super.ssl->offload == H2O_SOCKET_SSL_OFFLOAD_TBD); |
484 | | |
485 | | /* Postpone the decision, when we are still in the early stages of the connection, as we want to use userspace TLS for |
486 | | * generating small TLS records. TODO: integrate with TLS record size calculation logic. */ |
487 | | if (sock->super.bytes_written < 65536) |
488 | | return; |
489 | | |
490 | | /* load the key to the kernel */ |
491 | | struct { |
492 | | uint8_t key[PTLS_MAX_SECRET_SIZE]; |
493 | | uint8_t iv[PTLS_MAX_DIGEST_SIZE]; |
494 | | uint64_t seq; |
495 | | union { |
496 | | struct tls12_crypto_info_aes_gcm_128 aesgcm128; |
497 | | struct tls12_crypto_info_aes_gcm_256 aesgcm256; |
498 | | } tx_params; |
499 | | size_t tx_params_size; |
500 | | } keys; |
501 | | |
502 | | /* at the moment, only TLS/1.3 connections using aes-gcm is supported */ |
503 | | if (sock->super.ssl->ptls == NULL) |
504 | | goto Fail; |
505 | | ptls_cipher_suite_t *cipher = ptls_get_cipher(sock->super.ssl->ptls); |
506 | | switch (cipher->id) { |
507 | | case PTLS_CIPHER_SUITE_AES_128_GCM_SHA256: |
508 | | case PTLS_CIPHER_SUITE_AES_256_GCM_SHA384: |
509 | | break; |
510 | | default: |
511 | | goto Fail; |
512 | | } |
513 | | if (ptls_get_traffic_keys(sock->super.ssl->ptls, 1, keys.key, keys.iv, &keys.seq) != 0) |
514 | | goto Fail; |
515 | | keys.seq = htobe64(keys.seq); /* converted to big endian ASAP */ |
516 | | |
517 | | #define SETUP_TX_PARAMS(target, type) \ |
518 | | do { \ |
519 | | keys.tx_params.target.info.version = TLS_1_3_VERSION; \ |
520 | | keys.tx_params.target.info.cipher_type = type; \ |
521 | | H2O_BUILD_ASSERT(sizeof(keys.tx_params.target.key) == cipher->aead->key_size); \ |
522 | | memcpy(keys.tx_params.target.key, keys.key, cipher->aead->key_size); \ |
523 | | H2O_BUILD_ASSERT(cipher->aead->iv_size == 12); \ |
524 | | H2O_BUILD_ASSERT(sizeof(keys.tx_params.target.salt) == 4); \ |
525 | | memcpy(keys.tx_params.target.salt, keys.iv, 4); \ |
526 | | H2O_BUILD_ASSERT(sizeof(keys.tx_params.target.iv) == 8); \ |
527 | | memcpy(keys.tx_params.target.iv, keys.iv + 4, 8); \ |
528 | | H2O_BUILD_ASSERT(sizeof(keys.tx_params.target.rec_seq) == sizeof(keys.seq)); \ |
529 | | memcpy(keys.tx_params.target.rec_seq, &keys.seq, sizeof(keys.seq)); \ |
530 | | keys.tx_params_size = sizeof(keys.tx_params.target); \ |
531 | | } while (0) |
532 | | switch (cipher->id) { |
533 | | case PTLS_CIPHER_SUITE_AES_128_GCM_SHA256: |
534 | | SETUP_TX_PARAMS(aesgcm128, TLS_CIPHER_AES_GCM_128); |
535 | | break; |
536 | | case PTLS_CIPHER_SUITE_AES_256_GCM_SHA384: |
537 | | SETUP_TX_PARAMS(aesgcm256, TLS_CIPHER_AES_GCM_256); |
538 | | break; |
539 | | default: |
540 | | goto Fail; |
541 | | } |
542 | | #undef SETUP_TX_PARAMS |
543 | | |
544 | | /* set to kernel */ |
545 | | if (setsockopt(sock->fd, SOL_TCP, TCP_ULP, "tls", sizeof("tls")) != 0) |
546 | | goto Fail; |
547 | | if (setsockopt(sock->fd, SOL_TLS, TLS_TX, &keys.tx_params, keys.tx_params_size) != 0) |
548 | | goto Fail; |
549 | | sock->super.ssl->offload = H2O_SOCKET_SSL_OFFLOAD_ON; |
550 | | |
551 | | Exit: |
552 | | ptls_clear_memory(&keys, sizeof(keys)); |
553 | | return; |
554 | | |
555 | | Fail: |
556 | | sock->super.ssl->offload = H2O_SOCKET_SSL_OFFLOAD_NONE; |
557 | | goto Exit; |
558 | | } |
559 | | #endif |
560 | | |
561 | | /** |
562 | | * `bufs` should be an array capable of storing `bufcnt + 1` objects, as we will be flattening `sendvec` at the end of `bufs` before |
563 | | * encryption; see `write_core`. |
564 | | */ |
565 | | static int do_write_with_sendvec(h2o_socket_t *_sock, h2o_iovec_t *bufs, size_t bufcnt, h2o_sendvec_t *sendvec) |
566 | 0 | { |
567 | 0 | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
568 | |
|
569 | 0 | assert(sendvec->callbacks->read_ != NULL); |
570 | 0 | assert(sock->sendvec.callbacks == NULL); |
571 | | |
572 | | /* If userspace TLS is used, rely on `read_` which is a mandatory callback. Otherwise, rely on `send_` if it is available. */ |
573 | 0 | if (sock->super.ssl != NULL) { |
574 | | #if H2O_USE_KTLS |
575 | | if (sock->super.ssl->offload == H2O_SOCKET_SSL_OFFLOAD_TBD) |
576 | | switch_to_ktls(sock); |
577 | | if (sock->super.ssl->offload == H2O_SOCKET_SSL_OFFLOAD_ON && sendvec->callbacks->send_ == NULL) |
578 | | return 0; |
579 | | #endif |
580 | 0 | } else { |
581 | 0 | if (sendvec->callbacks->send_ == NULL) |
582 | 0 | return 0; |
583 | 0 | } |
584 | | |
585 | | /* handling writes with sendvec, here */ |
586 | 0 | sock->sendvec = *sendvec; |
587 | 0 | do_write(&sock->super, bufs, bufcnt); |
588 | |
|
589 | 0 | return 1; |
590 | 0 | } |
591 | | |
592 | | int h2o_socket_get_fd(h2o_socket_t *_sock) |
593 | 6.45k | { |
594 | 6.45k | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
595 | 6.45k | return sock->fd; |
596 | 6.45k | } |
597 | | |
598 | | void do_read_start(h2o_socket_t *_sock) |
599 | 37.8k | { |
600 | 37.8k | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
601 | | |
602 | 37.8k | link_to_statechanged(sock); |
603 | 37.8k | } |
604 | | |
605 | | void do_read_stop(h2o_socket_t *_sock) |
606 | 18.0k | { |
607 | 18.0k | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
608 | | |
609 | 18.0k | sock->_flags &= ~H2O_SOCKET_FLAG_IS_READ_READY; |
610 | 18.0k | link_to_statechanged(sock); |
611 | 18.0k | } |
612 | | |
613 | | void h2o_socket_dont_read(h2o_socket_t *_sock, int dont_read) |
614 | 0 | { |
615 | 0 | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
616 | |
|
617 | 0 | if (dont_read) { |
618 | 0 | sock->_flags |= H2O_SOCKET_FLAG_DONT_READ; |
619 | 0 | } else { |
620 | 0 | sock->_flags &= ~H2O_SOCKET_FLAG_DONT_READ; |
621 | 0 | } |
622 | 0 | } |
623 | | |
624 | | int do_export(h2o_socket_t *_sock, h2o_socket_export_t *info) |
625 | 0 | { |
626 | 0 | struct st_h2o_evloop_socket_t *sock = (void *)_sock; |
627 | |
|
628 | 0 | assert((sock->_flags & H2O_SOCKET_FLAG_IS_DISPOSED) == 0); |
629 | 0 | evloop_do_on_socket_export(sock); |
630 | 0 | sock->_flags = H2O_SOCKET_FLAG_IS_DISPOSED | (sock->_flags & H2O_SOCKET_FLAG__EPOLL_IS_REGISTERED); |
631 | |
|
632 | 0 | info->fd = sock->fd; |
633 | 0 | sock->fd = -1; |
634 | |
|
635 | 0 | return 0; |
636 | 0 | } |
637 | | |
638 | | h2o_socket_t *do_import(h2o_loop_t *loop, h2o_socket_export_t *info) |
639 | 0 | { |
640 | 0 | return h2o_evloop_socket_create(loop, info->fd, 0); |
641 | 0 | } |
642 | | |
643 | | h2o_loop_t *h2o_socket_get_loop(h2o_socket_t *_sock) |
644 | 6.45k | { |
645 | 6.45k | struct st_h2o_evloop_socket_t *sock = (void *)_sock; |
646 | 6.45k | return sock->loop; |
647 | 6.45k | } |
648 | | |
649 | | socklen_t get_sockname_uncached(h2o_socket_t *_sock, struct sockaddr *sa) |
650 | 0 | { |
651 | 0 | struct st_h2o_evloop_socket_t *sock = (void *)_sock; |
652 | 0 | socklen_t len = sizeof(struct sockaddr_storage); |
653 | 0 | if (getsockname(sock->fd, sa, &len) != 0) |
654 | 0 | return 0; |
655 | 0 | return len; |
656 | 0 | } |
657 | | |
658 | | socklen_t get_peername_uncached(h2o_socket_t *_sock, struct sockaddr *sa) |
659 | 1.53k | { |
660 | 1.53k | struct st_h2o_evloop_socket_t *sock = (void *)_sock; |
661 | 1.53k | socklen_t len = sizeof(struct sockaddr_storage); |
662 | 1.53k | if (getpeername(sock->fd, sa, &len) != 0) |
663 | 0 | return 0; |
664 | 1.53k | return len; |
665 | 1.53k | } |
666 | | |
667 | | static struct st_h2o_evloop_socket_t *create_socket(h2o_evloop_t *loop, int fd, int flags) |
668 | 22.1k | { |
669 | 22.1k | struct st_h2o_evloop_socket_t *sock; |
670 | | |
671 | 22.1k | sock = h2o_mem_alloc(sizeof(*sock)); |
672 | 22.1k | memset(sock, 0, sizeof(*sock)); |
673 | 22.1k | h2o_buffer_init(&sock->super.input, &h2o_socket_buffer_prototype); |
674 | 22.1k | sock->loop = loop; |
675 | 22.1k | sock->fd = fd; |
676 | 22.1k | sock->_flags = flags; |
677 | 22.1k | sock->max_read_size = h2o_evloop_socket_max_read_size; /* by default, we read up to 1MB at once */ |
678 | 22.1k | sock->_next_pending = sock; |
679 | 22.1k | sock->_next_statechanged = sock; |
680 | | |
681 | 22.1k | evloop_do_on_socket_create(sock); |
682 | | |
683 | 22.1k | return sock; |
684 | 22.1k | } |
685 | | |
686 | | /** |
687 | | * Sets TCP_NODELAY if the given file descriptor is likely to be a TCP socket. The intent of this function is to reduce number of |
688 | | * unnecessary system calls. Therefore, we skip setting TCP_NODELAY when it is certain that the socket is not a TCP socket, |
689 | | * otherwise call setsockopt. |
690 | | */ |
691 | | static void set_nodelay_if_likely_tcp(int fd, struct sockaddr *sa) |
692 | 22.1k | { |
693 | 22.1k | if (sa != NULL && !(sa->sa_family == AF_INET || sa->sa_family == AF_INET6)) |
694 | 3.98k | return; |
695 | | |
696 | 18.1k | int on = 1; |
697 | 18.1k | setsockopt(fd, IPPROTO_TCP, TCP_NODELAY, &on, sizeof(on)); |
698 | 18.1k | } |
699 | | |
700 | | h2o_socket_t *h2o_evloop_socket_create(h2o_evloop_t *loop, int fd, int flags) |
701 | 18.1k | { |
702 | | /* It is the reponsibility of the event loop to modify the properties of a socket for its use (e.g., set O_NONBLOCK). */ |
703 | 18.1k | fcntl(fd, F_SETFL, O_NONBLOCK); |
704 | 18.1k | set_nodelay_if_likely_tcp(fd, NULL); |
705 | | |
706 | 18.1k | return &create_socket(loop, fd, flags)->super; |
707 | 18.1k | } |
708 | | |
709 | | h2o_socket_t *h2o_evloop_socket_accept(h2o_socket_t *_listener) |
710 | 0 | { |
711 | 0 | struct st_h2o_evloop_socket_t *listener = (struct st_h2o_evloop_socket_t *)_listener; |
712 | 0 | int fd; |
713 | 0 | h2o_socket_t *sock; |
714 | 0 | union { |
715 | 0 | struct sockaddr sa; |
716 | 0 | struct sockaddr_in sin4; |
717 | 0 | struct sockaddr_in6 sin6; |
718 | 0 | } peeraddr; |
719 | 0 | socklen_t peeraddrlen = sizeof(peeraddr); |
720 | |
|
721 | 0 | #if H2O_USE_ACCEPT4 |
722 | 0 | if ((fd = accept4(listener->fd, &peeraddr.sa, &peeraddrlen, SOCK_NONBLOCK | SOCK_CLOEXEC)) == -1) |
723 | 0 | return NULL; |
724 | 0 | sock = &create_socket(listener->loop, fd, H2O_SOCKET_FLAG_IS_ACCEPTED_CONNECTION)->super; |
725 | | #else |
726 | | if ((fd = cloexec_accept(listener->fd, &peeraddr.sa, &peeraddrlen)) == -1) |
727 | | return NULL; |
728 | | fcntl(fd, F_SETFL, O_NONBLOCK); |
729 | | sock = &create_socket(listener->loop, fd, H2O_SOCKET_FLAG_IS_ACCEPTED_CONNECTION)->super; |
730 | | #endif |
731 | 0 | if (peeraddrlen <= sizeof(peeraddr)) { |
732 | 0 | h2o_socket_setpeername(sock, &peeraddr.sa, peeraddrlen); |
733 | 0 | } else { |
734 | 0 | peeraddr.sa.sa_family = AF_UNSPEC; |
735 | 0 | } |
736 | | |
737 | | /* note: even on linux, the accepted socket might not inherit TCP_NODELAY from the listening socket; see |
738 | | * https://github.com/h2o/h2o/pull/2542#issuecomment-760700859 */ |
739 | 0 | set_nodelay_if_likely_tcp(fd, &peeraddr.sa); |
740 | |
|
741 | 0 | ptls_log_init_conn_state(&sock->_log_state, ptls_openssl_random_bytes, 0, &peeraddr.sa); |
742 | |
|
743 | 0 | return sock; |
744 | 0 | } |
745 | | |
746 | | h2o_socket_t *h2o_socket_connect(h2o_loop_t *loop, struct sockaddr *addr, socklen_t addrlen, h2o_socket_cb cb, const char **err) |
747 | 3.98k | { |
748 | 3.98k | int fd, connect_ret; |
749 | 3.98k | struct st_h2o_evloop_socket_t *sock; |
750 | | |
751 | 3.98k | if ((fd = cloexec_socket(addr->sa_family, SOCK_STREAM, 0)) == -1) { |
752 | 0 | if (err != NULL) { |
753 | 0 | *err = h2o_socket_error_socket_fail; |
754 | 0 | } |
755 | 0 | return NULL; |
756 | 0 | } |
757 | 3.98k | fcntl(fd, F_SETFL, O_NONBLOCK); |
758 | | |
759 | 3.98k | if (!((connect_ret = connect(fd, addr, addrlen)) == 0 || errno == EINPROGRESS)) { |
760 | 0 | if (err != NULL) |
761 | 0 | *err = h2o_socket_get_error_string(errno, h2o_socket_error_conn_fail); |
762 | 0 | close(fd); |
763 | 0 | return NULL; |
764 | 0 | } |
765 | | |
766 | 3.98k | sock = create_socket(loop, fd, H2O_SOCKET_FLAG_IS_CONNECTING); |
767 | 3.98k | set_nodelay_if_likely_tcp(fd, addr); |
768 | | |
769 | 3.98k | if (connect_ret == 0) { |
770 | | /* connection has been established synchronously; notify the fact without going back to epoll */ |
771 | 3.98k | sock->_flags |= H2O_SOCKET_FLAG_IS_WRITE_NOTIFY | H2O_SOCKET_FLAG_IS_CONNECTING_CONNECTED; |
772 | 3.98k | sock->super._cb.write = cb; |
773 | 3.98k | link_to_pending(sock); |
774 | 3.98k | } else { |
775 | 0 | h2o_socket_notify_write(&sock->super, cb); |
776 | 0 | } |
777 | 3.98k | return &sock->super; |
778 | 3.98k | } |
779 | | |
780 | | void h2o_evloop_socket_set_max_read_size(h2o_socket_t *_sock, size_t max_size) |
781 | 4.70k | { |
782 | 4.70k | struct st_h2o_evloop_socket_t *sock = (void *)_sock; |
783 | 4.70k | sock->max_read_size = max_size; |
784 | 4.70k | } |
785 | | |
786 | | h2o_evloop_t *create_evloop(size_t sz) |
787 | 3 | { |
788 | 3 | h2o_evloop_t *loop = h2o_mem_alloc(sz); |
789 | | |
790 | 3 | memset(loop, 0, sz); |
791 | 3 | loop->_statechanged.tail_ref = &loop->_statechanged.head; |
792 | 3 | update_now(loop); |
793 | | /* 3 levels * 32-slots => 1 second goes into 2nd, becomes O(N) above approx. 31 seconds */ |
794 | 3 | loop->_timeouts = h2o_timerwheel_create(3, loop->_now_millisec); |
795 | | |
796 | 3 | return loop; |
797 | 3 | } |
798 | | |
799 | | void update_now(h2o_evloop_t *loop) |
800 | 3.34M | { |
801 | 3.34M | gettimeofday(&loop->_tv_at, NULL); |
802 | 3.34M | loop->_now_nanosec = ((uint64_t)loop->_tv_at.tv_sec * 1000000 + loop->_tv_at.tv_usec) * 1000; |
803 | 3.34M | loop->_now_millisec = loop->_now_nanosec / 1000000; |
804 | 3.34M | } |
805 | | |
806 | | int32_t adjust_max_wait(h2o_evloop_t *loop, int32_t max_wait) |
807 | 1.65M | { |
808 | 1.65M | uint64_t wake_at = h2o_timerwheel_get_wake_at(loop->_timeouts); |
809 | | |
810 | 1.65M | update_now(loop); |
811 | | |
812 | 1.65M | if (wake_at <= loop->_now_millisec) { |
813 | 934 | max_wait = 0; |
814 | 1.64M | } else { |
815 | 1.64M | uint64_t delta = wake_at - loop->_now_millisec; |
816 | 1.64M | if (delta < max_wait) |
817 | 14.8k | max_wait = (int32_t)delta; |
818 | 1.64M | } |
819 | | |
820 | 1.65M | return max_wait; |
821 | 1.65M | } |
822 | | |
823 | | void h2o_socket_notify_write(h2o_socket_t *_sock, h2o_socket_cb cb) |
824 | 119 | { |
825 | 119 | struct st_h2o_evloop_socket_t *sock = (struct st_h2o_evloop_socket_t *)_sock; |
826 | 119 | assert(sock->super._cb.write == NULL); |
827 | 119 | assert(sock->super._cb.write_flags == 0); |
828 | 119 | assert(sock->super._write_buf.cnt == 0); |
829 | 119 | assert(!has_pending_ssl_bytes(sock->super.ssl)); |
830 | | |
831 | 119 | sock->super._cb.write = cb; |
832 | 119 | link_to_statechanged(sock); |
833 | 119 | } |
834 | | |
835 | | static void run_socket(struct st_h2o_evloop_socket_t *sock) |
836 | 75.2k | { |
837 | 75.2k | if ((sock->_flags & H2O_SOCKET_FLAG_IS_DISPOSED) != 0) { |
838 | | /* is freed in updatestates phase */ |
839 | 8.94k | return; |
840 | 8.94k | } |
841 | | |
842 | 66.2k | if ((sock->_flags & H2O_SOCKET_FLAG_IS_READ_READY) != 0) { |
843 | 39.9k | sock->_flags &= ~H2O_SOCKET_FLAG_IS_READ_READY; |
844 | 39.9k | read_on_ready(sock); |
845 | 39.9k | } |
846 | | |
847 | 66.2k | if ((sock->_flags & H2O_SOCKET_FLAG_IS_WRITE_NOTIFY) != 0) { |
848 | 35.8k | const char *err = NULL; |
849 | 35.8k | assert(sock->super._cb.write != NULL); |
850 | 35.8k | sock->_flags &= ~H2O_SOCKET_FLAG_IS_WRITE_NOTIFY; |
851 | 35.8k | if (sock->super._write_buf.cnt != 0 || has_pending_ssl_bytes(sock->super.ssl) || sock->sendvec.callbacks != NULL) { |
852 | | /* error */ |
853 | 591 | err = h2o_socket_error_io; |
854 | 591 | sock->super._write_buf.cnt = 0; |
855 | 591 | if (has_pending_ssl_bytes(sock->super.ssl)) |
856 | 0 | dispose_ssl_output_buffer(sock->super.ssl); |
857 | 591 | sock->sendvec.callbacks = NULL; |
858 | 35.2k | } else if ((sock->_flags & H2O_SOCKET_FLAG_IS_CONNECTING) != 0) { |
859 | | /* completion of connect; determine error if we do not know whether the connection has been successfully estabilshed */ |
860 | 3.89k | if ((sock->_flags & H2O_SOCKET_FLAG_IS_CONNECTING_CONNECTED) == 0) { |
861 | 0 | int so_err = 0; |
862 | 0 | socklen_t l = sizeof(so_err); |
863 | 0 | so_err = 0; |
864 | 0 | if (getsockopt(sock->fd, SOL_SOCKET, SO_ERROR, &so_err, &l) != 0 || so_err != 0) |
865 | 0 | err = h2o_socket_get_error_string(so_err, h2o_socket_error_conn_fail); |
866 | 0 | } |
867 | 3.89k | sock->_flags &= ~(H2O_SOCKET_FLAG_IS_CONNECTING | H2O_SOCKET_FLAG_IS_CONNECTING_CONNECTED); |
868 | 3.89k | } |
869 | 35.8k | on_write_complete(&sock->super, err); |
870 | 35.8k | } |
871 | 66.2k | } |
872 | | |
873 | | static void run_pending(h2o_evloop_t *loop) |
874 | 1.67M | { |
875 | 1.67M | struct st_h2o_evloop_socket_t *sock; |
876 | | |
877 | 1.73M | while (loop->_pending_as_server != NULL || loop->_pending_as_client != NULL) { |
878 | 80.2k | while ((sock = loop->_pending_as_client) != NULL) { |
879 | 14.7k | loop->_pending_as_client = sock->_next_pending; |
880 | 14.7k | sock->_next_pending = sock; |
881 | 14.7k | run_socket(sock); |
882 | 14.7k | } |
883 | 65.4k | if ((sock = loop->_pending_as_server) != NULL) { |
884 | 60.4k | loop->_pending_as_server = sock->_next_pending; |
885 | 60.4k | sock->_next_pending = sock; |
886 | 60.4k | run_socket(sock); |
887 | 60.4k | } |
888 | 65.4k | } |
889 | 1.67M | } |
890 | | |
891 | | void h2o_evloop_destroy(h2o_evloop_t *loop) |
892 | 0 | { |
893 | 0 | struct st_h2o_evloop_socket_t *sock; |
894 | | |
895 | | /* timeouts are governed by the application and MUST be destroyed prior to destroying the loop */ |
896 | 0 | assert(h2o_timerwheel_get_wake_at(loop->_timeouts) == UINT64_MAX); |
897 | | |
898 | | /* dispose all socket */ |
899 | 0 | while ((sock = loop->_pending_as_client) != NULL) { |
900 | 0 | loop->_pending_as_client = sock->_next_pending; |
901 | 0 | sock->_next_pending = sock; |
902 | 0 | h2o_socket_close((h2o_socket_t *)sock); |
903 | 0 | } |
904 | 0 | while ((sock = loop->_pending_as_server) != NULL) { |
905 | 0 | loop->_pending_as_server = sock->_next_pending; |
906 | 0 | sock->_next_pending = sock; |
907 | 0 | h2o_socket_close((h2o_socket_t *)sock); |
908 | 0 | } |
909 | | |
910 | | /* now all socket are disposedand and placed in linked list statechanged |
911 | | * we can freeing memory in cycle by next_statechanged, |
912 | | */ |
913 | 0 | while ((sock = loop->_statechanged.head) != NULL) { |
914 | 0 | loop->_statechanged.head = sock->_next_statechanged; |
915 | 0 | free(sock); |
916 | 0 | } |
917 | | |
918 | | /* dispose backend-specific data */ |
919 | 0 | evloop_do_dispose(loop); |
920 | | |
921 | | /* lastly we need to free loop memory */ |
922 | 0 | h2o_timerwheel_destroy(loop->_timeouts); |
923 | 0 | free(loop); |
924 | 0 | } |
925 | | |
926 | | int h2o_evloop_run(h2o_evloop_t *loop, int32_t max_wait) |
927 | 1.65M | { |
928 | 1.65M | ++loop->run_count; |
929 | | |
930 | | /* Update socket states, poll, set readable flags, perform pending writes. */ |
931 | 1.65M | if (evloop_do_proceed(loop, max_wait) != 0) |
932 | 4 | return -1; |
933 | | |
934 | | /* Run the pending callbacks. */ |
935 | 1.65M | run_pending(loop); |
936 | | |
937 | | /* Run the expired timers at the same time invoking pending callbacks for every timer callback. This is an locality |
938 | | * optimization; handles things like timeout -> write -> on_write_complete for each object. |
939 | | * Expired timers are fetched and run at most 10 times, after which `h2o_evloop_run` returns even if there is a |
940 | | * pending immediate timer. By doing so, we guarantee that the server can make progress by polling the socket, doing |
941 | | * I/O, as well as running other operations coded in the caller of `h2s_evloop_run`, even if there is broken code |
942 | | * that registers an immediate timer perpetually. */ |
943 | 1.67M | for (int i = 0; i < 10; ++i) { |
944 | 1.67M | h2o_linklist_t expired; |
945 | 1.67M | h2o_linklist_init_anchor(&expired); |
946 | 1.67M | h2o_timerwheel_get_expired(loop->_timeouts, loop->_now_millisec, &expired); |
947 | 1.67M | if (h2o_linklist_is_empty(&expired)) |
948 | 1.65M | break; |
949 | 19.9k | do { |
950 | 19.9k | h2o_timerwheel_entry_t *timer = H2O_STRUCT_FROM_MEMBER(h2o_timerwheel_entry_t, _link, expired.next); |
951 | 19.9k | h2o_linklist_unlink(&timer->_link); |
952 | 19.9k | timer->cb(timer); |
953 | 19.9k | run_pending(loop); |
954 | 19.9k | } while (!h2o_linklist_is_empty(&expired)); |
955 | 19.9k | } |
956 | | |
957 | 1.65M | assert(loop->_pending_as_client == NULL); |
958 | 1.65M | assert(loop->_pending_as_server == NULL); |
959 | | |
960 | 1.65M | if (h2o_sliding_counter_is_running(&loop->exec_time_nanosec_counter)) { |
961 | 40.1k | update_now(loop); |
962 | 40.1k | h2o_sliding_counter_stop(&loop->exec_time_nanosec_counter, loop->_now_nanosec); |
963 | 40.1k | } |
964 | | |
965 | 1.65M | return 0; |
966 | 1.65M | } |