From c3eb5cd20235f643ac61ad7e5cf67965dd9d5beb Mon Sep 17 00:00:00 2001 From: Cloud Wu Date: Fri, 20 Jun 2014 02:49:48 +0800 Subject: [PATCH] After connecting, socket can send before connected. remove block connect api --- examples/config_log | 6 ++++-- examples/globallog.lua | 16 ++++++++++++++++ examples/main_log.lua | 1 - service-src/service_harbor.c | 24 +++++++++++++++--------- service-src/service_master.c | 3 ++- service/bootstrap.lua | 2 +- skynet-src/skynet_socket.c | 6 ------ skynet-src/skynet_socket.h | 1 - skynet-src/socket_server.c | 35 ++++++++++------------------------- skynet-src/socket_server.h | 2 -- 10 files changed, 48 insertions(+), 48 deletions(-) diff --git a/examples/config_log b/examples/config_log index 9261e435..42a8692b 100644 --- a/examples/config_log +++ b/examples/config_log @@ -3,8 +3,10 @@ mqueue = 256 cpath = "./cservice/?.so" logger = nil harbor = 2 -address = "127.0.0.1:2527" -master = "127.0.0.1:2013" +--address = "127.0.0.1:2527" +--master = "127.0.0.1:2013" +address = "172.16.100.201:2527" +master = "172.16.100.209:2013" start = "main_log" luaservice ="./service/?.lua;./test/?.lua;./examples/?.lua" snax = "./examples/?.lua;./test/?.lua" diff --git a/examples/globallog.lua b/examples/globallog.lua index ec431cce..9b4dce4c 100644 --- a/examples/globallog.lua +++ b/examples/globallog.lua @@ -1,5 +1,21 @@ local skynet = require "skynet" +skynet.register_protocol { + name = "text", + id = skynet.PTYPE_TEXT, + pack = function (...) + local n = select ("#" , ...) + if n == 0 then + return "" + elseif n == 1 then + return tostring(...) + else + return table.concat({...}," ") + end + end, + unpack = skynet.tostring +} + skynet.start(function() skynet.dispatch("text", function(session, address, text) print("[GLOBALLOG]", skynet.address(address),text) diff --git a/examples/main_log.lua b/examples/main_log.lua index 8a5194a0..f74b563f 100644 --- a/examples/main_log.lua +++ b/examples/main_log.lua @@ -2,7 +2,6 @@ local skynet = require "skynet" skynet.start(function() print("Log server start") - local service = skynet.newservice("service_mgr") skynet.monitor "simplemonitor" local log = skynet.newservice("globallog") skynet.exit() diff --git a/service-src/service_harbor.c b/service-src/service_harbor.c index aa6763a7..a9cfaec4 100644 --- a/service-src/service_harbor.c +++ b/service-src/service_harbor.c @@ -302,7 +302,7 @@ harbor_release(struct harbor *h) { } static int -_connect_to(struct harbor *h, const char *ipaddress, bool blocking) { +_connect_to(struct harbor *h, const char *ipaddress) { char * port = strchr(ipaddress,':'); if (port==NULL) { return -1; @@ -316,11 +316,7 @@ _connect_to(struct harbor *h, const char *ipaddress, bool blocking) { skynet_error(h->ctx, "Harbor(%d) connect to %s:%d", h->id, tmp, portid); - if (blocking) { - return skynet_socket_block_connect(h->ctx, tmp, portid); - } else { - return skynet_socket_connect(h->ctx, tmp, portid); - } + return skynet_socket_connect(h->ctx, tmp, portid); } static inline void @@ -399,7 +395,7 @@ _update_remote_address(struct harbor *h, int harbor_id, const char * ipaddr) { skynet_free(h->remote_addr[harbor_id]); h->remote_addr[harbor_id] = NULL; } - h->remote_fd[harbor_id] = _connect_to(h, ipaddr, false); + h->remote_fd[harbor_id] = _connect_to(h, ipaddr); response_close(h, harbor_id); } @@ -448,7 +444,9 @@ _request_master(struct harbor *h, const char name[GLOBALNAME_LENGTH], size_t i, to_bigendian(buffer, handle); memcpy(buffer+4,name,i); - _send_package(h->ctx, h->master_fd, buffer, 4+i); + if (h->master_fd >= 0) { + _send_package(h->ctx, h->master_fd, buffer, 4+i); + } } /* @@ -534,7 +532,14 @@ harbor_id(struct harbor *h, int fd) { static void close_harbor(struct harbor *h, int fd) { + if (fd == h->master_fd) { + skynet_socket_close(h->ctx, fd); + skynet_error(h->ctx, "Master disconnected"); + h->master_fd = -1; + return; + } int id = harbor_id(h,fd); + if (id == 0) return; skynet_error(h->ctx, "Harbor %d closed",id); @@ -551,6 +556,7 @@ open_harbor(struct harbor *h, int fd) { assert(h->connected[id] == false); monitor_clear(h, id); h->connected[id] = true; + skynet_error(h->ctx, "Harbor(%d) connected", id); } static void @@ -715,7 +721,7 @@ harbor_init(struct harbor *h, struct skynet_context *ctx, const char * args) { sscanf(args,"%s %s %d",master_addr, local_addr, &harbor_id); h->master_addr = skynet_strdup(master_addr); h->id = harbor_id; - h->master_fd = _connect_to(h, master_addr, true); + h->master_fd = _connect_to(h, master_addr); if (h->master_fd == -1) { fprintf(stderr, "Harbor: Connect to master failed\n"); exit(1); diff --git a/service-src/service_master.c b/service-src/service_master.c index b5d8a560..a4fb8b1b 100644 --- a/service-src/service_master.c +++ b/service-src/service_master.c @@ -121,6 +121,7 @@ _connect_to(struct master *m, int id) { int port = strtol(portstr+1,NULL,10); skynet_error(ctx, "Master connect to harbor(%d) %s:%d", id, tmp, port); m->remote_fd[id] = skynet_socket_connect(ctx, tmp, port); + m->connected[id] = true; } static inline void @@ -217,7 +218,7 @@ socket_id(struct master *m, int id) { static void on_connected(struct master *m, int id) { _broadcast(m, m->remote_addr[id], strlen(m->remote_addr[id]), id); - m->connected[id] = true; +// m->connected[id] = true; int i; for (i=1;islot[id % MAX_SOCKET]; if (s->type == SOCKET_TYPE_INVALID) { if (__sync_bool_compare_and_swap(&s->type, SOCKET_TYPE_INVALID, SOCKET_TYPE_RESERVE)) { + s->id = id; + s->fd = -1; return id; } else { // retry @@ -287,7 +289,7 @@ new_fd(struct socket_server *ss, int id, int fd, uintptr_t opaque, bool add) { // return -1 when connecting static int -open_socket(struct socket_server *ss, struct request_open * request, struct socket_message *result, bool blocking) { +open_socket(struct socket_server *ss, struct request_open * request, struct socket_message *result) { int id = request->id; result->opaque = request->opaque; result->id = id; @@ -316,18 +318,14 @@ open_socket(struct socket_server *ss, struct request_open * request, struct sock continue; } socket_keepalive(sock); - if (!blocking) { - sp_nonblocking(sock); - } + sp_nonblocking(sock); status = connect( sock, ai_ptr->ai_addr, ai_ptr->ai_addrlen); if ( status != 0 && errno != EINPROGRESS) { close(sock); sock = -1; continue; } - if (blocking) { - sp_nonblocking(sock); - } + sp_nonblocking(sock); break; } @@ -512,7 +510,7 @@ send_socket(struct socket_server *ss, struct request_send * request, struct sock return -1; } assert(s->type != SOCKET_TYPE_PLISTEN && s->type != SOCKET_TYPE_LISTEN); - if (send_buffer_empty(s)) { + if (send_buffer_empty(s) && s->type == SOCKET_TYPE_CONNECTED) { int n = write(s->fd, request->buffer, request->sz); if (n<0) { switch(errno) { @@ -687,7 +685,7 @@ ctrl_cmd(struct socket_server *ss, struct socket_message *result) { case 'K': return close_socket(ss,(struct request_close *)buffer, result); case 'O': - return open_socket(ss, (struct request_open *)buffer, result, false); + return open_socket(ss, (struct request_open *)buffer, result); case 'X': result->opaque = 0; result->id = 0; @@ -765,7 +763,9 @@ report_connect(struct socket_server *ss, struct socket *s, struct socket_message result->opaque = s->opaque; result->id = s->id; result->ud = 0; - sp_write(ss->event_fd, s->fd, s, false); + if (send_buffer_empty(s)) { + sp_write(ss->event_fd, s->fd, s, false); + } union sockaddr_all u; socklen_t slen = sizeof(u); if (getpeername(s->fd, &u.s, &slen) == 0) { @@ -939,19 +939,6 @@ socket_server_connect(struct socket_server *ss, uintptr_t opaque, const char * a return request.u.open.id; } -int -socket_server_block_connect(struct socket_server *ss, uintptr_t opaque, const char * addr, int port) { - struct request_package request; - struct socket_message result; - open_request(ss, &request, opaque, addr, port); - int ret = open_socket(ss, &request.u.open, &result, true); - if (ret == SOCKET_OPEN) { - return result.id; - } else { - return -1; - } -} - // return -1 when error int64_t socket_server_send(struct socket_server *ss, int id, const void * buffer, int sz) { @@ -959,7 +946,6 @@ socket_server_send(struct socket_server *ss, int id, const void * buffer, int sz if (s->id != id || s->type == SOCKET_TYPE_INVALID) { return -1; } - assert(s->type != SOCKET_TYPE_RESERVE); struct request_package request; request.u.send.id = id; @@ -976,7 +962,6 @@ socket_server_send_lowpriority(struct socket_server *ss, int id, const void * bu if (s->id != id || s->type == SOCKET_TYPE_INVALID) { return; } - assert(s->type != SOCKET_TYPE_RESERVE); struct request_package request; request.u.send.id = id; diff --git a/skynet-src/socket_server.h b/skynet-src/socket_server.h index 2fbb6939..ea15c0e1 100644 --- a/skynet-src/socket_server.h +++ b/skynet-src/socket_server.h @@ -36,6 +36,4 @@ int socket_server_listen(struct socket_server *, uintptr_t opaque, const char * int socket_server_connect(struct socket_server *, uintptr_t opaque, const char * addr, int port); int socket_server_bind(struct socket_server *, uintptr_t opaque, int fd); -int socket_server_block_connect(struct socket_server *, uintptr_t opaque, const char * addr, int port); - #endif