From 4d3810d8a1e6fa8332d1de6093ed7e167e348500 Mon Sep 17 00:00:00 2001 From: 0G-C <218736830+0G-C@users.noreply.github.com> Date: Tue, 21 Jul 2026 15:01:21 +0800 Subject: [PATCH 1/3] fix session close race leading to orphan Psub connections --- .../DefaultHealthCheckInstanceFactory.java | 1 + .../AbstractInstanceSessionManager.java | 107 +++++----- .../session/InstanceSessionManager.java | 2 + .../healthcheck/session/RedisSession.java | 75 +++++-- .../healthcheck/session/PsubRaceTest.java | 184 ++++++++++++++++++ ...AbstractConsoleInstanceSessionManager.java | 8 + 6 files changed, 321 insertions(+), 56 deletions(-) create mode 100644 redis/redis-checker/src/test/java/com/ctrip/xpipe/redis/checker/healthcheck/session/PsubRaceTest.java diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/impl/DefaultHealthCheckInstanceFactory.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/impl/DefaultHealthCheckInstanceFactory.java index 706b7dfbb8..772900f45d 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/impl/DefaultHealthCheckInstanceFactory.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/impl/DefaultHealthCheckInstanceFactory.java @@ -92,6 +92,7 @@ public DefaultHealthCheckInstanceFactory(CheckerConfig checkerConfig, HealthChec public void remove(RedisHealthCheckInstance instance) { Endpoint endpoint = instance.getEndpoint(); endpointFactory.remove(new HostPort(endpoint.getHost(), endpoint.getPort())); + redisSessionManager.removeSession(endpoint); stopCheck(instance); } diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java index 05d50582e8..f465740c70 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java @@ -2,6 +2,8 @@ import com.ctrip.xpipe.api.endpoint.Endpoint; import com.ctrip.xpipe.api.foundation.FoundationService; +import com.ctrip.xpipe.api.monitor.Task; +import com.ctrip.xpipe.api.monitor.TransactionMonitor; import com.ctrip.xpipe.concurrent.AbstractExceptionLogTask; import com.ctrip.xpipe.endpoint.HostPort; import com.ctrip.xpipe.pool.XpipeNettyClientKeyedObjectPool; @@ -64,85 +66,94 @@ public abstract class AbstractInstanceSessionManager implements InstanceSessionM private CheckerConfig config; @VisibleForTesting - public static long checkUnusedRedisDelaySeconds = 4; + public static long checkUnusedRedisDelaySeconds = 3600; + + @VisibleForTesting + public static long checkRedisDelaySeconds = 4; @PostConstruct - public void postConstruct(){ + public void postConstruct() { scheduled.scheduleAtFixedRate(new AbstractExceptionLogTask() { @Override - protected void doRun() throws Exception { + protected void doRun() { try { removeUnusedInstances(); } catch (Exception e) { logger.error("[removeUnusedInstances]", e); } + } + }, checkUnusedRedisDelaySeconds, checkUnusedRedisDelaySeconds, TimeUnit.SECONDS); + scheduled.scheduleAtFixedRate(new AbstractExceptionLogTask() { + @Override + protected void doRun() { for(RedisSession redisSession : sessions.values()){ - try{ + try { redisSession.check(); - }catch (Exception e){ + } catch (Exception e) { logger.error("[check]" + redisSession, e); } } } - }, checkUnusedRedisDelaySeconds, checkUnusedRedisDelaySeconds, TimeUnit.SECONDS); + }, checkRedisDelaySeconds, checkRedisDelaySeconds, TimeUnit.SECONDS); } @Override - public RedisSession findOrCreateSession(Endpoint endpoint) { + public synchronized RedisSession findOrCreateSession(Endpoint endpoint) { RedisSession session = sessions.get(endpoint); if (session == null) { - synchronized (this) { - session = sessions.get(endpoint); - if (session == null) { - session = new RedisSession(endpoint, scheduled, keyedObjectPool, config); - sessions.put(endpoint, session); - } - } + session = new RedisSession(endpoint, scheduled, keyedObjectPool, config); + sessions.put(endpoint, session); } return session; } @Override - public RedisSession findOrCreateSession(HostPort hostPort) { + public synchronized RedisSession findOrCreateSession(HostPort hostPort) { return findOrCreateSession(endpointFactory.getOrCreateEndpoint(hostPort)); } @VisibleForTesting - protected void removeUnusedInstances() { - Set currentStoredRedises = sessions.keySet(); - if(currentStoredRedises.isEmpty()) - return; - - Set redisInUse = getInUseInstances(); - if(redisInUse == null || redisInUse.isEmpty()) { - return; - } - List unusedRedises = new LinkedList<>(); - - for(Endpoint endpoint : currentStoredRedises) { - if(!redisInUse.contains(new HostPort(endpoint.getHost(), endpoint.getPort()))) { - unusedRedises.add(endpoint); - } - } + protected synchronized void removeUnusedInstances() { + TransactionMonitor.DEFAULT.logTransactionSwallowException( + "session.cleanup", "removeUnusedInstances", new Task() { + @Override + public void go() { + Set currentStoredRedises = sessions.keySet(); + if(currentStoredRedises.isEmpty()) + return; + + Set redisInUse = getInUseInstances(); + if(redisInUse == null || redisInUse.isEmpty()) { + return; + } + List unusedRedises = new LinkedList<>(); - if(unusedRedises.isEmpty()) { - return; - } - unusedRedises.forEach(endpoint -> { - RedisSession redisSession = sessions.getOrDefault(endpoint, null); - if(redisSession != null) { - logger.info("[removeUnusedRedises]Redis: {} not in use, remove from session manager", endpoint); - // add try logic to continue working on others - try { - redisSession.closeConnection(); - } catch (Exception ignore) { + for(Endpoint endpoint : currentStoredRedises) { + if(!redisInUse.contains(new HostPort(endpoint.getHost(), endpoint.getPort()))) { + unusedRedises.add(endpoint); + } + } + if(unusedRedises.isEmpty()) { + return; } - sessions.remove(endpoint); + unusedRedises.forEach(endpoint -> { + try { + logger.info("[removeUnusedRedises]Redis: {} not in use, remove from session manager", endpoint); + removeSession(endpoint); + } catch (Exception e) { + logger.warn("[removeUnusedRedises] close session {} failed", endpoint, e); + } + }); + } + + @Override + public java.util.Map getData() { + return null; } }); } @@ -150,13 +161,21 @@ protected void removeUnusedInstances() { protected abstract Set getInUseInstances(); + public synchronized boolean removeSession(Endpoint endpoint) { + RedisSession session = sessions.remove(endpoint); + if (session == null) return false; + session.close(); + return true; + } + + protected void closeAllConnections() { try { executors.execute(new Runnable() { @Override public void run() { for (RedisSession session : sessions.values()) { - session.closeConnection(); + session.close(); } } }); diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/InstanceSessionManager.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/InstanceSessionManager.java index ed2f233139..3573e9d7d4 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/InstanceSessionManager.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/InstanceSessionManager.java @@ -13,4 +13,6 @@ public interface InstanceSessionManager { RedisSession findOrCreateSession(Endpoint endpoint); RedisSession findOrCreateSession(HostPort hostPort); + + boolean removeSession(Endpoint endpoint); } diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/RedisSession.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/RedisSession.java index df5e2c7f19..eed432ccab 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/RedisSession.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/RedisSession.java @@ -21,6 +21,7 @@ import java.util.Map; import java.util.Set; import java.util.concurrent.*; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer; import java.util.function.Supplier; @@ -31,7 +32,7 @@ *

* Dec 1, 2016 2:28:43 PM */ -public class RedisSession { +public class RedisSession implements java.io.Closeable { private static Logger logger = LoggerFactory.getLogger(RedisSession.class); @@ -41,6 +42,8 @@ public class RedisSession { private Endpoint endpoint; + private AtomicBoolean isClosed = new AtomicBoolean(false); + private ConcurrentMap, PubSubConnectionWrapper> subscribConns = new ConcurrentHashMap<>(); private ScheduledExecutorService scheduled; @@ -67,7 +70,7 @@ public RedisSession() { } public void check() { - + makeSureOpen(); for (Map.Entry, PubSubConnectionWrapper> entry : subscribConns.entrySet()) { Set channel = entry.getKey(); @@ -94,7 +97,7 @@ public void check() { } public synchronized void closeSubscribedChannel(String... channel) { - + makeSureOpen(); PubSubConnectionWrapper pubSubConnectionWrapper = subscribConns.get(Sets.newHashSet(channel)); if (pubSubConnectionWrapper != null) { logger.debug("[closeSubscribedChannel]{}, {}", endpoint, channel); @@ -104,18 +107,22 @@ public synchronized void closeSubscribedChannel(String... channel) { } public synchronized void subscribeIfAbsent(SubscribeCallback callback, String... channel) { + makeSureOpen(); subscribeIfAbsent(callback, () -> new SubscribeCommand(clientPool, scheduled, commandTimeOut, channel), channel); } public synchronized void crdtsubscribeIfAbsent(SubscribeCallback callback, String... channel) { + makeSureOpen(); subscribeIfAbsent(callback, () -> new CRDTSubscribeCommand(clientPool, scheduled, commandTimeOut, channel), channel); } public synchronized void psubscribeIfAbsent(SubscribeCallback callback, String... channel) { + makeSureOpen(); subscribeIfAbsent(callback, () -> new PsubscribeCommand(clientPool, scheduled, commandTimeOut, channel), channel); } private synchronized void subscribeIfAbsent(SubscribeCallback callback, Supplier subCommandSupplier, String... channel) { + makeSureOpen(); PubSubConnectionWrapper pubSubConnectionWrapper = subscribConns.get(Sets.newHashSet(channel)); if (pubSubConnectionWrapper == null || pubSubConnectionWrapper.shouldCreateNewSession()) { if(pubSubConnectionWrapper != null) { @@ -138,6 +145,7 @@ public void operationComplete(CommandFuture commandFuture) throws Except } public synchronized void publish(String channel, String message) { + makeSureOpen(); PublishCommand pubCommand = new PublishCommand(clientPool, scheduled, commandTimeOut, channel, message); silentCommand(pubCommand); @@ -152,6 +160,7 @@ public void operationComplete(CommandFuture commandFuture) throws Except } public synchronized void crdtpublish(String channel, String message) { + makeSureOpen(); CRDTPublishCommand pubCommand = new CRDTPublishCommand(clientPool, scheduled, commandTimeOut, channel, message); silentCommand(pubCommand); @@ -166,6 +175,7 @@ public void operationComplete(CommandFuture commandFuture) throws Except } public CommandFuture ping(final PingCallback callback) { + makeSureOpen(); // if connect has been established PingCommand pingCommand = new PingCommand(clientPool, scheduled, commandTimeOut); silentCommand(pingCommand); @@ -184,6 +194,7 @@ public void operationComplete(CommandFuture commandFuture) throws Except } public CommandFuture role(RollCallback callback) { + makeSureOpen(); RoleCommand command = new RoleCommand(clientPool, commandTimeOut, false, scheduled); silentCommand(command); command.execute().addListener(new CommandFutureListener() { @@ -200,6 +211,7 @@ public void operationComplete(CommandFuture commandFuture) throws Exceptio } public void configRewrite(BiConsumer consumer) { + makeSureOpen(); ConfigRewrite command = new ConfigRewrite(clientPool, scheduled, commandTimeOut); silentCommand(command); command.execute().addListener(new CommandFutureListener() { @@ -216,7 +228,7 @@ public void operationComplete(CommandFuture commandFuture) throws Except } public String roleSync() throws InterruptedException, ExecutionException, TimeoutException { - + makeSureOpen(); RoleCommand command = new RoleCommand(clientPool, waitResultSeconds * 1000, true, scheduled); silentCommand(command); return command.execute().get().getServerRole().name(); @@ -240,23 +252,25 @@ public void operationComplete(CommandFuture commandFuture) throws Exception { } public CommandFuture expireSize(Callbackable callback) { + makeSureOpen(); ExpireSizeCommand command = new ExpireSizeCommand(clientPool, scheduled, commandTimeOut); return addHookAndExecute(command, callback); } public CommandFuture tombstoneSize(Callbackable callback) { + makeSureOpen(); TombstoneSizeCommand command = new TombstoneSizeCommand(clientPool, scheduled, commandTimeOut); return addHookAndExecute(command, callback); } public CommandFuture info(final String infoSection, Callbackable callback) { - + makeSureOpen(); InfoCommand command = new InfoCommand(clientPool, infoSection, scheduled, commandTimeOut); return addHookAndExecute(command, callback); } public CommandFuture crdtInfo(final String infoSection, Callbackable callback) { - + makeSureOpen(); InfoCommand command = new CRDTInfoCommand(clientPool, infoSection, scheduled, commandTimeOut); return addHookAndExecute(command, callback); } @@ -284,6 +298,7 @@ public CommandFuture crdtInfoReplication(Callbackable callback) } public void isDiskLessSync(Callbackable callback) { + makeSureOpen(); ConfigGetCommand.ConfigGetDisklessSync command = new ConfigGetCommand.ConfigGetDisklessSync(clientPool, scheduled, commandTimeOut); silentCommand(command); command.execute().addListener(new CommandFutureListener() { @@ -300,6 +315,7 @@ public void operationComplete(CommandFuture commandFuture) throws Excep } public void ConfigGet(Callbackable callback, String args) { + makeSureOpen(); ConfigGetCommand.ConfigGetAnyCommand command = new ConfigGetCommand.ConfigGetAnyCommand(clientPool, scheduled, args); silentCommand(command); @@ -316,6 +332,7 @@ public void operationComplete(CommandFuture commandFuture) throws Except } public void CRDTConfigGet(Callbackable callback, String args) { + makeSureOpen(); CRDTConfigGetCommand command = new CRDTConfigGetCommand(clientPool, scheduled, args); silentCommand(command); @@ -334,12 +351,14 @@ public void operationComplete(CommandFuture commandFuture) throws Except public InfoResultExtractor syncInfo(InfoCommand.INFO_TYPE infoType) throws InterruptedException, ExecutionException, TimeoutException { + makeSureOpen(); InfoCommand infoCommand = new InfoCommand(clientPool, infoType, scheduled); String info = infoCommand.execute().get(2000, TimeUnit.MILLISECONDS); return new InfoResultExtractor(info); } public InfoResultExtractor syncCRDTInfo(InfoCommand.INFO_TYPE infoType) throws InterruptedException, ExecutionException, TimeoutException { + makeSureOpen(); CRDTInfoCommand command = new CRDTInfoCommand(clientPool, infoType, scheduled); String info = command.execute().get(2000, TimeUnit.MILLISECONDS); return new InfoResultExtractor(info); @@ -347,6 +366,7 @@ public InfoResultExtractor syncCRDTInfo(InfoCommand.INFO_TYPE infoType) throws I public CommandFuture getRedisReplInfo() { + makeSureOpen(); InfoReplicationCommand command = new InfoReplicationCommand(clientPool, scheduled, commandTimeOut); silentCommand(command); return command.execute(); @@ -456,19 +476,50 @@ public boolean shouldCreateNewSession() { } } - public void closeConnection() { - for(PubSubConnectionWrapper connectionWrapper : subscribConns.values()) { - try { - connectionWrapper.closeAndClean(); - } catch (Exception ignore) {} + @Override + public void close() { + if (!cmpAndSetClosed()) { + logger.info("[close][{}] already closed, skip", endpoint); + return; + } + synchronized (this) { + for(PubSubConnectionWrapper connectionWrapper : subscribConns.values()) { + try { + connectionWrapper.closeAndClean(); + } catch (Exception ignore) {} + } } try { clientPool.clear(); } catch (Throwable th) { - logger.info("[closeConnection][{}] fail", endpoint, th); + logger.info("[close][{}] fail", endpoint, th); + } + } + + /** + * @deprecated use {@link #close()} instead + */ + @Deprecated + public void closeConnection() { + close(); + } + + + private boolean cmpAndSetClosed() { + return isClosed.compareAndSet(false, true); + } + + + private void makeSureOpen() { + if (isClosed.get()) { + throw new IllegalStateException("[RedisSession][closed] " + endpoint); } } + public boolean isClosed() { + return isClosed.get(); + } + public RedisSession setKeyedNettyClientPool(XpipeNettyClientKeyedObjectPool keyedNettyClientPool) { this.clientPool = keyedNettyClientPool.getKeyPool(endpoint); return this; diff --git a/redis/redis-checker/src/test/java/com/ctrip/xpipe/redis/checker/healthcheck/session/PsubRaceTest.java b/redis/redis-checker/src/test/java/com/ctrip/xpipe/redis/checker/healthcheck/session/PsubRaceTest.java new file mode 100644 index 0000000000..925792cee9 --- /dev/null +++ b/redis/redis-checker/src/test/java/com/ctrip/xpipe/redis/checker/healthcheck/session/PsubRaceTest.java @@ -0,0 +1,184 @@ +package com.ctrip.xpipe.redis.checker.healthcheck.session; + +import com.ctrip.xpipe.AbstractTest; +import com.ctrip.xpipe.api.endpoint.Endpoint; +import com.ctrip.xpipe.concurrent.AbstractExceptionLogTask; +import com.ctrip.xpipe.endpoint.DefaultEndPoint; +import com.ctrip.xpipe.netty.commands.NettyClient; +import com.ctrip.xpipe.pool.XpipeNettyClientKeyedObjectPool; +import com.ctrip.xpipe.redis.checker.TestConfig; +import com.ctrip.xpipe.redis.checker.config.CheckerConfig; +import com.ctrip.xpipe.simpleserver.AbstractIoAction; +import com.ctrip.xpipe.simpleserver.IoActionFactory; +import com.ctrip.xpipe.simpleserver.Server; +import org.apache.commons.pool2.impl.GenericObjectPool; +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.net.Socket; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +/** + * 真实复现:通过最上层的定时调度,两条线程完全独立地跑,不固定任何顺序。 + * + * 线程A: AbstractInstanceSessionManager + * 每 DELAY_A ms 调 removeUnusedInstances + * → 把 endpoint 从"在用列表"移除 → sessions.remove → closeConnection + * → 再放回"在用列表" → findOrCreateSession → 新 session → psub + * + * 线程B: PsubAction + * 每 DELAY_B ms 调 doTask + * → instance.getRedisSession().psubscribeIfAbsent() + * + * 两线程共用同一个 endpoint 的池。 + * 每次 A 移除了 endpoint、closeConnection,B 还在用旧 session 的引用做 psub。 + * 旧 session 产生孤儿、新 session 又 borrow——Active 积累。 + * + * 本测试不安排顺序、不 await、不 CountDownLatch——让调度器自己决定。 + */ +public class PsubRaceTest extends AbstractTest { + + private Server fakeRedis; + private XpipeNettyClientKeyedObjectPool pool; + private Endpoint endpoint; + private CheckerConfig config; + private GenericObjectPool underlying; + private volatile RedisSession latestSession; + private volatile boolean running; + private int initialActive; + + @Before + public void setup() throws Exception { + System.setProperty("DEFAULT_REDIS_COMMAND_TIME_OUT_SECONDS", "600"); + fakeRedis = startServer(randomPort(), psubIoActionFactory()); + endpoint = new DefaultEndPoint("127.0.0.1", fakeRedis.getPort()); + pool = getXpipeNettyClientKeyedObjectPool(); + config = new TestConfig(); + underlying = (GenericObjectPool) pool.getObjectPool(endpoint); + initialActive = underlying.getNumActive(); + } + + @After + public void teardown() throws Exception { + if (fakeRedis != null) fakeRedis.stop(); + } + + private IoActionFactory psubIoActionFactory() { + return socket -> new AbstractIoAction(socket) { + private boolean responded = false; + @Override protected Object doRead(InputStream ins) throws IOException { + byte[] buf = new byte[4096]; int len = ins.read(buf); + if (len < 0) return null; + byte[] data = new byte[len]; + System.arraycopy(buf, 0, data, 0, len); return data; + } + @Override protected void doWrite(OutputStream ous, Object readResult) throws IOException { + if (readResult == null || responded) return; + responded = true; + ous.write("*3\r\n$9\r\npsubscribe\r\n$5\r\nxpipe*\r\n:1\r\n".getBytes()); ous.flush(); + } + }; + } + + /** + * 真实复现:两条完全独立的定时任务线程,通过连接池耦合。 + * + * 线程A (manager 模拟): + * ① 从"在用"列表移除 endpoint → removeUnusedInstances + * ② 如果 currentSession 存在, closeConnection + sessions.remove + * ③ endpoint 重新加入"在用"列表 → findOrCreateSession → 新 session + * ④ 新 session.psubscribeIfAbsent → 从池 borrow + * + * 线程B (PsubAction 模拟): + * ① 取 latestSession + * ② 调 psubscribeIfAbsent → subscribConns 可能已被线程A清空 + * ③ 发现空了 → 建新 Psubscribe → borrow → 不归还 + * + * 两线程没有同步、没有 await、没有栅栏。 + * 完全依赖线程调度器自然产生的竞态窗口来制造孤儿。 + * 跟生产环境一样:元信息抖动越频繁、线程调度交错越多,孤儿越多。 + */ + @Test + public void testRealRaceWithScheduledExecutors() throws Exception { + final int intervalManager = 30; // 线程A间隔 (ms) + final int intervalPsub = 60; // 线程B间隔 (ms)——更慢,避免每次都赶上 clean + final int durationMs = 30000; // 跑 30 秒——让孤儿有机会累积 + running = true; + + // 先创建 sessionA(模拟初始 instance 的 session,固定给线程B用) + final RedisSession psubSession = new RedisSession(endpoint, scheduled, pool, config); + psubSession.psubscribeIfAbsent(new NoopCallback(), "xpipe*"); + // 线程A 持有的 session,会被替换(模拟 manager 换 instance) + latestSession = psubSession; + System.out.println("[start] initial Active=" + underlying.getNumActive()); + + ScheduledExecutorService executor = new ScheduledThreadPoolExecutor(2); + + // === 线程A: 模拟 AbstractInstanceSessionManager(元信息变更) === + executor.scheduleAtFixedRate(() -> { + try { + // 模拟 removeUnusedInstances: + // 杀 psubSession(线程B 的 session)——生产里会清理所有关联 session + psubSession.close(); + + // 让线程B 产生孤儿:psubSession.subscribConns 已被清空 + // 等线程B 下一次跑时,它看到的是空 subscribConns → 建新 Psubscribe → orphan + + // 同时模拟 manager 创建新 session(新 instance) + RedisSession newSession = new RedisSession(endpoint, scheduled, pool, config); + newSession.psubscribeIfAbsent(new NoopCallback(), "xpipe*"); + latestSession = newSession; + } catch (Exception e) { } + }, 0, 50, TimeUnit.MILLISECONDS); + + // === 线程B: 模拟 PsubAction(固定的 session,不跟着换) === + executor.scheduleAtFixedRate(() -> { + try { + // 线程B 持有一开始绑定的 session 引用,不会因为 manager 操作而改变 + // 即使 manager closed it + 移除了,线程B 还是在用这个老 session + psubSession.psubscribeIfAbsent(new NoopCallback(), "xpipe*"); + } catch (Exception e) { } + }, 0, 30, TimeUnit.MILLISECONDS); + + // 让两线程自由竞态跑 durationMs 毫秒 + Thread.sleep(durationMs); + running = false; + executor.shutdownNow(); + executor.awaitTermination(2, TimeUnit.SECONDS); + + int finalActive = underlying.getNumActive(); + int activeDelta = finalActive - initialActive; + String msg = "[result] " + durationMs + "ms 初始Active=" + initialActive + + " 最终Active=" + finalActive + " 差值=" + activeDelta + + " Idle=" + underlying.getNumIdle() + " max=" + underlying.getMaxTotal(); + System.err.println(msg); + java.io.FileWriter fw = new java.io.FileWriter("target/race-final-result.txt"); + fw.write(msg); fw.close(); + + // 诊断日志留底,不设置严格断言(timing 敏感) + System.out.println(msg); + + latestSession.close(); + } + + @Test + public void testPsubNeverReturnsWithoutRace() throws Exception { + RedisSession session = new RedisSession(endpoint, scheduled, pool, config); + session.psubscribeIfAbsent(new NoopCallback(), "xpipe*"); + Assert.assertEquals(1, underlying.getNumActive()); + Assert.assertEquals(0, underlying.getNumIdle()); + session.close(); + } + + private static class NoopCallback implements RedisSession.SubscribeCallback { + @Override public void message(String channel, String message) { } + @Override public void fail(Throwable e) { } + } +} diff --git a/redis/redis-console/src/main/java/com/ctrip/xpipe/redis/console/healthcheck/session/AbstractConsoleInstanceSessionManager.java b/redis/redis-console/src/main/java/com/ctrip/xpipe/redis/console/healthcheck/session/AbstractConsoleInstanceSessionManager.java index 6471dc21bc..677196e12b 100644 --- a/redis/redis-console/src/main/java/com/ctrip/xpipe/redis/console/healthcheck/session/AbstractConsoleInstanceSessionManager.java +++ b/redis/redis-console/src/main/java/com/ctrip/xpipe/redis/console/healthcheck/session/AbstractConsoleInstanceSessionManager.java @@ -134,6 +134,14 @@ protected void removeUnusedInstances() { protected abstract Set getInUseInstances(); + @Override + public boolean removeSession(Endpoint endpoint) { + RedisSession session = sessions.remove(endpoint); + if (session == null) return false; + session.closeConnection(); + return true; + } + protected void closeAllConnections() { try { executors.execute(new Runnable() { From ef3f247c67e5e62063e7f9d696af71ec90fe2c3f Mon Sep 17 00:00:00 2001 From: 0G-C <218736830+0G-C@users.noreply.github.com> Date: Thu, 23 Jul 2026 15:01:42 +0800 Subject: [PATCH 2/3] make session remove interval configurable via QConfig --- .../xpipe/redis/checker/alert/ALERT_TYPE.java | 16 +++++++++ .../redis/checker/config/CheckerConfig.java | 2 ++ .../checker/config/impl/CommonConfigBean.java | 6 ++++ .../AbstractInstanceSessionManager.java | 35 ++++++++++++------- .../ctrip/xpipe/redis/checker/TestConfig.java | 5 +++ .../config/impl/DefaultConsoleConfig.java | 5 +++ .../console/health/RemoveUnusedRedisTest.java | 2 +- 7 files changed, 58 insertions(+), 13 deletions(-) diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/alert/ALERT_TYPE.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/alert/ALERT_TYPE.java index 0448bf10d4..ec96236c15 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/alert/ALERT_TYPE.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/alert/ALERT_TYPE.java @@ -853,6 +853,22 @@ public DetailDesc detailDesc() { public ALERT_LEVEL getAlertLevel() { return ALERT_LEVEL.HIGH; } + }, + CHECKER_SESSION_MANAGER_FAIL("checker session manager fail", EMAIL_XPIPE_ADMIN) { + @Override + public boolean urgent() { return false; } + + @Override + public boolean reportRecovery() { return false; } + + @Override + public DetailDesc detailDesc() { + return new DetailDesc("session manager task start fail", + "session removeUnusedTask start fail, 可能会导致连接泄露"); + } + + @Override + public ALERT_LEVEL getAlertLevel() { return ALERT_LEVEL.HIGH; } }; diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/config/CheckerConfig.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/config/CheckerConfig.java index b01afacfc8..33121fb659 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/config/CheckerConfig.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/config/CheckerConfig.java @@ -128,4 +128,6 @@ default boolean supportSentinelBeacon(String clusterName) { boolean shouldComputeExtraInHash(); + int getSessionRemoveUnusedDelayMillis(); + } diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/config/impl/CommonConfigBean.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/config/impl/CommonConfigBean.java index db0328b3fe..9efc1ca5b6 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/config/impl/CommonConfigBean.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/config/impl/CommonConfigBean.java @@ -220,4 +220,10 @@ public long getAbnormalClusterStatusMonitorIntervalMilli() { DEFAULT_ABNORMAL_CLUSTER_STATUS_MONITOR_INTERVAL_MILLI); } + public static final String KEY_SESSION_REMOVE_UNUSED_DELAY_MILLIS = "checker.session.remove.unused.delay.millis"; + + public int getSessionRemoveUnusedDelayMillis() { + return getIntProperty(KEY_SESSION_REMOVE_UNUSED_DELAY_MILLIS, 3600000); + } + } diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java index f465740c70..b19504140b 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java @@ -8,10 +8,13 @@ import com.ctrip.xpipe.endpoint.HostPort; import com.ctrip.xpipe.pool.XpipeNettyClientKeyedObjectPool; import com.ctrip.xpipe.redis.checker.CheckerConsoleService; +import com.ctrip.xpipe.redis.checker.alert.ALERT_TYPE; +import com.ctrip.xpipe.redis.checker.alert.AlertManager; import com.ctrip.xpipe.redis.checker.config.CheckerConfig; import com.ctrip.xpipe.redis.checker.healthcheck.impl.HealthCheckEndpointFactory; import com.ctrip.xpipe.redis.core.meta.MetaCache; import com.ctrip.xpipe.utils.VisibleForTesting; +import com.ctrip.xpipe.utils.job.DynamicDelayPeriodTask; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; @@ -65,25 +68,26 @@ public abstract class AbstractInstanceSessionManager implements InstanceSessionM @Autowired private CheckerConfig config; - @VisibleForTesting - public static long checkUnusedRedisDelaySeconds = 3600; + @Autowired + private AlertManager alertManager; @VisibleForTesting public static long checkRedisDelaySeconds = 4; + private DynamicDelayPeriodTask removeUnusedTask; + @PostConstruct public void postConstruct() { - scheduled.scheduleAtFixedRate(new AbstractExceptionLogTask() { - @Override - protected void doRun() { - try { - removeUnusedInstances(); - } catch (Exception e) { - logger.error("[removeUnusedInstances]", e); - } - } - }, checkUnusedRedisDelaySeconds, checkUnusedRedisDelaySeconds, TimeUnit.SECONDS); + this.removeUnusedTask = new DynamicDelayPeriodTask("RemoveUnusedInstances", + this::removeUnusedInstances, config::getSessionRemoveUnusedDelayMillis, scheduled); + try { + removeUnusedTask.start(); + } catch (Exception e) { + logger.error("[postConstruct] start removeUnusedTask fail", e); + alertManager.alert(null, null, null, ALERT_TYPE.CHECKER_SESSION_MANAGER_FAIL, + "session removeUnusedTask start fail: " + e.getMessage()); + } scheduled.scheduleAtFixedRate(new AbstractExceptionLogTask() { @Override @@ -186,6 +190,13 @@ public void run() { @PreDestroy public void preDestroy(){ + if (removeUnusedTask != null) { + try { + removeUnusedTask.stop(); + } catch (Exception e) { + logger.error("[preDestroy] stop removeUnusedTask fail", e); + } + } closeAllConnections(); } diff --git a/redis/redis-checker/src/test/java/com/ctrip/xpipe/redis/checker/TestConfig.java b/redis/redis-checker/src/test/java/com/ctrip/xpipe/redis/checker/TestConfig.java index 40ed16019c..f2aab8f294 100644 --- a/redis/redis-checker/src/test/java/com/ctrip/xpipe/redis/checker/TestConfig.java +++ b/redis/redis-checker/src/test/java/com/ctrip/xpipe/redis/checker/TestConfig.java @@ -334,4 +334,9 @@ public boolean checkBeaconLastModifyTime() { public boolean shouldComputeExtraInHash() { return false; } + + @Override + public int getSessionRemoveUnusedDelayMillis() { + return 3600000; + } } diff --git a/redis/redis-console/src/main/java/com/ctrip/xpipe/redis/console/config/impl/DefaultConsoleConfig.java b/redis/redis-console/src/main/java/com/ctrip/xpipe/redis/console/config/impl/DefaultConsoleConfig.java index 55bb672de6..60f79d60cb 100644 --- a/redis/redis-console/src/main/java/com/ctrip/xpipe/redis/console/config/impl/DefaultConsoleConfig.java +++ b/redis/redis-console/src/main/java/com/ctrip/xpipe/redis/console/config/impl/DefaultConsoleConfig.java @@ -671,6 +671,11 @@ public boolean shouldComputeExtraInHash() { return commonConfigBean.shouldComputeExtraInHash(); } + @Override + public int getSessionRemoveUnusedDelayMillis() { + return commonConfigBean.getSessionRemoveUnusedDelayMillis(); + } + @Override public long getCheckIsolateInterval() { return consoleConfigBean.getIsolateCheckIntervalMilli(); diff --git a/redis/redis-console/src/test/java/com/ctrip/xpipe/redis/console/health/RemoveUnusedRedisTest.java b/redis/redis-console/src/test/java/com/ctrip/xpipe/redis/console/health/RemoveUnusedRedisTest.java index f662ffca8e..6c9d0f1f39 100644 --- a/redis/redis-console/src/test/java/com/ctrip/xpipe/redis/console/health/RemoveUnusedRedisTest.java +++ b/redis/redis-console/src/test/java/com/ctrip/xpipe/redis/console/health/RemoveUnusedRedisTest.java @@ -57,7 +57,7 @@ public class RemoveUnusedRedisTest extends AbstractConsoleDbTest { @Before public void beforeRemoveUnusedRedisTest() throws Exception { MockitoAnnotations.initMocks(this); - DefaultRedisSessionManager.checkUnusedRedisDelaySeconds = 2; + manager.setConfig(checkerConfig); // random port to avoid port conflict port = randomPort(); From 6c0800bfbd61b291c680a14f97f984b4b5456972 Mon Sep 17 00:00:00 2001 From: 0G-C <218736830+0G-C@users.noreply.github.com> Date: Thu, 23 Jul 2026 17:33:35 +0800 Subject: [PATCH 3/3] use logAlertEvent for session removeUnusedTask start failure --- .../xpipe/redis/checker/alert/ALERT_TYPE.java | 16 ---------------- .../session/AbstractInstanceSessionManager.java | 9 ++------- 2 files changed, 2 insertions(+), 23 deletions(-) diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/alert/ALERT_TYPE.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/alert/ALERT_TYPE.java index ec96236c15..0448bf10d4 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/alert/ALERT_TYPE.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/alert/ALERT_TYPE.java @@ -853,22 +853,6 @@ public DetailDesc detailDesc() { public ALERT_LEVEL getAlertLevel() { return ALERT_LEVEL.HIGH; } - }, - CHECKER_SESSION_MANAGER_FAIL("checker session manager fail", EMAIL_XPIPE_ADMIN) { - @Override - public boolean urgent() { return false; } - - @Override - public boolean reportRecovery() { return false; } - - @Override - public DetailDesc detailDesc() { - return new DetailDesc("session manager task start fail", - "session removeUnusedTask start fail, 可能会导致连接泄露"); - } - - @Override - public ALERT_LEVEL getAlertLevel() { return ALERT_LEVEL.HIGH; } }; diff --git a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java index b19504140b..d9043ff668 100644 --- a/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java +++ b/redis/redis-checker/src/main/java/com/ctrip/xpipe/redis/checker/healthcheck/session/AbstractInstanceSessionManager.java @@ -2,14 +2,13 @@ import com.ctrip.xpipe.api.endpoint.Endpoint; import com.ctrip.xpipe.api.foundation.FoundationService; +import com.ctrip.xpipe.api.monitor.EventMonitor; import com.ctrip.xpipe.api.monitor.Task; import com.ctrip.xpipe.api.monitor.TransactionMonitor; import com.ctrip.xpipe.concurrent.AbstractExceptionLogTask; import com.ctrip.xpipe.endpoint.HostPort; import com.ctrip.xpipe.pool.XpipeNettyClientKeyedObjectPool; import com.ctrip.xpipe.redis.checker.CheckerConsoleService; -import com.ctrip.xpipe.redis.checker.alert.ALERT_TYPE; -import com.ctrip.xpipe.redis.checker.alert.AlertManager; import com.ctrip.xpipe.redis.checker.config.CheckerConfig; import com.ctrip.xpipe.redis.checker.healthcheck.impl.HealthCheckEndpointFactory; import com.ctrip.xpipe.redis.core.meta.MetaCache; @@ -68,9 +67,6 @@ public abstract class AbstractInstanceSessionManager implements InstanceSessionM @Autowired private CheckerConfig config; - @Autowired - private AlertManager alertManager; - @VisibleForTesting public static long checkRedisDelaySeconds = 4; @@ -85,8 +81,7 @@ public void postConstruct() { removeUnusedTask.start(); } catch (Exception e) { logger.error("[postConstruct] start removeUnusedTask fail", e); - alertManager.alert(null, null, null, ALERT_TYPE.CHECKER_SESSION_MANAGER_FAIL, - "session removeUnusedTask start fail: " + e.getMessage()); + EventMonitor.DEFAULT.logAlertEvent("session-removeUnused-startFail"); } scheduled.scheduleAtFixedRate(new AbstractExceptionLogTask() {