From eeddf25ffcddb3e62a203e712c9f2717614c1f41 Mon Sep 17 00:00:00 2001 From: "linguantian.lgt" Date: Tue, 18 Aug 2026 14:58:35 +0800 Subject: [PATCH] fix: isolate suspect server state per client Synchronize suspect cleanup with roster refresh and make stale cleanup idempotent. --- .../alipay/oceanbase/rpc/ObTableClient.java | 9 +- .../rpc/bolt/transport/ObTableConnection.java | 5 +- .../location/model/RouteTableRefresher.java | 143 ++++++++++---- .../rpc/location/model/TableRoster.java | 8 +- .../rpc/location/model/TableRoute.java | 18 +- .../alipay/oceanbase/rpc/table/ObTable.java | 19 +- .../model/RouteTableRefresherTest.java | 181 ++++++++++++++++++ 7 files changed, 331 insertions(+), 52 deletions(-) create mode 100644 src/test/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresherTest.java diff --git a/src/main/java/com/alipay/oceanbase/rpc/ObTableClient.java b/src/main/java/com/alipay/oceanbase/rpc/ObTableClient.java index 0532be52..bae57eaf 100644 --- a/src/main/java/com/alipay/oceanbase/rpc/ObTableClient.java +++ b/src/main/java/com/alipay/oceanbase/rpc/ObTableClient.java @@ -944,7 +944,8 @@ public ObTable addTable(ObServerAddr addr){ logger.info("server from response not exist in route cache, server ip {}, port {} , execute add Table.", addr.getIp(), addr.getSvrPort()); ObTable obTable = new ObTable.Builder(addr.getIp(), addr.getSvrPort()) // .setLoginInfo(tenantName, userName, password, database, getClientType(runningMode)) // - .setProperties(getProperties()).setObServerAddr(addr).build(); + .setProperties(getProperties()).setObServerAddr(addr) + .setFailureHandler(tableRoute::reportObServerFailure).build(); tableRoster.put(addr, obTable); return obTable; } catch (Exception e) { @@ -983,14 +984,12 @@ public Row transformToRow(String tableName, Object[] rowkey) throws Exception { } public void dealWithRpcTimeoutForSingleTablet(ObServerAddr addr, String tableName, long tabletId) throws Exception { - RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(addr); - RouteTableRefresher.addIntoSuspectIPs(suspectAddr); + tableRoute.reportObServerFailure(addr); tableRoute.refreshPartitionLocation(tableName, tabletId, null); } public void dealWithRpcTimeoutForBatchTablet(ObServerAddr addr, String tableName) throws Exception { - RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(addr); - RouteTableRefresher.addIntoSuspectIPs(suspectAddr); + tableRoute.reportObServerFailure(addr); tableRoute.refreshTabletLocationBatch(tableName); } diff --git a/src/main/java/com/alipay/oceanbase/rpc/bolt/transport/ObTableConnection.java b/src/main/java/com/alipay/oceanbase/rpc/bolt/transport/ObTableConnection.java index ce2286c1..134472f3 100644 --- a/src/main/java/com/alipay/oceanbase/rpc/bolt/transport/ObTableConnection.java +++ b/src/main/java/com/alipay/oceanbase/rpc/bolt/transport/ObTableConnection.java @@ -20,7 +20,6 @@ import com.alipay.oceanbase.rpc.ObGlobal; import com.alipay.oceanbase.rpc.exception.*; import com.alipay.oceanbase.rpc.location.LocationUtil; -import com.alipay.oceanbase.rpc.location.model.RouteTableRefresher; import com.alipay.oceanbase.rpc.protocol.payload.impl.login.ObTableLoginRequest; import com.alipay.oceanbase.rpc.protocol.payload.impl.login.ObTableLoginResult; import com.alipay.oceanbase.rpc.table.ObTable; @@ -119,9 +118,7 @@ private boolean connect() throws Exception { if (tries >= maxTryTimes) { if (!obTable.isOdpMode() && obTable.getObServerAddr() != null) { - RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer( - obTable.getObServerAddr()); - RouteTableRefresher.addIntoSuspectIPs(suspectAddr); + obTable.reportConnectionFailure(); } LOGGER.warn("connect failed after max " + maxTryTimes + " tries " + TraceUtil.formatIpPort(obTable)); diff --git a/src/main/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresher.java b/src/main/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresher.java index 315f06be..cafcce75 100644 --- a/src/main/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresher.java +++ b/src/main/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresher.java @@ -22,13 +22,13 @@ import java.sql.Statement; import java.util.*; import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; import com.alipay.oceanbase.rpc.ObTableClient; import com.alipay.oceanbase.rpc.exception.ObTableEntryRefreshException; import com.alipay.oceanbase.rpc.exception.ObTableTryLockTimeoutException; -import com.alipay.oceanbase.rpc.exception.ObTableUnexpectedException; import com.alipay.oceanbase.rpc.location.LocationUtil; import com.alipay.oceanbase.rpc.table.ObTable; import org.slf4j.Logger; @@ -49,11 +49,15 @@ public class RouteTableRefresher { private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(2); - private final static ConcurrentHashMap suspectLocks = new ConcurrentHashMap<>(); // ObServer -> access lock + private final ConcurrentHashMap suspectLocks = new ConcurrentHashMap<>(); // ObServer -> access lock - private final static ConcurrentHashMap suspectServers = new ConcurrentHashMap<>(); // ObServer -> information structure + private final ConcurrentHashMap suspectServers = new ConcurrentHashMap<>(); // ObServer -> information structure - private final static HashMap serverLastAccessTimestamps = new HashMap<>(); // ObServer -> last access timestamp + private final ConcurrentHashMap serverLastAccessTimestamps = new ConcurrentHashMap<>(); // ObServer -> last access timestamp + + private final Set activeServers = ConcurrentHashMap.newKeySet(); + + private final AtomicBoolean closed = new AtomicBoolean(false); public RouteTableRefresher(ObTableClient tableClient, ObUserAuth sysUA) { this.tableClient = tableClient; @@ -70,6 +74,9 @@ public void start() { } public void close() { + if (!closed.compareAndSet(false, true)) { + return; + } try { scheduler.shutdown(); // wait at most 1 seconds to close the scheduler @@ -79,6 +86,39 @@ public void close() { } catch (InterruptedException e) { logger.warn("scheduler await for terminate interrupted: {}.", e.getMessage()); scheduler.shutdownNow(); + Thread.currentThread().interrupt(); + } finally { + suspectServers.clear(); + suspectLocks.clear(); + serverLastAccessTimestamps.clear(); + activeServers.clear(); + } + } + + /** + * Reconcile the keep-alive state with the latest authoritative tenant roster. + * Servers removed from the roster must not remain in, or be re-added to, the suspect set. + */ + public void refreshActiveServers(Collection servers) { + if (closed.get()) { + return; + } + Set newServers = new HashSet<>(); + if (servers != null) { + newServers.addAll(servers); + } + activeServers.retainAll(newServers); + activeServers.addAll(newServers); + + for (ObServerAddr addr : suspectServers.keySet()) { + if (!newServers.contains(addr)) { + removeFromSuspectIPs(addr); + } + } + for (ObServerAddr addr : serverLastAccessTimestamps.keySet()) { + if (!newServers.contains(addr)) { + serverLastAccessTimestamps.remove(addr); + } } } @@ -127,6 +167,10 @@ private void doRsListCheck() { private void doCheckAliveTask() { for (Map.Entry entry : suspectServers.entrySet()) { try { + if (!activeServers.contains(entry.getKey())) { + removeFromSuspectIPs(entry.getKey()); + continue; + } checkAlive(entry.getKey()); } catch (Exception e) { // silence resolving @@ -161,7 +205,7 @@ private void checkAlive(ObServerAddr addr) { if (t instanceof SQLException) { // occurred during query calcFailureOrClearCache(addr); - } if (t instanceof ObTableEntryRefreshException) { + } else if (t instanceof ObTableEntryRefreshException) { // occurred during connection construction ObTableEntryRefreshException e = (ObTableEntryRefreshException) t; if (e.isConnectInactive()) { @@ -192,11 +236,22 @@ private void checkAlive(ObServerAddr addr) { } } - public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedException { + public void addIntoSuspectIPs(ObServerAddr addr) { + if (addr == null) { + return; + } + addIntoSuspectIPs(new SuspectObServer(addr)); + } + + private void addIntoSuspectIPs(SuspectObServer server) { if (server == null || server.getAddr() == null) { return; } ObServerAddr addr = server.getAddr(); + if (closed.get() || !activeServers.contains(addr)) { + logger.debug("ignore suspect report for inactive server: {}", addr); + return; + } if (suspectServers.get(addr) != null) { // already in the list, directly return return; @@ -218,6 +273,9 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE // already in the list, directly break break; } + if (closed.get() || !activeServers.contains(addr)) { + break; + } Long lastServerAccessTs = serverLastAccessTimestamps.get(addr); if (lastServerAccessTs != null) { long interval = System.currentTimeMillis() - lastServerAccessTs; @@ -235,6 +293,11 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE ++retryTimes; logger.warn("wait to try lock to timeout 1s when add observer into suspect ips, server: {}, tryTimes: {}", addr.toString(), retryTimes, e); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + logger.debug("interrupted while adding observer into suspect ips, server: {}", + addr); + break; } } // end while } finally { @@ -247,43 +310,33 @@ public static void addIntoSuspectIPs(SuspectObServer server) throws InterruptedE private void removeFromSuspectIPs(ObServerAddr addr) { Lock lock = suspectLocks.get(addr); if (lock == null) { - // lock must have been added before remove - throw new ObTableUnexpectedException(String.format("ObServer [%s:%d] need to be add into suspect ips before remove", - addr.getIp(), addr.getSvrPort())); + suspectServers.remove(addr); + logger.debug("suspect server has already been removed: {}", addr); + return; } boolean acquired = false; try { - int retryTimes = 0; - while (true) { - try { - acquired = lock.tryLock(1, TimeUnit.SECONDS); - if (!acquired) { - throw new ObTableTryLockTimeoutException("try to get suspect server lock timeout, timeout: 1s"); - } - // no need to remove lock - SuspectObServer server = suspectServers.remove(addr); - if (server != null) { - int failure = server.getFailure(); - if (failure < failureLimit) { - ObTable obTable = tableClient.getTableRoute().getTableRoster().getTable(addr); - if (obTable != null && !obTable.isValid()) { - obTable.setValid(); - } - } + acquired = lock.tryLock(1, TimeUnit.SECONDS); + if (!acquired) { + logger.debug("defer suspect removal because lock is busy, server: {}", addr); + return; + } + // Keep the lock and cooldown entry until this refresher closes. This prevents a + // concurrent add from using a different lock and preserves the existing cooldown. + SuspectObServer server = suspectServers.remove(addr); + if (server != null) { + int failure = server.getFailure(); + if (failure < failureLimit && activeServers.contains(addr)) { + ObTable obTable = tableClient.getTableRoute().getTableRoster().getTable(addr); + if (obTable != null && !obTable.isValid()) { + obTable.setValid(); } - logger.debug("removed server from suspect list: {}", addr); - break; - } catch (ObTableTryLockTimeoutException e) { - // if try lock timeout, need to retry - ++retryTimes; - logger.warn("wait to try lock to timeout when add observer into suspect ips, server: {}, tryTimes: {}", - addr.toString(), retryTimes, e); - } catch (InterruptedException e) { - // do not throw exception to user layer - // next background task will continue to remove it - logger.warn("waiting to get lock while interrupted by other threads", e); } } + logger.debug("removed server from suspect list: {}", addr); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + logger.debug("interrupted while removing observer from suspect ips, server: {}", addr); } finally { if (acquired) { lock.unlock(); @@ -294,6 +347,10 @@ private void removeFromSuspectIPs(ObServerAddr addr) { private void calcFailureOrClearCache(ObServerAddr addr) { TableRoute tableRoute = tableClient.getTableRoute(); SuspectObServer server = suspectServers.get(addr); + if (server == null) { + logger.debug("skip failure calculation for removed suspect server: {}", addr); + return; + } server.incrementFailure(); int failure = server.getFailure(); if (failure >= failureLimit) { @@ -304,6 +361,18 @@ private void calcFailureOrClearCache(ObServerAddr addr) { addr, failure); } + int suspectServerCount() { + return suspectServers.size(); + } + + boolean containsSuspectServer(ObServerAddr addr) { + return suspectServers.containsKey(addr); + } + + int lastAccessTimestampCount() { + return serverLastAccessTimestamps.size(); + } + public static class SuspectObServer { private final ObServerAddr addr; private final long accessTimestamp; diff --git a/src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoster.java b/src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoster.java index e803f282..2c3dfe80 100644 --- a/src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoster.java +++ b/src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoster.java @@ -19,6 +19,7 @@ import java.util.*; import java.util.concurrent.ConcurrentHashMap; +import java.util.function.Consumer; import com.alipay.oceanbase.rpc.ObTableClient; import com.alipay.oceanbase.rpc.exception.ObTableCloseException; @@ -36,6 +37,7 @@ public class TableRoster { private Properties properties = new Properties(); private Map tableConfigs = new HashMap<>(); private ObTableClientType clientType; + private Consumer failureHandler; /* * ServerAddr(all) -> ObTableConnection */ @@ -65,6 +67,9 @@ public void setProperties(Properties properties) { public void setTableConfigs(Map tableConfigs) { this.tableConfigs = tableConfigs; } + public void setFailureHandler(Consumer failureHandler) { + this.failureHandler = failureHandler; + } public ObTable getTable(ObServerAddr addr) { return tables.get(addr); } @@ -101,7 +106,8 @@ public List refreshTablesAndGetNewServers(List ne ObTable obTable = new ObTable.Builder(addr.getIp(), addr.getSvrPort()) // .setLoginInfo(tenantName, userName, password, database, clientType) // - .setProperties(properties).setConfigs(tableConfigs).setObServerAddr(addr).build(); + .setProperties(properties).setConfigs(tableConfigs).setObServerAddr(addr) + .setFailureHandler(failureHandler).build(); ObTable oldObTable = tables.putIfAbsent(addr, obTable); logger.warn("add new table addr, {}", addr.toString()); if (oldObTable != null) { // maybe create two ob table concurrently, close current ob table diff --git a/src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoute.java b/src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoute.java index a1e7c616..81ff049f 100644 --- a/src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoute.java +++ b/src/main/java/com/alipay/oceanbase/rpc/location/model/TableRoute.java @@ -64,7 +64,7 @@ public class TableRoute { private IndexLocations indexLocations = null; // global index location private TableGroupCache tableGroupCache = null; private OdpInfo odpInfo = null; - private RouteTableRefresher routeRefresher = null; + private volatile RouteTableRefresher routeRefresher = null; private long lastRefreshMetadataTimestamp = -1; public final Lock refreshTableRosterLock = new ReentrantLock(); @@ -348,7 +348,8 @@ public void initRoster(TableEntryKey rootServerKey, boolean initialized, tableClient.getPassword(), tableClient.getDatabase(), tableClient.getClientType(runningMode)) .setProperties(tableClient.getProperties()) - .setConfigs(tableClient.getTableConfigs()).setObServerAddr(addr).build(); + .setConfigs(tableClient.getTableConfigs()).setObServerAddr(addr) + .setFailureHandler(this::reportObServerFailure).build(); addr2Table.put(addr, obTable); servers.add(addr); } catch (Exception e) { @@ -363,6 +364,7 @@ public void initRoster(TableEntryKey rootServerKey, boolean initialized, tableClient.getUserName(), tableClient.getPassword(), tableClient.getDatabase(), tableClient.getClientType(runningMode), tableClient.getProperties(), tableClient.getTableConfigs()); + this.tableRoster.setFailureHandler(this::reportObServerFailure); this.tableRoster.setTables(addr2Table); this.serverRoster.reset(servers); @@ -410,9 +412,17 @@ public void initRoster(TableEntryKey rootServerKey, boolean initialized, public void launchRouteRefresher() { routeRefresher = new RouteTableRefresher(tableClient, sysUA); + routeRefresher.refreshActiveServers(serverRoster.getMembers()); routeRefresher.start(); } + public void reportObServerFailure(ObServerAddr addr) { + RouteTableRefresher refresher = routeRefresher; + if (refresher != null) { + refresher.addIntoSuspectIPs(addr); + } + } + public void removeObServer(ObServerAddr addr) { logger.debug("remove useless table addr, {}", addr.toString()); ConcurrentHashMap tables = this.tableRoster.getTables(); @@ -480,6 +490,10 @@ public void refreshRosterByRsList(List newRsList) throws Exception // update new ob table and get new server address List servers = tableRoster.refreshTablesAndGetNewServers(replicaLocations); serverRoster.reset(servers); + RouteTableRefresher refresher = routeRefresher; + if (refresher != null) { + refresher.refreshActiveServers(servers); + } // 2. Get Server LDC info for weak read consistency. success = false; diff --git a/src/main/java/com/alipay/oceanbase/rpc/table/ObTable.java b/src/main/java/com/alipay/oceanbase/rpc/table/ObTable.java index 807f730c..fd943ab7 100644 --- a/src/main/java/com/alipay/oceanbase/rpc/table/ObTable.java +++ b/src/main/java/com/alipay/oceanbase/rpc/table/ObTable.java @@ -26,7 +26,6 @@ import com.alipay.oceanbase.rpc.exception.*; import com.alipay.oceanbase.rpc.filter.ObTableFilter; import com.alipay.oceanbase.rpc.location.model.ObServerAddr; -import com.alipay.oceanbase.rpc.location.model.RouteTableRefresher; import com.alipay.oceanbase.rpc.mutation.*; import com.alipay.oceanbase.rpc.protocol.payload.ObPayload; import com.alipay.oceanbase.rpc.protocol.payload.impl.execute.*; @@ -51,6 +50,7 @@ import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.ReentrantLock; +import java.util.function.Consumer; import static com.alipay.oceanbase.rpc.property.Property.*; @@ -68,6 +68,7 @@ public class ObTable extends AbstractObTable implements Lifecycle { private ObTableRemoting realClient; private ObTableConnectionPool connectionPool; private ObServerAddr addr; // just used in background keep-alive + private Consumer failureHandler; private ObTableServerCapacity serverCapacity = new ObTableServerCapacity(); @@ -526,8 +527,13 @@ public ObPayload executeWithConnection(final ObPayload request, private void dealWithReconnectFailForObTableConnection() throws InterruptedException { setDirty(); - RouteTableRefresher.SuspectObServer suspectAddr = new RouteTableRefresher.SuspectObServer(addr); - RouteTableRefresher.addIntoSuspectIPs(suspectAddr); + reportConnectionFailure(); + } + + public void reportConnectionFailure() { + if (failureHandler != null && addr != null) { + failureHandler.accept(addr); + } } private void checkObTableOperationResult(String ip, int port, Object result) { @@ -747,6 +753,7 @@ public static class Builder { private String password; private String database; private ObServerAddr addr = null; // only used in background keep-alive + private Consumer failureHandler; ObTableClientType clientType; private Properties properties = new Properties(); @@ -806,6 +813,11 @@ public Builder setObServerAddr(ObServerAddr addr) { return this; } + public Builder setFailureHandler(Consumer failureHandler) { + this.failureHandler = failureHandler; + return this; + } + /* * Build. */ @@ -822,6 +834,7 @@ public ObTable build() throws Exception { obTable.setClientType(clientType); obTable.setIsOdpMode(isOdpMode); obTable.setObServerAddr(addr); + obTable.failureHandler = failureHandler; obTable.init(); diff --git a/src/test/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresherTest.java b/src/test/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresherTest.java new file mode 100644 index 00000000..d82db833 --- /dev/null +++ b/src/test/java/com/alipay/oceanbase/rpc/location/model/RouteTableRefresherTest.java @@ -0,0 +1,181 @@ +/*- + * #%L + * com.oceanbase:obkv-table-client + * %% + * Copyright (C) 2021 - 2026 OceanBase + * %% + * OBKV Table Client Framework is licensed under Mulan PSL v2. + * You can use this software according to the terms and conditions of the Mulan PSL v2. + * You may obtain a copy of Mulan PSL v2 at: + * http://license.coscl.org.cn/MulanPSL2 + * THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND, + * EITHER EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT, + * MERCHANTABILITY OR FIT FOR A PARTICULAR PURPOSE. + * See the Mulan PSL v2 for more details. + * #L% + */ +package com.alipay.oceanbase.rpc.location.model; + +import com.alipay.oceanbase.rpc.ObTableClient; +import org.junit.After; +import org.junit.Test; + +import java.lang.reflect.Method; +import java.util.Collections; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class RouteTableRefresherTest { + + private RouteTableRefresher refresherA; + private RouteTableRefresher refresherB; + + @After + public void tearDown() { + if (refresherA != null) { + refresherA.close(); + } + if (refresherB != null) { + refresherB.close(); + } + } + + @Test + public void suspectStateIsIsolatedPerClient() { + ObServerAddr addr = server("127.0.0.1"); + refresherA = refresherFor(new RecordingTableRoute()); + refresherB = refresherFor(new RecordingTableRoute()); + refresherA.refreshActiveServers(Collections.singleton(addr)); + refresherB.refreshActiveServers(Collections.singleton(addr)); + + refresherA.addIntoSuspectIPs(addr); + + assertTrue(refresherA.containsSuspectServer(addr)); + assertFalse(refresherB.containsSuspectServer(addr)); + assertEquals(1, refresherA.suspectServerCount()); + assertEquals(0, refresherB.suspectServerCount()); + } + + @Test + public void rosterRetirementClearsStateAndRejectsStaleReports() { + ObServerAddr addr = server("127.0.0.2"); + refresherA = refresherFor(new RecordingTableRoute()); + refresherA.refreshActiveServers(Collections.singleton(addr)); + refresherA.addIntoSuspectIPs(addr); + assertEquals(1, refresherA.lastAccessTimestampCount()); + + refresherA.refreshActiveServers(Collections.emptyList()); + refresherA.addIntoSuspectIPs(addr); + + assertFalse(refresherA.containsSuspectServer(addr)); + assertEquals(0, refresherA.lastAccessTimestampCount()); + } + + @Test + public void staleFailureAfterRosterRemovalIsIdempotent() throws Exception { + ObServerAddr addr = server("127.0.0.3"); + RecordingTableRoute route = new RecordingTableRoute(); + refresherA = refresherFor(route); + refresherA.refreshActiveServers(Collections.singleton(addr)); + refresherA.addIntoSuspectIPs(addr); + refresherA.refreshActiveServers(Collections.emptyList()); + + invokePrivate(refresherA, "calcFailureOrClearCache", addr); + invokePrivate(refresherA, "removeFromSuspectIPs", addr); + invokePrivate(refresherA, "removeFromSuspectIPs", addr); + + assertEquals(0, route.removeCount.get()); + assertEquals(0, refresherA.suspectServerCount()); + } + + @Test + public void failureLimitEvictsServerOnlyOnce() throws Exception { + ObServerAddr addr = server("127.0.0.4"); + RecordingTableRoute route = new RecordingTableRoute(); + refresherA = refresherFor(route); + refresherA.refreshActiveServers(Collections.singleton(addr)); + refresherA.addIntoSuspectIPs(addr); + + invokePrivate(refresherA, "calcFailureOrClearCache", addr); + invokePrivate(refresherA, "calcFailureOrClearCache", addr); + invokePrivate(refresherA, "calcFailureOrClearCache", addr); + invokePrivate(refresherA, "calcFailureOrClearCache", addr); + + assertEquals(1, route.removeCount.get()); + assertEquals(0, refresherA.suspectServerCount()); + } + + @Test + public void closeClearsInstanceState() { + ObServerAddr addr = server("127.0.0.5"); + refresherA = refresherFor(new RecordingTableRoute()); + refresherA.refreshActiveServers(Collections.singleton(addr)); + refresherA.addIntoSuspectIPs(addr); + + refresherA.close(); + + assertEquals(0, refresherA.suspectServerCount()); + assertEquals(0, refresherA.lastAccessTimestampCount()); + } + + @Test + public void closingOneClientDoesNotClearAnotherClientState() { + ObServerAddr addr = server("127.0.0.6"); + refresherA = refresherFor(new RecordingTableRoute()); + refresherB = refresherFor(new RecordingTableRoute()); + refresherA.refreshActiveServers(Collections.singleton(addr)); + refresherB.refreshActiveServers(Collections.singleton(addr)); + refresherA.addIntoSuspectIPs(addr); + refresherB.addIntoSuspectIPs(addr); + + refresherA.close(); + + assertEquals(0, refresherA.suspectServerCount()); + assertEquals(1, refresherB.suspectServerCount()); + } + + private RouteTableRefresher refresherFor(TableRoute route) { + return new RouteTableRefresher(new TestObTableClient(route), new ObUserAuth("root@sys", + "unused")); + } + + private ObServerAddr server(String ip) { + return new ObServerAddr(ip, 2881, 2882); + } + + private void invokePrivate(RouteTableRefresher refresher, String methodName, ObServerAddr addr) + throws Exception { + Method method = RouteTableRefresher.class.getDeclaredMethod(methodName, ObServerAddr.class); + method.setAccessible(true); + method.invoke(refresher, addr); + } + + private static final class TestObTableClient extends ObTableClient { + private final TableRoute route; + + private TestObTableClient(TableRoute route) { + this.route = route; + } + + @Override + public TableRoute getTableRoute() { + return route; + } + } + + private static final class RecordingTableRoute extends TableRoute { + private final AtomicInteger removeCount = new AtomicInteger(); + + private RecordingTableRoute() { + super(null, null); + } + + @Override + public void removeObServer(ObServerAddr addr) { + removeCount.incrementAndGet(); + } + } +}