summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorruki <[email protected]>2025-11-08 00:44:35 +0800
committerruki <[email protected]>2025-11-08 00:44:35 +0800
commit45f03066b9738062aaa512de7c1c9f0e7c096ec4 (patch)
treeecf48a18381af6002997f96b4b7bec047993088f
parent041d3951b8817b0bc51d1de39b4c7a5da8ba6e36 (diff)
use pipe_event instead of event
-rw-r--r--xmake/core/base/private/async_task.lua104
-rw-r--r--xmake/core/base/private/pipe_event.lua144
-rw-r--r--xmake/core/base/thread.lua36
3 files changed, 231 insertions, 53 deletions
diff --git a/xmake/core/base/private/async_task.lua b/xmake/core/base/private/async_task.lua
index fbb2d8d0b..ca176816e 100644
--- a/xmake/core/base/private/async_task.lua
+++ b/xmake/core/base/private/async_task.lua
@@ -27,6 +27,7 @@ local utils = require("base/utils")
local thread = require("base/thread")
local option = require("base/option")
local path = require("base/path")
+local pipe_event = require("base/private/pipe_event")
-- the task status
local is_stopped = false
@@ -37,8 +38,7 @@ local task_event = nil
local task_queue = nil
local task_mutex = nil
--- object pool for event and sharedata
-local event_pool = {}
+-- object pool for sharedata
local sharedata_pool = {}
function async_task._absolute_dirs(searchdirs)
@@ -54,19 +54,13 @@ function async_task._absolute_dirs(searchdirs)
return dirs
end
--- get event from pool or create new one
function async_task._get_event()
- local event = table.remove(event_pool)
- if not event then
- event = thread.event()
- end
- return event
+ return pipe_event.new("async_task")
end
--- return event to pool
function async_task._put_event(event)
if event then
- table.insert(event_pool, event)
+ event:close()
end
end
@@ -102,10 +96,14 @@ function async_task._loop(event, queue, mutex, is_stopped, is_diagnosis)
local function _restore_thread_objects(cmd)
if cmd.event_data then
- cmd.event = thread._deserialize_object(cmd.event_data)
+ local event, errors = thread._deserialize_object(cmd.event_data)
+ assert(event, errors or "failed to deserialize event")
+ cmd.event = event
end
if cmd.result_data then
- cmd.result = thread._deserialize_object(cmd.result_data)
+ local result, errors = thread._deserialize_object(cmd.result_data)
+ assert(result, errors or "failed to deserialize result")
+ cmd.result = result
end
end
@@ -263,55 +261,71 @@ end
-- post task and wait for result
function async_task._post_task(cmd, is_detach, return_data)
- local cmd_event, cmd_result
+ local cmd_event
+ local cmd_result
- -- create event and result for non-detach mode
+ -- create pipe and result for non-detach mode
if not is_detach then
cmd_event = async_task._get_event()
+ if not cmd_event then
+ return false, "failed to acquire event"
+ end
cmd_result = async_task._get_sharedata()
-
- -- serialize thread objects for passing to worker thread
- cmd.event_data = thread._serialize_object(cmd_event)
cmd.result_data = thread._serialize_object(cmd_result)
+ if not cmd.result_data then
+ async_task._put_sharedata(cmd_result)
+ async_task._put_event(cmd_event)
+ return false, "failed to serialize sharedata"
+ end
+ cmd.event_data = thread._serialize_object(cmd_event)
+ if not cmd.event_data then
+ async_task._put_sharedata(cmd_result)
+ async_task._put_event(cmd_event)
+ return false, "failed to serialize event"
+ end
end
task_mutex:lock()
task_queue:push(cmd)
- local queue_size = task_queue:size()
task_mutex:unlock()
+ task_event:post()
+
if is_detach then
- -- We cache some tasks before executing them to avoid frequent thread switching.
- if queue_size > 10 then
- task_event:post()
- end
return true
- else
- -- wait for completion
- task_event:post()
- cmd_event:wait(-1)
- local result = cmd_result:get()
- async_task._put_event(cmd_event)
- async_task._put_sharedata(cmd_result)
- if result and result.ok then
- if return_data then
- local data = result.data
- if type(data) == "table" then
- return data, #data
- elseif data ~= nil then
- return data
- else
- return nil, 0
- end
+ end
+
+ local wait_ok, wait_errors = cmd_event:wait(-1)
+
+ local result
+ if wait_ok then
+ result = cmd_result:get()
+ end
+ async_task._put_sharedata(cmd_result)
+ async_task._put_event(cmd_event)
+
+ if not wait_ok then
+ return false, wait_errors or "wait event failed"
+ end
+
+ if result and result.ok then
+ if return_data then
+ local data = result.data
+ if type(data) == "table" then
+ return data, #data
+ elseif data ~= nil then
+ return data
else
- return true
- end
- else
- if return_data then
return nil, 0
- else
- return false, result and result.errors or "unknown error"
end
+ else
+ return true
+ end
+ else
+ if return_data then
+ return nil, 0
+ else
+ return false, result and result.errors or "unknown error"
end
end
end
diff --git a/xmake/core/base/private/pipe_event.lua b/xmake/core/base/private/pipe_event.lua
new file mode 100644
index 000000000..20f1373fb
--- /dev/null
+++ b/xmake/core/base/private/pipe_event.lua
@@ -0,0 +1,144 @@
+--!A cross-platform build utility based on Lua
+--
+-- Licensed under the Apache License, Version 2.0 (the "License");
+-- you may not use this file except in compliance with the License.
+-- You may obtain a copy of the License at
+--
+-- http://www.apache.org/licenses/LICENSE-2.0
+--
+-- Unless required by applicable law or agreed to in writing, software
+-- distributed under the License is distributed on an "AS IS" BASIS,
+-- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+-- See the License for the specific language governing permissions and
+-- limitations under the License.
+--
+-- Copyright (C) 2015-present, Xmake Open Source Community.
+--
+-- @author ruki
+-- @file pipe_event.lua
+--
+
+-- define module
+local pipe_event = pipe_event or {}
+local _instance = _instance or {}
+
+-- load modules
+local pipe = require("base/pipe")
+local libc = require("base/libc")
+local bytes = require("base/bytes")
+local table = require("base/table")
+
+function _instance.new(name)
+ local event = table.inherit(_instance)
+ event._PIPE_EVENT = true
+ event._BUFFER = bytes(2)
+ event._NAME = name or "pipe_event"
+ local reader, writer, errors = pipe.openpair("BA")
+ if not reader or not writer then
+ if reader then reader:close() end
+ if writer then writer:close() end
+ return nil, errors or "failed to open pipe"
+ end
+ event._READER = reader
+ event._WRITER = writer
+ event._WRITER_PTR = nil
+ return event
+end
+
+function _instance:name()
+ return self._NAME
+end
+
+function _instance:post()
+ local writer = self._WRITER
+ if not writer then
+ return false, "pipe event writer closed"
+ end
+ local ok, errors = writer:write("1")
+ if ok < 0 then
+ return false, errors or "pipe event post failed"
+ end
+ writer:close()
+ self._WRITER = nil
+ self._WRITER_PTR = nil
+ return true
+end
+
+function _instance:wait(timeout)
+ if not self._READER then
+ return false, "pipe event reader closed"
+ end
+ local events, errors = self._READER:wait(pipe.EV_READ, timeout or -1)
+ if events < 0 then
+ return false, errors
+ end
+ local read, read_errors = self._READER:read(self._BUFFER, 1)
+ if read < 0 then
+ return false, read_errors
+ end
+ return read
+end
+
+function _instance:close()
+ if self._READER then
+ self._READER:close()
+ end
+ if self._WRITER then
+ self._WRITER:close()
+ end
+ self._READER = nil
+ self._WRITER = nil
+ self._WRITER_PTR = nil
+end
+
+-- return pipe cdata for serialization (writer pointer is stable for passing across threads)
+function _instance:cdata()
+ if self._WRITER then
+ return self._WRITER:cdata()
+ end
+ return self._WRITER_PTR
+end
+
+function _instance:__gc()
+ self:close()
+end
+
+function _instance:_serialize()
+ if not self._WRITER and self._WRITER_PTR then
+ return {ptr = self._WRITER_PTR, name = self:name()}
+ end
+ if not self._WRITER then
+ return nil
+ end
+ local ptr = libc.dataptr(self._WRITER:cdata(), {ffi = false})
+ if not ptr then
+ return nil
+ end
+ self._WRITER._PIPE = nil
+ self._WRITER = nil
+ self._WRITER_PTR = ptr
+ return {ptr = ptr, name = self:name()}
+end
+
+function _instance:_deserialize(data)
+ if not data or not data.ptr then
+ return false, "invalid pipe event data"
+ end
+ self:close()
+ local writer = pipe.new(libc.ptraddr(data.ptr, {ffi = false}))
+ if not writer then
+ return false, "invalid pipe pointer"
+ end
+ self._NAME = data.name or self._NAME or "pipe_event"
+ self._WRITER = writer
+ self._WRITER_PTR = data.ptr
+ return true
+end
+
+function pipe_event.new(name)
+ return _instance.new(name)
+end
+
+return pipe_event
+
+
diff --git a/xmake/core/base/thread.lua b/xmake/core/base/thread.lua
index 66e6ba2af..26ac672c3 100644
--- a/xmake/core/base/thread.lua
+++ b/xmake/core/base/thread.lua
@@ -28,14 +28,15 @@ local _queue = _queue or {}
local _sharedata = _sharedata or {}
-- load modules
-local io = require("base/io")
-local libc = require("base/libc")
-local pipe = require("base/pipe")
-local bytes = require("base/bytes")
-local table = require("base/table")
-local string = require("base/string")
-local scheduler = require("base/scheduler")
-local sandbox = require("sandbox/sandbox")
+local io = require("base/io")
+local libc = require("base/libc")
+local pipe = require("base/pipe")
+local pipe_event = require("base/private/pipe_event")
+local bytes = require("base/bytes")
+local table = require("base/table")
+local string = require("base/string")
+local scheduler = require("base/scheduler")
+local sandbox = require("sandbox/sandbox")
-- the thread status
thread.STATUS_READY = 1
@@ -810,6 +811,15 @@ function thread._serialize_object(obj)
result.sharedata = true
result.name = obj:name()
result.caddr = libc.dataptr(obj:cdata(), {ffi = false})
+ elseif obj._PIPE_EVENT then
+ local data = obj:_serialize()
+ if not data then
+ return nil
+ end
+ result.pipe_event = true
+ result.ptr = data.ptr
+ result.name = data.name
+ result.caddr = data.ptr
else
return nil
end
@@ -836,6 +846,16 @@ function thread._deserialize_object(data)
return _queue.new(data.name, cdata)
elseif data.sharedata then
return _sharedata.new(data.name, cdata)
+ elseif data.pipe_event then
+ local event = pipe_event.new(data.name)
+ if not event then
+ return nil, "failed to create pipe event"
+ end
+ local ok, errors = event:_deserialize({ptr = data.ptr, name = data.name})
+ if not ok then
+ return nil, errors
+ end
+ return event
end
return nil