|
4 | 4 | #include <fluent-bit/flb_mem.h> |
5 | 5 | #include <fluent-bit/flb_parser.h> |
6 | 6 | #include <fluent-bit/flb_error.h> |
| 7 | +#include <fluent-bit/flb_io.h> |
7 | 8 | #include <fluent-bit/flb_network.h> |
8 | 9 | #include <fluent-bit/flb_socket.h> |
9 | 10 | #include <fluent-bit/flb_time.h> |
10 | 11 |
|
11 | 12 | #include <time.h> |
| 13 | + |
| 14 | +#ifndef FLB_SYSTEM_WINDOWS |
| 15 | +#include <pthread.h> |
| 16 | +#include <signal.h> |
| 17 | +#include <sys/socket.h> |
| 18 | +#include <unistd.h> |
| 19 | +#endif |
| 20 | + |
12 | 21 | #include "flb_tests_internal.h" |
13 | 22 |
|
14 | 23 | #define TEST_HOSTv4 "127.0.0.1" |
@@ -181,9 +190,238 @@ void test_ipv6_bracketed_listen() |
181 | 190 | } |
182 | 191 | } |
183 | 192 |
|
| 193 | +#ifndef FLB_SYSTEM_WINDOWS |
| 194 | + |
| 195 | +#define TEST_WRITE_SIZE (16 * 1024) |
| 196 | +#define TEST_WRITE_MAX_ELAPSED_MILLISECONDS 900 |
| 197 | + |
| 198 | +struct socket_reader_context { |
| 199 | + int fd; |
| 200 | + int close_peer; |
| 201 | + size_t bytes_read; |
| 202 | +}; |
| 203 | + |
| 204 | +static void *socket_reader(void *data) |
| 205 | +{ |
| 206 | + char buffer[4096]; |
| 207 | + ssize_t bytes_read; |
| 208 | + struct socket_reader_context *context; |
| 209 | + |
| 210 | + context = data; |
| 211 | + |
| 212 | + /* |
| 213 | + * Give the writer time to fill the small send buffer and enter the |
| 214 | + * writability wait before draining the peer. |
| 215 | + */ |
| 216 | + flb_time_msleep(20); |
| 217 | + |
| 218 | + if (context->close_peer == FLB_TRUE) { |
| 219 | + flb_socket_close(context->fd); |
| 220 | + return NULL; |
| 221 | + } |
| 222 | + |
| 223 | + while (context->bytes_read < TEST_WRITE_SIZE) { |
| 224 | + bytes_read = recv(context->fd, buffer, sizeof(buffer), 0); |
| 225 | + |
| 226 | + if (bytes_read > 0) { |
| 227 | + context->bytes_read += bytes_read; |
| 228 | + } |
| 229 | + else if (bytes_read < 0 && errno == EINTR) { |
| 230 | + continue; |
| 231 | + } |
| 232 | + else { |
| 233 | + break; |
| 234 | + } |
| 235 | + } |
| 236 | + |
| 237 | + return NULL; |
| 238 | +} |
| 239 | + |
| 240 | +static long elapsed_milliseconds(struct timespec *start, struct timespec *end) |
| 241 | +{ |
| 242 | + return (end->tv_sec - start->tv_sec) * 1000 + |
| 243 | + (end->tv_nsec - start->tv_nsec) / 1000000; |
| 244 | +} |
| 245 | + |
| 246 | +void test_nonblocking_socket_write_waits_for_writability() |
| 247 | +{ |
| 248 | + char *buffer; |
| 249 | + int result; |
| 250 | + int send_buffer_size; |
| 251 | + int sockets[2]; |
| 252 | + long elapsed; |
| 253 | + size_t bytes_written; |
| 254 | + pthread_t reader_thread; |
| 255 | + struct timespec start; |
| 256 | + struct timespec end; |
| 257 | + struct socket_reader_context reader_context; |
| 258 | + |
| 259 | + result = socketpair(AF_UNIX, SOCK_STREAM, 0, sockets); |
| 260 | + if (!TEST_CHECK(result == 0)) { |
| 261 | + return; |
| 262 | + } |
| 263 | + |
| 264 | + send_buffer_size = 4096; |
| 265 | + result = setsockopt(sockets[0], |
| 266 | + SOL_SOCKET, |
| 267 | + SO_SNDBUF, |
| 268 | + &send_buffer_size, |
| 269 | + sizeof(send_buffer_size)); |
| 270 | + if (!TEST_CHECK(result == 0)) { |
| 271 | + flb_socket_close(sockets[0]); |
| 272 | + flb_socket_close(sockets[1]); |
| 273 | + return; |
| 274 | + } |
| 275 | + |
| 276 | + result = flb_net_socket_nonblocking(sockets[0]); |
| 277 | + if (!TEST_CHECK(result == 0)) { |
| 278 | + flb_socket_close(sockets[0]); |
| 279 | + flb_socket_close(sockets[1]); |
| 280 | + return; |
| 281 | + } |
| 282 | + |
| 283 | + buffer = flb_malloc(TEST_WRITE_SIZE); |
| 284 | + if (!TEST_CHECK(buffer != NULL)) { |
| 285 | + flb_socket_close(sockets[0]); |
| 286 | + flb_socket_close(sockets[1]); |
| 287 | + return; |
| 288 | + } |
| 289 | + memset(buffer, 'x', TEST_WRITE_SIZE); |
| 290 | + |
| 291 | + reader_context.fd = sockets[1]; |
| 292 | + reader_context.close_peer = FLB_FALSE; |
| 293 | + reader_context.bytes_read = 0; |
| 294 | + |
| 295 | + result = pthread_create(&reader_thread, NULL, socket_reader, &reader_context); |
| 296 | + if (!TEST_CHECK(result == 0)) { |
| 297 | + flb_free(buffer); |
| 298 | + flb_socket_close(sockets[0]); |
| 299 | + flb_socket_close(sockets[1]); |
| 300 | + return; |
| 301 | + } |
| 302 | + |
| 303 | + clock_gettime(CLOCK_MONOTONIC, &start); |
| 304 | + result = flb_io_fd_write(sockets[0], |
| 305 | + buffer, |
| 306 | + TEST_WRITE_SIZE, |
| 307 | + &bytes_written); |
| 308 | + clock_gettime(CLOCK_MONOTONIC, &end); |
| 309 | + |
| 310 | + shutdown(sockets[0], SHUT_WR); |
| 311 | + pthread_join(reader_thread, NULL); |
| 312 | + |
| 313 | + elapsed = elapsed_milliseconds(&start, &end); |
| 314 | + |
| 315 | + TEST_CHECK(result == TEST_WRITE_SIZE); |
| 316 | + TEST_CHECK(bytes_written == TEST_WRITE_SIZE); |
| 317 | + TEST_CHECK(reader_context.bytes_read == TEST_WRITE_SIZE); |
| 318 | + TEST_CHECK(elapsed < TEST_WRITE_MAX_ELAPSED_MILLISECONDS); |
| 319 | + |
| 320 | + flb_free(buffer); |
| 321 | + flb_socket_close(sockets[0]); |
| 322 | + flb_socket_close(sockets[1]); |
| 323 | +} |
| 324 | + |
| 325 | +void test_nonblocking_socket_write_propagates_peer_close() |
| 326 | +{ |
| 327 | + char *buffer; |
| 328 | + int result; |
| 329 | + int send_buffer_size; |
| 330 | + int socket_error; |
| 331 | + int sockets[2]; |
| 332 | + int write_result; |
| 333 | + size_t bytes_written; |
| 334 | + pthread_t reader_thread; |
| 335 | + struct sigaction ignore_action; |
| 336 | + struct sigaction previous_action; |
| 337 | + struct socket_reader_context reader_context; |
| 338 | + |
| 339 | + result = socketpair(AF_UNIX, SOCK_STREAM, 0, sockets); |
| 340 | + if (!TEST_CHECK(result == 0)) { |
| 341 | + return; |
| 342 | + } |
| 343 | + |
| 344 | + send_buffer_size = 4096; |
| 345 | + result = setsockopt(sockets[0], |
| 346 | + SOL_SOCKET, |
| 347 | + SO_SNDBUF, |
| 348 | + &send_buffer_size, |
| 349 | + sizeof(send_buffer_size)); |
| 350 | + if (!TEST_CHECK(result == 0)) { |
| 351 | + flb_socket_close(sockets[0]); |
| 352 | + flb_socket_close(sockets[1]); |
| 353 | + return; |
| 354 | + } |
| 355 | + |
| 356 | + result = flb_net_socket_nonblocking(sockets[0]); |
| 357 | + if (!TEST_CHECK(result == 0)) { |
| 358 | + flb_socket_close(sockets[0]); |
| 359 | + flb_socket_close(sockets[1]); |
| 360 | + return; |
| 361 | + } |
| 362 | + |
| 363 | + buffer = flb_malloc(TEST_WRITE_SIZE); |
| 364 | + if (!TEST_CHECK(buffer != NULL)) { |
| 365 | + flb_socket_close(sockets[0]); |
| 366 | + flb_socket_close(sockets[1]); |
| 367 | + return; |
| 368 | + } |
| 369 | + memset(buffer, 'x', TEST_WRITE_SIZE); |
| 370 | + |
| 371 | + reader_context.fd = sockets[1]; |
| 372 | + reader_context.close_peer = FLB_TRUE; |
| 373 | + reader_context.bytes_read = 0; |
| 374 | + |
| 375 | + result = pthread_create(&reader_thread, NULL, socket_reader, &reader_context); |
| 376 | + if (!TEST_CHECK(result == 0)) { |
| 377 | + flb_free(buffer); |
| 378 | + flb_socket_close(sockets[0]); |
| 379 | + flb_socket_close(sockets[1]); |
| 380 | + return; |
| 381 | + } |
| 382 | + |
| 383 | + memset(&ignore_action, 0, sizeof(ignore_action)); |
| 384 | + ignore_action.sa_handler = SIG_IGN; |
| 385 | + sigemptyset(&ignore_action.sa_mask); |
| 386 | + result = sigaction(SIGPIPE, &ignore_action, &previous_action); |
| 387 | + if (!TEST_CHECK(result == 0)) { |
| 388 | + pthread_join(reader_thread, NULL); |
| 389 | + flb_free(buffer); |
| 390 | + flb_socket_close(sockets[0]); |
| 391 | + return; |
| 392 | + } |
| 393 | + |
| 394 | + errno = 0; |
| 395 | + write_result = flb_io_fd_write(sockets[0], |
| 396 | + buffer, |
| 397 | + TEST_WRITE_SIZE, |
| 398 | + &bytes_written); |
| 399 | + socket_error = errno; |
| 400 | + |
| 401 | + result = sigaction(SIGPIPE, &previous_action, NULL); |
| 402 | + pthread_join(reader_thread, NULL); |
| 403 | + |
| 404 | + TEST_CHECK(result == 0); |
| 405 | + TEST_CHECK(write_result == -1); |
| 406 | + TEST_CHECK(socket_error == EPIPE || |
| 407 | + socket_error == ECONNRESET || |
| 408 | + socket_error == ENOTCONN); |
| 409 | + |
| 410 | + flb_free(buffer); |
| 411 | + flb_socket_close(sockets[0]); |
| 412 | +} |
| 413 | + |
| 414 | +#endif |
| 415 | + |
184 | 416 | TEST_LIST = { |
185 | 417 | { "ipv4_client_server", test_ipv4_client_server}, |
186 | 418 | { "ipv6_client_server", test_ipv6_client_server}, |
187 | 419 | { "ipv6_bracketed_listen", test_ipv6_bracketed_listen}, |
| 420 | +#ifndef FLB_SYSTEM_WINDOWS |
| 421 | + { "nonblocking_socket_write_waits_for_writability", |
| 422 | + test_nonblocking_socket_write_waits_for_writability}, |
| 423 | + { "nonblocking_socket_write_propagates_peer_close", |
| 424 | + test_nonblocking_socket_write_propagates_peer_close}, |
| 425 | +#endif |
188 | 426 | { 0 } |
189 | 427 | }; |
0 commit comments