mirror of
https://github.com/rayaman/multi.git
synced 2026-09-05 07:27:35 -04:00
thread.hold(proxy.conn)
This commit is contained in:
@@ -35,16 +35,27 @@ function multi:chop(obj)
|
||||
if type(v) == "function" then
|
||||
table.insert(list, i)
|
||||
elseif type(v) == "table" and v.Type == multi.CONNECTOR then
|
||||
table.insert(list, {i, multi:newProxy(multi:chop(v)):init()})
|
||||
-- local stc = "stc_"..list[0].."_"..i
|
||||
-- list[-1][#list[-1] + 1] = {i, stc}
|
||||
-- list[#list+1] = i
|
||||
-- obj[stc] = multi:newSystemThreadedConnection(stc):init()
|
||||
-- obj["_"..i.."_"] = function(...)
|
||||
-- return obj[stc](...)
|
||||
-- end
|
||||
v.getThreadID = function() -- Special function we are adding
|
||||
return THREAD_ID
|
||||
end
|
||||
|
||||
v.getUniqueName = function(self)
|
||||
return self.__link_name
|
||||
end
|
||||
|
||||
local l = multi:chop(v)
|
||||
v.__link_name = l[0]
|
||||
v.__name = i
|
||||
|
||||
table.insert(list, {i, multi:newProxy(l):init()})
|
||||
end
|
||||
end
|
||||
table.insert(list, "isConnection")
|
||||
if obj.Type == multi.CONNECTOR then
|
||||
obj.isConnection = function() return true end
|
||||
else
|
||||
obj.isConnection = function() return false end
|
||||
end
|
||||
return list
|
||||
end
|
||||
|
||||
@@ -97,15 +108,14 @@ function multi:newProxy(list)
|
||||
THREAD = multi.integration.THREAD
|
||||
self.send = THREAD.waitFor(self.name.."_S")
|
||||
self.recv = THREAD.waitFor(self.name.."_R")
|
||||
self.Type = multi.PROXY
|
||||
for _,v in pairs(self.funcs) do
|
||||
if type(v) == "table" then
|
||||
-- We got a connection
|
||||
v[2]:init()
|
||||
|
||||
--setmetatable(v[2],getmetatable(multi:newConnection()))
|
||||
|
||||
self[v[1]] = v[2]
|
||||
v[2].Parent = self
|
||||
else
|
||||
lastObj = self
|
||||
self[v] = thread:newFunction(function(self,...)
|
||||
if self == me then
|
||||
me.send:push({v, true, ...})
|
||||
@@ -119,7 +129,6 @@ function multi:newProxy(list)
|
||||
table.remove(data, 1)
|
||||
for i=1,#data do
|
||||
if type(data[i]) == "table" and data[i]._self_ref_ then
|
||||
-- So if we get a self return as a return, we should return the proxy!
|
||||
data[i] = me
|
||||
end
|
||||
end
|
||||
@@ -136,9 +145,40 @@ function multi:newProxy(list)
|
||||
return c
|
||||
end
|
||||
|
||||
function multi:newSystemThreadedProcessor(name, cores)
|
||||
multi.PROXY = "proxy"
|
||||
|
||||
local name = name or "STP_"..multi.randomString(4) -- set a random name if none was given.
|
||||
local targets = {}
|
||||
|
||||
local nFunc = 0
|
||||
function multi:newTargetedFunction(ID, proc, name, func, holup) -- This registers with the queue
|
||||
if type(name)=="function" then
|
||||
holup = func
|
||||
func = name
|
||||
name = "JQ_TFunc_"..nFunc
|
||||
end
|
||||
nFunc = nFunc + 1
|
||||
proc.jobqueue:registerFunction(name, func)
|
||||
return thread:newFunction(function(...)
|
||||
local id = proc:pushJob(ID, name, ...)
|
||||
local link
|
||||
local rets
|
||||
link = proc.jobqueue.OnJobCompleted(function(jid,...)
|
||||
if id==jid then
|
||||
rets = {...}
|
||||
end
|
||||
end)
|
||||
return thread.hold(function()
|
||||
if rets then
|
||||
return multi.unpack(rets) or multi.NIL
|
||||
end
|
||||
end)
|
||||
end, holup), name
|
||||
end
|
||||
|
||||
local jid = -1
|
||||
function multi:newSystemThreadedProcessor(cores)
|
||||
|
||||
local name = "STP_"..multi.randomString(4) -- set a random name if none was given.
|
||||
|
||||
local autoscale = autoscale or false -- Will scale up the number of cores that the process uses.
|
||||
local c = {}
|
||||
@@ -154,54 +194,123 @@ function multi:newSystemThreadedProcessor(name, cores)
|
||||
c.OnObjectCreated = multi:newConnection()
|
||||
c.parent = self
|
||||
c.jobqueue = multi:newSystemThreadedJobQueue(c.cores)
|
||||
|
||||
c.targetedQueue = multi:newSystemThreadedQueue(name.."_target"):init()
|
||||
|
||||
c.jobqueue:registerFunction("enable_targets",function(name)
|
||||
local multi, thread = require("multi"):init()
|
||||
local qname = THREAD_NAME .. "_t_queue"
|
||||
local targetedQueue = THREAD.waitFor(name):init()
|
||||
local tjq = multi:newSystemThreadedQueue(qname):init()
|
||||
targetedQueue:push({tonumber(THREAD_ID), qname})
|
||||
multi:newThread("TargetedJobHandler", function()
|
||||
local queueReturn = _G["__QR"]
|
||||
while true do
|
||||
local dat = thread.hold(function()
|
||||
return tjq:pop()
|
||||
end)
|
||||
if dat then
|
||||
thread:newThread("test",function()
|
||||
local name = table.remove(dat, 1)
|
||||
local jid = table.remove(dat, 1)
|
||||
local args = table.remove(dat, 1)
|
||||
queueReturn:push{jid, _G[name](multi.unpack(args)), queue}
|
||||
end).OnError(multi.error)
|
||||
end
|
||||
end
|
||||
end).OnError(multi.error)
|
||||
end)
|
||||
|
||||
function c:pushJob(ID, name, ...)
|
||||
targets[ID]:push{name, jid, {...}}
|
||||
jid = jid - 1
|
||||
return jid + 1
|
||||
end
|
||||
|
||||
c.jobqueue:doToAll(function(name)
|
||||
enable_targets(name)
|
||||
end, name.."_target")
|
||||
|
||||
local count = 0
|
||||
while count < c.cores do
|
||||
local dat = c.targetedQueue:pop()
|
||||
if dat then
|
||||
targets[dat[1]] = multi.integration.THREAD.waitFor(dat[2]):init()
|
||||
count = count + 1
|
||||
end
|
||||
end
|
||||
|
||||
c.jobqueue:registerFunction("packObj",function(obj)
|
||||
local multi, thread = require("multi"):init()
|
||||
obj.getThreadID = function() -- Special function we are adding
|
||||
return THREAD_ID
|
||||
end
|
||||
|
||||
obj.getUniqueName = function(self)
|
||||
return self.__link_name
|
||||
end
|
||||
|
||||
local list = multi:chop(obj)
|
||||
obj.__link_name = list[0]
|
||||
|
||||
local proxy = multi:newProxy(list):init()
|
||||
|
||||
return proxy
|
||||
end)
|
||||
|
||||
c.spawnThread = c.jobqueue:newFunction("__spawnThread__", function(name, func, ...)
|
||||
local multi, thread = require("multi"):init()
|
||||
local proxy = multi:newProxy(multi:chop(thread:newThread(name, func, ...))):init()
|
||||
return proxy
|
||||
local obj = thread:newThread(name, func, ...)
|
||||
return packObj(obj)
|
||||
end, true)
|
||||
|
||||
c.spawnTask = c.jobqueue:newFunction("__spawnTask__", function(obj, func, ...)
|
||||
local multi, thread = require("multi"):init()
|
||||
local obj = multi[obj](multi, func, ...)
|
||||
local proxy = multi:newProxy(multi:chop(obj)):init()
|
||||
return proxy
|
||||
return packObj(obj)
|
||||
end, true)
|
||||
|
||||
function c:newLoop(func, notime)
|
||||
return self.spawnTask("newLoop", func, notime):init()
|
||||
proxy = self.spawnTask("newLoop", func, notime):init()
|
||||
proxy.__proc = self
|
||||
return proxy
|
||||
end
|
||||
|
||||
function c:newTLoop(func, time)
|
||||
return self.spawnTask("newTLoop", func, time):init()
|
||||
proxy = self.spawnTask("newTLoop", func, time):init()
|
||||
proxy.__proc = self
|
||||
return proxy
|
||||
end
|
||||
|
||||
function c:newUpdater(skip, func)
|
||||
return self.spawnTask("newUpdater", func, notime):init()
|
||||
proxy = self.spawnTask("newUpdater", func, notime):init()
|
||||
proxy.__proc = self
|
||||
return proxy
|
||||
end
|
||||
|
||||
function c:newEvent(task, func)
|
||||
return self.spawnTask("newEvent", task, func):init()
|
||||
proxy = self.spawnTask("newEvent", task, func):init()
|
||||
proxy.__proc = self
|
||||
return proxy
|
||||
end
|
||||
|
||||
function c:newAlarm(set, func)
|
||||
return self.spawnTask("newAlarm", set, func):init()
|
||||
proxy = self.spawnTask("newAlarm", set, func):init()
|
||||
proxy.__proc = self
|
||||
return proxy
|
||||
end
|
||||
|
||||
function c:newStep(start, reset, count, skip)
|
||||
return self.spawnTask("newStep", start, reset, count, skip):init()
|
||||
proxy = self.spawnTask("newStep", start, reset, count, skip):init()
|
||||
proxy.__proc = self
|
||||
return proxy
|
||||
end
|
||||
|
||||
function c:newTStep(start ,reset, count, set)
|
||||
return self.spawnTask("newTStep", start, reset, count, set):init()
|
||||
proxy = self.spawnTask("newTStep", start, reset, count, set):init()
|
||||
proxy.__proc = self
|
||||
return proxy
|
||||
end
|
||||
|
||||
c.OnObjectCreated(function(proc, obj)
|
||||
if not(obj.Type == multi.UPDATER or obj.Type == multi.LOOP) then
|
||||
return multi.error("Invalid type!")
|
||||
end
|
||||
end)
|
||||
|
||||
function c:getHandler()
|
||||
-- Not needed
|
||||
end
|
||||
@@ -219,7 +328,9 @@ function multi:newSystemThreadedProcessor(name, cores)
|
||||
end
|
||||
|
||||
function c:newThread(name, func, ...)
|
||||
return self.spawnThread(name, func, ...):init()
|
||||
proxy = self.spawnThread(name, func, ...):init()
|
||||
proxy.__proc = self
|
||||
return proxy
|
||||
end
|
||||
|
||||
function c:newFunction(func, holdme)
|
||||
@@ -249,3 +360,42 @@ function multi:newSystemThreadedProcessor(name, cores)
|
||||
return c
|
||||
end
|
||||
|
||||
-- Modify thread.hold to handle proxies
|
||||
local thread_ref = thread.hold
|
||||
function thread.hold(n, opt)
|
||||
if type(n) == "table" and n.Type == multi.PROXY and n.isConnection() then
|
||||
local ready = false
|
||||
local args
|
||||
local id = n.getThreadID()
|
||||
local name = n:getUniqueName()
|
||||
local func = multi:newTargetedFunction(id, n.Parent.__proc, "conn_"..multi.randomString(8), function(_name)
|
||||
local multi, thread = require("multi"):init()
|
||||
local obj = _G[_name]
|
||||
local rets = {thread.hold(obj)}
|
||||
for i,v in pairs(rets) do
|
||||
if v.Type then
|
||||
rets[i] = {_self_ref_ = "parent"}
|
||||
end
|
||||
end
|
||||
return unpack(rets)
|
||||
end)
|
||||
func(name).OnReturn(function(...)
|
||||
ready = true
|
||||
args = {...}
|
||||
end)
|
||||
local ret = {thread_ref(function()
|
||||
if ready then
|
||||
return multi.unpack(args) or multi.NIL
|
||||
end
|
||||
end, opt)}
|
||||
for i,v in pairs(ret) do
|
||||
if type(v) == "table" and v._self_ref_ == "parent" then
|
||||
ret[i] = n.Parent
|
||||
end
|
||||
end
|
||||
return unpack(ret)
|
||||
else
|
||||
return thread_ref(n, opt)
|
||||
end
|
||||
end
|
||||
|
||||
|
||||
Reference in New Issue
Block a user