xref: /unit/src/nxt_port.c (revision 1131:ec7d924d8dfb)
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