diff options
| author | ruki <[email protected]> | 2025-11-08 00:44:35 +0800 |
|---|---|---|
| committer | ruki <[email protected]> | 2025-11-08 00:44:35 +0800 |
| commit | 45f03066b9738062aaa512de7c1c9f0e7c096ec4 (patch) | |
| tree | ecf48a18381af6002997f96b4b7bec047993088f | |
| parent | 041d3951b8817b0bc51d1de39b4c7a5da8ba6e36 (diff) | |
use pipe_event instead of event
| -rw-r--r-- | xmake/core/base/private/async_task.lua | 104 | ||||
| -rw-r--r-- | xmake/core/base/private/pipe_event.lua | 144 | ||||
| -rw-r--r-- | xmake/core/base/thread.lua | 36 |
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 |
