summaryrefslogtreecommitdiff
path: root/xmake/modules/private
diff options
context:
space:
mode:
authorruki <[email protected]>2023-09-12 23:07:33 +0800
committerruki <[email protected]>2023-09-12 23:07:33 +0800
commitedef4b17accd31b07ec0079d305e8f1be8d4fff9 (patch)
treea19d6dc6f760ff4dfedbabe17f0206abcab0cc65 /xmake/modules/private
parentf3f566ae6ffc126153af87c4056f29cfb4579105 (diff)
mark async.runjobs as public api
Diffstat (limited to 'xmake/modules/private')
-rw-r--r--xmake/modules/private/action/build/object.lua2
-rw-r--r--xmake/modules/private/action/require/check.lua2
-rw-r--r--xmake/modules/private/action/require/impl/install_packages.lua2
-rw-r--r--xmake/modules/private/async/runjobs.lua276
-rw-r--r--xmake/modules/private/check/checkers/clang/tidy.lua2
5 files changed, 7 insertions, 277 deletions
diff --git a/xmake/modules/private/action/build/object.lua b/xmake/modules/private/action/build/object.lua
index 5b8188d49..ea6f3035a 100644
--- a/xmake/modules/private/action/build/object.lua
+++ b/xmake/modules/private/action/build/object.lua
@@ -24,7 +24,7 @@ import("core.theme.theme")
import("core.tool.compiler")
import("core.project.depend")
import("private.cache.build_cache")
-import("private.async.runjobs")
+import("async.runjobs")
import("utils.progress")
import("private.service.distcc_build.client", {alias = "distcc_build_client"})
diff --git a/xmake/modules/private/action/require/check.lua b/xmake/modules/private/action/require/check.lua
index cf1e3000d..44851d29d 100644
--- a/xmake/modules/private/action/require/check.lua
+++ b/xmake/modules/private/action/require/check.lua
@@ -22,7 +22,7 @@
import("core.base.option")
import("core.base.task")
import("lib.detect.find_tool")
-import("private.async.runjobs")
+import("async.runjobs")
import("private.action.require.impl.package")
import("private.action.require.impl.repository")
import("private.action.require.impl.environment")
diff --git a/xmake/modules/private/action/require/impl/install_packages.lua b/xmake/modules/private/action/require/impl/install_packages.lua
index 22dcbfe48..6fd659182 100644
--- a/xmake/modules/private/action/require/impl/install_packages.lua
+++ b/xmake/modules/private/action/require/impl/install_packages.lua
@@ -24,7 +24,7 @@ import("core.base.hashset")
import("core.base.scheduler")
import("core.project.project")
import("core.base.tty")
-import("private.async.runjobs")
+import("async.runjobs")
import("utils.progress")
import("net.fasturl")
import("private.action.require.impl.package")
diff --git a/xmake/modules/private/async/runjobs.lua b/xmake/modules/private/async/runjobs.lua
index 98ad63747..ab6776b82 100644
--- a/xmake/modules/private/async/runjobs.lua
+++ b/xmake/modules/private/async/runjobs.lua
@@ -19,280 +19,10 @@
--
-- imports
-import("core.base.scheduler")
-import("utils.progress")
-
--- print back characters
-function _print_backchars(backnum)
- if backnum > 0 then
- local str = ('\b'):rep(backnum) .. (' '):rep(backnum) .. ('\b'):rep(backnum)
- if #str > 0 then
- printf(str)
- end
- end
-end
+import("async.runjobs")
-- asynchronous run jobs
---
--- e.g.
--- runjobs("test", function (index) print("hello") end, {total = 100, comax = 6, timeout = 1000, on_timer = function (running_jobs_indices) end})
--- runjobs("test", function () os.sleep(10000) end, { progress = true })
--- runjobs("test", function () os.sleep(10000) end, { progress = { chars = {'/','\'} } }) -- see module utils.progress
---
--- local jobs = jobpool.new()
--- local root = jobs:addjob("job/root", function (idx, total)
--- print(idx, total)
--- end)
--- for i = 1, 3 do
--- local job = jobs:addjob("job/" .. i, function (idx, total)
--- print(idx, total)
--- end, {rootjob = root})
--- end
--- runjobs("test", jobs, {comax = 6, timeout = 1000, on_timer = function (running_jobs_indices) end})
---
--- distributed build:
--- runjobs("test", jobs, {comax = 6, distcc = distcc_build_client.singleton()}
---
function main(name, jobs, opt)
-
- -- init options
- op = opt or {}
- local total = opt.total or (type(jobs) == "table" and jobs:size()) or 1
- local comax = opt.comax or math.min(total, 4)
- local distcc = opt.distcc
- local timeout = opt.timeout or 500
- local group_name = name
- local jobs_cb = type(jobs) == "function" and jobs or nil
- assert(timeout < 60000, "runjobs: invalid timeout!")
- assert(jobs, "runjobs: no jobs!")
-
- -- show waiting tips?
- local showprogress = io.isatty() and (opt.progress or opt.showtips) -- we need to hide wait characters if is not a tty
- local progress_helper
- local backnum = 0
- if showprogress then
- local opt = nil
- if type(showprogress) == 'table' then opt = showprogress end
- progress_helper = progress.new(nil, opt)
- end
-
- -- isolate environments
- local is_isolated = false
- local co_running = scheduler.co_running()
- if co_running and opt.isolate then
- is_isolated = co_running:is_isolated()
- co_running:isolate(true)
- end
-
- -- run timer
- local stop = false
- local running_jobs_indices = {}
- local group_timer
- if opt.on_timer then
- group_timer = group_name .. "/timer"
- scheduler.co_group_begin(group_timer, function (co_group)
- scheduler.co_start_withopt({name = name .. "/timer", isolate = opt.isolate}, function ()
- while not stop do
- os.sleep(timeout)
- if not stop then
- local indices
- if running_jobs_indices then
- indices = table.keys(running_jobs_indices)
- end
- opt.on_timer(indices)
- end
- end
- end)
- end)
- elseif showprogress then
- group_timer = group_name .. "/timer"
- scheduler.co_group_begin(group_timer, function (co_group)
- scheduler.co_start_withopt({name = name .. "/tips", isolate = opt.isolate}, function ()
- while not stop do
- os.sleep(timeout)
- if not stop then
-
- -- show waitchars
- local tips = nil
- local waitobjs = scheduler.co_group_waitobjs(group_name)
- if waitobjs:size() > 0 then
- local names = {}
- for _, obj in waitobjs:keys() do
- if obj:otype() == scheduler.OT_PROC then
- table.insert(names, obj:name())
- elseif obj:otype() == scheduler.OT_SOCK then
- table.insert(names, "sock")
- elseif obj:otype() == scheduler.OT_PIPE then
- table.insert(names, "pipe")
- end
- end
- names = table.unique(names)
- if #names > 0 then
- names = table.concat(names, ",")
- if #names > 16 then
- names = names:sub(1, 16) .. ".."
- end
- tips = string.format("(%d/%s)", waitobjs:size(), names)
- end
- end
-
- -- print back characters
- progress_helper:clear()
- _print_backchars(backnum)
-
- if tips then
- cprintf("${dim}%s${clear} ", tips)
- backnum = #tips + 1
- end
- progress_helper:write()
- end
- end
- end)
- end)
- end
-
- -- run jobs
- local index = 0
- local count = 0
- local count_as_index = opt.count_as_index
- local priority_prev = 0
- local priority_curr = 0
- local job_pending = nil
- 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 jobname
- local distccjob = false
- if not jobs_cb then
-
- -- get job priority
- local job, priority
- if job_pending then
- job = job_pending
- priority = priority_prev
- else
- job, priority = jobs:pop()
- end
- if not job then
- break
- end
-
- -- priority changed? we need to wait all running jobs exited
- priority_curr = priority or priority_prev
- assert(priority_curr >= priority_prev, "runjobs: invalid priority(%d < %d)!", priority_curr, priority_prev)
- if priority_curr > priority_prev then
- job_pending = job
- 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 distcc then
- local co_running = scheduler.co_running()
- if co_running then
- co_running:data_set("distcc.distccjob", distccjob)
- end
- end
- running_jobs_indices[i] = i
- if jobfunc then
- if opt.curdir then
- os.cd(opt.curdir)
- end
- jobfunc(count_as_index and count or i, total)
- count = count + 1
- end
- running_jobs_indices[i] = nil
- end,
- catch
- {
- function (errors)
-
- -- stop timer and disable show waitchars first
- stop = true
-
- -- remove wait charactor
- if showprogress then
- _print_backchars(backnum)
- progress_helper:stop()
- end
-
- -- do exit callback
- if opt.on_exit then
- opt.on_exit(errors)
- end
-
- -- re-throw this errors and abort scheduler
- raise(errors)
- end
- }
- }
- end, index)
- end
- end)
-
- -- only need one job exited if be same priority
- if priority_curr == priority_prev then
- scheduler.co_group_wait(group_name, {limit = 1})
- else
- -- need to wait all running jobs exited first if be different priority
- scheduler.co_group_wait(group_name)
- priority_prev = priority_curr
- end
- end
-
- -- wait all jobs exited
- scheduler.co_group_wait(group_name)
-
- -- wait timer job exited
- if group_timer then
- stop = true
- scheduler.co_group_wait(group_timer)
- end
-
- -- restore isolated environments
- if co_running and opt.isolate then
- co_running:isolate(is_isolated)
- end
-
- -- remove wait charactor
- if showprogress then
- _print_backchars(backnum)
- progress_helper:stop()
- end
-
- -- do exit callback
- if opt.on_exit then
- opt.on_exit()
- end
+ runjobs(name, jobs, opt)
+ wprint("please use import(\"async.runjobs\") instead of private.async.runjobs!")
end
diff --git a/xmake/modules/private/check/checkers/clang/tidy.lua b/xmake/modules/private/check/checkers/clang/tidy.lua
index 9d735e546..e37893918 100644
--- a/xmake/modules/private/check/checkers/clang/tidy.lua
+++ b/xmake/modules/private/check/checkers/clang/tidy.lua
@@ -24,7 +24,7 @@ import("core.base.task")
import("core.project.config")
import("core.project.project")
import("lib.detect.find_tool")
-import("private.async.runjobs")
+import("async.runjobs")
import("private.action.require.impl.packagenv")
import("private.action.require.impl.install_packages")