mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-24 12:20:41 +00:00
change skynet_send api (use uint32_t instead of string)
This commit is contained in:
@@ -10,12 +10,13 @@
|
||||
|
||||
struct worker {
|
||||
int init;
|
||||
char * name;
|
||||
uint32_t address;
|
||||
};
|
||||
|
||||
struct broker {
|
||||
int init;
|
||||
int id;
|
||||
uint32_t launcher;
|
||||
char * name;
|
||||
struct worker w[DEFAULT_NUMBER];
|
||||
};
|
||||
@@ -29,41 +30,42 @@ broker_create(void) {
|
||||
|
||||
void
|
||||
broker_release(struct broker * b) {
|
||||
int i;
|
||||
for (i=0;i<DEFAULT_NUMBER;i++) {
|
||||
free(b->w[i].name);
|
||||
}
|
||||
free(b->name);
|
||||
free(b);
|
||||
}
|
||||
|
||||
static void
|
||||
_init(struct broker *b, int session, const char * msg, size_t sz) {
|
||||
_init(struct broker *b, int session, uint32_t address) {
|
||||
assert(session > 0 && session <= DEFAULT_NUMBER);
|
||||
assert(msg);
|
||||
int id = session - 1;
|
||||
assert(b->w[id].init == 0);
|
||||
b->w[id].name = malloc(sz+1);
|
||||
memcpy(b->w[id].name, msg, sz);
|
||||
b->w[id].name[sz] = '\0';
|
||||
b->w[id].address = address;
|
||||
b->w[id].init = 1;
|
||||
++b->init;
|
||||
}
|
||||
|
||||
static void
|
||||
_forward(struct broker *b, struct skynet_context * context) {
|
||||
skynet_forward(context, b->w[b->id].name);
|
||||
skynet_forward(context, b->w[b->id].address);
|
||||
b->id = (b->id + 1) % DEFAULT_NUMBER;
|
||||
}
|
||||
|
||||
static int
|
||||
_cb(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) {
|
||||
_cb(struct skynet_context * context, void * ud, int session, uint32_t source, const void * msg, size_t sz) {
|
||||
struct broker * b = ud;
|
||||
if (b->init < DEFAULT_NUMBER) {
|
||||
_init(b, session, msg, sz);
|
||||
if (source != b->launcher)
|
||||
return 0;
|
||||
assert(sz == 9);
|
||||
char addr[10];
|
||||
memcpy(addr, msg, 9);
|
||||
addr[9] = '\0';
|
||||
uint32_t address = strtoul(addr+1, NULL, 16);
|
||||
assert(address != 0);
|
||||
_init(b, session, address);
|
||||
if (b->init == DEFAULT_NUMBER) {
|
||||
skynet_command(context, "REG", b->name);
|
||||
skynet_send(context, NULL, LAUNCHER, 0, NULL, 0, 0);
|
||||
skynet_send(context, 0, b->launcher, 0, NULL, 0, 0);
|
||||
}
|
||||
} else {
|
||||
_forward(b, context);
|
||||
@@ -75,6 +77,12 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr,
|
||||
|
||||
int
|
||||
broker_init(struct broker *b, struct skynet_context *ctx, const char * args) {
|
||||
b->launcher = skynet_queryname(ctx, LAUNCHER);
|
||||
if (b->launcher == 0) {
|
||||
skynet_error(ctx, "Can't query %s", LAUNCHER);
|
||||
return 1;
|
||||
}
|
||||
|
||||
char * service = strchr(args,' ');
|
||||
if (service == NULL) {
|
||||
return 1;
|
||||
@@ -92,7 +100,7 @@ broker_init(struct broker *b, struct skynet_context *ctx, const char * args) {
|
||||
if (len == 0)
|
||||
return 1;
|
||||
for (i=0;i<DEFAULT_NUMBER;i++) {
|
||||
int id = skynet_send(ctx, NULL, LAUNCHER , -1, service , len, 0);
|
||||
int id = skynet_send(ctx, 0, b->launcher , -1, service , len, 0);
|
||||
assert(id > 0 && id <= DEFAULT_NUMBER);
|
||||
}
|
||||
|
||||
|
||||
@@ -9,7 +9,7 @@
|
||||
#include <string.h>
|
||||
|
||||
static int
|
||||
_cb(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) {
|
||||
_cb(struct skynet_context * context, void * ud, int session, uint32_t source, const void * msg, size_t sz) {
|
||||
assert(sz <= 65535);
|
||||
int fd = (int)(intptr_t)ud;
|
||||
|
||||
@@ -31,7 +31,7 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr,
|
||||
}
|
||||
}
|
||||
if (err < 0) {
|
||||
skynet_error(context, "Client socket error : Drop message from %s session = %d", addr, session);
|
||||
skynet_error(context, "Client socket error : Drop message from %x session = %d", source, session);
|
||||
return 0;
|
||||
}
|
||||
assert(err == sz +2);
|
||||
|
||||
@@ -291,17 +291,22 @@ _send_package(int fd, const void * buffer, size_t sz) {
|
||||
|
||||
static int
|
||||
_send_remote(int fd, const char * buffer, size_t sz, struct remote_message_header * cookie) {
|
||||
struct iovec part[2];
|
||||
part[0].iov_base = (char *)buffer;
|
||||
part[0].iov_len = sz;
|
||||
uint16_t sz_header = htons(sz+sizeof(*cookie));
|
||||
struct iovec part[3];
|
||||
|
||||
part[0].iov_base = &sz_header;
|
||||
part[0].iov_len = 2;
|
||||
|
||||
part[1].iov_base = (char *)buffer;
|
||||
part[1].iov_len = sz;
|
||||
|
||||
uint32_t header[3];
|
||||
_header_to_message(cookie, header);
|
||||
|
||||
part[1].iov_base = header;
|
||||
part[1].iov_len = sizeof(header);
|
||||
part[2].iov_base = header;
|
||||
part[2].iov_len = sizeof(header);
|
||||
for (;;) {
|
||||
int err = writev(fd, part, 2);
|
||||
int err = writev(fd, part, 3);
|
||||
if (err < 0) {
|
||||
switch (errno) {
|
||||
case EAGAIN:
|
||||
@@ -309,7 +314,7 @@ _send_remote(int fd, const char * buffer, size_t sz, struct remote_message_heade
|
||||
continue;
|
||||
}
|
||||
}
|
||||
if (err != sz+sizeof(*cookie)) {
|
||||
if (err != sz+sizeof(*cookie)+2) {
|
||||
return 1;
|
||||
}
|
||||
return 0;
|
||||
@@ -405,28 +410,13 @@ _request_master(struct harbor *h, struct skynet_context * context, const char na
|
||||
n bytes string (name)
|
||||
*/
|
||||
|
||||
static void
|
||||
_id_to_hex(char *str, uint32_t id) {
|
||||
int i;
|
||||
static char hex[16] = { '0','1','2','3','4','5','6','7','8','9','A','B','C','D','E','F' };
|
||||
str[0] = ':';
|
||||
for (i=0;i<8;i++) {
|
||||
str[i+1] = hex[(id >> ((7-i) * 4))&0xf];
|
||||
}
|
||||
str[9] = '\0';
|
||||
}
|
||||
|
||||
static int
|
||||
_remote_send_handle(struct harbor *h, struct skynet_context * context, uint32_t source, uint32_t destination, int session, const char * msg, size_t sz) {
|
||||
int harbor_id = destination >> HANDLE_REMOTE_SHIFT;
|
||||
assert(harbor_id != 0);
|
||||
if (harbor_id == h->id) {
|
||||
// local message
|
||||
char srcstr[10];
|
||||
char desstr[10];
|
||||
_id_to_hex(srcstr, source);
|
||||
_id_to_hex(desstr, destination);
|
||||
skynet_send(context, srcstr, desstr , session, (void *)msg, sz, DONTCOPY);
|
||||
skynet_send(context, source, destination , session, (void *)msg, sz, DONTCOPY);
|
||||
return 1;
|
||||
}
|
||||
|
||||
@@ -491,7 +481,7 @@ _report_local_address(struct harbor *h, struct skynet_context * context, const c
|
||||
}
|
||||
|
||||
static int
|
||||
_mainloop(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) {
|
||||
_mainloop(struct skynet_context * context, void * ud, int session, uint32_t source, const void * msg, size_t sz) {
|
||||
struct harbor * h = ud;
|
||||
if (session == SESSION_CLIENT) {
|
||||
const char * cookie = msg;
|
||||
@@ -517,11 +507,7 @@ _mainloop(struct skynet_context * context, void * ud, int session, const char *
|
||||
_update_remote_name(h, context, msg, header.destination);
|
||||
}
|
||||
} else {
|
||||
char srcstr[10];
|
||||
char desstr[10];
|
||||
_id_to_hex(srcstr, header.source);
|
||||
_id_to_hex(desstr, header.destination);
|
||||
skynet_send(context, srcstr, desstr, (int)header.session, (void *)msg, sz-12, DONTCOPY);
|
||||
skynet_send(context, header.source, header.destination, (int)header.session, (void *)msg, sz-12, DONTCOPY);
|
||||
return 1;
|
||||
}
|
||||
} else {
|
||||
@@ -531,13 +517,12 @@ _mainloop(struct skynet_context * context, void * ud, int session, const char *
|
||||
return 0;
|
||||
}
|
||||
assert(sz == sizeof(*rmsg));
|
||||
uint32_t source_handle = strtoul(addr+1, NULL, 16);
|
||||
if (rmsg->destination.handle == 0) {
|
||||
if (_remote_send_name(h, context, source_handle , rmsg->destination.name, session, rmsg->message, rmsg->sz)) {
|
||||
if (_remote_send_name(h, context, source , rmsg->destination.name, session, rmsg->message, rmsg->sz)) {
|
||||
return 0;
|
||||
}
|
||||
} else {
|
||||
if (_remote_send_handle(h, context, source_handle , rmsg->destination.handle, session, rmsg->message, rmsg->sz)) {
|
||||
if (_remote_send_handle(h, context, source , rmsg->destination.handle, session, rmsg->message, rmsg->sz)) {
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
@@ -563,7 +548,6 @@ harbor_init(struct harbor *h, struct skynet_context *ctx, const char * args) {
|
||||
h->master_addr = strdup(master_addr);
|
||||
h->master_fd = master_fd;
|
||||
|
||||
const char * self_addr = skynet_command(ctx, "REG", NULL);
|
||||
char tmp[128];
|
||||
sprintf(tmp,"gate ! %s %d 0",local_addr,REMOTE_MAX);
|
||||
const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp);
|
||||
@@ -571,9 +555,15 @@ harbor_init(struct harbor *h, struct skynet_context *ctx, const char * args) {
|
||||
skynet_error(ctx, "Harbor : launch gate failed");
|
||||
return 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;
|
||||
}
|
||||
const char * self_addr = skynet_command(ctx, "REG", NULL);
|
||||
int n = sprintf(tmp,"broker %s",self_addr);
|
||||
skynet_send(ctx, NULL, gate_addr, 0, tmp, n, 0);
|
||||
skynet_send(ctx, NULL, gate_addr, 0, "start", 5, 0);
|
||||
skynet_send(ctx, 0, gate, 0, tmp, n, 0);
|
||||
skynet_send(ctx, 0, gate, 0, "start", 5, 0);
|
||||
|
||||
h->id = harbor_id;
|
||||
skynet_callback(ctx, h, _mainloop);
|
||||
|
||||
@@ -231,7 +231,7 @@ _update_address(struct skynet_context * context, struct master *m, int harbor_id
|
||||
*/
|
||||
|
||||
static int
|
||||
_mainloop(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) {
|
||||
_mainloop(struct skynet_context * context, void * ud, int session, uint32_t source, const void * msg, size_t sz) {
|
||||
struct master *m = ud;
|
||||
assert(session == SESSION_CLIENT);
|
||||
uint32_t handle = 0;
|
||||
@@ -256,15 +256,20 @@ int
|
||||
master_init(struct master *m, struct skynet_context *ctx, const char * args) {
|
||||
char tmp[strlen(args) + 32];
|
||||
sprintf(tmp,"gate ! %s %d 0",args,REMOTE_MAX);
|
||||
const char * self_addr = skynet_command(ctx, "REG", NULL);
|
||||
const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp);
|
||||
if (gate_addr == NULL) {
|
||||
skynet_error(ctx, "Master : launch gate failed");
|
||||
return 1;
|
||||
}
|
||||
uint32_t gate = strtoul(gate_addr+1, NULL, 16);
|
||||
if (gate == 0) {
|
||||
skynet_error(ctx, "Master : launch gate invalid %s", gate_addr);
|
||||
return 1;
|
||||
}
|
||||
const char * self_addr = skynet_command(ctx, "REG", NULL);
|
||||
int n = sprintf(tmp,"broker %s",self_addr);
|
||||
skynet_send(ctx, NULL, gate_addr, 0, tmp, n, 0);
|
||||
skynet_send(ctx, NULL, gate_addr, 0, "start", 5, 0);
|
||||
skynet_send(ctx, 0, gate, 0, tmp, n, 0);
|
||||
skynet_send(ctx, 0, gate, 0, "start", 5, 0);
|
||||
|
||||
skynet_callback(ctx, m, _mainloop);
|
||||
|
||||
|
||||
@@ -16,9 +16,8 @@ multicast_release(struct skynet_multicast_group *g) {
|
||||
}
|
||||
|
||||
static int
|
||||
_maincb(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) {
|
||||
_maincb(struct skynet_context * context, void * ud, int session, uint32_t source, const void * msg, size_t sz) {
|
||||
struct skynet_multicast_group *g = ud;
|
||||
uint32_t source = strtoul(addr+1, NULL, 16);
|
||||
if (source == 0) {
|
||||
char cmd = '\0';
|
||||
uint32_t handle = 0;
|
||||
|
||||
Reference in New Issue
Block a user