mirror of
https://github.com/cloudwu/skynet.git
synced 2026-07-25 12:43:09 +00:00
skynet.blockcall concurrent bugfix
This commit is contained in:
@@ -202,8 +202,8 @@ end
|
|||||||
|
|
||||||
function skynet.blockcall(addr, typename , ...)
|
function skynet.blockcall(addr, typename , ...)
|
||||||
local p = proto[typename]
|
local p = proto[typename]
|
||||||
|
c.command("LOCK")
|
||||||
local session = c.send(addr, p.id , nil , p.pack(...))
|
local session = c.send(addr, p.id , nil , p.pack(...))
|
||||||
c.command("LOCK",tostring(session))
|
|
||||||
return p.unpack(coroutine.yield("CALL", session))
|
return p.unpack(coroutine.yield("CALL", session))
|
||||||
end
|
end
|
||||||
|
|
||||||
|
|||||||
@@ -12,6 +12,9 @@
|
|||||||
#define DEFAULT_QUEUE_SIZE 64;
|
#define DEFAULT_QUEUE_SIZE 64;
|
||||||
#define MAX_GLOBAL_MQ 0x10000
|
#define MAX_GLOBAL_MQ 0x10000
|
||||||
|
|
||||||
|
#define MQ_IN_GLOBAL 1
|
||||||
|
#define MQ_LOCKED 2
|
||||||
|
|
||||||
struct message_queue {
|
struct message_queue {
|
||||||
uint32_t handle;
|
uint32_t handle;
|
||||||
int cap;
|
int cap;
|
||||||
@@ -81,7 +84,7 @@ skynet_mq_create(uint32_t handle) {
|
|||||||
q->head = 0;
|
q->head = 0;
|
||||||
q->tail = 0;
|
q->tail = 0;
|
||||||
q->lock = 0;
|
q->lock = 0;
|
||||||
q->in_global = 1;
|
q->in_global = MQ_IN_GLOBAL;
|
||||||
q->release = 0;
|
q->release = 0;
|
||||||
q->lock_session = 0;
|
q->lock_session = 0;
|
||||||
q->queue = malloc(sizeof(struct skynet_message) * q->cap);
|
q->queue = malloc(sizeof(struct skynet_message) * q->cap);
|
||||||
@@ -153,9 +156,14 @@ _pushhead(struct message_queue *q, struct skynet_message *message) {
|
|||||||
q->queue[head] = *message;
|
q->queue[head] = *message;
|
||||||
q->head = head;
|
q->head = head;
|
||||||
|
|
||||||
// this api use in push a unlock message, so the in_global flags must be 1 , but the q is not exist in global queue.
|
// this api use in push a unlock message, so the in_global flags must not be 0 ,
|
||||||
assert(q->in_global);
|
// but the q is not exist in global queue.
|
||||||
skynet_globalmq_push(q);
|
if (q->in_global == MQ_IN_GLOBAL) {
|
||||||
|
skynet_globalmq_push(q);
|
||||||
|
} else {
|
||||||
|
assert(q->in_global == MQ_LOCKED);
|
||||||
|
}
|
||||||
|
q->lock_session = 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
void
|
void
|
||||||
@@ -165,7 +173,6 @@ skynet_mq_push(struct message_queue *q, struct skynet_message *message) {
|
|||||||
|
|
||||||
if (q->lock_session !=0 && message->session == q->lock_session) {
|
if (q->lock_session !=0 && message->session == q->lock_session) {
|
||||||
_pushhead(q,message);
|
_pushhead(q,message);
|
||||||
q->lock_session = 0;
|
|
||||||
} else {
|
} else {
|
||||||
q->queue[q->tail] = *message;
|
q->queue[q->tail] = *message;
|
||||||
if (++ q->tail >= q->cap) {
|
if (++ q->tail >= q->cap) {
|
||||||
@@ -178,7 +185,7 @@ skynet_mq_push(struct message_queue *q, struct skynet_message *message) {
|
|||||||
|
|
||||||
if (q->lock_session == 0) {
|
if (q->lock_session == 0) {
|
||||||
if (q->in_global == 0) {
|
if (q->in_global == 0) {
|
||||||
q->in_global = 1;
|
q->in_global = MQ_IN_GLOBAL;
|
||||||
skynet_globalmq_push(q);
|
skynet_globalmq_push(q);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -191,6 +198,8 @@ void
|
|||||||
skynet_mq_lock(struct message_queue *q, int session) {
|
skynet_mq_lock(struct message_queue *q, int session) {
|
||||||
LOCK(q)
|
LOCK(q)
|
||||||
assert(q->lock_session == 0);
|
assert(q->lock_session == 0);
|
||||||
|
assert(q->in_global == MQ_IN_GLOBAL);
|
||||||
|
q->in_global = MQ_LOCKED;
|
||||||
q->lock_session = session;
|
q->lock_session = session;
|
||||||
UNLOCK(q)
|
UNLOCK(q)
|
||||||
}
|
}
|
||||||
@@ -214,8 +223,19 @@ skynet_mq_force_push(struct message_queue * queue) {
|
|||||||
void
|
void
|
||||||
skynet_mq_pushglobal(struct message_queue *queue) {
|
skynet_mq_pushglobal(struct message_queue *queue) {
|
||||||
assert(queue->in_global);
|
assert(queue->in_global);
|
||||||
if (queue->lock_session == 0) {
|
if (queue->in_global == MQ_LOCKED) {
|
||||||
skynet_globalmq_push(queue);
|
// lock message queue just now, unlock it.
|
||||||
|
LOCK(queue)
|
||||||
|
queue->in_global = MQ_IN_GLOBAL;
|
||||||
|
if (queue->lock_session == 0) {
|
||||||
|
skynet_globalmq_push(queue);
|
||||||
|
}
|
||||||
|
UNLOCK(queue)
|
||||||
|
} else {
|
||||||
|
queue->in_global = MQ_IN_GLOBAL;
|
||||||
|
if (queue->lock_session == 0) {
|
||||||
|
skynet_globalmq_push(queue);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -364,9 +364,7 @@ skynet_command(struct skynet_context * context, const char * cmd , const char *
|
|||||||
if (context->init == false) {
|
if (context->init == false) {
|
||||||
return NULL;
|
return NULL;
|
||||||
}
|
}
|
||||||
int session = strtol(param, NULL, 10);
|
skynet_mq_lock(context->queue, context->session_id+1);
|
||||||
assert(session);
|
|
||||||
skynet_mq_lock(context->queue, session);
|
|
||||||
return NULL;
|
return NULL;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user