From 13d0ec0a9cafbbd6d8a09b9f050575b0eb915f28 Mon Sep 17 00:00:00 2001 From: Joshua T Corbin Date: Thu, 22 Oct 2015 14:36:28 -0700 Subject: [PATCH 1/7] ServiceDispatchHandler: factor out .ensurePeerConnected --- service-proxy.js | 17 +++++++++++------ 1 file changed, 11 insertions(+), 6 deletions(-) diff --git a/service-proxy.js b/service-proxy.js index aa0d05d..35e5e7b 100644 --- a/service-proxy.js +++ b/service-proxy.js @@ -472,9 +472,7 @@ function refreshServicePeer(serviceName, hostPort) { // The old way: fully connect every egress to all affine peers. self.addPeerIndex(serviceName, hostPort); var peer = self.getServicePeer(serviceName, hostPort); - if (!peer.isConnected('out')) { - peer.connectTo(); - } + self.ensurePeerConnected(peer, 'service peer refresh'); }; ServiceDispatchHandler.prototype.addPeerIndex = @@ -494,6 +492,15 @@ function deletePeerIndex(serviceName, hostPort) { deleteIndexEntry(self.knownPeers, hostPort, serviceName); }; +ServiceDispatchHandler.prototype.ensurePeerConnected = +function ensurePeerConnected(peer, reason) { + if (peer.isConnected('out')) { + return; + } + + peer.connectTo(); +}; + ServiceDispatchHandler.prototype.computePartialRange = function computePartialRange(serviceName, hostPort) { var self = this; @@ -600,9 +607,7 @@ function refreshServicePeerPartially(serviceName, hostPort) { for (var i = 0; i < range.affineWorkers.length; i++) { peer = self._getServicePeer(chan, range.affineWorkers[i]); - if (!peer.isConnected('out')) { - peer.connectTo(); - } + self.ensurePeerConnected(peer, 'service peer affinity refresh'); } // TODO Drop peers that no longer have affinity for this service, such From 9f438bf90d911ab5f5b7bec68d4289edd5e7990a Mon Sep 17 00:00:00 2001 From: Joshua T Corbin Date: Mon, 26 Oct 2015 14:30:18 -0700 Subject: [PATCH 2/7] ServiceDispatchHandler#refreshServicePeerPartially: improve info change log --- service-proxy.js | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/service-proxy.js b/service-proxy.js index 35e5e7b..8b4e4f7 100644 --- a/service-proxy.js +++ b/service-proxy.js @@ -596,9 +596,9 @@ function refreshServicePeerPartially(serviceName, hostPort) { })); } - self.logger.info('Refreshing service peer affinity', self.extendLogInfo({ + self.logger.info('implementing affinity change', self.extendLogInfo({ serviceName: serviceName, - serviceHostPort: hostPort, + newPeer: hostPort, partialRange: range })); From 2807527d69d35d9589b3be7089ecd38b7caac78d Mon Sep 17 00:00:00 2001 From: Joshua T Corbin Date: Mon, 26 Oct 2015 14:34:50 -0700 Subject: [PATCH 3/7] ServiceDispatchHandler#refreshServicePeerPartially: separate affineWorkers from toConnect --- service-proxy.js | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) diff --git a/service-proxy.js b/service-proxy.js index 8b4e4f7..776c974 100644 --- a/service-proxy.js +++ b/service-proxy.js @@ -594,19 +594,29 @@ function refreshServicePeerPartially(serviceName, hostPort) { advertisingPeer: hostPort, partialRange: range })); + + } + + var toConnect = []; + var i; + var worker; + for (i = 0; i < range.affineWorkers.length; i++) { + worker = range.affineWorkers[i]; + toConnect.push(worker); } self.logger.info('implementing affinity change', self.extendLogInfo({ serviceName: serviceName, newPeer: hostPort, - partialRange: range + partialRange: range, + toConnect: toConnect })); self.addPeerIndex(serviceName, hostPort); self._getServicePeer(chan, hostPort); - for (var i = 0; i < range.affineWorkers.length; i++) { - peer = self._getServicePeer(chan, range.affineWorkers[i]); + for (i = 0; i < toConnect.length; i++) { + peer = self._getServicePeer(chan, toConnect[i]); self.ensurePeerConnected(peer, 'service peer affinity refresh'); } From 1233ccd885f1a406423cac5bfb4735f2b65f16de Mon Sep 17 00:00:00 2001 From: Joshua T Corbin Date: Thu, 22 Oct 2015 14:27:45 -0700 Subject: [PATCH 4/7] ServiceDispatchHandler: store lastRefresh rather than bool in knownPeers --- service-proxy.js | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/service-proxy.js b/service-proxy.js index 776c974..de4e3db 100644 --- a/service-proxy.js +++ b/service-proxy.js @@ -103,8 +103,8 @@ function ServiceDispatchHandler(options) { * hostPort :: string * lastRefresh :: number // timestamp * exitServices :: Map - * peersToReap :: Map - * knownPeers :: Map + * peersToReap :: Map + * knownPeers :: Map */ self.exitServices = Object.create(null); self.peersToReap = Object.create(null); @@ -465,24 +465,24 @@ function refreshServicePeer(serviceName, hostPort) { self.exitServices[serviceName] = now; if (self.partialAffinityEnabled) { - self.refreshServicePeerPartially(serviceName, hostPort); + self.refreshServicePeerPartially(serviceName, hostPort, now); return; } // The old way: fully connect every egress to all affine peers. - self.addPeerIndex(serviceName, hostPort); + self.addPeerIndex(serviceName, hostPort, now); var peer = self.getServicePeer(serviceName, hostPort); self.ensurePeerConnected(peer, 'service peer refresh'); }; ServiceDispatchHandler.prototype.addPeerIndex = -function addPeerIndex(serviceName, hostPort) { +function addPeerIndex(serviceName, hostPort, now) { var self = this; // Unmark recently seen peers, so they don't get reaped deleteIndexEntry(self.peersToReap, hostPort, serviceName); // Mark known peers, so they are candidates for future reaping - addIndexEntry(self.knownPeers, hostPort, serviceName, true); + addIndexEntry(self.knownPeers, hostPort, serviceName, now); }; ServiceDispatchHandler.prototype.deletePeerIndex = @@ -564,7 +564,7 @@ function computePartialRange(serviceName, hostPort) { }; ServiceDispatchHandler.prototype.refreshServicePeerPartially = -function refreshServicePeerPartially(serviceName, hostPort) { +function refreshServicePeerPartially(serviceName, hostPort, now) { var self = this; // guaranteed non-null by refreshServicePeer above; we call this only so @@ -612,7 +612,7 @@ function refreshServicePeerPartially(serviceName, hostPort) { toConnect: toConnect })); - self.addPeerIndex(serviceName, hostPort); + self.addPeerIndex(serviceName, hostPort, now); self._getServicePeer(chan, hostPort); for (i = 0; i < toConnect.length; i++) { From f4540a6eded13e60aff1fe92c98313debd1a338a Mon Sep 17 00:00:00 2001 From: Joshua T Corbin Date: Mon, 26 Oct 2015 14:37:25 -0700 Subject: [PATCH 5/7] ServiceDispatchHandler: add .connectedServicePeers index --- service-proxy.js | 39 +++++++++++++++++++++++++++------------ 1 file changed, 27 insertions(+), 12 deletions(-) diff --git a/service-proxy.js b/service-proxy.js index de4e3db..7382f91 100644 --- a/service-proxy.js +++ b/service-proxy.js @@ -99,14 +99,16 @@ function ServiceDispatchHandler(options) { /* service peer state data structures * - * serviceName :: string - * hostPort :: string - * lastRefresh :: number // timestamp - * exitServices :: Map - * peersToReap :: Map - * knownPeers :: Map + * serviceName :: string + * hostPort :: string + * lastRefresh :: number // timestamp + * exitServices :: Map + * peersToReap :: Map + * knownPeers :: Map + * connectedServicePeers :: Map> */ self.exitServices = Object.create(null); + self.connectedServicePeers = Object.create(null); self.peersToReap = Object.create(null); self.knownPeers = Object.create(null); @@ -470,15 +472,21 @@ function refreshServicePeer(serviceName, hostPort) { } // The old way: fully connect every egress to all affine peers. - self.addPeerIndex(serviceName, hostPort, now); + self.addPeerIndex(serviceName, hostPort, true, now); var peer = self.getServicePeer(serviceName, hostPort); - self.ensurePeerConnected(peer, 'service peer refresh'); + self.ensurePeerConnected(serviceName, peer, 'service peer refresh', now); }; ServiceDispatchHandler.prototype.addPeerIndex = -function addPeerIndex(serviceName, hostPort, now) { +function addPeerIndex(serviceName, hostPort, connected, now) { var self = this; + if (connected) { + addIndexEntry(self.connectedServicePeers, serviceName, hostPort, now); + } else { + deleteIndexEntry(self.connectedServicePeers, serviceName, hostPort); + } + // Unmark recently seen peers, so they don't get reaped deleteIndexEntry(self.peersToReap, hostPort, serviceName); // Mark known peers, so they are candidates for future reaping @@ -489,11 +497,16 @@ ServiceDispatchHandler.prototype.deletePeerIndex = function deletePeerIndex(serviceName, hostPort) { var self = this; + deleteIndexEntry(self.connectedServicePeers, serviceName, hostPort); deleteIndexEntry(self.knownPeers, hostPort, serviceName); }; ServiceDispatchHandler.prototype.ensurePeerConnected = -function ensurePeerConnected(peer, reason) { +function ensurePeerConnected(serviceName, peer, reason, now) { + var self = this; + + addIndexEntry(self.connectedServicePeers, serviceName, peer.hostPort, now); + if (peer.isConnected('out')) { return; } @@ -598,10 +611,12 @@ function refreshServicePeerPartially(serviceName, hostPort, now) { } var toConnect = []; + var isAffine = {}; var i; var worker; for (i = 0; i < range.affineWorkers.length; i++) { worker = range.affineWorkers[i]; + isAffine[worker] = true; toConnect.push(worker); } @@ -612,12 +627,12 @@ function refreshServicePeerPartially(serviceName, hostPort, now) { toConnect: toConnect })); - self.addPeerIndex(serviceName, hostPort, now); + self.addPeerIndex(serviceName, hostPort, !!isAffine[hostPort], now); self._getServicePeer(chan, hostPort); for (i = 0; i < toConnect.length; i++) { peer = self._getServicePeer(chan, toConnect[i]); - self.ensurePeerConnected(peer, 'service peer affinity refresh'); + self.ensurePeerConnected(serviceName, peer, 'service peer affinity refresh', now); } // TODO Drop peers that no longer have affinity for this service, such From a28ab8fbc9c2b90daa5d0981229513f1154b642d Mon Sep 17 00:00:00 2001 From: Joshua T Corbin Date: Mon, 26 Oct 2015 14:43:25 -0700 Subject: [PATCH 6/7] ServiceDispatchHandler#refreshServicePeerPartially: prune the toConnect collection --- service-proxy.js | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/service-proxy.js b/service-proxy.js index 7382f91..9ee3968 100644 --- a/service-proxy.js +++ b/service-proxy.js @@ -585,6 +585,7 @@ function refreshServicePeerPartially(serviceName, hostPort, now) { var chan = self.getServiceChannel(serviceName, false); var peer = chan.peers.get(hostPort); + var connectedPeers = self.connectedServicePeers[serviceName]; if (!peer) { peer = self._getServicePeer(chan, hostPort); @@ -607,7 +608,6 @@ function refreshServicePeerPartially(serviceName, hostPort, now) { advertisingPeer: hostPort, partialRange: range })); - } var toConnect = []; @@ -617,7 +617,9 @@ function refreshServicePeerPartially(serviceName, hostPort, now) { for (i = 0; i < range.affineWorkers.length; i++) { worker = range.affineWorkers[i]; isAffine[worker] = true; - toConnect.push(worker); + if (!connectedPeers || !connectedPeers[worker]) { + toConnect.push(worker); + } } self.logger.info('implementing affinity change', self.extendLogInfo({ From 1a2b2cc7b0aa3c89e7dcb49f95cd352ab79e282a Mon Sep 17 00:00:00 2001 From: Joshua T Corbin Date: Mon, 26 Oct 2015 14:44:22 -0700 Subject: [PATCH 7/7] ServiceDispatchHandler#refreshServicePeerPartially: differentiate affinity refresh from affinity change --- service-proxy.js | 20 +++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/service-proxy.js b/service-proxy.js index 9ee3968..dcedf2a 100644 --- a/service-proxy.js +++ b/service-proxy.js @@ -587,10 +587,24 @@ function refreshServicePeerPartially(serviceName, hostPort, now) { var peer = chan.peers.get(hostPort); var connectedPeers = self.connectedServicePeers[serviceName]; - if (!peer) { - peer = self._getServicePeer(chan, hostPort); + // simply freshen if not new + if (peer) { + var connected = connectedPeers && connectedPeers[hostPort]; + self.addPeerIndex(serviceName, peer.hostPort, connected, now); + if (connected) { + self.ensurePeerConnected(serviceName, peer, 'service peer affinity refresh', now); + } + + self.logger.info('refreshed peer partially', self.extendLogInfo({ + serviceName: serviceName, + serviceHostPort: hostPort, + isConnected: connected + })); + return; } + peer = self._getServicePeer(chan, hostPort); + var range = self.computePartialRange(serviceName, hostPort); if (range.relayIndex < 0) { self.logger.warn('Relay could not find itself in the affinity set for service', self.extendLogInfo({ @@ -634,7 +648,7 @@ function refreshServicePeerPartially(serviceName, hostPort, now) { for (i = 0; i < toConnect.length; i++) { peer = self._getServicePeer(chan, toConnect[i]); - self.ensurePeerConnected(serviceName, peer, 'service peer affinity refresh', now); + self.ensurePeerConnected(serviceName, peer, 'service peer affinity change', now); } // TODO Drop peers that no longer have affinity for this service, such