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