From 84cb4bb4c18b55d3a3ca1d8b4269f40fb9150477 Mon Sep 17 00:00:00 2001 From: Cloud Wu Date: Tue, 25 Nov 2014 16:04:26 +0800 Subject: [PATCH 1/7] use explicit conversion --- skynet-src/socket_server.c | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/skynet-src/socket_server.c b/skynet-src/socket_server.c index 4010055e..5622e4a8 100644 --- a/skynet-src/socket_server.c +++ b/skynet-src/socket_server.c @@ -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; From 2b13eb250dacb6ac4c9cb88e92b6b6d2778f0b5c Mon Sep 17 00:00:00 2001 From: Cloud Wu Date: Wed, 26 Nov 2014 21:18:49 +0800 Subject: [PATCH 2/7] update sproto for big-endian --- lualib-src/sproto/lsproto.c | 10 +++++----- lualib-src/sproto/sproto.c | 20 ++++++++++---------- 2 files changed, 15 insertions(+), 15 deletions(-) diff --git a/lualib-src/sproto/lsproto.c b/lualib-src/sproto/lsproto.c index af86e389..e5378713 100644 --- a/lualib-src/sproto/lsproto.c +++ b/lualib-src/sproto/lsproto.c @@ -87,7 +87,7 @@ struct encode_ud { int deep; }; -static int +static int encode(void *ud, const char *tagname, int type, int index, struct sproto_type *st, void *value, int length) { struct encode_ud *self = ud; lua_State *L = self->L; @@ -189,7 +189,7 @@ expand_buffer(lua_State *L, int osz, int nsz) { lightuserdata sproto_type table source - return string + return string */ static int lencode(lua_State *L) { @@ -229,7 +229,7 @@ struct decode_ud { int deep; }; -static int +static int decode(void *ud, const char *tagname, int type, int index, struct sproto_type *st, void *value, int length) { struct decode_ud * self = ud; lua_State *L = self->L; @@ -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; } diff --git a/lualib-src/sproto/sproto.c b/lualib-src/sproto/sproto.c index 24d36492..67e0d0d0 100644 --- a/lualib-src/sproto/sproto.c +++ b/lualib-src/sproto/sproto.c @@ -500,7 +500,7 @@ sproto_dump(struct sproto *s) { } // query -int +int sproto_prototag(struct sproto *sp, const char * name) { int i; for (i=0;iprotocol_n;i++) { @@ -529,7 +529,7 @@ query_proto(struct sproto *sp, int tag) { return NULL; } -struct sproto_type * +struct sproto_type * sproto_protoquery(struct sproto *sp, int proto, int what) { struct protocol * p; if (what <0 || what >1) { @@ -542,7 +542,7 @@ sproto_protoquery(struct sproto *sp, int proto, int what) { return NULL; } -const char * +const char * sproto_protoname(struct sproto *sp, int proto) { struct protocol * p = query_proto(sp, proto); if (p) { @@ -551,7 +551,7 @@ sproto_protoname(struct sproto *sp, int proto) { return NULL; } -struct sproto_type * +struct sproto_type * sproto_type(struct sproto *sp, const char * type_name) { int i; for (i=0;itype_n;i++) { @@ -806,7 +806,7 @@ encode_array(sproto_callback cb, void *ud, struct field *f, uint8_t *data, int s return fill_size(data, sz); } -int +int sproto_encode(struct sproto_type *st, void * buffer, int size, sproto_callback cb, void *ud) { uint8_t * header = buffer; uint8_t * data; @@ -830,7 +830,7 @@ sproto_encode(struct sproto_type *st, void * buffer, int size, sproto_callback c sz = encode_array(cb,ud, f, data, size); } else { switch(type) { - case SPROTO_TINTEGER: + case SPROTO_TINTEGER: case SPROTO_TBOOLEAN: { union { uint64_t u64; @@ -971,7 +971,7 @@ decode_array(sproto_callback cb, void *ud, struct field *f, uint8_t * stream) { } case SPROTO_TBOOLEAN: for (i=0;iname, SPROTO_TBOOLEAN, i+1, NULL, &value, sizeof(value)); } break; @@ -1050,7 +1050,7 @@ sproto_decode(struct sproto_type *st, const void * data, int size, sproto_callba } break; } - case SPROTO_TSTRING: + case SPROTO_TSTRING: case SPROTO_TSTRUCT: { uint32_t sz = todword(currentdata); if (cb(ud, f->name, f->type, 0, f->st, currentdata+SIZEOF_LENGTH, sz)) @@ -1124,7 +1124,7 @@ write_ff(const uint8_t * src, uint8_t * des, int n) { } } -int +int sproto_pack(const void * srcv, int srcsz, void * bufferv, int bufsz) { uint8_t tmp[8]; int i; @@ -1181,7 +1181,7 @@ sproto_pack(const void * srcv, int srcsz, void * bufferv, int bufsz) { return size; } -int +int sproto_unpack(const void * srcv, int srcsz, void * bufferv, int bufsz) { const uint8_t * src = srcv; uint8_t * buffer = bufferv; From 6f6039c13601d244e12a2190ed6200d248838879 Mon Sep 17 00:00:00 2001 From: Cloud Wu Date: Wed, 3 Dec 2014 11:06:09 +0800 Subject: [PATCH 3/7] include sys/socket.h for freebsd --- lualib-src/lua-socket.c | 1 + 1 file changed, 1 insertion(+) diff --git a/lualib-src/lua-socket.c b/lualib-src/lua-socket.c index df81ce17..eabb8620 100644 --- a/lualib-src/lua-socket.c +++ b/lualib-src/lua-socket.c @@ -9,6 +9,7 @@ #include #include +#include #include #include "skynet_socket.h" From a0d2c7172cf0bd9bb71037a9cb2a76be978ac315 Mon Sep 17 00:00:00 2001 From: xjdrew Date: Thu, 4 Dec 2014 15:50:19 +0800 Subject: [PATCH 4/7] =?UTF-8?q?=E4=BD=BF=E7=94=A8spinlock=EF=BC=8C?= =?UTF-8?q?=E9=81=BF=E5=85=8Dcas=E4=B8=8D=E8=83=BD=E5=85=85=E5=88=86?= =?UTF-8?q?=E8=B0=83=E5=BA=A6testdeadloop=E8=BF=99=E6=A0=B7=E7=9A=84?= =?UTF-8?q?=E7=94=A8=E4=BE=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- skynet-src/skynet_mq.c | 69 ++++++++++++------------------------------ test/testdeadloop.lua | 10 ++++++ 2 files changed, 30 insertions(+), 49 deletions(-) create mode 100644 test/testdeadloop.lua diff --git a/skynet-src/skynet_mq.c b/skynet-src/skynet_mq.c index 8cc8f9da..51f47ef1 100644 --- a/skynet-src/skynet_mq.c +++ b/skynet-src/skynet_mq.c @@ -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,33 @@ 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 - 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; + LOCK(q) + assert(queue->next == NULL); + 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; - } - - 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); + 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; } + mq->next = NULL; } - - 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; - } + UNLOCK(q) return mq; } @@ -243,8 +216,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; } diff --git a/test/testdeadloop.lua b/test/testdeadloop.lua new file mode 100644 index 00000000..215700e2 --- /dev/null +++ b/test/testdeadloop.lua @@ -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) From 2da2f121c8f8524984792775769f38c36b735aa0 Mon Sep 17 00:00:00 2001 From: xjdrew Date: Thu, 4 Dec 2014 16:16:17 +0800 Subject: [PATCH 5/7] delete unneccessary comment --- skynet-src/skynet_mq.c | 1 - 1 file changed, 1 deletion(-) diff --git a/skynet-src/skynet_mq.c b/skynet-src/skynet_mq.c index 51f47ef1..098cb210 100644 --- a/skynet-src/skynet_mq.c +++ b/skynet-src/skynet_mq.c @@ -48,7 +48,6 @@ void skynet_globalmq_push(struct message_queue * queue) { struct global_queue *q= Q; - // only one thread can set the slot (change q->queue[tail] from NULL to queue) LOCK(q) assert(queue->next == NULL); if(q->tail) { From fe9640e2dd6224b1f7257b08e806a5db3e019b5c Mon Sep 17 00:00:00 2001 From: dpull Date: Thu, 4 Dec 2014 20:10:24 +0800 Subject: [PATCH 6/7] =?UTF-8?q?1=E3=80=81mongo=5Fcollection:createIndex=20?= =?UTF-8?q?=E5=88=9B=E5=BB=BA=E7=B4=A2=E5=BC=95=202=E3=80=81mongo=5Fcollec?= =?UTF-8?q?tion:safe=5Finsert=20=E5=8F=AF=E4=BB=A5=E5=88=A4=E6=96=AD?= =?UTF-8?q?=E8=BF=94=E5=9B=9E=E5=80=BC=E7=9A=84insert=203=E3=80=81mongo=5F?= =?UTF-8?q?collection:findAndModify=20=E6=9F=A5=E8=AF=A2=E5=B9=B6=E4=BF=AE?= =?UTF-8?q?=E6=94=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- lualib/mongo.lua | 48 +++++++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 47 insertions(+), 1 deletion(-) diff --git a/lualib/mongo.lua b/lualib/mongo.lua index 458e1a6c..38ac20dd 100644 --- a/lualib/mongo.lua +++ b/lualib/mongo.lua @@ -227,7 +227,11 @@ function mongo_collection:insert(doc) sock:request(pack) end -function mongo_collection:batch_insert(docs) +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 docs[i]._id = bson.objectid() @@ -276,6 +280,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 From 176e4df90c4265ec85833eab5f4f6dd47bf08126 Mon Sep 17 00:00:00 2001 From: Cloud Wu Date: Mon, 8 Dec 2014 10:53:14 +0800 Subject: [PATCH 7/7] ready for v0.9.2 --- HISTORY.md | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/HISTORY.md b/HISTORY.md index e8835119..4fefeb98 100644 --- a/HISTORY.md +++ b/HISTORY.md @@ -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