new harbor

This commit is contained in:
云风
2012-08-17 11:50:45 +08:00
parent b06d380219
commit 4f6826de40
25 changed files with 331 additions and 1010 deletions

View File

@@ -1,4 +1,7 @@
.PHONY : all clean install
.PHONY : all clean
CFLAGS = -g -Wall
SHARED = -fPIC --shared
all : \
skynet \
@@ -14,7 +17,8 @@ lualib/socket.so \
lualib/lpeg.so \
lualib/protobuf.so \
lualib/int64.so \
skynet-master
service/master.so \
service/harbor.so
skynet : \
skynet-src/skynet_main.c \
@@ -26,36 +30,41 @@ skynet-src/skynet_start.c \
skynet-src/skynet_timer.c \
skynet-src/skynet_error.c \
skynet-src/skynet_harbor.c \
skynet-src/skynet_env.c \
master/master.c
gcc -Wall -g -Wl,-E -o $@ $^ -Iskynet-src -Imaster -lpthread -ldl -lrt -Wl,-E -llua -lm -lzmq
skynet-src/skynet_env.c
gcc $(CFLAGS) -Wl,-E -o $@ $^ -Iskynet-src -lpthread -ldl -lrt -Wl,-E -llua -lm
service/master.so : service-src/service_master.c
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src
service/harbor.so : service-src/service_harbor.c
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src
service/logger.so : skynet-src/skynet_logger.c
gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src
service/snlua.so : service-src/service_lua.c
gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src
service/gate.so : gate/mread.c gate/ringbuffer.c gate/main.c
gcc -Wall -g -fPIC --shared -o $@ $^ -Igate -Iskynet-src
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Igate -Iskynet-src
lualib/skynet.so : lualib-src/lua-skynet.c lua-serialize/serialize.c
gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src -Ilua-serialize
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src -Ilua-serialize
service/client.so : service-src/service_client.c
gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src
service/connection.so : connection/connection.c connection/main.c
gcc -Wall -g -fPIC --shared -o $@ $^ -Iconnection -Iskynet-src
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src -Iconnection
lualib/socket.so : connection/lua-socket.c
gcc -Wall -g -fPIC --shared -o $@ $^ -Iconnection -Iskynet-src
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src -Iconnection
service/broker.so : service-src/service_broker.c
gcc -Wall -g -fPIC --shared $^ -o $@ -Iskynet-src
gcc $(CFLAGS) $(SHARED) $^ -o $@ -Iskynet-src
lualib/lpeg.so : lpeg/lpeg.c
gcc -Wall -O2 -fPIC --shared $^ -o $@ -Ilpeg
gcc $(CFLAGS) $(SHARED) -O2 $^ -o $@ -Ilpeg
PROTOBUFSRC = \
lua-protobuf/context.c \
@@ -73,45 +82,14 @@ PROTOBUFSRC = \
lua-protobuf/decode.c
lualib/protobuf.so : $(PROTOBUFSRC) lua-protobuf/pbc-lua.c
gcc -O2 -Wall --shared -fPIC -o $@ $^
gcc $(CFLAGS) $(SHARED) -O2 $^ -o $@
lualib/int64.so : lua-int64/int64.c
gcc -O2 -Wall --shared -fPIC -o $@ $^
gcc $(CFLAGS) $(SHARED) -O2 $^ -o $@
client : client-src/client.c
gcc -Wall -g $^ -o $@ -lpthread
skynet-master : master/master.c master/main.c
gcc -g -Wall -Imaster -o $@ $^ -lzmq
gcc $(CFLAGS) $^ -o $@ -lpthread
clean :
rm skynet client skynet-master lualib/*.so service/*.so
$(SKYNET_PATH) :
mkdir $@
install-libs : | $(SKYNET_PATH)
cp -R lualib $(SKYNET_PATH)
install-bins : | $(SKYNET_PATH)
cp skynet $(SKYNET_PATH)
$(SKYNET_PATH)/service : | $(SKYNET_PATH)
mkdir $@
install-services : | $(SKYNET_PATH)/service
cp service/connection.so $(SKYNET_PATH)/service/
cp service/logger.so $(SKYNET_PATH)/service/
cp service/gate.so $(SKYNET_PATH)/service/
cp service/broker.so $(SKYNET_PATH)/service/
cp service/snlua.so $(SKYNET_PATH)/service/
cp service/client.so $(SKYNET_PATH)/service/
cp service/launcher.lua $(SKYNET_PATH)/service/
cp service/redis-cli.lua $(SKYNET_PATH)/service/
install : install-libs install-bins install-services
rm skynet client lualib/*.so service/*.so

View File

@@ -1,6 +1,6 @@
## Build
Install zeromq 2.2 and lua 5.2 first.
Install lua 5.2 first.
```
make

6
config
View File

@@ -3,10 +3,10 @@ thread = 8
mqueue = 256
logger = nil
harbor = 1
address = "tcp://127.0.0.1:2525"
master = "tcp://127.0.0.1:2012"
address = "127.0.0.1:2525"
master = "127.0.0.1:2012"
start = "main"
standalone = true
standalone = "0.0.0.0:2012"
luaservice = root.."service/?.lua;"..root.."service/?/init.lua"
cpath = root.."service/?.so"
protopath = root.."proto"

View File

@@ -3,6 +3,6 @@ mqueue = 256
cpath = "./service/?.so"
logger = nil
harbor = 2
address = "tcp://127.0.0.1:2526"
master = "tcp://127.0.0.1:2012"
address = "127.0.0.1:2526"
master = "127.0.0.1:2012"
start = "main_log"

View File

@@ -124,11 +124,11 @@ _poll(struct connection_server * server) {
}
}
static void
static int
_main(struct skynet_context * ctx, void * ud, int session, const char * uid, const void * msg, size_t sz) {
if (msg == NULL) {
_poll(ud);
return;
return 0;
}
const char * param = (const char *)msg + 4;
if (memcmp(msg, "ADD ", 4)==0) {
@@ -136,7 +136,7 @@ _main(struct skynet_context * ctx, void * ud, int session, const char * uid, con
int fd = strtol(param, &endptr, 10);
if (endptr == NULL) {
skynet_error(ctx, "[connection] Invalid ADD command from %s (session = %d)", uid, session);
return;
return 0;
}
int addr_sz = sz - (endptr - (char *)msg);
char * addr = malloc(addr_sz);
@@ -148,12 +148,14 @@ _main(struct skynet_context * ctx, void * ud, int session, const char * uid, con
int fd = strtol(param, &endptr, 10);
if (endptr == NULL) {
skynet_error(ctx, "[connection] Invalid DEL command from %s (session = %d)", uid, session);
return;
return 0;
}
_del(ud, fd);
} else {
skynet_error(ctx, "[connection] Invalid command from %s (session = %d)", uid, session);
}
return 0;
}
int

View File

@@ -122,6 +122,9 @@ _ctrl(struct skynet_context * ctx, struct gate * g, const void * msg, int sz) {
static void
_report(struct gate *g, struct skynet_context * ctx, const char * data, ...) {
if (g->watchdog == NULL) {
return;
}
va_list ap;
va_start(ap, data);
char tmp[1024];
@@ -134,14 +137,14 @@ _report(struct gate *g, struct skynet_context * ctx, const char * data, ...) {
static void
_forward(struct skynet_context * ctx,struct gate *g, int uid, void * data, size_t len) {
if (g->broker) {
skynet_send(ctx, NULL, g->broker, 0x7fffffff, data, len, 0);
skynet_send(ctx, NULL, g->broker, SESSION_CLIENT, data, len, 0);
return;
}
struct connection * agent = _id_to_agent(g,uid);
if (agent->agent) {
// todo: client package has not session , send 0x7fffffff
skynet_send(ctx, agent->client, agent->agent, 0x7fffffff, data, len, 0);
} else {
skynet_send(ctx, agent->client, agent->agent, SESSION_CLIENT, data, len, 0);
} else if (g->watchdog) {
char * tmp = malloc(len + 32);
int n = snprintf(tmp,len+32,"%d data ",uid);
memcpy(tmp+n,data,len);
@@ -178,12 +181,12 @@ _remove_id(struct gate *g, int uid) {
}
}
static void
static int
_cb(struct skynet_context * ctx, void * ud, int session, const char * uid, const void * msg, size_t sz) {
struct gate *g = ud;
if (msg) {
_ctrl(ctx, g , msg , (int)sz);
return;
return 0;
}
struct mread_pool * m = g->pool;
int connection_id = mread_poll(m,100); // timeout : 100ms
@@ -223,6 +226,7 @@ _cb(struct skynet_context * ctx, void * ud, int session, const char * uid, const
_break:
skynet_command(ctx, "TIMEOUT", "0");
}
return 0;
}
int
@@ -230,18 +234,41 @@ gate_init(struct gate *g , struct skynet_context * ctx, char * parm) {
int port = 0;
int max = 0;
int buffer = 0;
char watchdog[strlen(parm)+1];
int n = sscanf(parm, "%s %d %d %d",watchdog, &port,&max,&buffer);
int sz = strlen(parm)+1;
char watchdog[sz];
char binding[sz];
int n = sscanf(parm, "%s %s %d %d",watchdog, binding,&max,&buffer);
if (n!=4) {
skynet_error(ctx, "Invalid gate parm %s",parm);
return 1;
}
struct mread_pool * pool = mread_create(port, max, buffer);
char * portstr = strchr(binding,':');
uint32_t addr = INADDR_ANY;
if (portstr == NULL) {
port = strtol(binding, NULL, 10);
if (port <= 0) {
skynet_error(ctx, "Invalid gate address %s",parm);
return 1;
}
} else {
port = strtol(portstr + 1, NULL, 10);
if (port <= 0) {
skynet_error(ctx, "Invalid gate address %s",parm);
return 1;
}
portstr[0] = '\0';
addr=inet_addr(binding);
}
struct mread_pool * pool = mread_create(addr, port, max, buffer);
if (pool == NULL) {
skynet_error(ctx, "Create gate %s failed",parm);
return 1;
}
g->watchdog = strdup(watchdog);
if (watchdog[0] == '!') {
g->watchdog = NULL;
} else {
g->watchdog = strdup(watchdog);
}
g->pool = pool;
int cap = 1;
while (cap < max) {

View File

@@ -93,7 +93,7 @@ _set_nonblocking(int fd)
}
struct mread_pool *
mread_create(int port , int max , int buffer_size) {
mread_create(uint32_t addr, int port , int max , int buffer_size) {
int listen_fd = socket(AF_INET, SOCK_STREAM, 0);
if (listen_fd == -1) {
return NULL;
@@ -109,8 +109,8 @@ mread_create(int port , int max , int buffer_size) {
memset(&my_addr, 0, sizeof(struct sockaddr_in));
my_addr.sin_family = AF_INET;
my_addr.sin_port = htons(port);
my_addr.sin_addr.s_addr = htonl(INADDR_ANY); // INADDR_LOOPBACK
// printf("MREAD bind %s:%u\n",inet_ntoa(my_addr.sin_addr),ntohs(my_addr.sin_port));
my_addr.sin_addr.s_addr = addr;
printf("MREAD bind %s:%u\n",inet_ntoa(my_addr.sin_addr),ntohs(my_addr.sin_port));
if (bind(listen_fd, (struct sockaddr *)&my_addr, sizeof(struct sockaddr)) == -1) {
close(listen_fd);
return NULL;

View File

@@ -1,9 +1,11 @@
#ifndef MREAD_H
#define MREAD_H
#include <stdint.h>
struct mread_pool;
struct mread_pool * mread_create(int port , int max , int buffer);
struct mread_pool * mread_create(uint32_t addr, int port , int max , int buffer);
void mread_close(struct mread_pool *m);
int mread_poll(struct mread_pool *m , int timeout);

View File

@@ -5,24 +5,20 @@
#include <lauxlib.h>
#include <stdlib.h>
#include <string.h>
#include <assert.h>
static int
traceback (lua_State *L) {
const char *msg = lua_tostring(L, 1);
if (msg) {
luaL_traceback(L, L, msg, 1);
} else {
lua_pushliteral(L, "(no error message)");
}
return 1;
}
static void
static int
_cb(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) {
lua_State *L = ud;
lua_rawgetp(L, LUA_REGISTRYINDEX, traceback);
int trace = lua_gettop(L);
lua_rawgetp(L, LUA_REGISTRYINDEX, _cb);
int trace = 1;
int top = lua_gettop(L);
if (top == 1) {
lua_rawgetp(L, LUA_REGISTRYINDEX, _cb);
} else {
assert(top == 2);
lua_pushvalue(L,2);
}
int r;
if (msg == NULL) {
if (addr == NULL) {
@@ -41,7 +37,7 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr,
r = lua_pcall(L, 4, 0 , trace);
}
if (r == LUA_OK)
return;
return 0;
const char * self = skynet_command(context, "REG", NULL);
switch (r) {
case LUA_ERRRUN:
@@ -59,6 +55,8 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr,
};
lua_pop(L,1);
return 0;
}
static int
@@ -201,9 +199,6 @@ luaopen_skynet_c(lua_State *L) {
{ "unpack", _luaseri_unpack },
{ NULL, NULL },
};
lua_pushcfunction(L, traceback);
lua_rawsetp(L, LUA_REGISTRYINDEX, traceback);
luaL_newlibtable(L,l);
lua_getfield(L, LUA_REGISTRYINDEX, "skynet_context");

View File

@@ -56,7 +56,7 @@ _forward(struct broker *b, struct skynet_context * context) {
b->id = (b->id + 1) % DEFAULT_NUMBER;
}
static void
static int
_cb(struct skynet_context * context, void * ud, int session, const char * addr, const void * msg, size_t sz) {
struct broker * b = ud;
if (b->init < DEFAULT_NUMBER) {
@@ -68,6 +68,8 @@ _cb(struct skynet_context * context, void * ud, int session, const char * addr,
} else {
_forward(b, context);
}
return 0;
}

View File

@@ -8,7 +8,7 @@
#include <stdio.h>
#include <string.h>
static void
static int
_cb(struct skynet_context * context, void * ud, int session, const char * addr, 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,
}
}
assert(err == sz +2);
return;
return 0;
}
}

View File

@@ -4,12 +4,12 @@ local client = ...
local session_id = 0
skynet.filter(function (session, address , msg, sz)
if session == 0x7fffffff then
print("client message",skynet.tostring(msg,sz))
skynet.send("LOG", "client message :" .. skynet.tostring(msg,sz))
-- It's client, there is no session
session_id = session_id + 1
session = - session_id
else
print("skynet message",msg,sz)
skynet.send("LOG", "skynet message")
end
return session, address , msg, sz
end, function (msg,sz)

View File

@@ -2,7 +2,7 @@ local skynet = require "skynet"
print("Log server start")
local log = skynet.launch("snlua","globallog.lua")
local log = skynet.launch("snlua","globallog")
print("log",log)
skynet.exit()

View File

@@ -5,6 +5,7 @@
#include <stdint.h>
#define DONTCOPY 1
#define SESSION_CLIENT 0x7fffffff
struct skynet_context;
@@ -12,8 +13,9 @@ void skynet_error(struct skynet_context * context, const char *msg, ...);
const char * skynet_command(struct skynet_context * context, const char * cmd , const char * parm);
int skynet_send(struct skynet_context * context, const char * source, const char * addr , int session, void * msg, size_t sz, int flags);
void skynet_forward(struct skynet_context *, const char * addr);
typedef void (*skynet_cb)(struct skynet_context * context, void *ud, int session, const char * addr , const void * msg, size_t sz);
void skynet_forward(struct skynet_context *, const char * destination);
typedef int (*skynet_cb)(struct skynet_context * context, void *ud, int session, const char * addr , const void * msg, size_t sz);
void skynet_callback(struct skynet_context * context, void *ud, skynet_cb cb);
#endif

View File

@@ -3,10 +3,7 @@
#include <stdint.h>
// reserve high 8 bits for remote id
// see skynet_harbor.c REMOTE_MAX
#define HANDLE_MASK 0xffffff
#define HANDLE_REMOTE_SHIFT 24
#include "skynet_harbor.h"
struct skynet_context;

View File

@@ -1,674 +1,53 @@
#include "skynet_harbor.h"
#include "skynet_mq.h"
#include "skynet_handle.h"
#include "skynet_system.h"
#include "skynet_server.h"
#include "skynet.h"
#include <zmq.h>
#include <string.h>
#include <stdlib.h>
#include <stdio.h>
#include <assert.h>
#define HASH_SIZE 4096
#define DEFAULT_QUEUE_SIZE 1024
static struct skynet_context * REMOTE = 0;
static int HARBOR = 0;
// see skynet_handle.h for HANDLE_REMOTE_SHIFT
#define REMOTE_MAX 255
struct keyvalue {
struct keyvalue * next;
uint32_t hash;
char * key;
uint32_t value;
struct message_queue * queue;
};
struct hashmap {
struct keyvalue *node[HASH_SIZE];
};
struct remote_header {
uint32_t source;
uint32_t destination;
uint32_t session;
};
struct remote {
void *socket;
struct message_remote_queue *queue;
};
struct harbor {
void * zmq_context;
void * zmq_master_request;
void * zmq_local;
void * zmq_queue_notice;
int notice_event;
struct hashmap *map;
struct remote remote[REMOTE_MAX];
struct message_remote_queue *queue;
int harbor;
int lock;
};
static struct harbor *Z = NULL;
// todo: optimize for little endian system
static inline void
buffer_to_remote_header(uint8_t *buffer, struct remote_header *header) {
header->source = buffer[0] | buffer[1] << 8 | buffer[2] << 16 | buffer[3] << 24;
header->destination = buffer[4] | buffer[5] << 8 | buffer[6] << 16 | buffer[7] << 24;
header->session = buffer[8] | buffer[9] << 8 | buffer[10] << 16 | buffer[11] << 24;
void
skynet_harbor_send(struct remote_message *rmsg, uint32_t source, int session) {
skynet_context_send(REMOTE, rmsg, sizeof(*rmsg), source, session);
}
static inline void
remote_header_to_buffer(struct remote_header *header, uint8_t *buffer) {
buffer[0] = header->source & 0xff;
buffer[1] = (header->source >> 8) & 0xff;
buffer[2] = (header->source >>16) & 0xff;
buffer[3] = (header->source >>24)& 0xff;
buffer[4] = (header->destination) & 0xff;
buffer[5] = (header->destination >>8) & 0xff;
buffer[6] = (header->destination >>16) & 0xff;
buffer[7] = (header->destination >>24) & 0xff;
buffer[8] = (header->session) & 0xff;
buffer[9] = (header->session >>8) & 0xff;
buffer[10] = (header->session >>16) & 0xff;
buffer[11] = (header->session >>24) & 0xff;
}
static uint32_t
calc_hash(const char *name) {
void
skynet_harbor_register(struct remote_name *rname) {
int i;
uint32_t h = 0;
for (i=0;name[i];i++) {
h = h ^ ((h<<5)+(h>>2)+(uint8_t)name[i]);
}
h ^= i;
return h;
}
static inline void
_lock() {
while (__sync_lock_test_and_set(&Z->lock,1)) {}
}
static inline void
_unlock() {
__sync_lock_release(&Z->lock);
}
static struct hashmap *
_hash_new(void) {
struct hashmap * hash = malloc(sizeof(*hash));
memset(hash, 0, sizeof(*hash));
return hash;
}
static struct keyvalue *
_hash_search(struct hashmap * hash, const char * key) {
uint32_t h = calc_hash(key);
struct keyvalue * n = hash->node[h & (HASH_SIZE-1)];
while (n) {
if (n->hash == h && strcmp(n->key, key) == 0) {
return n;
}
n = n->next;
}
return NULL;
}
static void
_hash_insert(struct hashmap * hash, const char * key, uint32_t handle, struct message_queue *queue) {
uint32_t h = calc_hash(key);
struct keyvalue * node = malloc(sizeof(*node));
node->next = hash->node[h & (HASH_SIZE-1)];
node->hash = h;
node->key = strdup(key);
node->value = handle;
node->queue = queue;
hash->node[h & (HASH_SIZE-1)] = node;
}
// thread safe function
static void
send_notice() {
if (__sync_lock_test_and_set(&Z->notice_event,1)) {
// already send notice
return;
}
static __thread void * queue_notice = NULL;
if (queue_notice == NULL) {
void * pub = zmq_socket(Z->zmq_context, ZMQ_PUSH);
int r = zmq_connect(pub , "inproc://notice");
assert(r==0);
queue_notice = pub;
}
zmq_msg_t dummy;
zmq_msg_init(&dummy);
zmq_send(queue_notice,&dummy,0);
zmq_msg_close(&dummy);
}
// thread safe function
void
skynet_harbor_send(const char *name, uint32_t destination, struct skynet_message * msg) {
if (name == NULL) {
assert(destination!=0);
int remote_id = destination >> HANDLE_REMOTE_SHIFT;
assert(remote_id > 0 && remote_id <= REMOTE_MAX);
struct skynet_remote_message message;
message.destination = destination;
message.message = *msg;
skynet_remotemq_push(Z->queue, &message);
send_notice();
} else {
_lock();
struct keyvalue * node = _hash_search(Z->map, name);
if (node) {
uint32_t dest = node->value;
_unlock();
if (dest == 0) {
// push message to unknown name service queue
skynet_mq_push(node->queue, msg);
} else {
if (!skynet_harbor_message_isremote(dest)) {
// local message
if (skynet_context_push(dest, msg)) {
skynet_error(NULL, "Drop local message from %u to %s",msg->source, name);
}
return;
}
struct skynet_remote_message message;
message.destination = dest;
message.message = *msg;
skynet_remotemq_push(Z->queue,&message);
send_notice();
}
} else {
// never seen name before
struct message_queue * queue = skynet_mq_create(0);
skynet_mq_push(queue, msg);
_hash_insert(Z->map, name, 0, queue);
_unlock();
// 0 for query
skynet_harbor_register(name,0);
int number = 1;
for (i=0;i<GLOBALNAME_LENGTH;i++) {
char c = rname->name[i];
if (!(c >= '0' && c <='9')) {
number = 0;
break;
}
}
}
// thread safe function
//queue a register message (destination = 0)
void
skynet_harbor_register(const char *name, uint32_t handle) {
struct skynet_remote_message msg;
msg.destination = SKYNET_SYSTEM_NAME;
msg.message.source = handle;
msg.message.data = strdup(name);
msg.message.sz = 0;
skynet_remotemq_push(Z->queue,&msg);
send_notice();
}
// Always in main harbor thread
static void
_register_name(const char *name, uint32_t addr) {
_lock();
struct keyvalue * node = _hash_search(Z->map, name);
if (node) {
if (node->value) {
node->value = addr;
assert(node->queue == NULL);
} else {
node->value = addr;
}
} else {
_hash_insert(Z->map, name, addr, NULL);
}
if (addr == 0) {
_unlock();
return;
}
struct skynet_message msg;
struct message_queue * queue = node ? node->queue : NULL;
if (queue) {
if (skynet_harbor_message_isremote(addr)) {
while (!skynet_mq_pop(queue, &msg)) {
struct skynet_remote_message message;
message.destination = addr;
message.message = msg;
skynet_remotemq_push(Z->queue, &message);
send_notice();
}
} else {
while (!skynet_mq_pop(queue, &msg)) {
if (skynet_context_push(addr,&msg)) {
skynet_error(NULL,"Drop local message from %u to %s",msg.source,name);
}
}
}
node->queue = NULL;
}
_unlock();
if (queue) {
skynet_mq_release(queue);
}
}
// Always in main harbor thread
static void
_remote_harbor_update(int harbor_id, const char * addr) {
struct remote * r = &Z->remote[harbor_id-1];
void *socket = zmq_socket( Z->zmq_context, ZMQ_PUSH);
int rc = zmq_connect(socket, addr);
if (rc<0) {
skynet_error(NULL, "Can't connect to harbor %d %s",harbor_id,addr);
zmq_close(socket);
socket = NULL;
}
if (socket) {
void *old_socket = r->socket;
if (old_socket) {
zmq_close(old_socket);
}
struct message_remote_queue * queue = r->queue;
if (queue) {
struct skynet_remote_message msg;
while (!skynet_remotemq_pop(queue, &msg)) {
skynet_remotemq_push(Z->queue, &msg);
}
skynet_remotemq_release(queue);
r->queue = NULL;
}
r->socket = socket;
}
}
static void
_report_zmq_error(int rc) {
if (rc) {
fprintf(stderr, "zmq error : %s\n",zmq_strerror(errno));
exit(1);
}
}
static int
_isdecimal(int c) {
return c>='0' && c<='9';
}
// Name-updating protocols:
//
// 1) harbor_id=harbor_address
// 2) context_name=context_handle
static int
_split_name(uint8_t *buf, int len, int *np) {
uint8_t *sep;
if (len > 0 && _isdecimal(buf[0])) {
int i=0;
int n=0;
do {
n = n*10 + (buf[i]-'0');
} while(++i<len && _isdecimal(buf[i]));
if (i < len && buf[i] == '=') {
buf[i] = '\0';
*np = n;
return i;
}
} else if ((sep = memchr(buf, '=', len)) != NULL) {
*sep = '\0';
return (int)(sep-buf);
}
return -1;
}
// Always in main harbor thread
static void
_name_update() {
zmq_msg_t content;
zmq_msg_init(&content);
int rc = zmq_recv(Z->zmq_local,&content,0);
_report_zmq_error(rc);
int sz = zmq_msg_size(&content);
uint8_t * buffer = zmq_msg_data(&content);
int n = 0;
int i = _split_name(buffer, sz, &n);
if (i == -1) {
char tmp[sz+1];
memcpy(tmp,buffer,sz);
tmp[sz] = '\0';
skynet_error(NULL, "Invalid master update [%s]",tmp);
zmq_msg_close(&content);
return;
}
char tmp[sz-i];
memcpy(tmp,buffer+i+1,sz-i-1);
tmp[sz-i-1]='\0';
if (n>0 && n <= REMOTE_MAX) {
_remote_harbor_update(n, tmp);
} else {
uint32_t source = strtoul(tmp,NULL,16);
if (source == 0) {
skynet_error(NULL, "Invalid master update [%s=%s]",(const char *)buffer,tmp);
} else {
_register_name((const char *)buffer, source);
}
}
zmq_msg_close(&content);
}
// Always in main harbor thread
static void
remote_query_harbor(int harbor_id) {
char tmp[32];
int sz = sprintf(tmp,"%d",harbor_id);
zmq_msg_t request;
zmq_msg_init_size(&request,sz);
memcpy(zmq_msg_data(&request),tmp,sz);
zmq_send(Z->zmq_master_request, &request, 0);
zmq_msg_close(&request);
zmq_msg_t reply;
zmq_msg_init(&reply);
int rc = zmq_recv(Z->zmq_master_request, &reply, 0);
_report_zmq_error(rc);
sz = zmq_msg_size(&reply);
char tmp2[sz+1];
memcpy(tmp2,zmq_msg_data(&reply),sz);
tmp2[sz] = '\0';
if (sz == 0) {
skynet_error(NULL,"Request harbor %d failed",harbor_id);
} else {
_remote_harbor_update(harbor_id, tmp2);
}
zmq_msg_close(&reply);
}
// Always in main harbor thread
static void
_remote_register_name(const char *name, uint32_t source) {
char tmp[strlen(name) + 20];
int sz = sprintf(tmp,"%s=%X",name,source);
zmq_msg_t msg;
zmq_msg_init_size(&msg,sz);
memcpy(zmq_msg_data(&msg), tmp , sz);
zmq_send(Z->zmq_master_request, &msg,0);
zmq_msg_close(&msg);
zmq_msg_init(&msg);
int rc = zmq_recv(Z->zmq_master_request, &msg,0);
_report_zmq_error(rc);
zmq_msg_close(&msg);
}
// Always in main harbor thread
static void
_remote_query_name(const char *name) {
int sz = strlen(name);
zmq_msg_t msg;
zmq_msg_init_size(&msg,sz);
memcpy(zmq_msg_data(&msg), name , sz);
zmq_send(Z->zmq_master_request, &msg,0);
zmq_msg_close(&msg);
zmq_msg_init(&msg);
int rc = zmq_recv(Z->zmq_master_request, &msg,0);
_report_zmq_error(rc);
sz = zmq_msg_size(&msg);
char tmp[sz+1];
memcpy(tmp, zmq_msg_data(&msg),sz);
tmp[sz] = '\0';
uint32_t addr = strtoul(tmp,NULL,16);
_register_name(name,addr);
zmq_msg_close(&msg);
}
// Always in main harbor thread
static void
free_message(void *data, void *hint) {
free(data);
}
static void
remote_socket_send(void * socket, struct skynet_remote_message *msg) {
struct remote_header rh;
rh.source = msg->message.source;
rh.destination = msg->destination;
rh.session = msg->message.session;
zmq_msg_t part;
zmq_msg_init_size(&part,sizeof(struct remote_header));
uint8_t * buffer = zmq_msg_data(&part);
remote_header_to_buffer(&rh,buffer);
zmq_send(socket, &part, ZMQ_SNDMORE);
zmq_msg_close(&part);
zmq_msg_init_data(&part,msg->message.data,msg->message.sz,free_message,NULL);
zmq_send(socket, &part, 0);
zmq_msg_close(&part);
}
// Always in main harbor thread
// remote message has two part
// when part one is nil (size == 0), part two is name update
// Or part one is source:destination (8 bytes little endian), part two is a binary block for message
static void
_remote_recv() {
zmq_msg_t header;
zmq_msg_init(&header);
int rc = zmq_recv(Z->zmq_local,&header,0);
_report_zmq_error(rc);
size_t s = zmq_msg_size(&header);
if (s!=sizeof(struct remote_header)) {
// s should be 0
if (s>0) {
char tmp[s+1];
memcpy(tmp, zmq_msg_data(&header),s);
tmp[s] = '\0';
skynet_error(NULL,"Invalid master header [%s]",tmp);
}
_name_update();
return;
}
uint8_t * buffer = zmq_msg_data(&header);
struct remote_header rh;
buffer_to_remote_header(buffer, &rh);
zmq_close(&header);
zmq_msg_t * data = malloc(sizeof(zmq_msg_t));
zmq_msg_init(data);
rc = zmq_recv(Z->zmq_local,data,0);
_report_zmq_error(rc);
struct skynet_message msg;
msg.session = (int)rh.session;
msg.source = rh.source;
msg.data = data;
msg.sz = zmq_msg_size(data);
// push remote message to local message queue
if (skynet_context_push(rh.destination, &msg)) {
zmq_msg_close(data);
free(data);
skynet_error(NULL, "Drop remote message from %u to %u",rh.source, rh.destination);
}
}
// Always in main harbor thread
static void
_remote_send() {
struct skynet_remote_message msg;
while (!skynet_remotemq_pop(Z->queue,&msg)) {
_goback:
if (msg.destination == SKYNET_SYSTEM_NAME) {
// register name
char * name = msg.message.data;
if (msg.message.source) {
_remote_register_name(name, msg.message.source);
} else {
_remote_query_name(name);
}
free(name);
} else {
int harbor_id = (msg.destination >> HANDLE_REMOTE_SHIFT);
assert(harbor_id > 0);
struct remote * r = &Z->remote[harbor_id-1];
if (r->socket == NULL) {
if (r->queue == NULL) {
r->queue = skynet_remotemq_create();
skynet_remotemq_push(r->queue, &msg);
remote_query_harbor(harbor_id);
} else {
skynet_remotemq_push(r->queue, &msg);
}
} else {
remote_socket_send(r->socket, &msg);
}
}
}
__sync_lock_release(&Z->notice_event);
// double check
if (!skynet_remotemq_pop(Z->queue,&msg)) {
__sync_lock_test_and_set(&Z->notice_event, 1);
goto _goback;
}
}
// Main harbor thread
void *
skynet_harbor_dispatch_thread(void *ud) {
zmq_pollitem_t items[2];
items[0].socket = Z->zmq_queue_notice;
items[0].events = ZMQ_POLLIN;
items[1].socket = Z->zmq_local;
items[1].events = ZMQ_POLLIN;
for (;;) {
zmq_poll(items,2,-1);
if (items[0].revents) {
zmq_msg_t msg;
zmq_msg_init(&msg);
int rc = zmq_recv(Z->zmq_queue_notice,&msg,0);
_report_zmq_error(rc);
zmq_msg_close(&msg);
_remote_send();
}
if (items[1].revents) {
_remote_recv();
}
}
}
// Call only at init
static void
register_harbor(void *request, const char *local, int harbor) {
char tmp[1024];
sprintf(tmp,"%d",harbor);
size_t sz = strlen(tmp);
zmq_msg_t req;
zmq_msg_init_size (&req , sz);
memcpy(zmq_msg_data(&req),tmp,sz);
zmq_send (request, &req, 0);
zmq_msg_close (&req);
zmq_msg_t reply;
zmq_msg_init (&reply);
int rc = zmq_recv(request, &reply, 0);
_report_zmq_error(rc);
sz = zmq_msg_size (&reply);
if (sz > 0) {
memcpy(tmp,zmq_msg_data(&reply),sz);
tmp[sz] = '\0';
fprintf(stderr, "Harbor %d is already registered by %s\n", harbor, tmp);
exit(1);
}
zmq_msg_close (&reply);
sprintf(tmp,"%d=%s",harbor,local);
sz = strlen(tmp);
zmq_msg_init_size (&req , sz);
memcpy(zmq_msg_data(&req),tmp,sz);
zmq_send (request, &req, 0);
zmq_msg_close (&req);
zmq_msg_init (&reply);
rc = zmq_recv (request, &reply, 0);
_report_zmq_error(rc);
zmq_msg_close (&reply);
}
void
skynet_harbor_init(void * context, const char * master, const char *local, int harbor) {
if (harbor <=0 || harbor>255 || strlen(local) > 512) {
fprintf(stderr,"Invalid harbor id\n");
exit(1);
}
void *request = zmq_socket (context, ZMQ_REQ);
int r = zmq_connect(request, master);
if (r<0) {
fprintf(stderr, "Can't connect to master: %s\n",master);
exit(1);
}
void *harbor_socket = zmq_socket(context, ZMQ_PULL);
r = zmq_bind(harbor_socket, local);
if (r<0) {
fprintf(stderr, "Can't bind to local : %s\n",local);
exit(1);
}
register_harbor(request,local, harbor);
printf("Start harbor on : %s\n",local);
struct harbor * h = malloc(sizeof(*h));
memset(h, 0, sizeof(*h));
h->zmq_context = context;
h->zmq_master_request = request;
h->zmq_local = harbor_socket;
h->map = _hash_new();
h->harbor = harbor;
h->queue = skynet_remotemq_create();
h->zmq_queue_notice = zmq_socket(context, ZMQ_PULL);
r = zmq_bind(h->zmq_queue_notice, "inproc://notice");
assert(r==0);
Z = h;
}
// thread safe api
void *
skynet_harbor_message_open(struct skynet_message * message) {
return zmq_msg_data(message->data);
}
void
skynet_harbor_message_close(struct skynet_message * message) {
zmq_msg_close(message->data);
assert(number == 0);
skynet_context_send(REMOTE, rname, sizeof(*rname), 0, 0);
}
int
skynet_harbor_message_isremote(uint32_t handle) {
int harbor_id = handle >> HANDLE_REMOTE_SHIFT;
return !(harbor_id == 0 || harbor_id == Z->harbor);
return !(handle & HARBOR);
}
void
skynet_harbor_init(int harbor) {
HARBOR = harbor << HANDLE_REMOTE_SHIFT;
}
int
skynet_harbor_start(const char * master, const char *local) {
size_t sz = strlen(master) + strlen(local) + 32;
char args[sz];
sprintf(args, "%s %s %d",master,local,HARBOR >> HANDLE_REMOTE_SHIFT);
struct skynet_context * inst = skynet_context_new("harbor",args);
if (inst == NULL) {
return 1;
}
REMOTE = inst;
return 0;
}

View File

@@ -2,20 +2,30 @@
#define SKYNET_HARBOR_H
#include <stdint.h>
#include <stdlib.h>
struct skynet_message;
#define GLOBALNAME_LENGTH 16
#define REMOTE_MAX 256
void skynet_harbor_send(const char *name, uint32_t destination, struct skynet_message * message);
void skynet_harbor_register(const char *name, uint32_t handle);
// reserve high 8 bits for remote id
#define HANDLE_MASK 0xffffff
#define HANDLE_REMOTE_SHIFT 24
// remote message is diffrent from local message.
// We must use these api to open and close message , see skynet_server.c
struct remote_name {
char name[GLOBALNAME_LENGTH];
uint32_t handle;
};
struct remote_message {
struct remote_name destination;
const void * message;
size_t sz;
};
void skynet_harbor_send(struct remote_message *rmsg, uint32_t source, int session);
void skynet_harbor_register(struct remote_name *rname);
int skynet_harbor_message_isremote(uint32_t handle);
void * skynet_harbor_message_open(struct skynet_message * message);
void skynet_harbor_message_close(struct skynet_message * message);
// harbor worker thread
void * skynet_harbor_dispatch_thread(void *ud);
void skynet_harbor_init(void * context, const char * master, const char *local, int harbor);
void skynet_harbor_init(int harbor);
int skynet_harbor_start(const char * master, const char *local);
#endif

View File

@@ -10,7 +10,7 @@ struct skynet_config {
const char * master;
const char * local;
const char * start;
int standalone;
const char * standalone;
};
void skynet_start(struct skynet_config * config);

View File

@@ -24,12 +24,14 @@ logger_release(struct logger * inst) {
free(inst);
}
static void
static int
_logger(struct skynet_context * context, void *ud, int session, const char * uid, const void * msg, size_t sz) {
struct logger * inst = ud;
fprintf(inst->handle, "[%s] ",uid);
fwrite(msg, sz , 1, inst->handle);
fprintf(inst->handle, "\n");
return 0;
}
int

View File

@@ -20,6 +20,7 @@ optint(const char *key, int opt) {
return strtol(str, NULL, 10);
}
/*
static int
optboolean(const char *key, int opt) {
const char * str = skynet_getenv(key);
@@ -29,7 +30,7 @@ optboolean(const char *key, int opt) {
}
return strcmp(str,"true")==0;
}
*/
static const char *
optstring(const char *key,const char * opt) {
const char * str = skynet_getenv(key);
@@ -104,10 +105,10 @@ main(int argc, char *argv[]) {
config.module_path = optstring("cpath","./service/?.so");
config.logger = optstring("logger",NULL);
config.harbor = optint("harbor", 1);
config.master = optstring("master","tcp://127.0.0.1:2012");
config.master = optstring("master","127.0.0.1:2012");
config.start = optstring("start","main.lua");
config.local = optstring("address","tcp://127.0.0.1:2525");
config.standalone = optboolean("standalone",0);
config.local = optstring("address","127.0.0.1:2525");
config.standalone = optstring("standalone",NULL);
lua_close(L);

View File

@@ -13,6 +13,7 @@ struct message_queue {
int head;
int tail;
int lock;
int in_global;
struct skynet_message *queue;
};
@@ -24,14 +25,6 @@ struct global_queue {
struct message_queue ** queue;
};
struct message_remote_queue {
int cap;
int head;
int tail;
int lock;
struct skynet_remote_message *queue;
};
static struct global_queue *Q = NULL;
static inline void
@@ -42,7 +35,7 @@ _lock_global_queue() {
#define LOCK(q) while (__sync_lock_test_and_set(&(q)->lock,1)) {}
#define UNLOCK(q) __sync_lock_release(&(q)->lock);
void
static void
skynet_globalmq_push(struct message_queue * queue) {
struct global_queue *q= Q;
LOCK(q)
@@ -95,6 +88,7 @@ skynet_mq_create(uint32_t handle) {
q->head = 0;
q->tail = 0;
q->lock = 0;
q->in_global = 1;
q->queue = malloc(sizeof(struct skynet_message) * q->cap);
return q;
@@ -114,7 +108,7 @@ skynet_mq_handle(struct message_queue *q) {
int
skynet_mq_pop(struct message_queue *q, struct skynet_message *message) {
int ret = -1;
int ret = 1;
LOCK(q)
if (q->head != q->tail) {
@@ -124,6 +118,10 @@ skynet_mq_pop(struct message_queue *q, struct skynet_message *message) {
q->head = 0;
}
}
if (ret) {
q->in_global = 0;
}
UNLOCK(q)
@@ -134,23 +132,30 @@ void
skynet_mq_push(struct message_queue *q, struct skynet_message *message) {
LOCK(q)
q->queue[q->tail] = *message;
if (++ q->tail >= q->cap) {
q->tail = 0;
if (message) {
q->queue[q->tail] = *message;
if (++ q->tail >= q->cap) {
q->tail = 0;
}
if (q->head == q->tail) {
struct skynet_message *new_queue = malloc(sizeof(struct skynet_message) * q->cap * 2);
int i;
for (i=0;i<q->cap;i++) {
new_queue[i] = q->queue[(q->head + i) % q->cap];
}
q->head = 0;
q->tail = q->cap;
q->cap *= 2;
free(q->queue);
q->queue = new_queue;
}
}
if (q->head == q->tail) {
struct skynet_message *new_queue = malloc(sizeof(struct skynet_message) * q->cap * 2);
int i;
for (i=0;i<q->cap;i++) {
new_queue[i] = q->queue[(q->head + i) % q->cap];
}
q->head = 0;
q->tail = q->cap;
q->cap *= 2;
free(q->queue);
q->queue = new_queue;
if (q->in_global == 0) {
skynet_globalmq_push(q);
q->in_global = 1;
}
UNLOCK(q)
@@ -170,68 +175,8 @@ skynet_mq_init(int n) {
Q=q;
}
// remote message queue
struct message_remote_queue *
skynet_remotemq_create(void) {
struct message_remote_queue *q = malloc(sizeof(*q));
q->cap = DEFAULT_QUEUE_SIZE;
q->head = 0;
q->tail = 0;
q->lock = 0;
q->queue = malloc(sizeof(struct skynet_remote_message) * q->cap);
return q;
}
void
skynet_remotemq_release(struct message_remote_queue *q) {
free(q->queue);
free(q);
}
int
skynet_remotemq_pop(struct message_remote_queue *q, struct skynet_remote_message *message) {
int ret = -1;
LOCK(q)
if (q->head != q->tail) {
*message = q->queue[q->head];
ret = 0;
if ( ++ q->head >= q->cap) {
q->head = 0;
}
}
UNLOCK(q)
return ret;
}
void
skynet_remotemq_push(struct message_remote_queue *q, struct skynet_remote_message *message) {
assert(message->destination != 0);
LOCK(q)
q->queue[q->tail] = *message;
if (++ q->tail >= q->cap) {
q->tail = 0;
}
if (q->head == q->tail) {
struct skynet_remote_message *new_queue = malloc(sizeof(struct skynet_remote_message) * q->cap * 2);
int i;
for (i=0;i<q->cap;i++) {
new_queue[i] = q->queue[(q->head + i) % q->cap];
}
q->head = 0;
q->tail = q->cap;
q->cap *= 2;
free(q->queue);
q->queue = new_queue;
}
UNLOCK(q)
skynet_mq_force_push(struct message_queue * queue) {
assert(queue->in_global);
skynet_globalmq_push(queue);
}

View File

@@ -13,7 +13,6 @@ struct skynet_message {
struct message_queue;
void skynet_globalmq_push(struct message_queue *);
struct message_queue * skynet_globalmq_pop(void);
struct message_queue * skynet_mq_create(uint32_t handle);
@@ -23,19 +22,7 @@ uint32_t skynet_mq_handle(struct message_queue *);
// 0 for success
int skynet_mq_pop(struct message_queue *q, struct skynet_message *message);
void skynet_mq_push(struct message_queue *q, struct skynet_message *message);
struct skynet_remote_message {
uint32_t destination;
struct skynet_message message;
};
struct message_remote_queue;
struct message_remote_queue * skynet_remotemq_create(void);
void skynet_remotemq_release(struct message_remote_queue *);
int skynet_remotemq_pop(struct message_remote_queue *q, struct skynet_remote_message *message);
void skynet_remotemq_push(struct message_remote_queue *q, struct skynet_remote_message *message);
void skynet_mq_force_push(struct message_queue *q);
void skynet_mq_init(int cap);

View File

@@ -25,10 +25,9 @@ struct skynet_context {
void * cb_ud;
skynet_cb cb;
int session_id;
int in_global_queue;
int init;
uint32_t forward;
char * forward_address;
char forward_address[GLOBALNAME_LENGTH];
struct message_queue *queue;
};
@@ -36,10 +35,11 @@ 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] = hex[(id >> ((7-i) * 4))&0xf];
str[i+1] = hex[(id >> ((7-i) * 4))&0xf];
}
str[8] = '\0';
str[9] = '\0';
}
struct skynet_context *
@@ -58,17 +58,13 @@ skynet_context_new(const char * name, const char *param) {
ctx->ref = 2;
ctx->cb = NULL;
ctx->cb_ud = NULL;
// a trick, global loop can't be dispatch in init process
ctx->in_global_queue = 1;
ctx->forward = 0;
ctx->forward_address = NULL;
ctx->session_id = 0;
ctx->forward_address[0] = '\0';
ctx->init = 0;
ctx->handle = skynet_handle_register(ctx);
char * uid = ctx->handle_name;
uid[0] = ':';
_id_to_hex(uid+1, ctx->handle);
_id_to_hex(uid, ctx->handle);
struct message_queue * queue = ctx->queue = skynet_mq_create(ctx->handle);
// init function maybe use ctx->handle, so it must init at last
@@ -77,15 +73,10 @@ skynet_context_new(const char * name, const char *param) {
struct skynet_context * ret = skynet_context_release(ctx);
if (ret) {
ctx->init = 1;
skynet_globalmq_push(ctx->queue);
return ret;
} else {
// because of ctx->in_global_queue == 1 , so we should release queue here.
skynet_mq_release(queue);
}
return NULL;
skynet_mq_force_push(queue);
return ret;
} else {
ctx->in_global_queue = 0;
skynet_context_release(ctx);
skynet_handle_retire(ctx->handle);
return NULL;
@@ -95,7 +86,7 @@ skynet_context_new(const char * name, const char *param) {
int
skynet_context_newsession(struct skynet_context *ctx) {
int session = ++ctx->session_id;
if (session >= 0x7fffffff) {
if (session >= SESSION_CLIENT) {
ctx->session_id = 1;
return 1;
}
@@ -111,9 +102,7 @@ skynet_context_grab(struct skynet_context *ctx) {
static void
_delete_context(struct skynet_context *ctx) {
skynet_module_instance_release(ctx->mod, ctx->instance);
if (!ctx->in_global_queue) {
skynet_mq_release(ctx->queue);
}
skynet_mq_push(ctx->queue , NULL);
free(ctx);
}
@@ -126,13 +115,29 @@ skynet_context_release(struct skynet_context *ctx) {
return ctx;
}
int
skynet_context_push(uint32_t handle, struct skynet_message *message) {
struct skynet_context * ctx = skynet_handle_grab(handle);
if (ctx == NULL) {
return -1;
}
skynet_mq_push(ctx->queue, message);
skynet_context_release(ctx);
return 0;
}
static int
_forwarding(struct skynet_context *ctx, struct skynet_message *msg) {
if (ctx->forward) {
uint32_t des = ctx->forward;
ctx->forward = 0;
if (skynet_harbor_message_isremote(des)) {
skynet_harbor_send(NULL, des, msg);
struct remote_message * rmsg = malloc(sizeof(*rmsg));
rmsg->destination.handle = des;
rmsg->message = msg->data;
rmsg->sz = msg->sz;
skynet_harbor_send(rmsg, msg->source, msg->session);
} else {
if (skynet_context_push(des, msg)) {
free(msg->data);
@@ -142,10 +147,14 @@ _forwarding(struct skynet_context *ctx, struct skynet_message *msg) {
}
return 1;
}
if (ctx->forward_address) {
skynet_harbor_send(ctx->forward_address, 0, msg);
free(ctx->forward_address);
ctx->forward_address = NULL;
if (ctx->forward_address[0]) {
struct remote_message * rmsg = malloc(sizeof(*rmsg));
memcpy(rmsg->destination.name, ctx->forward_address,GLOBALNAME_LENGTH);
rmsg->destination.handle = 0;
rmsg->message = msg->data;
rmsg->sz = msg->sz;
skynet_harbor_send(rmsg, msg->source, msg->session);
ctx->forward_address[0] = '\0';
return 1;
}
return 0;
@@ -158,21 +167,10 @@ _dispatch_message(struct skynet_context *ctx, struct skynet_message *msg) {
ctx->cb(ctx, ctx->cb_ud, msg->session, NULL, msg->data, msg->sz);
} else {
char tmp[10];
tmp[0] = ':';
int not_delete;
_id_to_hex(tmp+1, msg->source);
if (skynet_harbor_message_isremote(msg->source)) {
void * data = skynet_harbor_message_open(msg);
ctx->cb(ctx, ctx->cb_ud, msg->session, tmp, data, msg->sz);
not_delete = _forwarding(ctx, msg);
if (!not_delete) {
skynet_harbor_message_close(msg);
}
} else {
ctx->cb(ctx, ctx->cb_ud, msg->session, tmp, msg->data, msg->sz);
not_delete = _forwarding(ctx, msg);
}
if (!not_delete) {
_id_to_hex(tmp, msg->source);
int reserve = ctx->cb(ctx, ctx->cb_ud, msg->session, tmp, msg->data, msg->sz);
reserve |= _forwarding(ctx, msg);
if (!reserve) {
free(msg->data);
}
}
@@ -185,9 +183,6 @@ _drop_queue(struct message_queue *q) {
int s = 0;
while(!skynet_mq_pop(q, &msg)) {
++s;
if (skynet_harbor_message_isremote(msg.source)) {
skynet_harbor_message_close(&msg);
}
free(msg.data);
}
skynet_mq_release(q);
@@ -211,19 +206,13 @@ skynet_context_message_dispatch(void) {
return 0;
}
assert(ctx->in_global_queue);
struct skynet_message msg;
if (skynet_mq_pop(q,&msg)) {
__sync_lock_release(&ctx->in_global_queue);
skynet_context_release(ctx);
return 0;
}
if (ctx->cb == NULL) {
if (skynet_harbor_message_isremote(msg.source)) {
skynet_harbor_message_close(&msg);
}
free(msg.data);
skynet_error(NULL, "Drop message from %x to %x without callback , size = %d",msg.source, handle, (int)msg.sz);
} else {
@@ -231,12 +220,23 @@ skynet_context_message_dispatch(void) {
}
assert(q == ctx->queue);
skynet_globalmq_push(q);
skynet_mq_force_push(q);
skynet_context_release(ctx);
return 0;
}
static void
_copy_name(char name[GLOBALNAME_LENGTH], const char * addr) {
int i;
for (i=0;i<GLOBALNAME_LENGTH && addr[i];i++) {
name[i] = addr[i];
}
for (;i<GLOBALNAME_LENGTH;i++) {
name[i] = '\0';
}
}
const char *
skynet_command(struct skynet_context * context, const char * cmd , const char * param) {
if (strcmp(cmd,"TIMEOUT") == 0) {
@@ -263,14 +263,10 @@ skynet_command(struct skynet_context * context, const char * cmd , const char *
return skynet_handle_namehandle(context->handle, param + 1);
} else {
assert(context->handle!=0);
int i;
for (i=0;i<param[i];i++) {
if (!(param[i] >= '0' && param[i] <= '9')) {
break;
}
}
assert(param[i]);
skynet_harbor_register(param, context->handle);
struct remote_name *rname = malloc(sizeof(*rname));
_copy_name(rname->name, param);
rname->handle = context->handle;
skynet_harbor_register(rname);
return NULL;
}
}
@@ -290,7 +286,10 @@ skynet_command(struct skynet_context * context, const char * cmd , const char *
if (name[0] == '.') {
return skynet_handle_namehandle(handle_id, name + 1);
} else {
skynet_harbor_register(name, handle_id);
struct remote_name *rname = malloc(sizeof(*rname));
_copy_name(rname->name, name);
rname->handle = handle_id;
skynet_harbor_register(rname);
}
return NULL;
}
@@ -327,8 +326,7 @@ skynet_command(struct skynet_context * context, const char * cmd , const char *
fprintf(stderr, "Launch %s %s failed\n",mod,args);
return NULL;
} else {
context->result[0] = ':';
_id_to_hex(context->result+1, inst->handle);
_id_to_hex(context->result, inst->handle);
return context->result;
}
}
@@ -366,7 +364,7 @@ skynet_command(struct skynet_context * context, const char * cmd , const char *
void
skynet_forward(struct skynet_context * context, const char * addr) {
uint32_t des = 0;
assert(context->forward == 0 && context->forward_address == NULL);
assert(context->forward == 0 && context->forward_address[0] == '\0');
if (addr[0] == ':') {
des = strtol(addr+1, NULL, 16);
assert(des != 0);
@@ -377,7 +375,7 @@ skynet_forward(struct skynet_context * context, const char * addr) {
return;
}
} else {
context->forward_address = strdup(addr);
_copy_name(context->forward_address,addr);
return;
}
context->forward = des;
@@ -385,9 +383,14 @@ skynet_forward(struct skynet_context * context, const char * addr) {
int
skynet_send(struct skynet_context * context, const char * source, const char * addr , int session, void * data, size_t sz, int flags) {
int session_id = session;
uint32_t source_handle;
if (source == NULL) {
source_handle = context->handle;
if (session < 0) {
session = skynet_context_newsession(context);
session_id = - session;
}
} else {
assert (source[0] == ':');
source_handle = strtoul(source+1, NULL, 16);
@@ -401,11 +404,6 @@ skynet_send(struct skynet_context * context, const char * source, const char * a
memcpy(msg, data, sz);
msg[sz] = '\0';
}
int session_id = session;
if (session < 0) {
session = skynet_context_newsession(context);
session_id = - session;
}
if (addr == NULL) {
return session;
}
@@ -420,29 +418,35 @@ skynet_send(struct skynet_context * context, const char * source, const char * a
return session;
}
} else {
struct skynet_message smsg;
smsg.source = source_handle;
smsg.session = session_id;
smsg.data = msg;
smsg.sz = sz;
skynet_harbor_send(addr, 0, &smsg);
struct remote_message * rmsg = malloc(sizeof(*rmsg));
_copy_name(rmsg->destination.name, addr);
rmsg->destination.handle = 0;
rmsg->message = msg;
rmsg->sz = sz;
skynet_harbor_send(rmsg, source_handle, session_id);
return session;
}
assert(des > 0);
struct skynet_message smsg;
smsg.source = source_handle;
smsg.session = session_id;
smsg.data = msg;
smsg.sz = sz;
if (skynet_harbor_message_isremote(des)) {
skynet_harbor_send(NULL, des, &smsg);
} else if (skynet_context_push(des, &smsg)) {
free(msg);
skynet_error(NULL, "Drop message from %x to %s (size=%d)", smsg.source, addr, (int)sz);
return -1;
struct remote_message * rmsg = malloc(sizeof(*rmsg));
rmsg->destination.handle = des;
rmsg->message = msg;
rmsg->sz = sz;
skynet_harbor_send(rmsg, source_handle, session_id);
} else {
struct skynet_message smsg;
smsg.source = source_handle;
smsg.session = session_id;
smsg.data = msg;
smsg.sz = sz;
if (skynet_context_push(des, &smsg)) {
free(msg);
skynet_error(NULL, "Drop message from %x to %s (size=%d)", smsg.source, addr, (int)sz);
return -1;
}
}
return session;
}
@@ -464,17 +468,13 @@ skynet_callback(struct skynet_context * context, void *ud, skynet_cb cb) {
context->cb_ud = ud;
}
int
skynet_context_push(uint32_t handle, struct skynet_message *message) {
struct skynet_context * ctx = skynet_handle_grab(handle);
if (ctx == NULL) {
return -1;
}
skynet_mq_push(ctx->queue, message);
if (__sync_lock_test_and_set(&ctx->in_global_queue, 1) == 0) {
skynet_globalmq_push(ctx->queue);
}
skynet_context_release(ctx);
void
skynet_context_send(struct skynet_context * ctx, void * msg, size_t sz, uint32_t source, int session) {
struct skynet_message smsg;
smsg.source = source;
smsg.session = session;
smsg.data = msg;
smsg.sz = sz;
return 0;
skynet_mq_push(ctx->queue, &smsg);
}

View File

@@ -2,6 +2,7 @@
#define SKYNET_SERVER_H
#include <stdint.h>
#include <stdlib.h>
struct skynet_context;
struct skynet_message;
@@ -12,6 +13,7 @@ struct skynet_context * skynet_context_release(struct skynet_context *);
uint32_t skynet_context_handle(struct skynet_context *);
void skynet_context_init(struct skynet_context *, uint32_t handle);
int skynet_context_push(uint32_t handle, struct skynet_message *message);
void skynet_context_send(struct skynet_context * context, void * msg, size_t sz, uint32_t source, int session);
int skynet_context_newsession(struct skynet_context *);
int skynet_context_message_dispatch(void); // return 1 when block

View File

@@ -5,12 +5,10 @@
#include "skynet_module.h"
#include "skynet_timer.h"
#include "skynet_harbor.h"
#include "skynet_master.h"
#include <pthread.h>
#include <unistd.h>
#include <assert.h>
#include <zmq.h>
static void *
_timer(void *p) {
@@ -33,59 +31,51 @@ _worker(void *p) {
static void
_start(int thread) {
pthread_t pid[thread+2];
pthread_t pid[thread+1];
pthread_create(&pid[0], NULL, _timer, NULL);
pthread_create(&pid[1], NULL, skynet_harbor_dispatch_thread, NULL);
int i;
for (i=2;i<thread+2;i++) {
for (i=1;i<thread+1;i++) {
pthread_create(&pid[i], NULL, _worker, NULL);
}
for (i=0;i<thread+2;i++) {
for (i=0;i<thread+1;i++) {
pthread_join(pid[i], NULL);
}
}
struct master_arg {
void * context;
const char * port;
};
static void *
_master_thread(void *ud) {
struct master_arg * args = ud;
skynet_master(args->context, args->port);
free(args);
return NULL;
}
static void
_start_master(void * context, const char * port) {
pthread_t pid;
struct master_arg * args = malloc(sizeof(*args));
args->context = context;
args->port = port;
pthread_create(&pid, NULL, _master_thread, args);
static int
_start_master(const char * master) {
struct skynet_context *ctx = skynet_context_new("master", master);
if (ctx == NULL)
return 1;
return 0;
}
void
skynet_start(struct skynet_config * config) {
void *context = zmq_init (1);
assert(context);
if (config->standalone) {
_start_master(context, config->master);
}
// harbor must be init first
skynet_harbor_init(context, config->master , config->local, config->harbor);
skynet_harbor_init(config->harbor);
skynet_handle_init(config->harbor);
skynet_mq_init(config->mqueue_size);
skynet_module_init(config->module_path);
skynet_timer_init();
if (config->standalone) {
if (_start_master(config->standalone)) {
return;
}
}
// harbor must be init first
if (skynet_harbor_start(config->master , config->local)) {
return;
}
struct skynet_context *ctx;
ctx = skynet_context_new("logger", config->logger);
assert(ctx);
if (ctx == NULL) {
return;
}
ctx = skynet_context_new("snlua", config->start);
_start(config->thread);