summaryrefslogtreecommitdiff
path: root/xmake/modules
diff options
context:
space:
mode:
authorruki <[email protected]>2022-06-29 00:53:40 +0800
committerruki <[email protected]>2022-06-29 00:53:40 +0800
commit7e2f4dc36f81a3bbd0c188591886e697c2320f5b (patch)
tree25945ddcad7654f536878e47c3d5a7e9733d48c7 /xmake/modules
parent3f131675ed167ae3eed2afc0bfb69b633a8080a6 (diff)
fix remote build
Diffstat (limited to 'xmake/modules')
-rw-r--r--xmake/modules/private/service/message.lua16
-rw-r--r--xmake/modules/private/service/remote_build/client.lua14
-rw-r--r--xmake/modules/private/service/remote_build/server_session.lua24
3 files changed, 48 insertions, 6 deletions
diff --git a/xmake/modules/private/service/message.lua b/xmake/modules/private/service/message.lua
index 7338b06ac..12792a129 100644
--- a/xmake/modules/private/service/message.lua
+++ b/xmake/modules/private/service/message.lua
@@ -37,6 +37,7 @@ message.CODE_PULL = 9 -- pull the given file from server
message.CODE_PUSH = 10 -- push the given file to server
message.CODE_FILEINFO = 11 -- get the given file info in server
message.CODE_EXISTINFO = 12 -- get exists info in server (use bloom filter)
+message.CODE_END = 13 -- end
-- init message
function message:init(body)
@@ -113,6 +114,11 @@ function message:is_existinfo()
return self:code() == message.CODE_EXISTINFO
end
+-- is end message?
+function message:is_end()
+ return self:code() == message.CODE_END
+end
+
-- get user authorization
function message:token()
return self:body().token
@@ -300,6 +306,16 @@ function new_existinfo(session_id, name, opt)
})
end
+-- new end message
+function new_end(session_id, opt)
+ opt = opt or {}
+ return _new({
+ code = message.CODE_END,
+ session_id = session_id,
+ token = opt.token
+ })
+end
+
function main(body)
return _new(body)
end
diff --git a/xmake/modules/private/service/remote_build/client.lua b/xmake/modules/private/service/remote_build/client.lua
index f76699790..fe964f054 100644
--- a/xmake/modules/private/service/remote_build/client.lua
+++ b/xmake/modules/private/service/remote_build/client.lua
@@ -274,7 +274,10 @@ function remote_build_client:runcmd(program, argv)
local stream = socket_stream(sock)
if stream:send_msg(message.new_runcmd(session_id, program, argv, {token = self:token()})) and stream:flush() then
local stdin_opt = {stop = false}
- scheduler.co_start(self._read_stdin, self, stream, stdin_opt)
+ local group_name = "remote_build/runcmd"
+ scheduler.co_group_begin(group_name, function (co_group)
+ scheduler.co_start(self._read_stdin, self, stream, stdin_opt)
+ end)
while true do
local msg = stream:recv_msg()
if msg then
@@ -291,6 +294,9 @@ function remote_build_client:runcmd(program, argv)
errors = string.format("recv output data(%d) failed!", msg:body().size)
break
end
+ elseif msg:is_end() then
+ ok = true
+ break
else
if msg:success() then
ok = true
@@ -304,6 +310,7 @@ function remote_build_client:runcmd(program, argv)
end
end
stdin_opt.stop = true
+ scheduler.co_group_wait(group_name)
end
if #leftstr > 0 then
cprint(leftstr)
@@ -503,9 +510,12 @@ function remote_build_client:_read_stdin(stream, opt)
os.sleep(500)
end
end
+ -- say bye
+ if stream:send_msg(message.new_end({token = self:token()})) then
+ stream:flush()
+ end
end
-
function remote_build_client:__tostring()
return "<remote_build_client>"
end
diff --git a/xmake/modules/private/service/remote_build/server_session.lua b/xmake/modules/private/service/remote_build/server_session.lua
index 340d9934d..d75e0668e 100644
--- a/xmake/modules/private/service/remote_build/server_session.lua
+++ b/xmake/modules/private/service/remote_build/server_session.lua
@@ -202,8 +202,11 @@ function server_session:runcmd(respmsg)
local stdout_rpipeopt = {rpipe = stdout_rpipe, stop = false}
-- read and write pipe
- scheduler.co_start(self._write_pipe, self, stdin_wpipeopt)
- scheduler.co_start(self._read_pipe, self, stdout_rpipeopt)
+ local group_name = "remote_build/runcmd"
+ scheduler.co_group_begin(group_name, function (co_group)
+ scheduler.co_start(self._write_pipe, self, stdin_wpipeopt)
+ scheduler.co_start(self._read_pipe, self, stdout_rpipeopt)
+ end)
-- run program
os.execv(program, argv, {curdir = self:sourcedir(), stdout = stdout_wpipe, stdin = stdin_rpipe, envs = {XMAKE_IN_SERVICE = "true"}})
@@ -213,6 +216,9 @@ function server_session:runcmd(respmsg)
stdin_wpipe:close()
stdout_rpipeopt.stop = true
stdout_wpipe:close()
+
+ -- wait pipes exits
+ scheduler.co_group_wait(group_name)
vprint("%s: run command ok", self)
end
@@ -268,7 +274,7 @@ function server_session:_ensure_sourcedir()
end
end
--- write data from pipe
+-- write data to pipe
function server_session:_write_pipe(opt)
local buff = bytes(256)
local wpipe = opt.wpipe
@@ -308,7 +314,7 @@ function server_session:_read_pipe(opt)
end
end
if not self:_send_data(data) then
- break;
+ break
end
elseif real == 0 then
if rpipe:wait(pipe.EV_READ, -1) < 0 then
@@ -322,6 +328,8 @@ function server_session:_read_pipe(opt)
if #leftstr > 0 then
cprint(leftstr)
end
+ -- say end to client
+ self:_send_end()
vprint("%s: %s: read data end", self, rpipe)
end
@@ -344,6 +352,14 @@ function server_session:_send_data(data)
end
end
+-- send end to stream
+function server_session:_send_end()
+ local stream = self:stream()
+ if stream:send_msg(message.new_end(self:id())) then
+ return stream:flush()
+ end
+end
+
-- recv syncfiles
function server_session:_recv_syncfiles(manifest, outputdir)
local stream = self:stream()