summaryrefslogtreecommitdiffhomepage
path: root/src/nxt_router.c
diff options
context:
space:
mode:
Diffstat (limited to '')
-rw-r--r--src/nxt_router.c279
1 files changed, 211 insertions, 68 deletions
diff --git a/src/nxt_router.c b/src/nxt_router.c
index 1dd7f8e5..7623ccbb 100644
--- a/src/nxt_router.c
+++ b/src/nxt_router.c
@@ -65,12 +65,14 @@ typedef struct {
typedef struct {
nxt_app_t *app;
nxt_router_temp_conf_t *temp_conf;
+ uint8_t proto; /* 1 bit */
} nxt_app_rpc_t;
typedef struct {
nxt_app_joint_t *app_joint;
uint32_t generation;
+ uint8_t proto; /* 1 bit */
} nxt_app_joint_rpc_t;
@@ -392,32 +394,52 @@ nxt_router_start_app_process_handler(nxt_task_t *task, nxt_port_t *port,
{
size_t size;
uint32_t stream;
- nxt_mp_t *mp;
nxt_int_t ret;
nxt_app_t *app;
nxt_buf_t *b;
- nxt_port_t *main_port;
+ nxt_port_t *dport;
nxt_runtime_t *rt;
nxt_app_joint_rpc_t *app_joint_rpc;
app = data;
- rt = task->thread->runtime;
- main_port = rt->port_by_type[NXT_PROCESS_MAIN];
+ nxt_thread_mutex_lock(&app->mutex);
- nxt_debug(task, "app '%V' %p start process", &app->name, app);
+ dport = app->proto_port;
- size = app->name.length + 1 + app->conf.length;
+ nxt_thread_mutex_unlock(&app->mutex);
- b = nxt_buf_mem_ts_alloc(task, task->thread->engine->mem_pool, size);
+ if (dport != NULL) {
+ nxt_debug(task, "app '%V' %p start process", &app->name, app);
- if (nxt_slow_path(b == NULL)) {
- goto failed;
- }
+ b = NULL;
- nxt_buf_cpystr(b, &app->name);
- *b->mem.free++ = '\0';
- nxt_buf_cpystr(b, &app->conf);
+ } else {
+ if (app->proto_port_requests > 0) {
+ nxt_debug(task, "app '%V' %p wait for prototype process",
+ &app->name, app);
+
+ app->proto_port_requests++;
+
+ goto skip;
+ }
+
+ nxt_debug(task, "app '%V' %p start prototype process", &app->name, app);
+
+ rt = task->thread->runtime;
+ dport = rt->port_by_type[NXT_PROCESS_MAIN];
+
+ size = app->name.length + 1 + app->conf.length;
+
+ b = nxt_buf_mem_alloc(task->thread->engine->mem_pool, size, 0);
+ if (nxt_slow_path(b == NULL)) {
+ goto failed;
+ }
+
+ nxt_buf_cpystr(b, &app->name);
+ *b->mem.free++ = '\0';
+ nxt_buf_cpystr(b, &app->conf);
+ }
app_joint_rpc = nxt_port_rpc_register_handler_ex(task, port,
nxt_router_app_port_ready,
@@ -429,7 +451,7 @@ nxt_router_start_app_process_handler(nxt_task_t *task, nxt_port_t *port,
stream = nxt_port_rpc_ex_stream(app_joint_rpc);
- ret = nxt_port_socket_write(task, main_port, NXT_PORT_MSG_START_PROCESS,
+ ret = nxt_port_socket_write(task, dport, NXT_PORT_MSG_START_PROCESS,
-1, stream, port->id, b);
if (nxt_slow_path(ret != NXT_OK)) {
nxt_port_rpc_cancel(task, port, stream);
@@ -439,26 +461,23 @@ nxt_router_start_app_process_handler(nxt_task_t *task, nxt_port_t *port,
app_joint_rpc->app_joint = app->joint;
app_joint_rpc->generation = app->generation;
+ app_joint_rpc->proto = (b != NULL);
- nxt_router_app_joint_use(task, app->joint, 1);
+ if (b != NULL) {
+ app->proto_port_requests++;
- nxt_router_app_use(task, app, -1);
+ b = NULL;
+ }
- return;
+ nxt_router_app_joint_use(task, app->joint, 1);
failed:
if (b != NULL) {
- mp = b->data;
- nxt_mp_free(mp, b);
- nxt_mp_release(mp);
+ nxt_mp_free(b->data, b);
}
- nxt_thread_mutex_lock(&app->mutex);
-
- app->pending_processes--;
-
- nxt_thread_mutex_unlock(&app->mutex);
+skip:
nxt_router_app_use(task, app, -1);
}
@@ -658,6 +677,12 @@ nxt_router_new_port_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg)
nxt_router_greet_controller(task, msg->u.new_port);
}
+ if (port != NULL && port->type == NXT_PROCESS_PROTOTYPE) {
+ nxt_port_rpc_handler(task, msg);
+
+ return;
+ }
+
if (port == NULL || port->type != NXT_PROCESS_APP) {
if (msg->port_msg.stream == 0) {
@@ -683,6 +708,8 @@ nxt_router_new_port_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg)
return;
}
+ nxt_debug(task, "new port id %d (%d)", port->id, port->type);
+
/*
* Port with "id == 0" is application 'main' port and it always
* should come with non-zero stream.
@@ -819,7 +846,8 @@ nxt_router_app_restart_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg)
nxt_app_t *app;
nxt_int_t ret;
nxt_str_t app_name;
- nxt_port_t *port, *reply_port, *shared_port, *old_shared_port;
+ nxt_port_t *reply_port, *shared_port, *old_shared_port;
+ nxt_port_t *proto_port;
nxt_port_msg_type_t reply;
reply_port = nxt_runtime_port_find(task->thread->runtime,
@@ -862,12 +890,15 @@ nxt_router_app_restart_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg)
nxt_thread_mutex_lock(&app->mutex);
- nxt_queue_each(port, &app->ports, nxt_port_t, app_link) {
+ proto_port = app->proto_port;
- (void) nxt_port_socket_write(task, port, NXT_PORT_MSG_QUIT, -1,
- 0, 0, NULL);
+ if (proto_port != NULL) {
+ nxt_debug(task, "send QUIT to prototype '%V' pid %PI", &app->name,
+ proto_port->pid);
- } nxt_queue_loop;
+ app->proto_port = NULL;
+ proto_port->app = NULL;
+ }
app->generation++;
@@ -883,6 +914,15 @@ nxt_router_app_restart_handler(nxt_task_t *task, nxt_port_recv_msg_t *msg)
nxt_port_close(task, old_shared_port);
nxt_port_use(task, old_shared_port, -1);
+ if (proto_port != NULL) {
+ (void) nxt_port_socket_write(task, proto_port, NXT_PORT_MSG_QUIT,
+ -1, 0, 0, NULL);
+
+ nxt_port_close(task, proto_port);
+
+ nxt_port_use(task, proto_port, -1);
+ }
+
reply = NXT_PORT_MSG_RPC_READY_LAST;
} else {
@@ -2735,54 +2775,67 @@ nxt_router_app_rpc_create(nxt_task_t *task,
uint32_t stream;
nxt_int_t ret;
nxt_buf_t *b;
- nxt_port_t *main_port, *router_port;
+ nxt_port_t *router_port, *dport;
nxt_runtime_t *rt;
nxt_app_rpc_t *rpc;
- rpc = nxt_mp_alloc(tmcf->mem_pool, sizeof(nxt_app_rpc_t));
- if (rpc == NULL) {
- goto fail;
- }
+ rt = task->thread->runtime;
- rpc->app = app;
- rpc->temp_conf = tmcf;
+ dport = app->proto_port;
- nxt_debug(task, "app '%V' prefork", &app->name);
+ if (dport == NULL) {
+ nxt_debug(task, "app '%V' prototype prefork", &app->name);
- size = app->name.length + 1 + app->conf.length;
+ size = app->name.length + 1 + app->conf.length;
- b = nxt_buf_mem_alloc(tmcf->mem_pool, size, 0);
- if (nxt_slow_path(b == NULL)) {
- goto fail;
- }
+ b = nxt_buf_mem_alloc(tmcf->mem_pool, size, 0);
+ if (nxt_slow_path(b == NULL)) {
+ goto fail;
+ }
- b->completion_handler = nxt_buf_dummy_completion;
+ b->completion_handler = nxt_buf_dummy_completion;
- nxt_buf_cpystr(b, &app->name);
- *b->mem.free++ = '\0';
- nxt_buf_cpystr(b, &app->conf);
+ nxt_buf_cpystr(b, &app->name);
+ *b->mem.free++ = '\0';
+ nxt_buf_cpystr(b, &app->conf);
+
+ dport = rt->port_by_type[NXT_PROCESS_MAIN];
+
+ } else {
+ nxt_debug(task, "app '%V' prefork", &app->name);
+
+ b = NULL;
+ }
- rt = task->thread->runtime;
- main_port = rt->port_by_type[NXT_PROCESS_MAIN];
router_port = rt->port_by_type[NXT_PROCESS_ROUTER];
- stream = nxt_port_rpc_register_handler(task, router_port,
+ rpc = nxt_port_rpc_register_handler_ex(task, router_port,
nxt_router_app_prefork_ready,
nxt_router_app_prefork_error,
- -1, rpc);
- if (nxt_slow_path(stream == 0)) {
+ sizeof(nxt_app_rpc_t));
+ if (nxt_slow_path(rpc == NULL)) {
goto fail;
}
- ret = nxt_port_socket_write(task, main_port, NXT_PORT_MSG_START_PROCESS,
- -1, stream, router_port->id, b);
+ rpc->app = app;
+ rpc->temp_conf = tmcf;
+ rpc->proto = (b != NULL);
+
+ stream = nxt_port_rpc_ex_stream(rpc);
+ ret = nxt_port_socket_write(task, dport,
+ NXT_PORT_MSG_START_PROCESS,
+ -1, stream, router_port->id, b);
if (nxt_slow_path(ret != NXT_OK)) {
nxt_port_rpc_cancel(task, router_port, stream);
goto fail;
}
- app->pending_processes++;
+ if (b == NULL) {
+ nxt_port_rpc_ex_set_peer(task, router_port, rpc, dport->pid);
+
+ app->pending_processes++;
+ }
return;
@@ -2807,9 +2860,24 @@ nxt_router_app_prefork_ready(nxt_task_t *task, nxt_port_recv_msg_t *msg,
port = msg->u.new_port;
nxt_assert(port != NULL);
- nxt_assert(port->type == NXT_PROCESS_APP);
nxt_assert(port->id == 0);
+ if (rpc->proto) {
+ nxt_assert(app->proto_port == NULL);
+ nxt_assert(port->type == NXT_PROCESS_PROTOTYPE);
+
+ nxt_port_inc_use(port);
+
+ app->proto_port = port;
+ port->app = app;
+
+ nxt_router_app_rpc_create(task, rpc->temp_conf, app);
+
+ return;
+ }
+
+ nxt_assert(port->type == NXT_PROCESS_APP);
+
port->app = app;
port->main_app_port = port;
@@ -2851,10 +2919,16 @@ nxt_router_app_prefork_error(nxt_task_t *task, nxt_port_recv_msg_t *msg,
app = rpc->app;
tmcf = rpc->temp_conf;
- nxt_log(task, NXT_LOG_WARN, "failed to start application \"%V\"",
- &app->name);
+ if (rpc->proto) {
+ nxt_log(task, NXT_LOG_WARN, "failed to start prototype \"%V\"",
+ &app->name);
- app->pending_processes--;
+ } else {
+ nxt_log(task, NXT_LOG_WARN, "failed to start application \"%V\"",
+ &app->name);
+
+ app->pending_processes--;
+ }
nxt_router_conf_error(task, tmcf);
}
@@ -4413,8 +4487,9 @@ static void
nxt_router_app_port_ready(nxt_task_t *task, nxt_port_recv_msg_t *msg,
void *data)
{
+ uint32_t n;
nxt_app_t *app;
- nxt_bool_t start_process;
+ nxt_bool_t start_process, restarted;
nxt_port_t *port;
nxt_app_joint_t *app_joint;
nxt_app_joint_rpc_t *app_joint_rpc;
@@ -4427,7 +4502,6 @@ nxt_router_app_port_ready(nxt_task_t *task, nxt_port_recv_msg_t *msg,
nxt_assert(app_joint != NULL);
nxt_assert(port != NULL);
- nxt_assert(port->type == NXT_PROCESS_APP);
nxt_assert(port->id == 0);
app = app_joint->app;
@@ -4444,11 +4518,51 @@ nxt_router_app_port_ready(nxt_task_t *task, nxt_port_recv_msg_t *msg,
nxt_thread_mutex_lock(&app->mutex);
+ restarted = (app->generation != app_joint_rpc->generation);
+
+ if (app_joint_rpc->proto) {
+ nxt_assert(app->proto_port == NULL);
+ nxt_assert(port->type == NXT_PROCESS_PROTOTYPE);
+
+ n = app->proto_port_requests;
+ app->proto_port_requests = 0;
+
+ if (nxt_slow_path(restarted)) {
+ nxt_thread_mutex_unlock(&app->mutex);
+
+ nxt_debug(task, "proto port ready for restarted app, send QUIT");
+
+ nxt_port_socket_write(task, port, NXT_PORT_MSG_QUIT, -1, 0, 0,
+ NULL);
+
+ } else {
+ port->app = app;
+ app->proto_port = port;
+
+ nxt_thread_mutex_unlock(&app->mutex);
+
+ nxt_port_use(task, port, 1);
+ }
+
+ port = task->thread->runtime->port_by_type[NXT_PROCESS_ROUTER];
+
+ while (n > 0) {
+ nxt_router_app_use(task, app, 1);
+
+ nxt_router_start_app_process_handler(task, port, app);
+
+ n--;
+ }
+
+ return;
+ }
+
+ nxt_assert(port->type == NXT_PROCESS_APP);
nxt_assert(app->pending_processes != 0);
app->pending_processes--;
- if (nxt_slow_path(app->generation != app_joint_rpc->generation)) {
+ if (nxt_slow_path(restarted)) {
nxt_debug(task, "new port ready for restarted app, send QUIT");
start_process = !task->thread->engine->shutdown
@@ -4591,7 +4705,6 @@ nxt_router_app_port_error(nxt_task_t *task, nxt_port_recv_msg_t *msg,
}
-
nxt_inline nxt_port_t *
nxt_router_app_get_port_for_quit(nxt_task_t *task, nxt_app_t *app)
{
@@ -4791,6 +4904,20 @@ nxt_router_app_port_close(nxt_task_t *task, nxt_port_t *port)
nxt_thread_mutex_lock(&app->mutex);
+ if (port == app->proto_port) {
+ app->proto_port = NULL;
+ port->app = NULL;
+
+ nxt_thread_mutex_unlock(&app->mutex);
+
+ nxt_debug(task, "app '%V' prototype pid %PI closed", &app->name,
+ port->pid);
+
+ nxt_port_use(task, port, -1);
+
+ return;
+ }
+
nxt_port_hash_remove(&app->port_hash, port);
app->port_hash_count--;
@@ -4979,7 +5106,7 @@ static void
nxt_router_free_app(nxt_task_t *task, void *obj, void *data)
{
nxt_app_t *app;
- nxt_port_t *port;
+ nxt_port_t *port, *proto_port;
nxt_app_joint_t *app_joint;
app_joint = obj;
@@ -4991,10 +5118,6 @@ nxt_router_free_app(nxt_task_t *task, void *obj, void *data)
break;
}
- nxt_debug(task, "send QUIT to app '%V' pid %PI", &app->name, port->pid);
-
- nxt_port_socket_write(task, port, NXT_PORT_MSG_QUIT, -1, 0, 0, NULL);
-
nxt_port_use(task, port, -1);
}
@@ -5015,8 +5138,28 @@ nxt_router_free_app(nxt_task_t *task, void *obj, void *data)
nxt_port_use(task, port, -1);
}
+ proto_port = app->proto_port;
+
+ if (proto_port != NULL) {
+ nxt_debug(task, "send QUIT to prototype '%V' pid %PI", &app->name,
+ proto_port->pid);
+
+ app->proto_port = NULL;
+ proto_port->app = NULL;
+ }
+
nxt_thread_mutex_unlock(&app->mutex);
+ if (proto_port != NULL) {
+ nxt_port_socket_write(task, proto_port, NXT_PORT_MSG_QUIT,
+ -1, 0, 0, NULL);
+
+ nxt_port_close(task, proto_port);
+
+ nxt_port_use(task, proto_port, -1);
+ }
+
+ nxt_assert(app->proto_port == NULL);
nxt_assert(app->processes == 0);
nxt_assert(app->active_requests == 0);
nxt_assert(app->port_hash_count == 0);