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