diff --git a/openframe-test-service-core/src/main/java/com/openframe/test/config/RedisConfig.java b/openframe-test-service-core/src/main/java/com/openframe/test/config/RedisConfig.java index c5d3137a8b..1de5a36af8 100644 --- a/openframe-test-service-core/src/main/java/com/openframe/test/config/RedisConfig.java +++ b/openframe-test-service-core/src/main/java/com/openframe/test/config/RedisConfig.java @@ -69,20 +69,28 @@ public static String getPassword() { } /** - * Whether the server runs in cluster mode. True for the in-cluster shard set every environment - * but dev still uses; dev points at a single Memorystore instance, where a cluster client fails - * on topology discovery because cluster mode is disabled server-side. + * Pins cluster mode instead of letting {@link com.openframe.test.data.redis.Redis} ask the server. + * An override for the case where the probe cannot be trusted; leaving it unset is the normal path. */ public static void setCluster(boolean enabled) { cluster = enabled; } - public static boolean isCluster() { + /** + * The pinned answer, or {@code null} when nobody pinned one and the server should be asked. + * + *

SaaS Redis is migrating to Memorystore for Valkey — a single node with cluster mode disabled — + * one environment at a time, so the answer differs per environment and changes under us as the + * migration proceeds. Detecting it costs one {@code CLUSTER INFO}, which both topologies answer, so + * an environment moves without anyone editing config, and this override exists only as an escape + * hatch. + */ + public static Boolean getConfiguredCluster() { if (cluster != null) { return cluster; } String env = System.getenv("REDIS_CLUSTER"); - return (env == null || env.trim().isEmpty()) || Boolean.parseBoolean(env); + return (env == null || env.trim().isEmpty()) ? null : Boolean.parseBoolean(env); } public static Set getClusterNodes() { diff --git a/openframe-test-service-core/src/main/java/com/openframe/test/data/redis/Redis.java b/openframe-test-service-core/src/main/java/com/openframe/test/data/redis/Redis.java index 8b5ba15273..ee8f8805ae 100644 --- a/openframe-test-service-core/src/main/java/com/openframe/test/data/redis/Redis.java +++ b/openframe-test-service-core/src/main/java/com/openframe/test/data/redis/Redis.java @@ -3,8 +3,9 @@ import com.openframe.test.config.RedisConfig; import lombok.extern.slf4j.Slf4j; import redis.clients.jedis.DefaultJedisClientConfig; +import redis.clients.jedis.HostAndPort; +import redis.clients.jedis.Jedis; import redis.clients.jedis.JedisClientConfig; -import redis.clients.jedis.JedisCluster; import redis.clients.jedis.JedisPooled; import redis.clients.jedis.UnifiedJedis; import redis.clients.jedis.params.ScanParams; @@ -22,10 +23,15 @@ import java.security.cert.CertificateFactory; import java.util.Collection; import java.util.List; +import java.util.Set; @Slf4j public class Redis { + /** What the server answered to CLUSTER INFO, remembered for the JVM. Null until first asked. */ + private static volatile Boolean detectedCluster; + + /** * Find the password-reset token for {@code email}. The auth-server stores it under the tenant-scoped, * hash-tagged key {@code of:{}:pwdreset:} with the email as the value. @@ -40,38 +46,117 @@ public class Redis { */ public static String getResetToken(String email) { String pattern = "of:{" + RedisConfig.getTenant() + "}:pwdreset:*"; - try (UnifiedJedis client = client()) { - ScanParams scanParams = new ScanParams().match(pattern).count(100); - String cursor = ScanParams.SCAN_POINTER_START; - do { - ScanResult scanResult = client.scan(cursor, scanParams); - List keys = scanResult.getResult(); - if (!keys.isEmpty()) { - // One slot for the whole batch, so this is a single round trip rather than a GET per key. - List emails = client.mget(keys.toArray(new String[0])); - for (int i = 0; i < keys.size(); i++) { - if (email.equals(emails.get(i))) { - return keys.get(i).split(":pwdreset:")[1]; + try { + JedisClientConfig config = clientConfig(); + if (clusterMode(config)) { + // SCAN is per-node: a cluster only ever reports the keys of the node answering it, so the + // seeds are walked one by one. This deliberately does not use JedisCluster — that needs + // CLUSTER SLOTS to succeed first, and the advertised addresses are not reachable from the + // test pod, which is how this lookup broke on qa (JedisClusterOperationException: could + // not initialize cluster slots cache) and took the password-reset case with it. + Set seeds = RedisConfig.getClusterNodes(); + int scanned = 0; + for (HostAndPort node : seeds) { + try (UnifiedJedis client = new JedisPooled(node, config)) { + String token = findToken(client, pattern, email); + scanned++; + if (token != null) { + return token; } + } catch (Exception e) { + // Node unreachable or holding none of the slots — try the next seed. + log.debug("Seed {} did not answer the password-reset scan: {}", node, e.toString()); } } - cursor = scanResult.getCursor(); - } while (!cursor.equals(ScanParams.SCAN_POINTER_START)); + if (scanned == 0) { + // "Scanned every node and the token is not there yet" and "could not scan anything" are + // the same null to the caller, and the caller polls on it until a timeout. Only one of + // those is worth waking someone for: a reset that never lands leaves the tenant-report + // account half-rotated, which a re-run cannot repair. + log.warn("None of the {} Redis seeds could be scanned for the password-reset token; " + + "the token may exist and be unreachable rather than absent", seeds.size()); + } + return null; + } + try (UnifiedJedis client = new JedisPooled(RedisConfig.getNode(), config)) { + return findToken(client, pattern, email); + } } catch (Exception e) { log.warn("Reading the password-reset token from Redis failed", e); + return null; } + } + + /** Walks one server's keyspace for the tenant's reset keys and returns the token whose value is the email. */ + private static String findToken(UnifiedJedis client, String pattern, String email) { + ScanParams scanParams = new ScanParams().match(pattern).count(100); + String cursor = ScanParams.SCAN_POINTER_START; + do { + ScanResult scanResult = client.scan(cursor, scanParams); + List keys = scanResult.getResult(); + if (!keys.isEmpty()) { + // The keys share the tenant hash tag, so one slot for the whole batch: a single round + // trip rather than a GET per key. + List emails = client.mget(keys.toArray(new String[0])); + for (int i = 0; i < keys.size(); i++) { + if (email.equals(emails.get(i))) { + return keys.get(i).split(":pwdreset:")[1]; + } + } + } + cursor = scanResult.getCursor(); + } while (!cursor.equals(ScanParams.SCAN_POINTER_START)); return null; } /** - * The client is built per call rather than cached: a caller polls at most a few dozen times, and a - * cached static client would trade those handshakes for a topology-staleness problem. + * Whether the server runs in cluster mode, asked once and remembered. + * + *

SaaS Redis is moving to Memorystore for Valkey — one node, cluster mode disabled, TLS and AUTH — + * one environment at a time, so this differs per environment and changes as the migration proceeds. + * {@code CLUSTER INFO} is answered by both topologies, so one round trip settles it and an + * environment migrates without anyone editing config. {@link RedisConfig#getConfiguredCluster()} + * still wins where someone pinned an answer. + * + *

A probe that cannot connect assumes a cluster for that one lookup - that is what every + * environment but dev is today, so it keeps the behaviour unchanged where the probe itself is the + * thing that is broken. Only an answer the server actually gave is remembered: this pod lives for + * days, and a probe lost to one dropped SYN must not pin a guess for all of them. */ - private static UnifiedJedis client() throws GeneralSecurityException, IOException { - JedisClientConfig config = clientConfig(); - return RedisConfig.isCluster() - ? new JedisCluster(RedisConfig.getClusterNodes(), config) - : new JedisPooled(RedisConfig.getNode(), config); + private static boolean clusterMode(JedisClientConfig config) { + Boolean pinned = RedisConfig.getConfiguredCluster(); + if (pinned != null) { + return pinned; + } + Boolean known = detectedCluster; + if (known != null) { + return known; + } + synchronized (Redis.class) { + Boolean answered = detectedCluster; + if (answered != null) { + return answered; + } + Boolean probed = probeCluster(config); + if (probed == null) { + return true; + } + detectedCluster = probed; + return probed; + } + } + + /** What the server says about itself, or {@code null} when it did not answer. */ + private static Boolean probeCluster(JedisClientConfig config) { + HostAndPort node = RedisConfig.getNode(); + try (Jedis jedis = new Jedis(node, config)) { + boolean enabled = jedis.clusterInfo().contains("cluster_enabled:1"); + log.info("Redis at {} reports cluster mode {}", node, enabled ? "enabled" : "disabled"); + return enabled; + } catch (Exception e) { + log.warn("Could not read CLUSTER INFO from {}; assuming a cluster for this lookup", node, e); + return null; + } } /**