From 83db25fe926b9fb86b49a5a5ecd12599b176a4f3 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E4=BA=91=E9=A3=8E?= Date: Wed, 21 Aug 2013 16:56:16 +0800 Subject: [PATCH] service_harbor use skynet_socket now --- service-src/service_harbor.c | 351 ++++++++++++++++++----------------- 1 file changed, 182 insertions(+), 169 deletions(-) diff --git a/service-src/service_harbor.c b/service-src/service_harbor.c index 3ff9c0e7..fd2537c6 100644 --- a/service-src/service_harbor.c +++ b/service-src/service_harbor.c @@ -1,24 +1,19 @@ #include "skynet.h" #include "skynet_harbor.h" +#include "skynet_socket.h" -#include -#include -#include -#include -#include -#include #include #include +#include #include #include #include -#include #define HASH_SIZE 4096 #define DEFAULT_QUEUE_SIZE 1024 struct msg { - char * buffer; + uint8_t * buffer; size_t size; }; @@ -52,11 +47,14 @@ struct remote_message_header { }; struct harbor { + struct skynet_context *ctx; + char * local_addr; int id; struct hashmap * map; int master_fd; char * master_addr; int remote_fd[REMOTE_MAX]; + bool connected[REMOTE_MAX]; char * remote_addr[REMOTE_MAX]; }; @@ -200,12 +198,14 @@ _hash_delete(struct hashmap *hash) { struct harbor * harbor_create(void) { struct harbor * h = malloc(sizeof(*h)); + h->ctx = NULL; h->id = 0; h->master_fd = -1; h->master_addr = NULL; int i; for (i=0;iremote_fd[i] = -1; + h->connected[i] = false; h->remote_addr[i] = NULL; } h->map = _hash_new(); @@ -214,14 +214,16 @@ harbor_create(void) { void harbor_release(struct harbor *h) { + struct skynet_context *ctx = h->ctx; if (h->master_fd >= 0) { - close(h->master_fd); + skynet_socket_close(ctx, h->master_fd); } free(h->master_addr); + free(h->local_addr); int i; for (i=0;iremote_fd[i] >= 0) { - close(h->remote_fd[i]); + skynet_socket_close(ctx, h->remote_fd[i]); free(h->remote_addr[i]); } } @@ -230,9 +232,7 @@ harbor_release(struct harbor *h) { } static int -_connect_to(struct skynet_context *ctx, const char *ipaddress) { - int fd = socket(AF_INET,SOCK_STREAM,0); - struct sockaddr_in my_addr; +_connect_to(struct harbor *h, const char *ipaddress) { char * port = strchr(ipaddress,':'); if (port==NULL) { return -1; @@ -242,114 +242,88 @@ _connect_to(struct skynet_context *ctx, const char *ipaddress) { memcpy(tmp,ipaddress,sz); tmp[sz] = '\0'; - my_addr.sin_addr.s_addr=inet_addr(tmp); - my_addr.sin_family=AF_INET; - my_addr.sin_port=htons(strtol(port+1,NULL,10)); + int portid = strtol(port+1, NULL,10); - int r = connect(fd,(struct sockaddr *)&my_addr,sizeof(struct sockaddr_in)); - - if (r == -1) { - close(fd); - skynet_error(ctx, "Connect to %s error", ipaddress); - return -1; - } - - return fd; + return skynet_socket_connect(h->ctx, tmp, portid); } static inline void -_header_to_message(const struct remote_message_header * header, uint32_t * message) { - message[0] = htonl(header->source); - message[1] = htonl(header->destination); - message[2] = htonl(header->session); +to_bigendian(uint8_t *buffer, uint32_t n) { + buffer[0] = (n >> 24) & 0xff; + buffer[1] = (n >> 16) & 0xff; + buffer[2] = (n >> 8) & 0xff; + buffer[3] = n & 0xff; +} + +static inline void +_header_to_message(const struct remote_message_header * header, uint8_t * message) { + to_bigendian(message , header->source); + to_bigendian(message+4 , header->destination); + to_bigendian(message+8 , header->session); +} + +static inline uint32_t +from_bigendian(uint32_t n) { + union { + uint32_t big; + uint8_t bytes[4]; + } u; + u.big = n; + return u.bytes[0] << 24 | u.bytes[1] << 16 | u.bytes[2] << 8 | u.bytes[3]; } static inline void _message_to_header(const uint32_t *message, struct remote_message_header *header) { - header->source = ntohl(message[0]); - header->destination = ntohl(message[1]); - header->session = ntohl(message[2]); + header->source = from_bigendian(message[0]); + header->destination = from_bigendian(message[1]); + header->session = from_bigendian(message[2]); } -static int -_send_package(int fd, const void * buffer, size_t sz) { - uint32_t header = htonl(sz); - struct iovec part[2]; - part[0].iov_base = &header; - part[0].iov_len = 4; - part[1].iov_base = (void*)buffer; - part[1].iov_len = sz; +static void +_send_package(struct skynet_context *ctx, int fd, const void * buffer, size_t sz) { + uint8_t * sendbuf = malloc(sz+4); + to_bigendian(sendbuf, sz); + memcpy(sendbuf+4, buffer, sz); - for (;;) { - int err = writev(fd, part, 2); - if (err < 0) { - switch (errno) { - case EAGAIN: - case EINTR: - continue; - } - } - if (err != sz+4) { - return 1; - } - return 0; - } -} - -static int -_send_remote(int fd, const char * buffer, size_t sz, struct remote_message_header * cookie) { - uint32_t sz_header = htonl(sz+sizeof(*cookie)); - struct iovec part[3]; - - part[0].iov_base = &sz_header; - part[0].iov_len = 4; - - part[1].iov_base = (char *)buffer; - part[1].iov_len = sz; - - uint32_t header[3]; - _header_to_message(cookie, header); - - part[2].iov_base = header; - part[2].iov_len = sizeof(header); - for (;;) { - int err = writev(fd, part, 3); - if (err < 0) { - switch (errno) { - case EAGAIN: - case EINTR: - continue; - } - } - if (err != sz+sizeof(*cookie)+4) { - return 1; - } - return 0; + if (skynet_socket_send(ctx, fd, sendbuf, sz+4)) { + skynet_error(ctx, "Send to %d error", fd); } } static void -_update_remote_address(struct skynet_context * context, struct harbor *h, int harbor_id, const char * ipaddr) { +_send_remote(struct skynet_context * ctx, int fd, const char * buffer, size_t sz, struct remote_message_header * cookie) { + uint32_t sz_header = sz+sizeof(*cookie); + uint8_t * sendbuf = malloc(sz_header+4); + to_bigendian(sendbuf, sz_header); + memcpy(sendbuf+4, buffer, sz); + _header_to_message(cookie, sendbuf+4+sz); + + if (skynet_socket_send(ctx, fd, sendbuf, sz_header+4)) { + skynet_error(ctx, "Remote send to %d error", fd); + } +} + +static void +_update_remote_address(struct harbor *h, int harbor_id, const char * ipaddr) { if (harbor_id == h->id) { return; } assert(harbor_id > 0 && harbor_id< REMOTE_MAX); + struct skynet_context * context = h->ctx; if (h->remote_fd[harbor_id] >=0) { - close(h->remote_fd[harbor_id]); + skynet_socket_close(context, h->remote_fd[harbor_id]); free(h->remote_addr[harbor_id]); h->remote_addr[harbor_id] = NULL; } - h->remote_fd[harbor_id] = _connect_to(context, ipaddr); - if (h->remote_fd[harbor_id] >= 0) { - free(h->remote_addr[harbor_id]); - h->remote_addr[harbor_id] = strdup(ipaddr); - } + h->remote_fd[harbor_id] = _connect_to(h, ipaddr); + h->connected[harbor_id] = false; } static void -_dispatch_queue(struct harbor *h, struct skynet_context * context, struct msg_queue * queue, uint32_t handle, const char name[GLOBALNAME_LENGTH] ) { +_dispatch_queue(struct harbor *h, struct msg_queue * queue, uint32_t handle, const char name[GLOBALNAME_LENGTH] ) { int harbor_id = handle >> HANDLE_REMOTE_SHIFT; assert(harbor_id != 0); + struct skynet_context * context = h->ctx; int fd = h->remote_fd[harbor_id]; if (fd < 0) { char tmp [GLOBALNAME_LENGTH+1]; @@ -360,53 +334,37 @@ _dispatch_queue(struct harbor *h, struct skynet_context * context, struct msg_qu } struct msg * m = _pop_queue(queue); while (m) { - struct remote_message_header * cookie = (struct remote_message_header *)(m->buffer + m->size - sizeof(*cookie)); - cookie->destination |= (handle & HANDLE_MASK); - _header_to_message(cookie, (uint32_t *)cookie); - int err = _send_package(fd, m->buffer, m->size); - if (err) { - close(fd); - h->remote_fd[harbor_id] = _connect_to(context, h->remote_addr[harbor_id]); - if (h->remote_fd[harbor_id] < 0) { - skynet_error(context, "Reconnect to harbor %d %s failed",harbor_id, h->remote_addr[harbor_id]); - return; - } - } + struct remote_message_header cookie; + uint8_t *ptr = m->buffer + m->size - sizeof(cookie); + memcpy(&cookie, ptr, sizeof(cookie)); + cookie.destination |= (handle & HANDLE_MASK); + _header_to_message(&cookie, ptr); + _send_package(context, fd, m->buffer, m->size); m = _pop_queue(queue); } } static void -_update_remote_name(struct harbor *h, struct skynet_context * context, const char name[GLOBALNAME_LENGTH], uint32_t handle) { +_update_remote_name(struct harbor *h, const char name[GLOBALNAME_LENGTH], uint32_t handle) { struct keyvalue * node = _hash_search(h->map, name); if (node == NULL) { node = _hash_insert(h->map, name); } node->value = handle; if (node->queue) { - _dispatch_queue(h, context, node->queue, handle, name); + _dispatch_queue(h, node->queue, handle, name); _release_queue(node->queue); node->queue = NULL; } } static void -_request_master(struct harbor *h, struct skynet_context * context, const char name[GLOBALNAME_LENGTH], size_t i, uint32_t handle) { - char buffer[4+i]; - handle = htonl(handle); - memcpy(buffer, &handle, 4); +_request_master(struct harbor *h, const char name[GLOBALNAME_LENGTH], size_t i, uint32_t handle) { + uint8_t buffer[4+i]; + to_bigendian(buffer, handle); memcpy(buffer+4,name,i); - int err = _send_package(h->master_fd, buffer, 4+i); - if (err) { - close(h->master_fd); - h->master_fd = _connect_to(context, h->master_addr); - if (h->master_fd < 0) { - skynet_error(context, "Reconnect to master server %s failed", h->master_addr); - return; - } - _send_package(h->master_fd, buffer, 4+i); - } + _send_package(h->ctx, h->master_fd, buffer, 4+i); } /* @@ -418,9 +376,10 @@ _request_master(struct harbor *h, struct skynet_context * context, const char na */ static int -_remote_send_handle(struct harbor *h, struct skynet_context * context, uint32_t source, uint32_t destination, int type, int session, const char * msg, size_t sz) { +_remote_send_handle(struct harbor *h, uint32_t source, uint32_t destination, int type, int session, const char * msg, size_t sz) { int harbor_id = destination >> HANDLE_REMOTE_SHIFT; assert(harbor_id != 0); + struct skynet_context * context = h->ctx; if (harbor_id == h->id) { // local message skynet_send(context, source, destination , type | PTYPE_TAG_DONTCOPY, session, (void *)msg, sz); @@ -428,42 +387,34 @@ _remote_send_handle(struct harbor *h, struct skynet_context * context, uint32_t } int fd = h->remote_fd[harbor_id]; - if (fd >= 0) { + if (fd >= 0 && h->connected[harbor_id]) { struct remote_message_header cookie; cookie.source = source; cookie.destination = (destination & HANDLE_MASK) | ((uint32_t)type << HANDLE_REMOTE_SHIFT); cookie.session = (uint32_t)session; - int err = _send_remote(fd, msg,sz,&cookie); - if (err) { - close(fd); - h->remote_fd[harbor_id] = _connect_to(context, h->remote_addr[harbor_id]); - if (h->remote_fd[harbor_id] < 0) { - skynet_error(context, "Reconnect to harbor %d : %s failed", harbor_id, h->remote_addr[harbor_id]); - return 0; - } - } + _send_remote(context, fd, msg,sz,&cookie); } else { - _request_master(h, context, NULL, 0, harbor_id); + _request_master(h, NULL, 0, harbor_id); skynet_error(context, "Drop message to harbor %d from %x to %x (session = %d, msgsz = %d)",harbor_id, source, destination,session,(int)sz); } return 0; } static void -_remote_register_name(struct harbor *h, struct skynet_context * context, const char name[GLOBALNAME_LENGTH], uint32_t handle) { +_remote_register_name(struct harbor *h, const char name[GLOBALNAME_LENGTH], uint32_t handle) { int i; for (i=0;imap, name); if (node == NULL) { node = _hash_insert(h->map, name); @@ -478,23 +429,67 @@ _remote_send_name(struct harbor *h, struct skynet_context * context, uint32_t so header.session = (uint32_t)session; _push_queue(node->queue, msg, sz, &header); // 0 for request - _remote_register_name(h, context, name, 0); + _remote_register_name(h, name, 0); return 1; } else { - return _remote_send_handle(h, context, source, node->value, type, session, msg, sz); + return _remote_send_handle(h, source, node->value, type, session, msg, sz); } } +static int +harbor_id(struct harbor *h, int fd) { + int i; + for (i=1;iremote_fd[i] == fd) + return i; + } + return 0; +} + static void -_report_local_address(struct harbor *h, struct skynet_context * context, const char * local_address, int harbor_id) { - size_t sz = strlen(local_address); - _request_master(h, context, local_address, sz, harbor_id); +close_harbor(struct harbor *h, int fd) { + int id = harbor_id(h,fd); + if (id == 0) + return; + skynet_error(h->ctx, "Harbor %d closed",id); + skynet_socket_close(h->ctx, fd); + h->remote_fd[id] = -1; + h->connected[id] = false; +} + +static void +open_harbor(struct harbor *h, int fd) { + int id = harbor_id(h,fd); + if (id == 0) + return; + assert(h->connected[id] == false); + h->connected[id] = true; } static int _mainloop(struct skynet_context * context, void * ud, int type, int session, uint32_t source, const void * msg, size_t sz) { struct harbor * h = ud; switch (type) { + case PTYPE_SOCKET: { + const struct skynet_socket_message * message = msg; + switch(message->type) { + case SKYNET_SOCKET_TYPE_DATA: + free(message->buffer); + skynet_error(context, "recv invalid socket message (size=%d)", message->ud); + break; + case SKYNET_SOCKET_TYPE_ACCEPT: + skynet_error(context, "recv invalid socket accept message"); + break; + case SKYNET_SOCKET_TYPE_ERROR: + case SKYNET_SOCKET_TYPE_CLOSE: + close_harbor(h, message->id); + break; + case SKYNET_SOCKET_TYPE_CONNECT: + open_harbor(h, message->id); + break; + } + return 0; + } case PTYPE_HARBOR: { // remote message in const char * cookie = msg; @@ -508,7 +503,7 @@ _mainloop(struct skynet_context * context, void * ud, int type, int session, uin char ip [sz - 11]; memcpy(ip, msg, sz-12); ip[sz-11] = '\0'; - _update_remote_address(context, h, header.destination, ip); + _update_remote_address(h, header.destination, ip); } else { // update global name if (sz - 12 > GLOBALNAME_LENGTH) { @@ -517,7 +512,7 @@ _mainloop(struct skynet_context * context, void * ud, int type, int session, uin name[sz-11] = '\0'; skynet_error(context, "Global name is too long %s", name); } - _update_remote_name(h, context, msg, header.destination); + _update_remote_name(h, msg, header.destination); } } else { uint32_t destination = header.destination; @@ -532,18 +527,18 @@ _mainloop(struct skynet_context * context, void * ud, int type, int session, uin // register name message const struct remote_message *rmsg = msg; assert (sz == sizeof(rmsg->destination)); - _remote_register_name(h, context, rmsg->destination.name, rmsg->destination.handle); + _remote_register_name(h, rmsg->destination.name, rmsg->destination.handle); return 0; } default: { // remote message out const struct remote_message *rmsg = msg; if (rmsg->destination.handle == 0) { - if (_remote_send_name(h, context, source , rmsg->destination.name, type, session, rmsg->message, rmsg->sz)) { + if (_remote_send_name(h, source , rmsg->destination.name, type, session, rmsg->message, rmsg->sz)) { return 0; } } else { - if (_remote_send_handle(h, context, source , rmsg->destination.handle, type, session, rmsg->message, rmsg->sz)) { + if (_remote_send_handle(h, source , rmsg->destination.handle, type, session, rmsg->message, rmsg->sz)) { return 0; } } @@ -553,44 +548,62 @@ _mainloop(struct skynet_context * context, void * ud, int type, int session, uin } } -int -harbor_init(struct harbor *h, struct skynet_context *ctx, const char * args) { - int sz = strlen(args)+1; - char master_addr[sz]; - char local_addr[sz]; - int harbor_id = 0; - sscanf(args,"%s %s %d",master_addr, local_addr, &harbor_id); - int master_fd = _connect_to(ctx, master_addr); - if (master_fd < 0) { - skynet_error(ctx, "Harbor : Connect to master %s faild",master_addr); - return 1; +static int +_connect_master(struct skynet_context * ctx, void * ud, int type, int session, uint32_t source, const void * msg, size_t sz) { + struct harbor *h = ud; + assert(type == PTYPE_SOCKET); + const struct skynet_socket_message * message = msg; + if (message->id != h->master_fd) { + if (message->type == SKYNET_SOCKET_TYPE_DATA) { + free(message->buffer); + } + skynet_error(ctx, "Invalid socket message incoming before connecting to master"); + return 0; } - printf("Connect to master %s\n",master_addr); - - h->master_addr = strdup(master_addr); - h->master_fd = master_fd; + if (message->type == SKYNET_SOCKET_TYPE_ERROR) { + fprintf(stderr, "Harbor: Conenct to master failed\n"); + exit(1); + } + assert(message->type == SKYNET_SOCKET_TYPE_CONNECT); char tmp[128]; - sprintf(tmp,"gate L ! %s %d %d 0",local_addr, PTYPE_HARBOR, REMOTE_MAX); + sprintf(tmp,"gate L ! %s %d %d 0",h->local_addr, PTYPE_HARBOR, REMOTE_MAX); const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp); if (gate_addr == NULL) { - skynet_error(ctx, "Harbor : launch gate failed"); - return 1; + fprintf(stderr, "Harbor : launch gate failed\n"); + exit(1); } uint32_t gate = strtoul(gate_addr+1 , NULL, 16); if (gate == 0) { - skynet_error(ctx, "Harbor : launch gate invalid %s", gate_addr); - return 1; + fprintf(stderr, "Harbor : launch gate invalid %s", gate_addr); + exit(1); } const char * self_addr = skynet_command(ctx, "REG", NULL); int n = sprintf(tmp,"broker %s",self_addr); skynet_send(ctx, 0, gate, PTYPE_TEXT, 0, tmp, n); skynet_send(ctx, 0, gate, PTYPE_TEXT, 0, "start", 5); - h->id = harbor_id; skynet_callback(ctx, h, _mainloop); - _report_local_address(h, ctx, local_addr, harbor_id); + _request_master(h, h->local_addr, strlen(h->local_addr), h->id); + + return 0; +} + +int +harbor_init(struct harbor *h, struct skynet_context *ctx, const char * args) { + h->ctx = ctx; + int sz = strlen(args)+1; + char master_addr[sz]; + char local_addr[sz]; + int harbor_id = 0; + sscanf(args,"%s %s %d",master_addr, local_addr, &harbor_id); + h->master_addr = strdup(master_addr); + h->master_fd = _connect_to(h, master_addr); + h->local_addr = strdup(local_addr); + h->id = harbor_id; + + skynet_callback(ctx, h, _connect_master); return 0; }