mirror of
https://github.com/rayaman/multi.git
synced 2026-09-05 07:27:35 -04:00
Updated to 1.10.0!
Check changes.md for what was done
This commit is contained in:
@@ -33,43 +33,50 @@ function love.run()
|
||||
if love.load then love.load(arg) end
|
||||
if love.timer then love.timer.step() end
|
||||
local dt = 0
|
||||
while true do
|
||||
-- Process events.
|
||||
if love.event then
|
||||
love.event.pump()
|
||||
for e,a,b,c,d in love.event.poll() do
|
||||
if e == "quit" then
|
||||
if not love.quit or not love.quit() then
|
||||
if love.audio then
|
||||
love.audio.stop()
|
||||
local breakme=false
|
||||
multi:newThread("MAIN-RUN",function()
|
||||
while true do
|
||||
-- Process events.
|
||||
if love.event then
|
||||
love.event.pump()
|
||||
for name, a,b,c,d,e,f in love.event.poll() do
|
||||
if name == "quit" then
|
||||
if not love.quit or not love.quit() then
|
||||
breakme=true
|
||||
thread.kill()
|
||||
break
|
||||
end
|
||||
return
|
||||
end
|
||||
love.handlers[name](a,b,c,d,e,f)
|
||||
end
|
||||
love.handlers[e](a,b,c,d)
|
||||
end
|
||||
if love.timer then
|
||||
love.timer.step()
|
||||
dt = love.timer.getDelta()
|
||||
end
|
||||
if love.update then love.update(dt) end
|
||||
if love.window and love.graphics and love.window.isCreated() then
|
||||
love.graphics.clear()
|
||||
love.graphics.origin()
|
||||
if love.draw then love.draw() end
|
||||
multi.dManager()
|
||||
love.graphics.setColor(255,255,255,255)
|
||||
if multi.draw then multi.draw() end
|
||||
love.graphics.present()
|
||||
end
|
||||
thread.sleep()
|
||||
end
|
||||
if love.timer then
|
||||
love.timer.step()
|
||||
dt = love.timer.getDelta()
|
||||
end
|
||||
if love.update then love.update(dt) end
|
||||
end)
|
||||
while not breakme do
|
||||
love.timer.sleep(.005)
|
||||
multi:uManager(dt)
|
||||
if multi.boost then
|
||||
for i=1,multi.boost-1 do
|
||||
multi:uManager(dt)
|
||||
end
|
||||
end
|
||||
multi:uManager(dt)
|
||||
if love.window and love.graphics and love.window.isCreated() then
|
||||
love.graphics.clear()
|
||||
love.graphics.origin()
|
||||
if love.draw then love.draw() end
|
||||
multi.dManager()
|
||||
love.graphics.setColor(255,255,255,255)
|
||||
if multi.draw then multi.draw() end
|
||||
love.graphics.present()
|
||||
end
|
||||
end
|
||||
return
|
||||
end
|
||||
multi.drawF={}
|
||||
function multi:dManager()
|
||||
@@ -81,34 +88,34 @@ function multi:onDraw(func,i)
|
||||
i=i or 1
|
||||
table.insert(self.drawF,i,func)
|
||||
end
|
||||
function multi:lManager()
|
||||
if love.event then
|
||||
love.event.pump()
|
||||
for e,a,b,c,d in love.event.poll() do
|
||||
if e == "quit" then
|
||||
if not love.quit or not love.quit() then
|
||||
if love.audio then
|
||||
love.audio.stop()
|
||||
end
|
||||
return nil
|
||||
end
|
||||
end
|
||||
love.handlers[e](a,b,c,d)
|
||||
end
|
||||
end
|
||||
if love.timer then
|
||||
love.timer.step()
|
||||
dt = love.timer.getDelta()
|
||||
end
|
||||
if love.update then love.update(dt) end
|
||||
multi:uManager(dt)
|
||||
if love.window and love.graphics and love.window.isCreated() then
|
||||
love.graphics.clear()
|
||||
love.graphics.origin()
|
||||
if love.draw then love.draw() end
|
||||
multi.dManager()
|
||||
love.graphics.setColor(255,255,255,255)
|
||||
if multi.draw then multi.draw() end
|
||||
love.graphics.present()
|
||||
end
|
||||
end
|
||||
--~ function multi:lManager()
|
||||
--~ if love.event then
|
||||
--~ love.event.pump()
|
||||
--~ for e,a,b,c,d in love.event.poll() do
|
||||
--~ if e == "quit" then
|
||||
--~ if not love.quit or not love.quit() then
|
||||
--~ if love.audio then
|
||||
--~ love.audio.stop()
|
||||
--~ end
|
||||
--~ return nil
|
||||
--~ end
|
||||
--~ end
|
||||
--~ love.handlers[e](a,b,c,d)
|
||||
--~ end
|
||||
--~ end
|
||||
--~ if love.timer then
|
||||
--~ love.timer.step()
|
||||
--~ dt = love.timer.getDelta()
|
||||
--~ end
|
||||
--~ if love.update then love.update(dt) end
|
||||
--~ multi:uManager(dt)
|
||||
--~ if love.window and love.graphics and love.window.isCreated() then
|
||||
--~ love.graphics.clear()
|
||||
--~ love.graphics.origin()
|
||||
--~ if love.draw then love.draw() end
|
||||
--~ multi.dManager()
|
||||
--~ love.graphics.setColor(255,255,255,255)
|
||||
--~ if multi.draw then multi.draw() end
|
||||
--~ love.graphics.present()
|
||||
--~ end
|
||||
--~ end
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -21,6 +21,7 @@ LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
|
||||
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
|
||||
SOFTWARE.
|
||||
]]
|
||||
package.path="?/init.lua;?.lua;"..package.path
|
||||
function os.getOS()
|
||||
if package.config:sub(1,1)=='\\' then
|
||||
return 'windows'
|
||||
@@ -32,9 +33,13 @@ end
|
||||
lanes=require("lanes").configure()
|
||||
--~ package.path="lua/?/init.lua;lua/?.lua;"..package.path
|
||||
require("multi") -- get it all and have it on all lanes
|
||||
isMainThread=true
|
||||
function multi:canSystemThread()
|
||||
return true
|
||||
end
|
||||
function multi:getPlatform()
|
||||
return "lanes"
|
||||
end
|
||||
local multi=multi
|
||||
-- Step 2 set up the linda objects
|
||||
local __GlobalLinda = lanes.linda() -- handles global stuff
|
||||
@@ -85,7 +90,10 @@ else
|
||||
THREAD.__CORES=tonumber(io.popen("nproc --all"):read("*n"))
|
||||
end
|
||||
function THREAD.kill() -- trigger the lane destruction
|
||||
-- coroutine.yield({"_kill_",":)"})
|
||||
error("Thread was killed!")
|
||||
end
|
||||
function THREAD.getName()
|
||||
return THREAD_NAME
|
||||
end
|
||||
--[[ Step 4 We need to get sleeping working to handle timing... We want idle wait, not busy wait
|
||||
Idle wait keeps the CPU running better where busy wait wastes CPU cycles... Lanes does not have a sleep method
|
||||
@@ -102,12 +110,17 @@ function THREAD.hold(n)
|
||||
repeat wait() until n()
|
||||
end
|
||||
-- Step 5 Basic Threads!
|
||||
function multi:newSystemThread(name,func)
|
||||
function multi:newSystemThread(name,func,...)
|
||||
local c={}
|
||||
local __self=c
|
||||
c.name=name
|
||||
c.Type="sthread"
|
||||
c.thread=lanes.gen("*", func)()
|
||||
local THREAD_NAME=name
|
||||
local function func2(...)
|
||||
_G["THREAD_NAME"]=THREAD_NAME
|
||||
func()
|
||||
end
|
||||
c.thread=lanes.gen("*", func2)(...)
|
||||
function c:kill()
|
||||
--self.status:Destroy()
|
||||
self.thread:cancel()
|
||||
|
||||
@@ -2,17 +2,22 @@ require("multi.compat.love2d")
|
||||
function multi:canSystemThread()
|
||||
return true
|
||||
end
|
||||
function multi:getPlatform()
|
||||
return "love2d"
|
||||
end
|
||||
multi.integration={}
|
||||
multi.integration.love2d={}
|
||||
multi.integration.love2d.ThreadBase=[[
|
||||
tab={...}
|
||||
__THREADNAME__=tab[2]
|
||||
__THREADID__=tab[1]
|
||||
__THREADID__=table.remove(tab,1)
|
||||
__THREADNAME__=table.remove(tab,1)
|
||||
require("love.filesystem")
|
||||
require("love.system")
|
||||
require("love.timer")
|
||||
require("love.image")
|
||||
require("multi")
|
||||
GLOBAL={}
|
||||
isMainThread=false
|
||||
setmetatable(GLOBAL,{
|
||||
__index=function(t,k)
|
||||
__sync__()
|
||||
@@ -31,7 +36,7 @@ setmetatable(GLOBAL,{
|
||||
function __sync__()
|
||||
local data=__mythread__:pop()
|
||||
while data do
|
||||
love.timer.sleep(.001)
|
||||
love.timer.sleep(.01)
|
||||
if type(data)=="string" then
|
||||
local cmd,tp,name,d=data:match("(%S-) (%S-) (%S-) (.+)")
|
||||
if name=="__DIEPLZ"..__THREADID__.."__" then
|
||||
@@ -136,6 +141,12 @@ function dump(func)
|
||||
end
|
||||
return table.concat(code)
|
||||
end
|
||||
function sThread.getName()
|
||||
return __THREADNAME__
|
||||
end
|
||||
function sThread.kill()
|
||||
error("Thread was killed!")
|
||||
end
|
||||
function sThread.set(name,val)
|
||||
GLOBAL[name]=val
|
||||
end
|
||||
@@ -155,17 +166,9 @@ end
|
||||
function sThread.hold(n)
|
||||
repeat __sync__() until n()
|
||||
end
|
||||
multi:newLoop(function(self)
|
||||
self:Pause()
|
||||
local ld=multi:getLoad()
|
||||
self:Resume()
|
||||
if ld<80 then
|
||||
love.timer.sleep(.01)
|
||||
end
|
||||
end)
|
||||
updater=multi:newUpdater()
|
||||
updater:OnUpdate(__sync__)
|
||||
func=loadDump([=[INSERT_USER_CODE]=])()
|
||||
func=loadDump([=[INSERT_USER_CODE]=])(unpack(tab))
|
||||
multi:mainloop()
|
||||
]]
|
||||
GLOBAL={} -- Allow main thread to interact with these objects as well
|
||||
@@ -187,6 +190,10 @@ setmetatable(GLOBAL,{
|
||||
})
|
||||
THREAD={} -- Allow main thread to interact with these objects as well
|
||||
multi.integration.love2d.mainChannel=love.thread.getChannel("__MainChan__")
|
||||
isMainThread=true
|
||||
function THREAD.getName()
|
||||
return __THREADNAME__
|
||||
end
|
||||
function ToStr(val, name, skipnewlines, depth)
|
||||
skipnewlines = skipnewlines or false
|
||||
depth = depth or 0
|
||||
@@ -265,12 +272,12 @@ local function randomString(n)
|
||||
end
|
||||
return str
|
||||
end
|
||||
function multi:newSystemThread(name,func) -- the main method
|
||||
function multi:newSystemThread(name,func,...) -- the main method
|
||||
local c={}
|
||||
c.name=name
|
||||
c.ID=c.name.."<ID|"..randomString(8)..">"
|
||||
c.thread=love.thread.newThread(multi.integration.love2d.ThreadBase:gsub("INSERT_USER_CODE",dump(func)))
|
||||
c.thread:start(c.ID,c.name)
|
||||
c.thread:start(c.ID,c.name,...)
|
||||
function c:kill()
|
||||
multi.integration.GLOBAL["__DIEPLZ"..self.ID.."__"]="__DIEPLZ"..self.ID.."__"
|
||||
end
|
||||
@@ -288,9 +295,7 @@ function THREAD.get(name)
|
||||
return GLOBAL[name]
|
||||
end
|
||||
function THREAD.waitFor(name)
|
||||
multi.OBJ_REF:Pause()
|
||||
repeat multi:lManager() until GLOBAL[name]
|
||||
multi.OBJ_REF:Resume()
|
||||
repeat multi:uManager() until GLOBAL[name]
|
||||
return GLOBAL[name]
|
||||
end
|
||||
function THREAD.getCores()
|
||||
@@ -300,9 +305,7 @@ function THREAD.sleep(n)
|
||||
love.timer.sleep(n)
|
||||
end
|
||||
function THREAD.hold(n)
|
||||
multi.OBJ_REF:Pause()
|
||||
repeat multi:lManager() until n()
|
||||
multi.OBJ_REF:Resume()
|
||||
repeat multi:uManager() until n()
|
||||
end
|
||||
__channels__={}
|
||||
multi.integration.GLOBAL=GLOBAL
|
||||
@@ -311,7 +314,6 @@ updater=multi:newUpdater()
|
||||
updater:OnUpdate(function(self)
|
||||
local data=multi.integration.love2d.mainChannel:pop()
|
||||
while data do
|
||||
--print("MAIN:",data)
|
||||
if type(data)=="string" then
|
||||
local cmd,tp,name,d=data:match("(%S-) (%S-) (%S-) (.+)")
|
||||
if cmd=="SYNC" then
|
||||
|
||||
@@ -0,0 +1,127 @@
|
||||
--[[
|
||||
MIT License
|
||||
|
||||
Copyright (c) 2017 Ryan Ward
|
||||
|
||||
Permission is hereby granted, free of charge, to any person obtaining a copy
|
||||
of this software and associated documentation files (the "Software"), to deal
|
||||
in the Software without restriction, including without limitation the rights
|
||||
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
|
||||
copies of the Software, and to permit persons to whom the Software is
|
||||
furnished to do so, subject to the following conditions:
|
||||
|
||||
The above copyright notice and this permission notice shall be included in all
|
||||
copies or substantial portions of the Software.
|
||||
|
||||
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
|
||||
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
|
||||
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
|
||||
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
|
||||
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
|
||||
SOFTWARE.
|
||||
]]
|
||||
|
||||
-- I DEMAND USAGE FOR LUVIT
|
||||
-- Cannot use discordia without my multitasking library (Which I love more that the luvit platform... then again i'm partial :P)
|
||||
package.path="?/init.lua;?.lua;"..package.path
|
||||
local function _INIT(luvitThread,timer)
|
||||
-- lots of this stuff should be able to stay the same
|
||||
function os.getOS()
|
||||
if package.config:sub(1,1)=='\\' then
|
||||
return 'windows'
|
||||
else
|
||||
return 'unix'
|
||||
end
|
||||
end
|
||||
-- Step 1 get setup threads on luvit... Sigh how do i even...
|
||||
require("multi")
|
||||
isMainThread=true
|
||||
function multi:canSystemThread()
|
||||
return true
|
||||
end
|
||||
function multi:getPlatform()
|
||||
return "luvit"
|
||||
end
|
||||
local multi=multi
|
||||
-- Step 2 set up the Global table... is this possible?
|
||||
local GLOBAL={}
|
||||
setmetatable(GLOBAL,{
|
||||
__index=function(t,k)
|
||||
--print("No Global table when using luvit integration!")
|
||||
return nil
|
||||
end,
|
||||
__newindex=function(t,k,v)
|
||||
--print("No Global table when using luvit integration!")
|
||||
end,
|
||||
})
|
||||
local THREAD={}
|
||||
function THREAD.set(name,val)
|
||||
--print("No Global table when using luvit integration!")
|
||||
end
|
||||
function THREAD.get(name)
|
||||
--print("No Global table when using luvit integration!")
|
||||
end
|
||||
local function randomString(n)
|
||||
local str = ''
|
||||
local strings = {'a','b','c','d','e','f','g','h','i','j','k','l','m','n','o','p','q','r','s','t','u','v','w','x','y','z','1','2','3','4','5','6','7','8','9','0','A','B','C','D','E','F','G','H','I','J','K','L','M','N','O','P','Q','R','S','T','U','V','W','X','Y','Z'}
|
||||
for i=1,n do
|
||||
str = str..''..strings[math.random(1,#strings)]
|
||||
end
|
||||
return str
|
||||
end
|
||||
function THREAD.waitFor(name)
|
||||
--print("No Global table when using luvit integration!")
|
||||
end
|
||||
function THREAD.testFor(name,val,sym)
|
||||
--print("No Global table when using luvit integration!")
|
||||
end
|
||||
function THREAD.getCores()
|
||||
return THREAD.__CORES
|
||||
end
|
||||
if os.getOS()=="windows" then
|
||||
THREAD.__CORES=tonumber(os.getenv("NUMBER_OF_PROCESSORS"))
|
||||
else
|
||||
THREAD.__CORES=tonumber(io.popen("nproc --all"):read("*n"))
|
||||
end
|
||||
function THREAD.kill() -- trigger the thread destruction
|
||||
error("Thread was Killed!")
|
||||
end
|
||||
-- hmmm if im cleaver I can get this to work... but since data passing isn't going to be a thing its probably not important
|
||||
function THREAD.sleep(n)
|
||||
--print("No Global table when using luvit integration!")
|
||||
end
|
||||
function THREAD.hold(n)
|
||||
--print("No Global table when using luvit integration!")
|
||||
end
|
||||
-- Step 5 Basic Threads!
|
||||
local function entry(path,name,func,...)
|
||||
local timer = require'timer'
|
||||
local luvitThread = require'thread'
|
||||
package.path=path
|
||||
loadstring(func)(...)
|
||||
end
|
||||
function multi:newSystemThread(name,func,...)
|
||||
local c={}
|
||||
local __self=c
|
||||
c.name=name
|
||||
c.Type="sthread"
|
||||
c.thread={}
|
||||
c.func=string.dump(func)
|
||||
function c:kill()
|
||||
-- print("No Global table when using luvit integration!")
|
||||
end
|
||||
luvitThread.start(entry,package.path,name,c.func,...)
|
||||
return c
|
||||
end
|
||||
print("Integrated Luvit!")
|
||||
multi.integration={} -- for module creators
|
||||
multi.integration.GLOBAL=GLOBAL
|
||||
multi.integration.THREAD=THREAD
|
||||
require("multi.integration.shared")
|
||||
-- Start the main mainloop... This allows you to process your multi objects, but the engine on the main thread will be limited to .001 or 1 millisecond sigh...
|
||||
local interval = timer.setInterval(1, function ()
|
||||
multi:uManager()
|
||||
end)
|
||||
end
|
||||
return {init=function(threadHandle,timerHandle) _INIT(threadHandle,timerHandle) return GLOBAL,THREAD end}
|
||||
@@ -21,35 +21,70 @@ LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
|
||||
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
|
||||
SOFTWARE.
|
||||
]]
|
||||
function multi.randomString(n)
|
||||
local str = ''
|
||||
local strings = {'a','b','c','d','e','f','g','h','i','j','k','l','m','n','o','p','q','r','s','t','u','v','w','x','y','z','1','2','3','4','5','6','7','8','9','0','A','B','C','D','E','F','G','H','I','J','K','L','M','N','O','P','Q','R','S','T','U','V','W','X','Y','Z'}
|
||||
for i=1,n do
|
||||
str = str..''..strings[math.random(1,#strings)]
|
||||
end
|
||||
return str
|
||||
end
|
||||
function multi:newSystemThreadedQueue(name) -- in love2d this will spawn a channel on both ends
|
||||
local c={} -- where we will store our object
|
||||
c.name=name -- set the name this is important for the love2d side
|
||||
if love then -- check love
|
||||
if love.thread then -- make sure we can use the threading module
|
||||
function c:init() -- create an init function so we can mimic on bith love2d and lanes
|
||||
function c:init() -- create an init function so we can mimic on both love2d and lanes
|
||||
self.chan=love.thread.getChannel(self.name) -- create channel by the name self.name
|
||||
function self:push(v) -- push to the channel
|
||||
self.chan:push({type(v),resolveData(v)})
|
||||
local tab
|
||||
if type(v)=="table" then
|
||||
tab = {}
|
||||
for i,c in pairs(v) do
|
||||
if type(c)=="function" then
|
||||
tab[i]="\1"..string.dump(c)
|
||||
else
|
||||
tab[i]=c
|
||||
end
|
||||
end
|
||||
self.chan:push(tab)
|
||||
else
|
||||
self.chan:push(c)
|
||||
end
|
||||
end
|
||||
function self:pop() -- pop from the channel
|
||||
local tab=self.chan:pop()
|
||||
--print(tab)
|
||||
if not tab then return end
|
||||
return resolveType(tab[1],tab[2])
|
||||
local v=self.chan:pop()
|
||||
if not v then return end
|
||||
if type(v)=="table" then
|
||||
tab = {}
|
||||
for i,c in pairs(v) do
|
||||
if type(c)=="string" then
|
||||
if c:sub(1,1)=="\1" then
|
||||
tab[i]=loadstring(c:sub(2,-1))
|
||||
else
|
||||
tab[i]=c
|
||||
end
|
||||
else
|
||||
tab[i]=c
|
||||
end
|
||||
end
|
||||
return tab
|
||||
else
|
||||
return self.chan:pop()
|
||||
end
|
||||
end
|
||||
function self:peek()
|
||||
local tab=self.chan:peek()
|
||||
--print(tab)
|
||||
if not tab then return end
|
||||
return resolveType(tab[1],tab[2])
|
||||
local v=self.chan:peek()
|
||||
if not v then return end
|
||||
if type(v)=="table" then
|
||||
tab = {}
|
||||
for i,c in pairs(v) do
|
||||
if type(c)=="string" then
|
||||
if c:sub(1,1)=="\1" then
|
||||
tab[i]=loadstring(c:sub(2,-1))
|
||||
else
|
||||
tab[i]=c
|
||||
end
|
||||
else
|
||||
tab[i]=c
|
||||
end
|
||||
end
|
||||
return tab
|
||||
else
|
||||
return self.chan:pop()
|
||||
end
|
||||
end
|
||||
GLOBAL[self.name]=self -- send the object to the thread through the global interface
|
||||
return self -- return the object
|
||||
@@ -59,7 +94,7 @@ function multi:newSystemThreadedQueue(name) -- in love2d this will spawn a chann
|
||||
error("Make sure you required the love.thread module!") -- tell the user if he/she didn't require said module
|
||||
end
|
||||
else
|
||||
c.linda=lanes.linda() -- lanes is a bit eaiser, create the linda on the main thread
|
||||
c.linda=lanes.linda() -- lanes is a bit easier, create the linda on the main thread
|
||||
function c:push(v) -- push to the queue
|
||||
self.linda:send("Q",v)
|
||||
end
|
||||
@@ -76,6 +111,100 @@ function multi:newSystemThreadedQueue(name) -- in love2d this will spawn a chann
|
||||
end
|
||||
return c
|
||||
end
|
||||
function multi:newSystemThreadedConnection(name,protect)
|
||||
local c={}
|
||||
c.name = name
|
||||
c.protect=protect
|
||||
local sThread=multi.integration.THREAD
|
||||
local GLOBAL=multi.integration.GLOBAL
|
||||
function c:init()
|
||||
require("multi")
|
||||
if multi:getPlatform()=="love2d" then
|
||||
GLOBAL=_G.GLOBAL
|
||||
sThread=_G.sThread
|
||||
end
|
||||
local conn = {}
|
||||
conn.name = self.name
|
||||
conn.count = 0
|
||||
if isMainThread then
|
||||
if GLOBAL[self.name.."THREADED_CONNQ"] then -- if this thing exists then lets grab it, we are doing something different here. instead of cleaning things up, we will gave a dedicated queue to manage things
|
||||
conn.queueCall = sThread.waitFor(self.name.."THREADED_CALLQ"):init()
|
||||
else
|
||||
conn.queueCall = multi:newSystemThreadedQueue(self.name.."THREADED_CALLQ"):init()
|
||||
end
|
||||
else
|
||||
require("multi") -- so things don't break, but also allows bi-directional connections to work
|
||||
conn.queueCall = sThread.waitFor(self.name.."THREADED_CALLQ"):init()
|
||||
end
|
||||
setmetatable(conn,{__call=function(self,...) return self:connect(...) end})
|
||||
conn.obj=multi:newConnection(self.protect)
|
||||
function conn:connect(func)
|
||||
return self.obj(func)
|
||||
end
|
||||
function conn:fConnect(func)
|
||||
return self.obj:fConnect(func)
|
||||
end
|
||||
function conn:holdUT(n)
|
||||
self.obj:holdUT(n)
|
||||
end
|
||||
function conn:Bind(t)
|
||||
self.obj:Bind(t)
|
||||
end
|
||||
function conn:Remove()
|
||||
self.obj:Remove()
|
||||
end
|
||||
function conn:getConnection(name,ingore)
|
||||
return self.obj:getConnection(name,ingore)
|
||||
end
|
||||
function conn:Fire(...)
|
||||
local args = {...}
|
||||
table.insert(args,1,multi.randomString(8))
|
||||
table.insert(args,1,self.name)
|
||||
table.insert(args,1,"F")
|
||||
self.queueCall:push(args)
|
||||
if self.trigger_self then
|
||||
self.obj:Fire(...)
|
||||
end
|
||||
end
|
||||
self.cleanup = .01
|
||||
function conn:SetCleanUpRate(n)
|
||||
self.cleanup=n or .01
|
||||
end
|
||||
conn.lastid=""
|
||||
conn.looper = multi:newLoop(function(self)
|
||||
local con = self.link
|
||||
local data = con.queueCall:peek()
|
||||
if not data then return end
|
||||
local id = data[3]
|
||||
if data[1]=="F" and data[2]==con.name and con.lastid~=id then
|
||||
con.lastid=id
|
||||
table.remove(data,1)-- Remove the first 3 elements
|
||||
table.remove(data,1)-- Remove the first 3 elements
|
||||
table.remove(data,1)-- Remove the first 3 elements
|
||||
con.obj:Fire(unpack(data))
|
||||
multi:newThread("Clean_UP",function()
|
||||
thread.sleep(con.cleanup)
|
||||
local dat = con.queueCall:peek()
|
||||
if not dat then return end
|
||||
table.remove(data,1)-- Remove the first 3 elements
|
||||
table.remove(data,1)-- Remove the first 3 elements
|
||||
table.remove(data,1)-- Remove the first 3 elements
|
||||
if dat[3]==id then
|
||||
con.queueCall:pop()
|
||||
end
|
||||
end)
|
||||
end
|
||||
end)
|
||||
conn.HoldUT=conn.holdUT
|
||||
conn.looper.link=conn
|
||||
conn.Connect=conn.connect
|
||||
conn.FConnect=conn.fConnect
|
||||
conn.GetConnection=conn.getConnection
|
||||
return conn
|
||||
end
|
||||
GLOBAL[name]=c
|
||||
return c
|
||||
end
|
||||
function multi:systemThreadedBenchmark(n,p)
|
||||
n=n or 1
|
||||
local cores=multi.integration.THREAD.getCores()
|
||||
@@ -120,88 +249,43 @@ function multi:systemThreadedBenchmark(n,p)
|
||||
end)
|
||||
return c
|
||||
end
|
||||
function multi:newSystemThreadedTable(name,n)
|
||||
local c={} -- where we will store our object
|
||||
c.name=name -- set the name this is important for the love2d side
|
||||
c.cores=n
|
||||
c.hasT={}
|
||||
if love then -- check love
|
||||
if love.thread then -- make sure we can use the threading module
|
||||
function c:init() -- create an init function so we can mimic on bith love2d and lanes
|
||||
self.tab={}
|
||||
self.chan=love.thread.getChannel(self.name) -- create channel by the name self.name
|
||||
function self:waitFor(name) -- pop from the channel
|
||||
repeat self:sync() until self[name]
|
||||
return self[name]
|
||||
end
|
||||
function self:sync()
|
||||
local data=self.chan:peek()
|
||||
if data then
|
||||
local cmd,tp,name,d=data:match("(%S-) (%S-) (%S-) (.+)")
|
||||
if not self.hasT[name] then
|
||||
if type(data)=="string" then
|
||||
if cmd=="SYNC" then
|
||||
self.tab[name]=resolveType(tp,d) -- this is defined in the loveManager.lua file
|
||||
self.hasT[name]=true
|
||||
end
|
||||
else
|
||||
self.tab[name]=data
|
||||
end
|
||||
self.chan:pop()
|
||||
end
|
||||
end
|
||||
end
|
||||
function self:reset(name)
|
||||
self.hasT[core]=nil
|
||||
end
|
||||
setmetatable(self,{
|
||||
__index=function(t,k)
|
||||
self:sync()
|
||||
return self.tab[k]
|
||||
end,
|
||||
__newindex=function(t,k,v)
|
||||
self:sync()
|
||||
self.tab[k]=v
|
||||
if type(v)=="userdata" then
|
||||
self.chan:push(v)
|
||||
else
|
||||
for i=1,self.cores do
|
||||
self.chan:push("SYNC "..type(v).." "..k.." "..resolveData(v)) -- this is defined in the loveManager.lua file
|
||||
end
|
||||
end
|
||||
end,
|
||||
})
|
||||
GLOBAL[self.name]=self -- send the object to the thread through the global interface
|
||||
return self -- return the object
|
||||
end
|
||||
return c
|
||||
function multi:newSystemThreadedTable(name)
|
||||
local c={}
|
||||
c.name=name -- set the name this is important for identifying what is what
|
||||
local sThread=multi.integration.THREAD
|
||||
local GLOBAL=multi.integration.GLOBAL
|
||||
function c:init() -- create an init function so we can mimic on both love2d and lanes
|
||||
if multi:getPlatform()=="love2d" then
|
||||
GLOBAL=_G.GLOBAL
|
||||
sThread=_G.sThread
|
||||
end
|
||||
local cc={}
|
||||
cc.tab={}
|
||||
if isMainThread then
|
||||
cc.conn = multi:newSystemThreadedConnection(self.name.."_Tabled_Connection"):init()
|
||||
else
|
||||
error("Make sure you required the love.thread module!") -- tell the user if he/she didn't require said module
|
||||
cc.conn = sThread.waitFor(self.name.."_Tabled_Connection"):init()
|
||||
end
|
||||
else
|
||||
c.linda=lanes.linda() -- lanes is a bit eaiser, create the linda on the main thread
|
||||
function c:waitFor(name)
|
||||
while self[name]==nil do
|
||||
-- Waiting
|
||||
end
|
||||
return self[name]
|
||||
function cc:waitFor(name)
|
||||
repeat multi:uManager() until tab[name]~=nil
|
||||
return tab[name]
|
||||
end
|
||||
function c:sync()
|
||||
return -- just so we match the love2d side
|
||||
end
|
||||
function c:init() -- set the metatable
|
||||
setmetatable(self,{
|
||||
__index=function(t,k)
|
||||
return self.linda:get(k)
|
||||
end,
|
||||
__newindex=function(t,k,v)
|
||||
self.linda:set(k,v)
|
||||
end,
|
||||
})
|
||||
return self
|
||||
end
|
||||
multi.integration.GLOBAL[name]=c -- send the object to the thread through the global interface
|
||||
local link = cc
|
||||
cc.conn(function(k,v)
|
||||
link.tab[k]=v
|
||||
end)
|
||||
setmetatable(cc,{
|
||||
__index=function(t,k)
|
||||
return t.tab[k]
|
||||
end,
|
||||
__newindex=function(t,k,v)
|
||||
t.tab[k]=v
|
||||
t.conn:Fire(k,v)
|
||||
end,
|
||||
})
|
||||
return cc
|
||||
end
|
||||
GLOBAL[c.name]=c
|
||||
return c
|
||||
end
|
||||
function multi:newSystemThreadedJobQueue(numOfCores)
|
||||
@@ -212,8 +296,7 @@ function multi:newSystemThreadedJobQueue(numOfCores)
|
||||
c.queueOUT=multi:newSystemThreadedQueue("THREADED_JQO"):init()
|
||||
c.queueALL=multi:newSystemThreadedQueue("THREADED_QALL"):init()
|
||||
c.REG=multi:newSystemThreadedQueue("THREADED_JQ_F_REG"):init()
|
||||
-- registerJob(name,func)
|
||||
-- pushJob(...)
|
||||
c.OnReady=multi:newConnection()
|
||||
function c:registerJob(name,func)
|
||||
for i=1,self.cores do
|
||||
self.REG:push({name,func})
|
||||
@@ -222,9 +305,10 @@ function multi:newSystemThreadedJobQueue(numOfCores)
|
||||
function c:pushJob(name,...)
|
||||
self.queueOUT:push({self.jobnum,name,...})
|
||||
self.jobnum=self.jobnum+1
|
||||
return self.jobnum-1
|
||||
end
|
||||
local GLOBAL=multi.integration.GLOBAL -- set up locals incase we are using lanes
|
||||
local sThread=multi.integration.THREAD -- set up locals incase we are using lanes
|
||||
local GLOBAL=multi.integration.GLOBAL -- set up locals in case we are using lanes
|
||||
local sThread=multi.integration.THREAD -- set up locals in case we are using lanes
|
||||
function c:doToAll(func)
|
||||
local TaskName=multi.randomString(16)
|
||||
for i=1,self.cores do
|
||||
@@ -232,17 +316,26 @@ function multi:newSystemThreadedJobQueue(numOfCores)
|
||||
end
|
||||
end
|
||||
function c:start()
|
||||
self:doToAll(function()
|
||||
_G["__started__"]=true
|
||||
SFunc()
|
||||
multi:newEvent(function()
|
||||
return self.ThreadsLoaded==true
|
||||
end):OnEvent(function(evnt)
|
||||
GLOBAL["THREADED_JQ"]=nil -- remove it
|
||||
GLOBAL["THREADED_JQO"]=nil -- remove it
|
||||
GLOBAL["THREADED_JQ_F_REG"]=nil -- remove it
|
||||
self:doToAll(function()
|
||||
_G["__started__"]=true
|
||||
SFunc()
|
||||
end)
|
||||
evnt:Destroy()
|
||||
end)
|
||||
end
|
||||
GLOBAL["__JQ_COUNT__"]=c.cores
|
||||
for i=1,c.cores do
|
||||
multi:newSystemThread("System Threaded Job Queue Worker Thread #"..i,function()
|
||||
multi:newSystemThread("System Threaded Job Queue Worker Thread #"..i,function(name,ind)
|
||||
require("multi")
|
||||
ThreadName=name
|
||||
__sleep__=.001
|
||||
if love then -- lets make sure we don't reference upvalues if using love2d
|
||||
if love then -- lets make sure we don't reference up-values if using love2d
|
||||
GLOBAL=_G.GLOBAL
|
||||
sThread=_G.sThread
|
||||
__sleep__=.1
|
||||
@@ -251,10 +344,6 @@ function multi:newSystemThreadedJobQueue(numOfCores)
|
||||
JQO=sThread.waitFor("THREADED_JQ"):init() -- Grab it
|
||||
REG=sThread.waitFor("THREADED_JQ_F_REG"):init() -- Grab it
|
||||
QALL=sThread.waitFor("THREADED_QALL"):init() -- Grab it
|
||||
sThread.sleep(.1) -- lets wait for things to work out
|
||||
GLOBAL["THREADED_JQ"]=nil -- remove it
|
||||
GLOBAL["THREADED_JQO"]=nil -- remove it
|
||||
GLOBAL["THREADED_JQ_F_REG"]=nil -- remove it
|
||||
QALLT={}
|
||||
FUNCS={}
|
||||
SFunc=multi:newFunction(function(self)
|
||||
@@ -287,10 +376,12 @@ function multi:newSystemThreadedJobQueue(numOfCores)
|
||||
return FUNCS[k]
|
||||
end
|
||||
})
|
||||
lastjob=os.clock()
|
||||
MainLoop=multi:newLoop(function(self)
|
||||
if __started__ then
|
||||
local job=JQI:pop()
|
||||
if job then
|
||||
lastjob=os.clock()
|
||||
local d=QALL:peek()
|
||||
if d then
|
||||
if not QALLT[d[1]] then
|
||||
@@ -311,17 +402,36 @@ function multi:newSystemThreadedJobQueue(numOfCores)
|
||||
end
|
||||
end
|
||||
end)
|
||||
multi:newThread("Idler",function()
|
||||
while true do
|
||||
if os.clock()-lastjob>1 then
|
||||
sThread.sleep(.1)
|
||||
end
|
||||
thread.sleep(.001)
|
||||
end
|
||||
end)
|
||||
JQO:push({"_THREADINIT_"})
|
||||
if not love then
|
||||
multi:mainloop()
|
||||
end
|
||||
end)
|
||||
end,"Thread<"..i..">",i)
|
||||
end
|
||||
c.OnJobCompleted=multi:newConnection()
|
||||
c.threadsResponded = 0
|
||||
c.updater=multi:newLoop(function(self)
|
||||
local data=self.link.queueIN:pop()
|
||||
while data do
|
||||
if data then
|
||||
self.link.OnJobCompleted:Fire(unpack(data))
|
||||
local a=unpack(data)
|
||||
if a=="_THREADINIT_" then
|
||||
self.link.threadsResponded=self.link.threadsResponded+1
|
||||
if self.link.threadsResponded==self.link.cores then
|
||||
self.link.ThreadsLoaded=true
|
||||
self.link.OnReady:Fire()
|
||||
end
|
||||
else
|
||||
self.link.OnJobCompleted:Fire(unpack(data))
|
||||
end
|
||||
end
|
||||
data=self.link.queueIN:pop()
|
||||
end
|
||||
|
||||
Reference in New Issue
Block a user