diff --git a/src/modules/thread/Channel.cpp b/src/modules/thread/Channel.cpp index 4998f4392..268e8be9c 100644 --- a/src/modules/thread/Channel.cpp +++ b/src/modules/thread/Channel.cpp @@ -22,6 +22,30 @@ #include #include +namespace +{ + union uslong + { + unsigned long u; + long i; + }; + + // target <= current, but semi-wrapsafe, one wrap, anyway + inline bool past(unsigned int target, unsigned int current) + { + if (target > current) + return false; + if (target == current) + return true; + + uslong t, c; + t.u = target; + c.u = current; + + return !(t.i < 0 && c.i > 0); + } +} + namespace love { namespace thread @@ -41,14 +65,14 @@ Channel *Channel::getChannel(const std::string &name) } Channel::Channel() - : named(false) + : named(false), sent(0), received(0) { mutex = newMutex(); cond = newConditional(); } Channel::Channel(const std::string &name) - : named(true), name(name) + : named(true), name(name), sent(0), received(0) { mutex = newMutex(); cond = newConditional(); @@ -68,10 +92,10 @@ Channel::~Channel() namedChannels.erase(name); } -void Channel::push(Variant *var) +unsigned long Channel::push(Variant *var) { if (!var) - return; + return 0; Lock l(mutex); var->retain(); // Keep a reference to ourselves @@ -79,7 +103,24 @@ void Channel::push(Variant *var) if (named && queue.empty()) retain(); queue.push(var); - cond->signal(); + cond->broadcast(); + + return ++sent; +} + +void Channel::supply(Variant *var) +{ + if (!var) + return; + + unsigned long id = push(var); + + mutex->lock(); + while (!past(id, received)) + { + cond->wait(mutex); + } + mutex->unlock(); } Variant *Channel::pop() @@ -91,6 +132,9 @@ Variant *Channel::pop() Variant *var = queue.front(); queue.pop(); + received++; + cond->broadcast(); + // Release our reference to ourselves // if we're empty and named. if (named && queue.empty()) diff --git a/src/modules/thread/Channel.h b/src/modules/thread/Channel.h index f3c58d1bd..d48b6fc76 100644 --- a/src/modules/thread/Channel.h +++ b/src/modules/thread/Channel.h @@ -43,14 +43,18 @@ private: std::string name; Channel(const std::string &name); + unsigned long sent; + unsigned long received; + public: Channel(); ~Channel(); static Channel *getChannel(const std::string &name); - void push(Variant *var); + unsigned long push(Variant *var); + void supply(Variant *var); // blocking push Variant *pop(); - Variant *demand(); + Variant *demand(); // blocking pop Variant *peek(); int count(); void clear(); diff --git a/src/modules/thread/wrap_Channel.cpp b/src/modules/thread/wrap_Channel.cpp index 4ad6d0fcf..23c40a3be 100644 --- a/src/modules/thread/wrap_Channel.cpp +++ b/src/modules/thread/wrap_Channel.cpp @@ -38,6 +38,15 @@ namespace thread return 0; } + int w_Channel_supply(lua_State *L) + { + Channel *c = luax_checkchannel(L, 1); + Variant *var = Variant::fromLua(L, 2); + c->supply(var); + var->release(); + return 0; + } + int w_Channel_pop(lua_State *L) { Channel *c = luax_checkchannel(L, 1); @@ -91,6 +100,7 @@ namespace thread static const luaL_Reg type_functions[] = { { "push", w_Channel_push }, + { "supply", w_Channel_supply }, { "pop", w_Channel_pop }, { "demand", w_Channel_demand }, { "peek", w_Channel_peek }, diff --git a/src/modules/thread/wrap_Channel.h b/src/modules/thread/wrap_Channel.h index 2f9f9332a..191aacfa9 100644 --- a/src/modules/thread/wrap_Channel.h +++ b/src/modules/thread/wrap_Channel.h @@ -30,6 +30,7 @@ namespace thread { Channel *luax_checkchannel(lua_State *L, int idx); int w_Channel_push(lua_State *L); + int w_Channel_supply(lua_State *L); int w_Channel_pop(lua_State *L); int w_Channel_demand(lua_State *L); int w_Channel_peek(lua_State *L);