summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorruki <[email protected]>2022-05-11 23:38:42 +0800
committerruki <[email protected]>2022-05-11 23:38:42 +0800
commitfc4f86e88a8fa170b3783ff82c57549ff6dcb210 (patch)
tree11a4fcedca2914aa16b325dd8159f28fdd56ac5d
parentcca1e7a2e77422b1bd94157ddbeba7e955d95117 (diff)
add distributed_ uild config
-rw-r--r--xmake/modules/private/service/config.lua14
-rw-r--r--xmake/modules/private/service/distributed_build/client.lua198
-rw-r--r--xmake/modules/private/service/distributed_build/server.lua127
-rw-r--r--xmake/modules/private/service/distributed_build/session.lua127
4 files changed, 466 insertions, 0 deletions
diff --git a/xmake/modules/private/service/config.lua b/xmake/modules/private/service/config.lua
index e2b19bdef..00fc61c92 100644
--- a/xmake/modules/private/service/config.lua
+++ b/xmake/modules/private/service/config.lua
@@ -58,7 +58,21 @@ function _generate_configfile()
-- with token authorization
token = token
}
+ },
+ distributed_build = {
+ server = {
+ listen = "0.0.0.0:9692",
+ workdir = path.join(servicedir, "distributed_build"),
+ },
+ client = {
+ -- without authorization: "127.0.0.1:9691"
+ -- with user authorization: "[email protected]:9691"
+ connect = "127.0.0.1:9692",
+ -- with token authorization
+ token = token
+ }
}
+
}
save(configs)
end
diff --git a/xmake/modules/private/service/distributed_build/client.lua b/xmake/modules/private/service/distributed_build/client.lua
new file mode 100644
index 000000000..da263c796
--- /dev/null
+++ b/xmake/modules/private/service/distributed_build/client.lua
@@ -0,0 +1,198 @@
+--!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, TBOOX Open Source Group.
+--
+-- @author ruki
+-- @file client.lua
+--
+
+-- imports
+import("core.base.bytes")
+import("core.base.base64")
+import("core.base.socket")
+import("core.base.option")
+import("core.base.scheduler")
+import("core.project.config", {alias = "project_config"})
+import("lib.detect.find_tool")
+import("private.service.config")
+import("private.service.message")
+import("private.service.client")
+import("private.service.stream", {alias = "socket_stream"})
+
+-- define module
+local distributed_build_client = distributed_build_client or client()
+local super = distributed_build_client:class()
+
+-- init client
+function distributed_build_client:init()
+ super.init(self)
+
+ -- init address
+ local address = assert(config.get("distributed_build.client.connect"), "config(distributed_build.client.connect): not found!")
+ super.address_set(self, address)
+
+ -- get project directory
+ local projectdir = os.projectdir()
+ local projectfile = os.projectfile()
+ if projectfile and os.isfile(projectfile) and projectdir then
+ self._PROJECTDIR = projectdir
+ self._WORKDIR = path.join(project_config.directory(), "distributed_build")
+ else
+ raise("we need enter a project directory with xmake.lua first!")
+ end
+
+ -- check environment
+ environment.check(false)
+end
+
+-- connect to the remote server
+function distributed_build_client:connect()
+ if self:is_connected() then
+ print("%s: has been connected!", self)
+ return
+ end
+
+ -- we need user authorization?
+ local token = config.get("distributed_build.client.token")
+ if not token and self:user() then
+
+ -- get user password
+ cprint("Please input user ${bright}%s${clear} password:", self:user())
+ io.flush()
+ local pass = (io.read() or ""):trim()
+ assert(pass ~= "", "password is empty!")
+
+ -- compute user authorization
+ token = base64.encode(self:user() .. ":" .. pass)
+ token = hash.md5(bytes(token))
+ end
+
+ -- do connect
+ local addr = self:addr()
+ local port = self:port()
+ local sock = assert(socket.connect(addr, port), "%s: server unreachable!", self)
+ local session_id = self:session_id()
+ local ok = false
+ local errors
+ print("%s: connect %s:%d ..", self, addr, port)
+ if sock then
+ local stream = socket_stream(sock)
+ if stream:send_msg(message.new_connect(session_id, {token = token})) and stream:flush() then
+ local msg = stream:recv_msg()
+ if msg then
+ vprint(msg:body())
+ if msg:success() then
+ ok = true
+ else
+ errors = msg:errors()
+ end
+ end
+ end
+ end
+ if ok then
+ print("%s: connected!", self)
+ else
+ print("%s: connect %s:%d failed, %s", self, addr, port, errors or "unknown")
+ end
+
+ -- update status
+ local status = self:status()
+ status.addr = addr
+ status.port = port
+ status.token = token
+ status.connected = ok
+ status.session_id = session_id
+ self:status_save()
+
+ -- sync files
+ if ok then
+ self:sync()
+ end
+end
+
+-- disconnect server
+function distributed_build_client:disconnect()
+ if not self:is_connected() then
+ print("%s: has been disconnected!", self)
+ return
+ end
+ local addr = self:addr()
+ local port = self:port()
+ local sock = socket.connect(addr, port)
+ local session_id = self:session_id()
+ local errors
+ local ok = false
+ print("%s: disconnect %s:%d ..", self, addr, port)
+ if sock then
+ local stream = socket_stream(sock)
+ if stream:send_msg(message.new_disconnect(session_id, {token = self:token()})) and stream:flush() then
+ local msg = stream:recv_msg()
+ if msg then
+ vprint(msg:body())
+ if msg:success() then
+ ok = true
+ else
+ errors = msg:errors()
+ end
+ end
+ end
+ else
+ -- server unreachable, but we still disconnect it.
+ wprint("%s: server unreachable!", self)
+ ok = true
+ end
+ if ok then
+ print("%s: disconnected!", self)
+ else
+ print("%s: disconnect %s:%d failed, %s", self, addr, port, errors or "unknown")
+ end
+
+ -- update status
+ local status = self:status()
+ status.token = nil
+ status.connected = not ok
+ self:status_save()
+end
+
+-- get class
+function distributed_build_client:class()
+ return distributed_build_client
+end
+
+-- is connected? we cannot depend on client:init when run action
+function is_connected()
+ -- the current process is in service? we cannot enable it
+ if os.getenv("XMAKE_IN_SERVICE") then
+ return false
+ end
+ local projectdir = os.projectdir()
+ local projectfile = os.projectfile()
+ if projectfile and os.isfile(projectfile) and projectdir then
+ local workdir = path.join(project_config.directory(), "distributed_build")
+ local statusfile = path.join(workdir, "status.txt")
+ if os.isfile(statusfile) then
+ local status = io.load(statusfile)
+ if status and status.connected then
+ return true
+ end
+ end
+ end
+end
+
+function main()
+ local instance = distributed_build_client()
+ instance:init()
+ return instance
+end
diff --git a/xmake/modules/private/service/distributed_build/server.lua b/xmake/modules/private/service/distributed_build/server.lua
new file mode 100644
index 000000000..88b3d689c
--- /dev/null
+++ b/xmake/modules/private/service/distributed_build/server.lua
@@ -0,0 +1,127 @@
+--!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, TBOOX Open Source Group.
+--
+-- @author ruki
+-- @file server.lua
+--
+
+-- imports
+import("private.service.config")
+import("private.service.message")
+import("private.service.server")
+import("private.service.stream", {alias = "socket_stream"})
+import("private.service.distributed_build.session", {alias = "server_session"})
+import("lib.detect.find_tool")
+
+-- define module
+local distributed_build_server = distributed_build_server or server()
+local super = distributed_build_server:class()
+
+-- init server
+function distributed_build_server:init(daemon)
+ super.init(self, daemon)
+ if self:daemon() then
+ config.load()
+ end
+
+ -- init address
+ local address = assert(config.get("distributed_build.server.listen"), "config(distributed_build.server.listen): not found!")
+ super.address_set(self, address)
+
+ -- init handler
+ super.handler_set(self, self._on_handle)
+
+ -- init sessions
+ self._SESSIONS = {}
+end
+
+-- get class
+function distributed_build_server:class()
+ return distributed_build_server
+end
+
+-- on handle message
+function distributed_build_server:_on_handle(stream, msg)
+ local session_id = msg:session_id()
+ local session = self:_session(session_id)
+ vprint("%s: %s: <session %s>: on handle message(%d)", self, stream:sock(), session_id, msg:code())
+ vprint(msg:body())
+ session:stream_set(stream)
+ local respmsg = msg:clone()
+ local session_errs
+ local session_ok = try
+ {
+ function()
+ if self:need_verfiy() then
+ local ok, errors = self:verify_user(msg:token(), stream:sock():peeraddr())
+ if not ok then
+ session_errs = errors
+ return false
+ end
+ end
+ if msg:is_connect() then
+ session:open()
+ elseif msg:is_disconnect() then
+ session:close()
+ self._SESSIONS[session_id] = nil
+ else
+ assert(session:is_connected(), "session has not been connected!")
+ -- TODO
+ end
+ return true
+ end,
+ catch
+ {
+ function (errors)
+ if errors then
+ session_errs = tostring(errors)
+ vprint(session_errs)
+ end
+ end
+ }
+ }
+ respmsg:status_set(session_ok)
+ if not session_ok and session_errs then
+ respmsg:errors_set(session_errs)
+ end
+ local ok = stream:send_msg(respmsg) and stream:flush()
+ vprint("%s: %s: <session %s>: send %s", self, stream:sock(), session_id, ok and "ok" or "failed")
+end
+
+-- get session
+function distributed_build_server:_session(session_id)
+ local session = self._SESSIONS[session_id]
+ if not session then
+ session = server_session(session_id)
+ self._SESSIONS[session_id] = session
+ end
+ return session
+end
+
+-- close session
+function distributed_build_server:_session_close(session_id)
+ self._SESSIONS[session_id] = nil
+end
+
+function distributed_build_server:__tostring()
+ return "<distributed_build_server>"
+end
+
+function main(daemon)
+ local instance = distributed_build_server()
+ instance:init(daemon ~= nil)
+ return instance
+end
diff --git a/xmake/modules/private/service/distributed_build/session.lua b/xmake/modules/private/service/distributed_build/session.lua
new file mode 100644
index 000000000..7d64948c6
--- /dev/null
+++ b/xmake/modules/private/service/distributed_build/session.lua
@@ -0,0 +1,127 @@
+--!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, TBOOX Open Source Group.
+--
+-- @author ruki
+-- @file session.lua
+--
+
+-- imports
+import("core.base.pipe")
+import("core.base.bytes")
+import("core.base.object")
+import("core.base.global")
+import("core.base.option")
+import("core.base.hashset")
+import("core.base.scheduler")
+import("private.service.config")
+import("private.service.message")
+
+-- define module
+local session = session or object()
+
+-- init session
+function session:init(session_id)
+ self._ID = session_id
+end
+
+-- get session id
+function session:id()
+ return self._ID
+end
+
+-- open session
+function session:open()
+ if self:is_connected() then
+ return
+ end
+
+ -- update status
+ local status = self:status()
+ status.connected = true
+ status.session_id = self:id()
+ self:status_save()
+end
+
+-- close session
+function session:close()
+ if not self:is_connected() then
+ return
+ end
+
+ -- update status
+ local status = self:status()
+ status.connected = false
+ status.session_id = self:id()
+ self:status_save()
+end
+
+-- set stream
+function session:stream_set(stream)
+ self._STREAM = stream
+end
+
+-- get stream
+function session:stream()
+ return self._STREAM
+end
+
+-- get work directory
+function session:workdir()
+ local workdir = config.get("distributed_build.server.workdir")
+ if not workdir then
+ workdir = path.join(global.directory(), "service", "distributed_build")
+ end
+ return path.join(workdir, "sessons", self:id())
+end
+
+-- is connected?
+function session:is_connected()
+ return self:status().connected
+end
+
+-- get the status
+function session:status()
+ local status = self._STATUS
+ local statusfile = self:statusfile()
+ if not status then
+ if os.isfile(statusfile) then
+ status = io.load(statusfile)
+ end
+ status = status or {}
+ self._STATUS = status
+ end
+ return status
+end
+
+-- save status
+function session:status_save()
+ io.save(self:statusfile(), self:status())
+end
+
+-- get status file
+function session:statusfile()
+ return path.join(self:workdir(), "status.txt")
+end
+
+function session:__tostring()
+ return string.format("<session %s>", self:id())
+end
+
+function main(session_id)
+ local instance = session()
+ instance:init(session_id)
+ return instance
+end