mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-25 12:43:09 +00:00
gate support user defined header size
This commit is contained in:
@@ -75,25 +75,38 @@ _write(lua_State *L) {
|
|||||||
static int
|
static int
|
||||||
_writeblock(lua_State *L) {
|
_writeblock(lua_State *L) {
|
||||||
int fd = luaL_checkinteger(L,1);
|
int fd = luaL_checkinteger(L,1);
|
||||||
int type = lua_type(L,2);
|
int header = luaL_checkinteger(L,2);
|
||||||
|
int type = lua_type(L,3);
|
||||||
const char * buffer = NULL;
|
const char * buffer = NULL;
|
||||||
size_t sz;
|
size_t sz;
|
||||||
if (type == LUA_TSTRING) {
|
if (type == LUA_TSTRING) {
|
||||||
buffer = lua_tolstring(L,2,&sz);
|
buffer = lua_tolstring(L,3,&sz);
|
||||||
} else if (type == LUA_TLIGHTUSERDATA) {
|
} else if (type == LUA_TLIGHTUSERDATA) {
|
||||||
buffer = lua_touserdata(L,2);
|
buffer = lua_touserdata(L,3);
|
||||||
sz = luaL_checkinteger(L,3);
|
sz = luaL_checkinteger(L,4);
|
||||||
}
|
}
|
||||||
|
|
||||||
if (sz > 65535) {
|
if (header == 2) {
|
||||||
luaL_error(L, "Too big package %d", (int)sz);
|
if (sz > 65535) {
|
||||||
|
luaL_error(L, "Too big package %d", (int)sz);
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
if (header != 4) {
|
||||||
|
luaL_error(L, "block header must be 2 or 4 bytes");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
struct iovec buf[2];
|
struct iovec buf[2];
|
||||||
// send big-endian header
|
if (header == 2) {
|
||||||
uint8_t head[2] = { sz >> 8 & 0xff , sz & 0xff };
|
// send big-endian header
|
||||||
buf[0].iov_base = head;
|
uint8_t head[2] = { sz >> 8 & 0xff , sz & 0xff };
|
||||||
buf[0].iov_len = 2;
|
buf[0].iov_base = head;
|
||||||
|
buf[0].iov_len = 2;
|
||||||
|
} else {
|
||||||
|
uint8_t head[4] = { sz >> 24 & 0xff, sz >> 16 & 0xff, sz >> 8 & 0xff , sz & 0xff };
|
||||||
|
buf[0].iov_base = head;
|
||||||
|
buf[0].iov_len = 4;
|
||||||
|
}
|
||||||
buf[1].iov_base = (void *)buffer;
|
buf[1].iov_base = (void *)buffer;
|
||||||
buf[1].iov_len = sz;
|
buf[1].iov_len = sz;
|
||||||
|
|
||||||
@@ -107,7 +120,7 @@ _writeblock(lua_State *L) {
|
|||||||
}
|
}
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
assert(err == sz +2);
|
assert(err == sz + header);
|
||||||
return 0;
|
return 0;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
20
gate/main.c
20
gate/main.c
@@ -25,6 +25,7 @@ struct gate {
|
|||||||
int cap;
|
int cap;
|
||||||
int max_connection;
|
int max_connection;
|
||||||
int client_tag;
|
int client_tag;
|
||||||
|
int header_size;
|
||||||
struct connection ** agent;
|
struct connection ** agent;
|
||||||
struct connection * map;
|
struct connection * map;
|
||||||
};
|
};
|
||||||
@@ -193,7 +194,7 @@ _cb(struct skynet_context * ctx, void * ud, int type, int session, uint32_t sour
|
|||||||
getpeername(fd, (struct sockaddr *)&remote_addr, &len);
|
getpeername(fd, (struct sockaddr *)&remote_addr, &len);
|
||||||
_report(g, ctx, "%d open %d %s:%u",id,fd,inet_ntoa(remote_addr.sin_addr),ntohs(remote_addr.sin_port));
|
_report(g, ctx, "%d open %d %s:%u",id,fd,inet_ntoa(remote_addr.sin_addr),ntohs(remote_addr.sin_port));
|
||||||
}
|
}
|
||||||
uint8_t * plen = mread_pull(m,2);
|
uint8_t * plen = mread_pull(m,g->header_size);
|
||||||
if (plen == NULL) {
|
if (plen == NULL) {
|
||||||
if (mread_closed(m)) {
|
if (mread_closed(m)) {
|
||||||
_remove_id(g,id);
|
_remove_id(g,id);
|
||||||
@@ -202,7 +203,12 @@ _cb(struct skynet_context * ctx, void * ud, int type, int session, uint32_t sour
|
|||||||
goto _break;
|
goto _break;
|
||||||
}
|
}
|
||||||
// big-endian
|
// big-endian
|
||||||
uint16_t len = plen[0] << 8 | plen[1];
|
uint16_t len ;
|
||||||
|
if (g->header_size == 2) {
|
||||||
|
len = plen[0] << 8 | plen[1];
|
||||||
|
} else {
|
||||||
|
len = plen[0] << 24 | plen[1] << 16 | plen[2] << 8 | plen[3];
|
||||||
|
}
|
||||||
|
|
||||||
void * data = mread_pull(m, len);
|
void * data = mread_pull(m, len);
|
||||||
if (data == NULL) {
|
if (data == NULL) {
|
||||||
@@ -230,11 +236,16 @@ gate_init(struct gate *g , struct skynet_context * ctx, char * parm) {
|
|||||||
char watchdog[sz];
|
char watchdog[sz];
|
||||||
char binding[sz];
|
char binding[sz];
|
||||||
int client_tag = 0;
|
int client_tag = 0;
|
||||||
int n = sscanf(parm, "%s %s %d %d %d",watchdog, binding,&client_tag , &max,&buffer);
|
char header;
|
||||||
if (n<3) {
|
int n = sscanf(parm, "%c %s %s %d %d %d",&header,watchdog, binding,&client_tag , &max,&buffer);
|
||||||
|
if (n<4) {
|
||||||
skynet_error(ctx, "Invalid gate parm %s",parm);
|
skynet_error(ctx, "Invalid gate parm %s",parm);
|
||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
|
if (header != 'S' && header !='L') {
|
||||||
|
skynet_error(ctx, "Invalid data header style");
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
if (client_tag == 0) {
|
if (client_tag == 0) {
|
||||||
client_tag = PTYPE_CLIENT;
|
client_tag = PTYPE_CLIENT;
|
||||||
}
|
}
|
||||||
@@ -279,6 +290,7 @@ gate_init(struct gate *g , struct skynet_context * ctx, char * parm) {
|
|||||||
g->max_connection = max;
|
g->max_connection = max;
|
||||||
g->id_index = 0;
|
g->id_index = 0;
|
||||||
g->client_tag = client_tag;
|
g->client_tag = client_tag;
|
||||||
|
g->header_size = header=='S' ? 2 : 4;
|
||||||
|
|
||||||
g->agent = malloc(cap * sizeof(struct connection *));
|
g->agent = malloc(cap * sizeof(struct connection *));
|
||||||
memset(g->agent, 0, cap * sizeof(struct connection *));
|
memset(g->agent, 0, cap * sizeof(struct connection *));
|
||||||
|
|||||||
@@ -270,10 +270,10 @@ _message_to_header(const uint32_t *message, struct remote_message_header *header
|
|||||||
|
|
||||||
static int
|
static int
|
||||||
_send_package(int fd, const void * buffer, size_t sz) {
|
_send_package(int fd, const void * buffer, size_t sz) {
|
||||||
uint16_t header = htons(sz);
|
uint32_t header = htonl(sz);
|
||||||
struct iovec part[2];
|
struct iovec part[2];
|
||||||
part[0].iov_base = &header;
|
part[0].iov_base = &header;
|
||||||
part[0].iov_len = 2;
|
part[0].iov_len = 4;
|
||||||
part[1].iov_base = (void*)buffer;
|
part[1].iov_base = (void*)buffer;
|
||||||
part[1].iov_len = sz;
|
part[1].iov_len = sz;
|
||||||
|
|
||||||
@@ -286,7 +286,7 @@ _send_package(int fd, const void * buffer, size_t sz) {
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (err != sz+2) {
|
if (err != sz+4) {
|
||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
return 0;
|
return 0;
|
||||||
@@ -295,11 +295,11 @@ _send_package(int fd, const void * buffer, size_t sz) {
|
|||||||
|
|
||||||
static int
|
static int
|
||||||
_send_remote(int fd, const char * buffer, size_t sz, struct remote_message_header * cookie) {
|
_send_remote(int fd, const char * buffer, size_t sz, struct remote_message_header * cookie) {
|
||||||
uint16_t sz_header = htons(sz+sizeof(*cookie));
|
uint32_t sz_header = htonl(sz+sizeof(*cookie));
|
||||||
struct iovec part[3];
|
struct iovec part[3];
|
||||||
|
|
||||||
part[0].iov_base = &sz_header;
|
part[0].iov_base = &sz_header;
|
||||||
part[0].iov_len = 2;
|
part[0].iov_len = 4;
|
||||||
|
|
||||||
part[1].iov_base = (char *)buffer;
|
part[1].iov_base = (char *)buffer;
|
||||||
part[1].iov_len = sz;
|
part[1].iov_len = sz;
|
||||||
@@ -318,7 +318,7 @@ _send_remote(int fd, const char * buffer, size_t sz, struct remote_message_heade
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (err != sz+sizeof(*cookie)+2) {
|
if (err != sz+sizeof(*cookie)+4) {
|
||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
return 0;
|
return 0;
|
||||||
@@ -568,7 +568,7 @@ harbor_init(struct harbor *h, struct skynet_context *ctx, const char * args) {
|
|||||||
h->master_fd = master_fd;
|
h->master_fd = master_fd;
|
||||||
|
|
||||||
char tmp[128];
|
char tmp[128];
|
||||||
sprintf(tmp,"gate ! %s %d %d 0",local_addr, PTYPE_HARBOR, REMOTE_MAX);
|
sprintf(tmp,"gate L ! %s %d %d 0",local_addr, PTYPE_HARBOR, REMOTE_MAX);
|
||||||
const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp);
|
const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp);
|
||||||
if (gate_addr == NULL) {
|
if (gate_addr == NULL) {
|
||||||
skynet_error(ctx, "Harbor : launch gate failed");
|
skynet_error(ctx, "Harbor : launch gate failed");
|
||||||
|
|||||||
@@ -131,18 +131,18 @@ _connect_to(const char *ipaddress) {
|
|||||||
|
|
||||||
static int
|
static int
|
||||||
_send_to(int fd, const void * buf, size_t sz, uint32_t handle) {
|
_send_to(int fd, const void * buf, size_t sz, uint32_t handle) {
|
||||||
char buffer[2 + sz + 12];
|
char buffer[4 + sz + 12];
|
||||||
uint16_t header = htons(sz+12);
|
uint32_t header = htonl(sz+12);
|
||||||
memcpy(buffer, &header, 2);
|
memcpy(buffer, &header, 4);
|
||||||
memcpy(buffer+2, buf, sz);
|
memcpy(buffer+4, buf, sz);
|
||||||
uint32_t u32 = 0;
|
uint32_t u32 = 0;
|
||||||
memcpy(buffer+2+sz,&u32,4);
|
memcpy(buffer+4+sz,&u32,4);
|
||||||
u32 = htonl(handle);
|
u32 = htonl(handle);
|
||||||
memcpy(buffer+2+sz+4,&u32,4);
|
memcpy(buffer+4+sz+4,&u32,4);
|
||||||
u32 = 0;
|
u32 = 0;
|
||||||
memcpy(buffer+2+sz+8,&u32,4);
|
memcpy(buffer+4+sz+8,&u32,4);
|
||||||
|
|
||||||
sz += 2 + 12;
|
sz += 4 + 12;
|
||||||
|
|
||||||
for (;;) {
|
for (;;) {
|
||||||
int err = send(fd, buffer, sz, 0);
|
int err = send(fd, buffer, sz, 0);
|
||||||
@@ -269,7 +269,7 @@ _mainloop(struct skynet_context * context, void * ud, int type, int session, uin
|
|||||||
int
|
int
|
||||||
master_init(struct master *m, struct skynet_context *ctx, const char * args) {
|
master_init(struct master *m, struct skynet_context *ctx, const char * args) {
|
||||||
char tmp[strlen(args) + 32];
|
char tmp[strlen(args) + 32];
|
||||||
sprintf(tmp,"gate ! %s %d %d 0",args,PTYPE_HARBOR,REMOTE_MAX);
|
sprintf(tmp,"gate L ! %s %d %d 0",args,PTYPE_HARBOR,REMOTE_MAX);
|
||||||
const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp);
|
const char * gate_addr = skynet_command(ctx, "LAUNCH", tmp);
|
||||||
if (gate_addr == NULL) {
|
if (gate_addr == NULL) {
|
||||||
skynet_error(ctx, "Master : launch gate failed");
|
skynet_error(ctx, "Master : launch gate failed");
|
||||||
|
|||||||
@@ -52,7 +52,7 @@ skynet.start(function()
|
|||||||
end
|
end
|
||||||
end)
|
end)
|
||||||
-- 0 for default client tag
|
-- 0 for default client tag
|
||||||
gate = skynet.launch("gate" , skynet.address(skynet.self()), port, 0, max_agent, buffer)
|
gate = skynet.launch("gate" , "S" , skynet.address(skynet.self()), port, 0, max_agent, buffer)
|
||||||
skynet.send(gate,"text", "start")
|
skynet.send(gate,"text", "start")
|
||||||
skynet.register(".watchdog")
|
skynet.register(".watchdog")
|
||||||
end)
|
end)
|
||||||
|
|||||||
Reference in New Issue
Block a user