mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-24 20:23:06 +00:00
@@ -1,4 +1,10 @@
|
||||
v0.9.0 (2014-11-17)
|
||||
v0.9.2 (2014-12-8)
|
||||
-----------
|
||||
* Simplify the message queue
|
||||
* Add create_index in mongo driver
|
||||
* Fix a bug in big-endian architecture (sproto)
|
||||
|
||||
v0.9.0 / v0.9.1 (2014-11-17)
|
||||
-----------
|
||||
* Add UDP support
|
||||
* Add IPv6 support
|
||||
|
||||
@@ -9,6 +9,7 @@
|
||||
#include <lua.h>
|
||||
#include <lauxlib.h>
|
||||
|
||||
#include <sys/socket.h>
|
||||
#include <arpa/inet.h>
|
||||
|
||||
#include "skynet_socket.h"
|
||||
|
||||
@@ -252,12 +252,12 @@ decode(void *ud, const char *tagname, int type, int index, struct sproto_type *s
|
||||
switch (type) {
|
||||
case SPROTO_TINTEGER: {
|
||||
// notice: in lua 5.2, 52bit integer support (not 64)
|
||||
lua_Integer v = *(lua_Integer *)value;
|
||||
lua_Integer v = *(uint64_t*)value;
|
||||
lua_pushinteger(L, v);
|
||||
break;
|
||||
}
|
||||
case SPROTO_TBOOLEAN: {
|
||||
int v = *(lua_Integer*)value;
|
||||
int v = *(uint64_t*)value;
|
||||
lua_pushboolean(L,v);
|
||||
break;
|
||||
}
|
||||
|
||||
@@ -971,7 +971,7 @@ decode_array(sproto_callback cb, void *ud, struct field *f, uint8_t * stream) {
|
||||
}
|
||||
case SPROTO_TBOOLEAN:
|
||||
for (i=0;i<sz;i++) {
|
||||
int value = stream[i];
|
||||
uint64_t value = stream[i];
|
||||
cb(ud, f->name, SPROTO_TBOOLEAN, i+1, NULL, &value, sizeof(value));
|
||||
}
|
||||
break;
|
||||
|
||||
@@ -229,6 +229,10 @@ function mongo_collection:insert(doc)
|
||||
sock:request(pack)
|
||||
end
|
||||
|
||||
function mongo_collection:safe_insert(doc)
|
||||
return self.database:runCommand("insert", self.name, "documents", {bson_encode(doc)})
|
||||
end
|
||||
|
||||
function mongo_collection:batch_insert(docs)
|
||||
for i=1,#docs do
|
||||
if docs[i]._id == nil then
|
||||
@@ -278,6 +282,48 @@ function mongo_collection:find(query, selector)
|
||||
} , cursor_meta)
|
||||
end
|
||||
|
||||
-- collection:createIndex({username = 1}, {unique = true})
|
||||
function mongo_collection:createIndex(keys, option)
|
||||
local name
|
||||
for k, v in pairs(keys) do
|
||||
assert(v == 1)
|
||||
name = (name == nil) and k or (name .. "_" .. k)
|
||||
end
|
||||
|
||||
local doc = {};
|
||||
doc.name = name
|
||||
doc.key = keys
|
||||
for k, v in pairs(option) do
|
||||
if v then
|
||||
doc[k] = true
|
||||
end
|
||||
end
|
||||
return self.database:runCommand("createIndexes", self.name, "indexes", {doc})
|
||||
end
|
||||
|
||||
mongo_collection.ensureIndex = mongo_collection.createIndex;
|
||||
|
||||
-- collection:findAndModify({query = {name = "userid"}, update = {["$inc"] = {nextid = 1}}, })
|
||||
-- keys, value type
|
||||
-- query, table
|
||||
-- sort, table
|
||||
-- remove, bool
|
||||
-- update, table
|
||||
-- new, bool
|
||||
-- fields, bool
|
||||
-- upsert, boolean
|
||||
function mongo_collection:findAndModify(doc)
|
||||
assert(doc.query)
|
||||
assert(doc.update or doc.remove)
|
||||
|
||||
local cmd = {"findAndModify", self.name};
|
||||
for k, v in pairs(doc) do
|
||||
table.insert(cmd, k)
|
||||
table.insert(cmd, v)
|
||||
end
|
||||
return self.database:runCommand(unpack(cmd))
|
||||
end
|
||||
|
||||
function mongo_cursor:hasNext()
|
||||
if self.__ptr == nil then
|
||||
if self.__document == nil then
|
||||
|
||||
@@ -32,12 +32,9 @@ struct message_queue {
|
||||
};
|
||||
|
||||
struct global_queue {
|
||||
uint32_t head;
|
||||
uint32_t tail;
|
||||
struct message_queue ** queue;
|
||||
// We use a separated flag array to ensure the mq is pushed.
|
||||
// See the comments below.
|
||||
struct message_queue *list;
|
||||
struct message_queue *head;
|
||||
struct message_queue *tail;
|
||||
int lock;
|
||||
};
|
||||
|
||||
static struct global_queue *Q = NULL;
|
||||
@@ -51,57 +48,32 @@ void
|
||||
skynet_globalmq_push(struct message_queue * queue) {
|
||||
struct global_queue *q= Q;
|
||||
|
||||
uint32_t tail = GP(__sync_fetch_and_add(&q->tail,1));
|
||||
|
||||
// only one thread can set the slot (change q->queue[tail] from NULL to queue)
|
||||
if (!__sync_bool_compare_and_swap(&q->queue[tail], NULL, queue)) {
|
||||
// The queue may full seldom, save queue in list
|
||||
LOCK(q)
|
||||
assert(queue->next == NULL);
|
||||
struct message_queue * last;
|
||||
do {
|
||||
last = q->list;
|
||||
queue->next = last;
|
||||
} while(!__sync_bool_compare_and_swap(&q->list, last, queue));
|
||||
|
||||
return;
|
||||
if(q->tail) {
|
||||
q->tail->next = queue;
|
||||
q->tail = queue;
|
||||
} else {
|
||||
q->head = q->tail = queue;
|
||||
}
|
||||
UNLOCK(q)
|
||||
}
|
||||
|
||||
struct message_queue *
|
||||
skynet_globalmq_pop() {
|
||||
struct global_queue *q = Q;
|
||||
uint32_t head = q->head;
|
||||
|
||||
if (head == q->tail) {
|
||||
// The queue is empty.
|
||||
return NULL;
|
||||
LOCK(q)
|
||||
struct message_queue *mq = q->head;
|
||||
if(mq) {
|
||||
q->head = mq->next;
|
||||
if(q->head == NULL) {
|
||||
assert(mq == q->tail);
|
||||
q->tail = NULL;
|
||||
}
|
||||
|
||||
uint32_t head_ptr = GP(head);
|
||||
|
||||
struct message_queue * list = q->list;
|
||||
if (list) {
|
||||
// If q->list is not empty, try to load it back to the queue
|
||||
struct message_queue *newhead = list->next;
|
||||
if (__sync_bool_compare_and_swap(&q->list, list, newhead)) {
|
||||
// try load list only once, if success , push it back to the queue.
|
||||
list->next = NULL;
|
||||
skynet_globalmq_push(list);
|
||||
}
|
||||
}
|
||||
|
||||
struct message_queue * mq = q->queue[head_ptr];
|
||||
if (mq == NULL) {
|
||||
// globalmq push not complete
|
||||
return NULL;
|
||||
}
|
||||
if (!__sync_bool_compare_and_swap(&q->head, head, head+1)) {
|
||||
return NULL;
|
||||
}
|
||||
// only one thread can get the slot (change q->queue[head_ptr] to NULL)
|
||||
if (!__sync_bool_compare_and_swap(&q->queue[head_ptr], mq, NULL)) {
|
||||
return NULL;
|
||||
mq->next = NULL;
|
||||
}
|
||||
UNLOCK(q)
|
||||
|
||||
return mq;
|
||||
}
|
||||
@@ -243,8 +215,6 @@ void
|
||||
skynet_mq_init() {
|
||||
struct global_queue *q = skynet_malloc(sizeof(*q));
|
||||
memset(q,0,sizeof(*q));
|
||||
q->queue = skynet_malloc(MAX_GLOBAL_MQ * sizeof(struct message_queue *));
|
||||
memset(q->queue, 0, sizeof(struct message_queue *) * MAX_GLOBAL_MQ);
|
||||
Q=q;
|
||||
}
|
||||
|
||||
|
||||
@@ -628,7 +628,7 @@ append_sendbuffer_(struct socket_server *ss, struct wb_list *s, struct request_s
|
||||
struct write_buffer * buf = MALLOC(size);
|
||||
struct send_object so;
|
||||
buf->userobject = send_object_init(ss, &so, request->buffer, request->sz);
|
||||
buf->ptr = so.buffer+n;
|
||||
buf->ptr = (char*)so.buffer+n;
|
||||
buf->sz = so.sz - n;
|
||||
buf->buffer = request->buffer;
|
||||
buf->next = NULL;
|
||||
|
||||
10
test/testdeadloop.lua
Normal file
10
test/testdeadloop.lua
Normal file
@@ -0,0 +1,10 @@
|
||||
local skynet = require "skynet"
|
||||
local function dead_loop()
|
||||
while true do
|
||||
skynet.sleep(0)
|
||||
end
|
||||
end
|
||||
|
||||
skynet.start(function()
|
||||
skynet.fork(dead_loop)
|
||||
end)
|
||||
Reference in New Issue
Block a user