diff options
| author | ruki <[email protected]> | 2022-05-21 22:55:22 +0800 |
|---|---|---|
| committer | ruki <[email protected]> | 2022-05-21 22:55:22 +0800 |
| commit | 593e92dc0691eaaf3e1217fcf47161dc856af05e (patch) | |
| tree | d8dc2099c3d00b5f2ce718c96e4090e80cfda033 /xmake/modules/private/service/remote_cache | |
| parent | bcde9c9f0f7b2f51e20ccf1a4fccb08dc0b2f9e4 (diff) | |
add existinfo
Diffstat (limited to 'xmake/modules/private/service/remote_cache')
3 files changed, 74 insertions, 0 deletions
diff --git a/xmake/modules/private/service/remote_cache/client.lua b/xmake/modules/private/service/remote_cache/client.lua index 98b715a68..969d3c654 100644 --- a/xmake/modules/private/service/remote_cache/client.lua +++ b/xmake/modules/private/service/remote_cache/client.lua @@ -25,6 +25,7 @@ import("core.base.socket") import("core.base.option") import("core.base.hashset") import("core.base.scheduler") +import("core.base.bloom_filter") import("core.project.config", {alias = "project_config"}) import("lib.detect.find_tool") import("private.service.client_config", {alias = "config"}) @@ -272,6 +273,47 @@ function remote_cache_client:cacheinfo(cachekey) return cacheinfo end +-- get the exist info of cache in server +function remote_cache_client:existinfo() + assert(self:is_connected(), "%s: has been not connected!", self) + local addr = self:addr() + local port = self:port() + local sock = assert(self:_sock_open(), "open socket failed!") + local session_id = self:session_id() + local errors + local existinfo + dprint("%s: get exist info in %s:%d ..", self, addr, port) + local stream = socket_stream(sock) + if stream:send_msg(message.new_existinfo(session_id, "objectfiles", {token = self:token()})) and stream:flush() then + local data = stream:recv_data() + if data then + local msg = stream:recv_msg() + if msg then + dprint(msg:body()) + if msg:success() then + local count = msg:body().count + if count and count > 0 then + local filter = bloom_filter.new() + filter:data_set(data) + existinfo = filter + end + else + errors = msg:errors() + end + end + else + errors = "recv exist info failed" + end + end + self:_sock_close(sock) + if existinfo then + dprint("%s: get exist info ok!", self) + else + dprint("%s: get exist info failed in %s:%d, %s", self, addr, port, errors or "unknown") + end + return existinfo +end + -- clean server files function remote_cache_client:clean() assert(self:is_connected(), "%s: has been not connected!", self) @@ -456,6 +498,7 @@ end function singleton() local instance = _g.singleton if not instance then + config.load() instance = new() _g.singleton = instance end diff --git a/xmake/modules/private/service/remote_cache/server.lua b/xmake/modules/private/service/remote_cache/server.lua index 42d0046ad..9ee83eb5b 100644 --- a/xmake/modules/private/service/remote_cache/server.lua +++ b/xmake/modules/private/service/remote_cache/server.lua @@ -92,6 +92,8 @@ function remote_cache_server:_on_handle(stream, msg) session:pull(respmsg) elseif msg:is_fileinfo() then session:fileinfo(respmsg) + elseif msg:is_existinfo() then + session:existinfo(respmsg) elseif msg:is_clean() then session:clean() end diff --git a/xmake/modules/private/service/remote_cache/server_session.lua b/xmake/modules/private/service/remote_cache/server_session.lua index 5a68877e6..38d86e55e 100644 --- a/xmake/modules/private/service/remote_cache/server_session.lua +++ b/xmake/modules/private/service/remote_cache/server_session.lua @@ -26,6 +26,7 @@ import("core.base.global") import("core.base.option") import("core.base.hashset") import("core.base.scheduler") +import("core.base.bloom_filter") import("private.service.server_config", {alias = "config"}) import("private.service.message") @@ -118,6 +119,34 @@ function server_session:fileinfo(respmsg) vprint("get cacheinfo(%s)", cachekey) end +-- get exist info +function server_session:existinfo(respmsg) + local body = respmsg:body() + local stream = self:stream() + local cachedir = self:cachedir() + local filter = bloom_filter.new() + local count = 0 + vprint("get existinfo(%s) ..", body.name) + for _, objectfile in ipairs(os.files(path.join(cachedir, "*", "*"))) do + local cachekey = path.basename(objectfile) + if cachekey then + filter:set(cachekey) + count = count + 1 + end + end + if count > 0 then + if not stream:send_data(filter:data(), {compress = true}) then + raise("send data failed!") + end + else + if not stream:send_emptydata() then + raise("send empty data failed!") + end + end + body.count = count + vprint("get existinfo(%s): %d ok", body.name, count) +end + -- clean files function server_session:clean() vprint("%s: clean files in %s ..", self, self:cachedir()) |
