diff options
| author | ruki <[email protected]> | 2022-06-29 00:53:40 +0800 |
|---|---|---|
| committer | ruki <[email protected]> | 2022-06-29 00:53:40 +0800 |
| commit | 7e2f4dc36f81a3bbd0c188591886e697c2320f5b (patch) | |
| tree | 25945ddcad7654f536878e47c3d5a7e9733d48c7 /xmake/modules | |
| parent | 3f131675ed167ae3eed2afc0bfb69b633a8080a6 (diff) | |
fix remote build
Diffstat (limited to 'xmake/modules')
| -rw-r--r-- | xmake/modules/private/service/message.lua | 16 | ||||
| -rw-r--r-- | xmake/modules/private/service/remote_build/client.lua | 14 | ||||
| -rw-r--r-- | xmake/modules/private/service/remote_build/server_session.lua | 24 |
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() |
