From 57f08426ee9f7bd44a8edc4231b94625e2b580d3 Mon Sep 17 00:00:00 2001 From: Cloud Wu Date: Thu, 2 Nov 2017 22:31:54 +0800 Subject: [PATCH] concurrent cluster link, fix issue #757 --- service/clusterd.lua | 23 +++++++++++++++++++++-- 1 file changed, 21 insertions(+), 2 deletions(-) diff --git a/service/clusterd.lua b/service/clusterd.lua index 151e83ac..559317ae 100644 --- a/service/clusterd.lua +++ b/service/clusterd.lua @@ -14,7 +14,18 @@ local function read_response(sock) return cluster.unpackresponse(msg) -- session, ok, data, padding end +local connecting = {} + local function open_channel(t, key) + local ct = connecting[key] + if ct then + local co = coroutine.running() + table.insert(ct, co) + skynet.wait(co) + return assert(ct.channel) + end + ct = {} + connecting[key] = ct local host, port = string.match(node_address[key], "([^:]+):(.*)$") local c = sc.channel { host = host, @@ -22,8 +33,16 @@ local function open_channel(t, key) response = read_response, nodelay = true, } - assert(c:connect(true)) - t[key] = c + local succ, err = pcall(c.connect, c, true) + if succ then + t[key] = c + ct.channel = c + end + connecting[key] = nil + for _, co in ipairs(ct) do + skynet.wakeup(co) + end + assert(succ, err) return c end