Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*
* <p>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<HostAndPort> getClusterNodes() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Comment thread
giokur marked this conversation as resolved.


/**
* Find the password-reset token for {@code email}. The auth-server stores it under the tenant-scoped,
* hash-tagged key {@code of:{<tenant>}:pwdreset:<token>} with the email as the value.
Expand All @@ -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<String> scanResult = client.scan(cursor, scanParams);
List<String> 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<String> 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<HostAndPort> 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<String> scanResult = client.scan(cursor, scanParams);
List<String> 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<String> 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.
*
* <p>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.
*
* <p>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;
}
}
Comment thread
giokur marked this conversation as resolved.

/** 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;
}
}

/**
Expand Down
Loading