xref: /unit/src/nxt_port.c (revision 753:0f0094dc0f51)
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