summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorruki <[email protected]>2025-10-05 18:15:13 +0800
committerruki <[email protected]>2025-10-05 18:15:13 +0800
commit0d3e6542776c0b90e24b594a022e8746c8ec2065 (patch)
tree44b72ebf448cb93b8dc9d533158c14349620f219
parent7b432414d56921bf663f1d7765c5807e51142fce (diff)
decrease coroutine count in runjobs
-rw-r--r--xmake/modules/async/runjobs.lua224
1 files changed, 92 insertions, 132 deletions
diff --git a/xmake/modules/async/runjobs.lua b/xmake/modules/async/runjobs.lua
index 1892d9cc5..cdcc2067d 100644
--- a/xmake/modules/async/runjobs.lua
+++ b/xmake/modules/async/runjobs.lua
@@ -49,11 +49,11 @@ function _init_progress(state, opt)
end
-- init progress wrapper
- state.count = 0
+ state.finished_count = 0
state.progress_factor = opt.progress_factor or 1.0
local progress_wrapper = {}
progress_wrapper.current = function ()
- return state.count
+ return state.finished_count
end
progress_wrapper.total = function ()
return state.total
@@ -61,7 +61,7 @@ function _init_progress(state, opt)
progress_wrapper.percent = function ()
local total = state.total
if total and total > 0 then
- return math.floor((state.count * state.progress_factor * 100) / total)
+ return math.floor((state.finished_count * state.progress_factor * 100) / total)
else
return 0
end
@@ -135,7 +135,81 @@ function _progress_loop(state)
end
-- comsume jobs
-function _comsume_jobs_loop()
+-- TODO distcc
+function _comsume_jobs_loop(state)
+ local jobs = state.jobs
+ local jobs_cb = state.jobs_cb
+ local total = state.total
+ local curdir = state.curdir
+ while state.finished_count < total and not state.stop do
+
+ -- get free job
+ local job
+ local job_func = jobs_cb
+ if not job_func then
+ job = jobs:getfree()
+ if job then
+ job_func = job.run
+ else
+ break
+ end
+ end
+
+ try
+ {
+ function ()
+ -- run job
+ local job_index = state.finished_count
+ state.running_jobs_indices[job_index] = job_index
+ if job_func then
+ if curdir then
+ os.cd(curdir)
+ end
+ state.finished_count = state.finished_count + 1
+ job_func(job_index, total, {progress = state.progress_wrapper})
+ end
+ state.running_jobs_indices[job_index] = nil
+ end,
+ catch
+ {
+ function (errors)
+ -- stop timer and disable show waitchars first
+ state.stop = true
+
+ -- remove wait charactor
+ if state.show_progress then
+ _print_backchars(state.backnum)
+ state.progress_helper:stop()
+ end
+
+ -- we need re-throw this errors outside scheduler
+ state.abort = true
+ if state.abort_errors == nil then
+ state.abort_errors = errors
+ end
+
+ -- kill all waited objects in this group
+ local waitobjs = scheduler.co_group_waitobjs(state.group_name)
+ if waitobjs:size() > 0 then
+ for _, obj in waitobjs:keys() do
+ -- TODO, kill pipe is not supported now
+ if obj.kill then
+ obj:kill()
+ end
+ end
+ end
+ end
+ },
+ finally
+ {
+ function ()
+ if job then
+ jobs:remove(job)
+ end
+ end
+ }
+ }
+ end
end
-- asynchronous run jobs
@@ -165,7 +239,7 @@ function main(name, jobs, opt)
-- init state
local state = {}
state.total = opt.total or (type(jobs) == "table" and jobs:size()) or 1
- state.comax = opt.comax or math.min(state.total, 4)
+ state.comax = opt.comax and tonumber(opt.comax) or math.min(state.total, 4)
state.distcc = opt.distcc
state.timeout = opt.timeout or 500
state.group_name = name
@@ -208,132 +282,18 @@ function main(name, jobs, opt)
end
-- run jobs
- local index = 0
- local abort = false
- local abort_errors
- local job_pending
- local total = state.total
- local comax = state.comax
- local distcc = state.distcc
- local group_name = state.group_name
- local jobs_cb = state.jobs_cb
- local jobs = state.jobs
- while index < total do
- scheduler.co_group_begin(group_name, function (co_group)
- local freemax = comax - #co_group
- local local_max = math.min(index + freemax, total)
- local total_max = local_max
- if distcc then
- total_max = math.min(index + freemax + distcc:freejobs(), total)
- end
- local jobfunc = jobs_cb
- while index < total_max do
-
- -- uses job pool?
- local job
- local jobname
- local distccjob = false
- if not jobs_cb then
-
- -- get free job
- job = job_pending and job_pending or jobs:getfree()
- if not job then
- break
- end
-
- -- we can only continue to run the job with distcc if local jobs are full
- if distcc and index >= local_max then
- if job.distcc then
- distccjob = true
- else
- job_pending = job
- break
- end
- end
-
- -- get run function
- jobfunc = job.run
- jobname = job.name
- job_pending = nil
- else
- jobname = tostring(index)
- end
-
- -- start this job
- index = index + 1
- scheduler.co_start_withopt({name = name .. '/' .. jobname, isolate = opt.isolate}, function(i)
- try
- {
- function()
- if state.stop then
- return
- end
- if distcc then
- local co_running = scheduler.co_running()
- if co_running then
- co_running:data_set("distcc.distccjob", distccjob)
- end
- end
- state.running_jobs_indices[i] = i
- if jobfunc then
- if opt.curdir then
- os.cd(opt.curdir)
- end
- state.count = state.count + 1
- jobfunc(i, total, {progress = state.progress_wrapper})
- end
- state.running_jobs_indices[i] = nil
- end,
- catch
- {
- function (errors)
-
- -- stop timer and disable show waitchars first
- state.stop = true
-
- -- remove wait charactor
- if state.show_progress then
- _print_backchars(state.backnum)
- state.progress_helper:stop()
- end
-
- -- we need re-throw this errors outside scheduler
- abort = true
- if abort_errors == nil then
- abort_errors = errors
- end
-
- -- kill all waited objects in this group
- local waitobjs = scheduler.co_group_waitobjs(group_name)
- if waitobjs:size() > 0 then
- for _, obj in waitobjs:keys() do
- -- TODO, kill pipe is not supported now
- if obj.kill then
- obj:kill()
- end
- end
- end
- end
- },
- finally
- {
- function ()
- if job then
- jobs:remove(job)
- end
- end
- }
- }
- end, index)
- end
- end)
-
- -- wait for free jobs
- scheduler.co_group_wait(group_name, {limit = 1})
- end
+ state.abort = false
+ state.abort_errors = nil
+ state.finished_count = 0
+ state.curdir = opt.curdir
+ scheduler.co_group_begin(state.group_name, function (co_group)
+ for id = 1, math.min(state.total, state.comax) do
+ scheduler.co_start_withopt({name = name .. '/' .. tostring(id), isolate = opt.isolate}, _comsume_jobs_loop, state)
+ end
+ end)
-- wait all jobs exited
- scheduler.co_group_wait(group_name)
+ scheduler.co_group_wait(state.group_name)
-- wait timer job exited
if group_timer then
@@ -354,7 +314,7 @@ function main(name, jobs, opt)
-- do exit callback
if opt.on_exit then
- opt.on_exit(abort_errors)
+ opt.on_exit(state.abort_errors)
end
-- re-throw abort errors
@@ -363,7 +323,7 @@ function main(name, jobs, opt)
-- because his causes a direct exit from the entire runloop and
-- a quick escape from nested try-catch blocks and coroutines groups.
-- so we can not catch runjobs errors, e.g. build fails
- if abort then
- raise(abort_errors)
+ if state.abort then
+ raise(state.abort_errors)
end
end