1 2 /* 3 * Copyright (C) Igor Sysoev 4 * Copyright (C) NGINX, Inc. 5 */ 6 7 #include <nxt_main.h> 8 #include <nxt_runtime.h> 9 #include <nxt_port.h> 10 #include <nxt_router.h> 11 12 13 static void nxt_port_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg); 14 15 static nxt_atomic_uint_t nxt_port_last_id = 1; 16 17 18 static void 19 nxt_port_mp_cleanup(nxt_task_t *task, void *obj, void *data) 20 { 21 nxt_mp_t *mp; 22 nxt_port_t *port; 23 24 port = obj; 25 mp = data; 26 27 nxt_assert(port->pair[0] == -1); 28 nxt_assert(port->pair[1] == -1); 29 30 nxt_assert(port->use_count == 0); 31 nxt_assert(port->app_link.next == NULL); 32 nxt_assert(port->idle_link.next == NULL); 33 34 nxt_assert(nxt_queue_is_empty(&port->messages)); 35 nxt_assert(nxt_lvlhsh_is_empty(&port->rpc_streams)); 36 nxt_assert(nxt_lvlhsh_is_empty(&port->rpc_peers)); 37 38 nxt_thread_mutex_destroy(&port->write_mutex); 39 40 nxt_mp_free(mp, port); 41 } 42 43 44 nxt_port_t * 45 nxt_port_new(nxt_task_t *task, nxt_port_id_t id, nxt_pid_t pid, 46 nxt_process_type_t type) 47 { 48 nxt_mp_t *mp; 49 nxt_port_t *port; 50 51 mp = nxt_mp_create(1024, 128, 256, 32); 52 53 if (nxt_slow_path(mp == NULL)) { 54 return NULL; 55 } 56 57 port = nxt_mp_zalloc(mp, sizeof(nxt_port_t)); 58 59 if (nxt_fast_path(port != NULL)) { 60 port->id = id; 61 port->pid = pid; 62 port->type = type; 63 port->mem_pool = mp; 64 port->use_count = 1; 65 66 nxt_mp_cleanup(mp, nxt_port_mp_cleanup, task, port, mp); 67 68 nxt_queue_init(&port->messages); 69 nxt_thread_mutex_create(&port->write_mutex); 70 nxt_queue_init(&port->pending_requests); 71 nxt_queue_init(&port->active_websockets); 72 73 } else { 74 nxt_mp_destroy(mp); 75 } 76 77 nxt_thread_log_debug("port %p %d:%d new, type %d", port, pid, id, type); 78 79 return port; 80 } 81 82 83 void 84 nxt_port_close(nxt_task_t *task, nxt_port_t *port) 85 { 86 nxt_debug(task, "port %p %d:%d close, type %d", port, port->pid, 87 port->id, port->type); 88 89 if (port->pair[0] != -1) { 90 nxt_port_rpc_close(task, port); 91 92 nxt_fd_close(port->pair[0]); 93 port->pair[0] = -1; 94 } 95 96 if (port->pair[1] != -1) { 97 nxt_fd_close(port->pair[1]); 98 port->pair[1] = -1; 99 100 if (port->app != NULL) { 101 nxt_router_app_port_close(task, port); 102 } 103 } 104 } 105 106 107 static void 108 nxt_port_release(nxt_task_t *task, nxt_port_t *port) 109 { 110 nxt_debug(task, "port %p %d:%d release, type %d", port, port->pid, 111 port->id, port->type); 112 113 port->app = NULL; 114 115 if (port->link.next != NULL) { 116 nxt_assert(port->process != NULL); 117 118 nxt_process_port_remove(port); 119 120 nxt_process_use(task, port->process, -1); 121 } 122 123 nxt_mp_release(port->mem_pool); 124 } 125 126 127 nxt_port_id_t 128 nxt_port_get_next_id() 129 { 130 return nxt_atomic_fetch_add(&nxt_port_last_id, 1); 131 } 132 133 134 void 135 nxt_port_reset_next_id() 136 { 137 nxt_port_last_id = 1; 138 } 139 140 141 void 142 nxt_port_enable(nxt_task_t *task, nxt_port_t *port, 143 nxt_port_handlers_t *handlers) 144 { 145 port->pid = nxt_pid; 146 port->handler = nxt_port_handler; 147 port->data = (nxt_port_handler_t *) (handlers); 148 149 nxt_port_read_enable(task, port); 150 } 151 152 153 static void 154 nxt_port_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 155 { 156 nxt_port_handler_t *handlers; 157 158 if (nxt_fast_path(msg->port_msg.type < NXT_PORT_MSG_MAX)) { 159 160 nxt_debug(task, "port %d: message type:%uD", 161 msg->port->socket.fd, msg->port_msg.type); 162 163 handlers = msg->port->data; 164 handlers[msg->port_msg.type](task, msg); 165 166 return; 167 } 168 169 nxt_alert(task, "port %d: unknown message type:%uD", 170 msg->port->socket.fd, msg->port_msg.type); 171 } 172 173 174 void 175 nxt_port_quit_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 176 { 177 nxt_runtime_quit(task, 0); 178 } 179 180 181 nxt_inline void 182 nxt_port_send_new_port(nxt_task_t *task, nxt_runtime_t *rt, 183 nxt_port_t *new_port, uint32_t stream) 184 { 185 nxt_port_t *port; 186 nxt_process_t *process; 187 188 nxt_debug(task, "new port %d for process %PI", 189 new_port->pair[1], new_port->pid); 190 191 nxt_runtime_process_each(rt, process) { 192 193 if (process->pid == new_port->pid || process->pid == nxt_pid) { 194 continue; 195 } 196 197 port = nxt_process_port_first(process); 198 199 if (nxt_proc_conn_matrix[port->type][new_port->type]) { 200 (void) nxt_port_send_port(task, port, new_port, stream); 201 } 202 203 } nxt_runtime_process_loop; 204 } 205 206 207 nxt_int_t 208 nxt_port_send_port(nxt_task_t *task, nxt_port_t *port, nxt_port_t *new_port, 209 uint32_t stream) 210 { 211 nxt_buf_t *b; 212 nxt_port_msg_new_port_t *msg; 213 214 b = nxt_buf_mem_ts_alloc(task, task->thread->engine->mem_pool, 215 sizeof(nxt_port_data_t)); 216 if (nxt_slow_path(b == NULL)) { 217 return NXT_ERROR; 218 } 219 220 nxt_debug(task, "send port %FD to process %PI", 221 new_port->pair[1], port->pid); 222 223 b->mem.free += sizeof(nxt_port_msg_new_port_t); 224 msg = (nxt_port_msg_new_port_t *) b->mem.pos; 225 226 msg->id = new_port->id; 227 msg->pid = new_port->pid; 228 msg->max_size = port->max_size; 229 msg->max_share = port->max_share; 230 msg->type = new_port->type; 231 232 return nxt_port_socket_write(task, port, NXT_PORT_MSG_NEW_PORT, 233 new_port->pair[1], stream, 0, b); 234 } 235 236 237 void 238 nxt_port_new_port_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 239 { 240 nxt_port_t *port; 241 nxt_process_t *process; 242 nxt_runtime_t *rt; 243 nxt_port_msg_new_port_t *new_port_msg; 244 245 rt = task->thread->runtime; 246 247 new_port_msg = (nxt_port_msg_new_port_t *) msg->buf->mem.pos; 248 249 /* TODO check b size and make plain */ 250 251 nxt_debug(task, "new port %d received for process %PI:%d", 252 msg->fd, new_port_msg->pid, new_port_msg->id); 253 254 port = nxt_runtime_port_find(rt, new_port_msg->pid, new_port_msg->id); 255 if (port != NULL) { 256 nxt_debug(task, "port %PI:%d already exists", new_port_msg->pid, 257 new_port_msg->id); 258 259 nxt_fd_close(msg->fd); 260 msg->fd = -1; 261 return; 262 } 263 264 process = nxt_runtime_process_get(rt, new_port_msg->pid); 265 if (nxt_slow_path(process == NULL)) { 266 return; 267 } 268 269 port = nxt_port_new(task, new_port_msg->id, new_port_msg->pid, 270 new_port_msg->type); 271 if (nxt_slow_path(port == NULL)) { 272 nxt_process_use(task, process, -1); 273 return; 274 } 275 276 nxt_process_port_add(task, process, port); 277 278 nxt_process_use(task, process, -1); 279 280 nxt_fd_nonblocking(task, msg->fd); 281 282 port->pair[0] = -1; 283 port->pair[1] = msg->fd; 284 port->max_size = new_port_msg->max_size; 285 port->max_share = new_port_msg->max_share; 286 287 port->socket.task = task; 288 289 nxt_runtime_port_add(task, port); 290 291 nxt_port_use(task, port, -1); 292 293 nxt_port_write_enable(task, port); 294 295 msg->u.new_port = port; 296 } 297 298 299 void 300 nxt_port_process_ready_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 301 { 302 nxt_port_t *port; 303 nxt_process_t *process; 304 nxt_runtime_t *rt; 305 306 rt = task->thread->runtime; 307 308 process = nxt_runtime_process_find(rt, msg->port_msg.pid); 309 if (nxt_slow_path(process == NULL)) { 310 return; 311 } 312 313 process->ready = 1; 314 315 nxt_assert(!nxt_queue_is_empty(&process->ports)); 316 317 port = nxt_process_port_first(process); 318 319 nxt_debug(task, "process %PI ready", msg->port_msg.pid); 320 321 nxt_port_send_new_port(task, rt, port, msg->port_msg.stream); 322 } 323 324 325 void 326 nxt_port_mmap_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 327 { 328 nxt_runtime_t *rt; 329 nxt_process_t *process; 330 331 rt = task->thread->runtime; 332 333 if (nxt_slow_path(msg->fd == -1)) { 334 nxt_log(task, NXT_LOG_WARN, "invalid fd passed with mmap message"); 335 336 return; 337 } 338 339 process = nxt_runtime_process_find(rt, msg->port_msg.pid); 340 if (nxt_slow_path(process == NULL)) { 341 nxt_log(task, NXT_LOG_WARN, "failed to get process #%PI", 342 msg->port_msg.pid); 343 344 goto fail_close; 345 } 346 347 nxt_port_incoming_port_mmap(task, process, msg->fd); 348 349 fail_close: 350 351 nxt_fd_close(msg->fd); 352 } 353 354 355 void 356 nxt_port_change_log_file(nxt_task_t *task, nxt_runtime_t *rt, nxt_uint_t slot, 357 nxt_fd_t fd) 358 { 359 nxt_buf_t *b; 360 nxt_port_t *port; 361 nxt_process_t *process; 362 363 nxt_debug(task, "change log file #%ui fd:%FD", slot, fd); 364 365 nxt_runtime_process_each(rt, process) { 366 367 if (nxt_pid == process->pid) { 368 continue; 369 } 370 371 port = nxt_process_port_first(process); 372 373 b = nxt_buf_mem_ts_alloc(task, task->thread->engine->mem_pool, 374 sizeof(nxt_port_data_t)); 375 if (nxt_slow_path(b == NULL)) { 376 continue; 377 } 378 379 *(nxt_uint_t *) b->mem.pos = slot; 380 b->mem.free += sizeof(nxt_uint_t); 381 382 (void) nxt_port_socket_write(task, port, NXT_PORT_MSG_CHANGE_FILE, 383 fd, 0, 0, b); 384 385 } nxt_runtime_process_loop; 386 } 387 388 389 void 390 nxt_port_change_log_file_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 391 { 392 nxt_buf_t *b; 393 nxt_uint_t slot; 394 nxt_file_t *log_file; 395 nxt_runtime_t *rt; 396 397 rt = task->thread->runtime; 398 399 b = msg->buf; 400 slot = *(nxt_uint_t *) b->mem.pos; 401 402 log_file = nxt_list_elt(rt->log_files, slot); 403 404 nxt_debug(task, "change log file %FD:%FD", msg->fd, log_file->fd); 405 406 /* 407 * The old log file descriptor must be closed at the moment when no 408 * other threads use it. dup2() allows to use the old file descriptor 409 * for new log file. This change is performed atomically in the kernel. 410 */ 411 if (nxt_file_redirect(log_file, msg->fd) == NXT_OK) { 412 if (slot == 0) { 413 (void) nxt_file_stderr(log_file); 414 } 415 } 416 } 417 418 419 void 420 nxt_port_data_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 421 { 422 size_t dump_size; 423 nxt_buf_t *b; 424 425 b = msg->buf; 426 dump_size = b->mem.free - b->mem.pos; 427 428 if (dump_size > 300) { 429 dump_size = 300; 430 } 431 432 nxt_debug(task, "data: %*s", dump_size, b->mem.pos); 433 } 434 435 436 void 437 nxt_port_remove_pid_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 438 { 439 nxt_buf_t *buf; 440 nxt_pid_t pid; 441 nxt_runtime_t *rt; 442 nxt_process_t *process; 443 444 buf = msg->buf; 445 446 nxt_assert(nxt_buf_used_size(buf) == sizeof(pid)); 447 448 nxt_memcpy(&pid, buf->mem.pos, sizeof(pid)); 449 450 msg->u.removed_pid = pid; 451 452 nxt_debug(task, "port remove pid %PI handler", pid); 453 454 rt = task->thread->runtime; 455 456 nxt_port_rpc_remove_peer(task, msg->port, pid); 457 458 process = nxt_runtime_process_find(rt, pid); 459 460 if (process) { 461 nxt_process_close_ports(task, process); 462 } 463 } 464 465 466 void 467 nxt_port_empty_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg) 468 { 469 nxt_debug(task, "port empty handler"); 470 } 471 472 473 typedef struct { 474 nxt_work_t work; 475 nxt_port_t *port; 476 nxt_port_post_handler_t handler; 477 } nxt_port_work_t; 478 479 480 static void 481 nxt_port_post_handler(nxt_task_t *task, void *obj, void *data) 482 { 483 nxt_port_t *port; 484 nxt_port_work_t *pw; 485 nxt_port_post_handler_t handler; 486 487 pw = obj; 488 port = pw->port; 489 handler = pw->handler; 490 491 nxt_free(pw); 492 493 handler(task, port, data); 494 495 nxt_port_use(task, port, -1); 496 } 497 498 499 nxt_int_t 500 nxt_port_post(nxt_task_t *task, nxt_port_t *port, 501 nxt_port_post_handler_t handler, void *data) 502 { 503 nxt_port_work_t *pw; 504 505 if (task->thread->engine == port->engine) { 506 handler(task, port, data); 507 508 return NXT_OK; 509 } 510 511 pw = nxt_zalloc(sizeof(nxt_port_work_t)); 512 513 if (nxt_slow_path(pw == NULL)) { 514 return NXT_ERROR; 515 } 516 517 nxt_atomic_fetch_add(&port->use_count, 1); 518 519 pw->work.handler = nxt_port_post_handler; 520 pw->work.task = &port->engine->task; 521 pw->work.obj = pw; 522 pw->work.data = data; 523 524 pw->port = port; 525 pw->handler = handler; 526 527 nxt_event_engine_post(port->engine, &pw->work); 528 529 return NXT_OK; 530 } 531 532 533 static void 534 nxt_port_release_handler(nxt_task_t *task, nxt_port_t *port, void *data) 535 { 536 /* no op */ 537 } 538 539 540 void 541 nxt_port_use(nxt_task_t *task, nxt_port_t *port, int i) 542 { 543 int c; 544 545 c = nxt_atomic_fetch_add(&port->use_count, i); 546 547 if (i < 0 && c == -i) { 548 549 if (task->thread->engine == port->engine) { 550 nxt_port_release(task, port); 551 552 return; 553 } 554 555 nxt_port_post(task, port, nxt_port_release_handler, NULL); 556 } 557 } 558