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