From ad1a47b8181bcef01b89fe11c6ea73412e046046 Mon Sep 17 00:00:00 2001 From: Eric Hare Date: Mon, 27 Jul 2026 08:28:42 -0700 Subject: [PATCH 1/5] Add Cassandra readiness health check --- .../CassandraConnectionHealthCheck.java | 170 ++++++++++++++ .../CassandraConnectionHealthCheckTest.java | 208 ++++++++++++++++++ .../v1/SessionEvictionIntegrationTest.java | 55 +++++ 3 files changed, 433 insertions(+) create mode 100644 src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java create mode 100644 src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java diff --git a/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java b/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java new file mode 100644 index 0000000000..67c1f06e2c --- /dev/null +++ b/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java @@ -0,0 +1,170 @@ +package io.stargate.sgv2.jsonapi.api.health; + +import com.datastax.oss.driver.api.core.AllNodesFailedException; +import com.datastax.oss.driver.api.core.connection.ClosedConnectionException; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.google.common.annotations.VisibleForTesting; +import io.stargate.sgv2.jsonapi.api.request.UserAgent; +import io.stargate.sgv2.jsonapi.api.request.tenant.Tenant; +import io.stargate.sgv2.jsonapi.api.request.tenant.TenantFactory; +import io.stargate.sgv2.jsonapi.config.DatabaseType; +import io.stargate.sgv2.jsonapi.config.OperationsConfig; +import io.stargate.sgv2.jsonapi.service.cqldriver.CQLSessionCache; +import io.stargate.sgv2.jsonapi.service.cqldriver.CqlCredentials; +import io.stargate.sgv2.jsonapi.service.cqldriver.CqlSessionCacheSupplier; +import jakarta.enterprise.context.ApplicationScoped; +import jakarta.inject.Inject; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.Base64; +import java.util.Objects; +import org.eclipse.microprofile.health.HealthCheck; +import org.eclipse.microprofile.health.HealthCheckResponse; +import org.eclipse.microprofile.health.Readiness; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * Health check that verifies Cassandra connectivity for the Data API readiness probe. + * + *

The check obtains the Cassandra session used by normal requests from the session cache, + * verifies that it is open, and executes a lightweight query against {@code system.local}. It is + * only active when the Data API is configured with a {@link DatabaseType#CASSANDRA} backend. + */ +@Readiness +@ApplicationScoped +public class CassandraConnectionHealthCheck implements HealthCheck { + + private static final Logger LOGGER = + LoggerFactory.getLogger(CassandraConnectionHealthCheck.class); + + private static final String HEALTH_CHECK_NAME = "cassandra-connection"; + private static final String HEALTH_CHECK_QUERY = "SELECT release_version FROM system.local"; + private static final Duration HEALTH_CHECK_TIMEOUT = Duration.ofSeconds(5); + private static final UserAgent HEALTH_CHECK_USER_AGENT = new UserAgent("DataAPI-HealthCheck/1.0"); + + private final CQLSessionCache sessionCache; + private final OperationsConfig operationsConfig; + private final Duration timeout; + private final String authToken; + + @Inject + public CassandraConnectionHealthCheck( + CqlSessionCacheSupplier sessionCacheSupplier, OperationsConfig operationsConfig) { + this( + Objects.requireNonNull(sessionCacheSupplier, "sessionCacheSupplier must not be null").get(), + operationsConfig, + HEALTH_CHECK_TIMEOUT); + } + + @VisibleForTesting + CassandraConnectionHealthCheck( + CQLSessionCache sessionCache, OperationsConfig operationsConfig, Duration timeout) { + this.sessionCache = Objects.requireNonNull(sessionCache, "sessionCache must not be null"); + this.operationsConfig = + Objects.requireNonNull(operationsConfig, "operationsConfig must not be null"); + this.timeout = Objects.requireNonNull(timeout, "timeout must not be null"); + this.authToken = + operationsConfig.databaseConfig().type() == DatabaseType.CASSANDRA + ? createAuthToken(operationsConfig.databaseConfig()) + : null; + } + + @Override + public HealthCheckResponse call() { + var responseBuilder = HealthCheckResponse.named(HEALTH_CHECK_NAME); + + if (operationsConfig.databaseConfig().type() != DatabaseType.CASSANDRA) { + return responseBuilder + .up() + .withData( + "reason", + "Cassandra connectivity check is not applicable for database type " + + operationsConfig.databaseConfig().type()) + .build(); + } + + var healthCheckTenant = TenantFactory.instance().create(null); + + try { + var session = + sessionCache + .getSession(healthCheckTenant, authToken, HEALTH_CHECK_USER_AGENT) + .await() + .atMost(timeout); + + if (session.isClosed()) { + LOGGER.warn("Cassandra session is closed during health check"); + evictSession(healthCheckTenant); + return responseBuilder.down().withData("reason", "Session is closed").build(); + } + + var statement = + SimpleStatement.builder(HEALTH_CHECK_QUERY) + .setTimeout(timeout) + .setConsistencyLevel(operationsConfig.queriesConfig().consistency().reads()) + .build(); + + var resultSet = session.execute(statement); + var row = resultSet.one(); + var version = row != null ? row.getString("release_version") : "unknown"; + + LOGGER.trace("Cassandra health check passed, version: {}", version); + + return responseBuilder + .up() + .withData("cassandra_version", version) + .withData("session_name", session.getName()) + .build(); + } catch (Exception e) { + if (isUnreliableSessionFailure(e)) { + evictSession(healthCheckTenant); + } + + LOGGER.error("Cassandra health check failed", e); + return responseBuilder + .down() + .withData("error", e.getClass().getSimpleName()) + .withData("message", e.getMessage() != null ? e.getMessage() : "Unknown error") + .build(); + } + } + + private static String createAuthToken(OperationsConfig.DatabaseConfig databaseConfig) { + return databaseConfig + .fixedToken() + .orElseGet( + () -> + CqlCredentials.USERNAME_PASSWORD_TOKEN_PREFIX + + encode(databaseConfig.userName()) + + ":" + + encode(databaseConfig.password())); + } + + private static String encode(String value) { + return Base64.getEncoder() + .encodeToString( + Objects.requireNonNull(value, "Cassandra credential must not be null") + .getBytes(StandardCharsets.UTF_8)); + } + + private static boolean isUnreliableSessionFailure(Throwable throwable) { + var current = throwable; + while (current != null) { + if (current instanceof AllNodesFailedException + || current instanceof ClosedConnectionException) { + return true; + } + current = current.getCause(); + } + return false; + } + + private void evictSession(Tenant tenant) { + try { + sessionCache.evictSession(tenant, authToken, HEALTH_CHECK_USER_AGENT); + } catch (Exception e) { + LOGGER.warn("Unable to evict the Cassandra session after a failed health check", e); + } + } +} diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java new file mode 100644 index 0000000000..b9898ad60d --- /dev/null +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java @@ -0,0 +1,208 @@ +package io.stargate.sgv2.jsonapi.api.health; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.*; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; +import com.datastax.oss.driver.api.core.connection.ClosedConnectionException; +import com.datastax.oss.driver.api.core.cql.ResultSet; +import com.datastax.oss.driver.api.core.cql.Row; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import io.smallrye.mutiny.Uni; +import io.stargate.sgv2.jsonapi.api.request.UserAgent; +import io.stargate.sgv2.jsonapi.api.request.tenant.Tenant; +import io.stargate.sgv2.jsonapi.api.request.tenant.TenantFactory; +import io.stargate.sgv2.jsonapi.config.DatabaseType; +import io.stargate.sgv2.jsonapi.config.OperationsConfig; +import io.stargate.sgv2.jsonapi.service.cqldriver.CQLSessionCache; +import io.stargate.sgv2.jsonapi.service.cqldriver.CqlCredentials; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.Base64; +import java.util.Optional; +import org.eclipse.microprofile.health.HealthCheckResponse; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +/** Tests for {@link CassandraConnectionHealthCheck}. */ +public class CassandraConnectionHealthCheckTest { + + private static final String FIXED_TOKEN = "fixed-token"; + private static final String USER_NAME = "test-user"; + private static final String PASSWORD = "test-password"; + + private CQLSessionCache sessionCache; + private OperationsConfig operationsConfig; + private OperationsConfig.DatabaseConfig databaseConfig; + private CqlSession session; + private CassandraConnectionHealthCheck healthCheck; + + @BeforeEach + public void setup() { + TenantFactory.initialize(DatabaseType.CASSANDRA); + + sessionCache = mock(CQLSessionCache.class); + session = mock(CqlSession.class); + + databaseConfig = mock(OperationsConfig.DatabaseConfig.class); + when(databaseConfig.type()).thenReturn(DatabaseType.CASSANDRA); + when(databaseConfig.fixedToken()).thenReturn(Optional.of(FIXED_TOKEN)); + when(databaseConfig.userName()).thenReturn(USER_NAME); + when(databaseConfig.password()).thenReturn(PASSWORD); + + var consistencyConfig = mock(OperationsConfig.QueriesConfig.ConsistencyConfig.class); + when(consistencyConfig.reads()).thenReturn(DefaultConsistencyLevel.LOCAL_QUORUM); + + var queriesConfig = mock(OperationsConfig.QueriesConfig.class); + when(queriesConfig.consistency()).thenReturn(consistencyConfig); + + operationsConfig = mock(OperationsConfig.class); + when(operationsConfig.databaseConfig()).thenReturn(databaseConfig); + when(operationsConfig.queriesConfig()).thenReturn(queriesConfig); + + healthCheck = + new CassandraConnectionHealthCheck(sessionCache, operationsConfig, Duration.ofMillis(100)); + } + + @AfterEach + public void cleanup() { + TenantFactory.reset(); + } + + @Test + public void successfulHealthCheck() { + var resultSet = mock(ResultSet.class); + var row = mock(Row.class); + + when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) + .thenReturn(Uni.createFrom().item(session)); + when(session.isClosed()).thenReturn(false); + when(session.execute(any(SimpleStatement.class))).thenReturn(resultSet); + when(resultSet.one()).thenReturn(row); + when(row.getString("release_version")).thenReturn("6.9.21"); + when(session.getName()).thenReturn("SINGLE-TENANT"); + + var response = healthCheck.call(); + + assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.UP); + assertThat(response.getData()) + .hasValueSatisfying(data -> assertThat(data).containsEntry("cassandra_version", "6.9.21")); + assertThat(response.getData()) + .hasValueSatisfying( + data -> assertThat(data).containsEntry("session_name", "SINGLE-TENANT")); + + var statementCaptor = ArgumentCaptor.forClass(SimpleStatement.class); + verify(session).execute(statementCaptor.capture()); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("SELECT release_version FROM system.local"); + assertThat(statementCaptor.getValue().getTimeout()).isEqualTo(Duration.ofMillis(100)); + assertThat(statementCaptor.getValue().getConsistencyLevel()) + .isEqualTo(DefaultConsistencyLevel.LOCAL_QUORUM); + } + + @Test + public void sessionAcquisitionFailureReportsDown() { + when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) + .thenReturn(Uni.createFrom().failure(new IllegalStateException("Cannot connect"))); + + var response = healthCheck.call(); + + assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); + assertThat(response.getData()) + .hasValueSatisfying( + data -> + assertThat(data) + .containsEntry("error", "IllegalStateException") + .containsEntry("message", "Cannot connect")); + } + + @Test + public void closedSessionReportsDownAndIsEvicted() { + when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) + .thenReturn(Uni.createFrom().item(session)); + when(session.isClosed()).thenReturn(true); + + var response = healthCheck.call(); + + assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); + assertThat(response.getData()) + .hasValueSatisfying(data -> assertThat(data).containsEntry("reason", "Session is closed")); + verify(sessionCache).evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + verify(session, never()).execute(any(SimpleStatement.class)); + } + + @Test + public void queryFailureReportsDownAndUnreliableSessionIsEvicted() { + when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) + .thenReturn(Uni.createFrom().item(session)); + when(session.isClosed()).thenReturn(false); + when(session.execute(any(SimpleStatement.class))) + .thenThrow(new ClosedConnectionException("Connection is closed")); + + var response = healthCheck.call(); + + assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); + assertThat(response.getData()) + .hasValueSatisfying( + data -> + assertThat(data) + .containsEntry("error", "ClosedConnectionException") + .containsEntry("message", "Connection is closed")); + verify(sessionCache).evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + } + + @Test + public void sessionAcquisitionTimeoutReportsDown() { + when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) + .thenReturn(Uni.createFrom().nothing()); + + var response = healthCheck.call(); + + assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); + verifyNoInteractions(session); + } + + @Test + public void configuredCredentialsAreUsedWhenFixedTokenIsNotSet() { + when(databaseConfig.fixedToken()).thenReturn(Optional.empty()); + healthCheck = + new CassandraConnectionHealthCheck(sessionCache, operationsConfig, Duration.ofMillis(100)); + + when(sessionCache.getSession(any(Tenant.class), any(String.class), any(UserAgent.class))) + .thenReturn(Uni.createFrom().item(session)); + when(session.isClosed()).thenReturn(true); + + healthCheck.call(); + + var expectedToken = + CqlCredentials.USERNAME_PASSWORD_TOKEN_PREFIX + + Base64.getEncoder().encodeToString(USER_NAME.getBytes(StandardCharsets.UTF_8)) + + ":" + + Base64.getEncoder().encodeToString(PASSWORD.getBytes(StandardCharsets.UTF_8)); + verify(sessionCache).getSession(any(Tenant.class), eq(expectedToken), any(UserAgent.class)); + } + + @Test + public void nonCassandraDatabaseDoesNotRunConnectivityCheck() { + when(databaseConfig.type()).thenReturn(DatabaseType.ASTRA); + healthCheck = + new CassandraConnectionHealthCheck(sessionCache, operationsConfig, Duration.ofMillis(100)); + + var response = healthCheck.call(); + + assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.UP); + assertThat(response.getData()) + .hasValueSatisfying( + data -> + assertThat(data) + .containsEntry( + "reason", + "Cassandra connectivity check is not applicable for database type ASTRA")); + verifyNoInteractions(sessionCache); + } +} diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java index 322dbd3711..663ed19f60 100644 --- a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java @@ -30,6 +30,8 @@ public class SessionEvictionIntegrationTest extends AbstractCollectionIntegratio private static final Logger LOGGER = LoggerFactory.getLogger(SessionEvictionIntegrationTest.class); + private static final String READINESS_PATH = "/stargate/health/ready"; + private static final String LIVENESS_PATH = "/stargate/health/live"; /** * Overridden to ensure we connect to the isolated container created for this test. @@ -69,6 +71,8 @@ public void testSessionEvictionOnAllNodesFailed() { .body("$", responseIsFindSuccess()) .body("data.document._id", is("before_crash")); + waitForReadinessStatus("UP", 30_000); + // 2. Stop the container to simulate DB failure // Use low-level dockerClient to stop the container without triggering Testcontainers' // cleanup/termination logic (which dbContainer.stop() would do). @@ -86,11 +90,15 @@ public void testSessionEvictionOnAllNodesFailed() { .body( "errors[0].errorCode", is(DatabaseException.Code.FAILED_TO_CONNECT_TO_DATABASE.name())); + waitForReadinessStatus("DOWN", 60_000); + given().when().get(LIVENESS_PATH).then().statusCode(200).body("status", is("UP")); + // 4. Restart the container to simulate recovery getDockerClient().startContainerCmd(getContainerId()).exec(); // 5. Wait for the database to become responsive again waitForDbRecovery(); + waitForReadinessStatus("UP", 120_000); // 6. Verify Session Recovery: check the data before crashing // Not to check that cassandra works, but to check that we are running the same container as @@ -225,6 +233,53 @@ private boolean isApiReady() { } } + /** + * Polls the readiness endpoint until both the overall health and Cassandra connectivity check + * have the expected status. + */ + private void waitForReadinessStatus(String expectedStatus, long timeoutMillis) { + var start = System.currentTimeMillis(); + var expectedStatusCode = "UP".equals(expectedStatus) ? 200 : 503; + Response lastResponse = null; + + while (System.currentTimeMillis() - start < timeoutMillis) { + try { + lastResponse = given().when().get(READINESS_PATH); + var jsonPath = lastResponse.jsonPath(); + var overallStatus = jsonPath.getString("status"); + var cassandraStatus = + jsonPath.getString("checks.find { it.name == 'cassandra-connection' }.status"); + + if (lastResponse.statusCode() == expectedStatusCode + && expectedStatus.equals(overallStatus) + && expectedStatus.equals(cassandraStatus)) { + return; + } + } catch (Exception e) { + LOGGER.debug("Readiness endpoint not in expected state yet", e); + } + + try { + Thread.sleep(1000); + } catch (InterruptedException ignored) { + Thread.currentThread().interrupt(); + break; + } + } + + var lastResponseDescription = + lastResponse == null + ? "no response" + : "HTTP " + lastResponse.statusCode() + ": " + lastResponse.asString(); + throw new RuntimeException( + "Readiness did not become " + + expectedStatus + + " within " + + timeoutMillis + + " ms. Last response: " + + lastResponseDescription); + } + /** Checks if Cassandra is up and normal by running "nodetool status" inside the container. */ private boolean isCassandraUp(DockerClient dockerClient, String containerId) { try { From 143a00683bae93dc58859c11066fa9726c363357 Mon Sep 17 00:00:00 2001 From: Eric Hare Date: Mon, 27 Jul 2026 09:54:05 -0700 Subject: [PATCH 2/5] Address Cassandra readiness review feedback --- CONFIGURATION.md | 22 ++++++++++ .../CassandraConnectionHealthCheck.java | 16 +++++-- .../sgv2/jsonapi/config/OperationsConfig.java | 8 ++-- .../CassandraConnectionHealthCheckTest.java | 43 +++++++++++++++++++ .../v1/SessionEvictionIntegrationTest.java | 2 +- 5 files changed, 83 insertions(+), 8 deletions(-) diff --git a/CONFIGURATION.md b/CONFIGURATION.md index b622e15f60..be02d9ba25 100644 --- a/CONFIGURATION.md +++ b/CONFIGURATION.md @@ -54,6 +54,28 @@ Other Quarkus properties that are specifically relevant for the service: | `stargate.jsonapi.operations.database-config.ddl-delay-millis` | `int` | `2000` | Delay between create table and create index to get the schema sync. | | `stargate.jsonapi.operations.vectorize-enabled` | `boolean` | `false` | Flag to enable server side vectorization. | +### Cassandra readiness + +When `stargate.jsonapi.operations.database-config.type` is `CASSANDRA`, the +`/stargate/health/ready` response includes a `cassandra-connection` check. The check obtains a +session through the application session cache and executes +`SELECT release_version FROM system.local`. + +If `stargate.jsonapi.operations.database-config.fixed-token` is configured, the readiness check +uses that token. Otherwise it connects with +`stargate.jsonapi.operations.database-config.user-name` and +`stargate.jsonapi.operations.database-config.password`, which both default to `cassandra`. These +credentials must be valid even when API clients supply different per-request credentials. The check +verifies connectivity with the configured default credentials; it does not validate every +request-specific credential. + +Readiness polling counts as session access and intentionally keeps the cached session active while +polling continues. Session acquisition and the validation query each have a five-second +timeout, so deployment probe timeouts should allow for both stages. The Helm chart defaults the +readiness probe timeout to ten seconds. + +For other database types, the check reports UP without accessing the Cassandra session cache. + ## Jsonapi metering configuration *Configuration for jsonapi metering, defined by [JsonApiMetricsConfig.java](io/stargate/sgv2/jsonapi/api/v1/metrics/JsonApiMetricsConfig.java).* diff --git a/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java b/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java index 67c1f06e2c..b5d0d2ccbc 100644 --- a/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java +++ b/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java @@ -27,9 +27,18 @@ /** * Health check that verifies Cassandra connectivity for the Data API readiness probe. * - *

The check obtains the Cassandra session used by normal requests from the session cache, - * verifies that it is open, and executes a lightweight query against {@code system.local}. It is - * only active when the Data API is configured with a {@link DatabaseType#CASSANDRA} backend. + *

The check obtains a session through the same session cache used by normal requests, verifies + * that it is open, and executes a lightweight query against {@code system.local}. It uses the fixed + * token when configured; otherwise it uses the configured Cassandra username and password. These + * default credentials verify base Cassandra connectivity, not the validity of every credential + * supplied on a request. + * + *

Readiness polling deliberately keeps the cached session active while polling continues. If + * requests use the same credentials, they share that session; otherwise the check maintains a + * dedicated session for the configured default credentials. + * + *

The database type is runtime configuration, so the bean remains registered for other database + * types. In those deployments it reports UP without accessing the Cassandra session cache. */ @Readiness @ApplicationScoped @@ -149,6 +158,7 @@ private static String encode(String value) { } private static boolean isUnreliableSessionFailure(Throwable throwable) { + // Session acquisition through Mutiny may wrap driver failures, so inspect the cause chain. var current = throwable; while (current != null) { if (current instanceof AllNodesFailedException diff --git a/src/main/java/io/stargate/sgv2/jsonapi/config/OperationsConfig.java b/src/main/java/io/stargate/sgv2/jsonapi/config/OperationsConfig.java index 981d89da89..63d56f7024 100644 --- a/src/main/java/io/stargate/sgv2/jsonapi/config/OperationsConfig.java +++ b/src/main/java/io/stargate/sgv2/jsonapi/config/OperationsConfig.java @@ -220,16 +220,16 @@ interface DatabaseConfig { DatabaseType type(); /** - * Username when connecting to cassandra database (when type is {@link DatabaseType#CASSANDRA}) - * and fixedToken is used + * Username used for Cassandra connections when fixedToken is configured, and by the Cassandra + * readiness check when fixedToken is not configured. */ @Nullable @WithDefault("cassandra") String userName(); /** - * Password when connecting to cassandra database (when type is {@link DatabaseType#CASSANDRA}) - * and fixedToken is used + * Password used for Cassandra connections when fixedToken is configured, and by the Cassandra + * readiness check when fixedToken is not configured. */ @Nullable @WithDefault("cassandra") diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java index b9898ad60d..582741d7c8 100644 --- a/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java @@ -5,6 +5,7 @@ import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.*; +import com.datastax.oss.driver.api.core.AllNodesFailedException; import com.datastax.oss.driver.api.core.CqlSession; import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; import com.datastax.oss.driver.api.core.connection.ClosedConnectionException; @@ -44,6 +45,7 @@ public class CassandraConnectionHealthCheckTest { @BeforeEach public void setup() { + TenantFactory.reset(); TenantFactory.initialize(DatabaseType.CASSANDRA); sessionCache = mock(CQLSessionCache.class); @@ -119,6 +121,8 @@ public void sessionAcquisitionFailureReportsDown() { assertThat(data) .containsEntry("error", "IllegalStateException") .containsEntry("message", "Cannot connect")); + verify(sessionCache, never()) + .evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); } @Test @@ -156,6 +160,43 @@ public void queryFailureReportsDownAndUnreliableSessionIsEvicted() { verify(sessionCache).evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); } + @Test + public void wrappedAllNodesFailedReportsDownAndUnreliableSessionIsEvicted() { + var allNodesFailed = mock(AllNodesFailedException.class); + + when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) + .thenReturn(Uni.createFrom().item(session)); + when(session.isClosed()).thenReturn(false); + when(session.execute(any(SimpleStatement.class))) + .thenThrow(new RuntimeException("Session failed", allNodesFailed)); + + var response = healthCheck.call(); + + assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); + assertThat(response.getData()) + .hasValueSatisfying( + data -> + assertThat(data) + .containsEntry("error", "RuntimeException") + .containsEntry("message", "Session failed")); + verify(sessionCache).evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + } + + @Test + public void queryFailureDoesNotEvictReliableSession() { + when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) + .thenReturn(Uni.createFrom().item(session)); + when(session.isClosed()).thenReturn(false); + when(session.execute(any(SimpleStatement.class))) + .thenThrow(new IllegalStateException("Invalid query")); + + var response = healthCheck.call(); + + assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); + verify(sessionCache, never()) + .evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + } + @Test public void sessionAcquisitionTimeoutReportsDown() { when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) @@ -164,6 +205,8 @@ public void sessionAcquisitionTimeoutReportsDown() { var response = healthCheck.call(); assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); + verify(sessionCache, never()) + .evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); verifyNoInteractions(session); } diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java index 663ed19f60..28ec1c2907 100644 --- a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java @@ -271,7 +271,7 @@ private void waitForReadinessStatus(String expectedStatus, long timeoutMillis) { lastResponse == null ? "no response" : "HTTP " + lastResponse.statusCode() + ": " + lastResponse.asString(); - throw new RuntimeException( + throw new AssertionError( "Readiness did not become " + expectedStatus + " within " From 55a232a4adb7a8476140670831f233566299f1ae Mon Sep 17 00:00:00 2001 From: Eric Hare Date: Mon, 27 Jul 2026 10:11:35 -0700 Subject: [PATCH 3/5] Reuse shared Cassandra test fixtures --- .../CassandraConnectionHealthCheckTest.java | 68 +++++++++++-------- 1 file changed, 40 insertions(+), 28 deletions(-) diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java index 582741d7c8..bbdf6e5ef2 100644 --- a/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java @@ -13,8 +13,8 @@ import com.datastax.oss.driver.api.core.cql.Row; import com.datastax.oss.driver.api.core.cql.SimpleStatement; import io.smallrye.mutiny.Uni; +import io.stargate.sgv2.jsonapi.TestConstants; import io.stargate.sgv2.jsonapi.api.request.UserAgent; -import io.stargate.sgv2.jsonapi.api.request.tenant.Tenant; import io.stargate.sgv2.jsonapi.api.request.tenant.TenantFactory; import io.stargate.sgv2.jsonapi.config.DatabaseType; import io.stargate.sgv2.jsonapi.config.OperationsConfig; @@ -33,10 +33,12 @@ /** Tests for {@link CassandraConnectionHealthCheck}. */ public class CassandraConnectionHealthCheckTest { - private static final String FIXED_TOKEN = "fixed-token"; + private static final UserAgent HEALTH_CHECK_USER_AGENT = new UserAgent("DataAPI-HealthCheck/1.0"); private static final String USER_NAME = "test-user"; private static final String PASSWORD = "test-password"; + private final TestConstants TEST_CONSTANTS = new TestConstants(); + private CQLSessionCache sessionCache; private OperationsConfig operationsConfig; private OperationsConfig.DatabaseConfig databaseConfig; @@ -53,7 +55,7 @@ public void setup() { databaseConfig = mock(OperationsConfig.DatabaseConfig.class); when(databaseConfig.type()).thenReturn(DatabaseType.CASSANDRA); - when(databaseConfig.fixedToken()).thenReturn(Optional.of(FIXED_TOKEN)); + when(databaseConfig.fixedToken()).thenReturn(Optional.of(TEST_CONSTANTS.AUTH_TOKEN)); when(databaseConfig.userName()).thenReturn(USER_NAME); when(databaseConfig.password()).thenReturn(PASSWORD); @@ -81,8 +83,7 @@ public void successfulHealthCheck() { var resultSet = mock(ResultSet.class); var row = mock(Row.class); - when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) - .thenReturn(Uni.createFrom().item(session)); + sessionRequestReturns(Uni.createFrom().item(session)); when(session.isClosed()).thenReturn(false); when(session.execute(any(SimpleStatement.class))).thenReturn(resultSet); when(resultSet.one()).thenReturn(row); @@ -109,8 +110,7 @@ public void successfulHealthCheck() { @Test public void sessionAcquisitionFailureReportsDown() { - when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) - .thenReturn(Uni.createFrom().failure(new IllegalStateException("Cannot connect"))); + sessionRequestReturns(Uni.createFrom().failure(new IllegalStateException("Cannot connect"))); var response = healthCheck.call(); @@ -121,14 +121,12 @@ public void sessionAcquisitionFailureReportsDown() { assertThat(data) .containsEntry("error", "IllegalStateException") .containsEntry("message", "Cannot connect")); - verify(sessionCache, never()) - .evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + verifySessionNotEvicted(); } @Test public void closedSessionReportsDownAndIsEvicted() { - when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) - .thenReturn(Uni.createFrom().item(session)); + sessionRequestReturns(Uni.createFrom().item(session)); when(session.isClosed()).thenReturn(true); var response = healthCheck.call(); @@ -136,14 +134,13 @@ public void closedSessionReportsDownAndIsEvicted() { assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); assertThat(response.getData()) .hasValueSatisfying(data -> assertThat(data).containsEntry("reason", "Session is closed")); - verify(sessionCache).evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + verifySessionEvicted(); verify(session, never()).execute(any(SimpleStatement.class)); } @Test public void queryFailureReportsDownAndUnreliableSessionIsEvicted() { - when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) - .thenReturn(Uni.createFrom().item(session)); + sessionRequestReturns(Uni.createFrom().item(session)); when(session.isClosed()).thenReturn(false); when(session.execute(any(SimpleStatement.class))) .thenThrow(new ClosedConnectionException("Connection is closed")); @@ -157,15 +154,14 @@ public void queryFailureReportsDownAndUnreliableSessionIsEvicted() { assertThat(data) .containsEntry("error", "ClosedConnectionException") .containsEntry("message", "Connection is closed")); - verify(sessionCache).evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + verifySessionEvicted(); } @Test public void wrappedAllNodesFailedReportsDownAndUnreliableSessionIsEvicted() { var allNodesFailed = mock(AllNodesFailedException.class); - when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) - .thenReturn(Uni.createFrom().item(session)); + sessionRequestReturns(Uni.createFrom().item(session)); when(session.isClosed()).thenReturn(false); when(session.execute(any(SimpleStatement.class))) .thenThrow(new RuntimeException("Session failed", allNodesFailed)); @@ -179,13 +175,12 @@ public void wrappedAllNodesFailedReportsDownAndUnreliableSessionIsEvicted() { assertThat(data) .containsEntry("error", "RuntimeException") .containsEntry("message", "Session failed")); - verify(sessionCache).evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + verifySessionEvicted(); } @Test public void queryFailureDoesNotEvictReliableSession() { - when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) - .thenReturn(Uni.createFrom().item(session)); + sessionRequestReturns(Uni.createFrom().item(session)); when(session.isClosed()).thenReturn(false); when(session.execute(any(SimpleStatement.class))) .thenThrow(new IllegalStateException("Invalid query")); @@ -193,20 +188,17 @@ public void queryFailureDoesNotEvictReliableSession() { var response = healthCheck.call(); assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); - verify(sessionCache, never()) - .evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + verifySessionNotEvicted(); } @Test public void sessionAcquisitionTimeoutReportsDown() { - when(sessionCache.getSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class))) - .thenReturn(Uni.createFrom().nothing()); + sessionRequestReturns(Uni.createFrom().nothing()); var response = healthCheck.call(); assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); - verify(sessionCache, never()) - .evictSession(any(Tenant.class), eq(FIXED_TOKEN), any(UserAgent.class)); + verifySessionNotEvicted(); verifyNoInteractions(session); } @@ -216,7 +208,8 @@ public void configuredCredentialsAreUsedWhenFixedTokenIsNotSet() { healthCheck = new CassandraConnectionHealthCheck(sessionCache, operationsConfig, Duration.ofMillis(100)); - when(sessionCache.getSession(any(Tenant.class), any(String.class), any(UserAgent.class))) + when(sessionCache.getSession( + eq(TEST_CONSTANTS.CASSANDRA_TENANT), any(String.class), eq(HEALTH_CHECK_USER_AGENT))) .thenReturn(Uni.createFrom().item(session)); when(session.isClosed()).thenReturn(true); @@ -227,7 +220,8 @@ public void configuredCredentialsAreUsedWhenFixedTokenIsNotSet() { + Base64.getEncoder().encodeToString(USER_NAME.getBytes(StandardCharsets.UTF_8)) + ":" + Base64.getEncoder().encodeToString(PASSWORD.getBytes(StandardCharsets.UTF_8)); - verify(sessionCache).getSession(any(Tenant.class), eq(expectedToken), any(UserAgent.class)); + verify(sessionCache) + .getSession(TEST_CONSTANTS.CASSANDRA_TENANT, expectedToken, HEALTH_CHECK_USER_AGENT); } @Test @@ -248,4 +242,22 @@ public void nonCassandraDatabaseDoesNotRunConnectivityCheck() { "Cassandra connectivity check is not applicable for database type ASTRA")); verifyNoInteractions(sessionCache); } + + private void sessionRequestReturns(Uni sessionResult) { + when(sessionCache.getSession( + TEST_CONSTANTS.CASSANDRA_TENANT, TEST_CONSTANTS.AUTH_TOKEN, HEALTH_CHECK_USER_AGENT)) + .thenReturn(sessionResult); + } + + private void verifySessionEvicted() { + verify(sessionCache) + .evictSession( + TEST_CONSTANTS.CASSANDRA_TENANT, TEST_CONSTANTS.AUTH_TOKEN, HEALTH_CHECK_USER_AGENT); + } + + private void verifySessionNotEvicted() { + verify(sessionCache, never()) + .evictSession( + TEST_CONSTANTS.CASSANDRA_TENANT, TEST_CONSTANTS.AUTH_TOKEN, HEALTH_CHECK_USER_AGENT); + } } From 861d862bb2fa1c7542238c65cc1275a3c046dc4f Mon Sep 17 00:00:00 2001 From: Eric Hare Date: Mon, 3 Aug 2026 20:20:52 -0700 Subject: [PATCH 4/5] chore: refactor to support astra readiness --- CONFIGURATION.md | 54 ++-- .../CassandraConnectionHealthCheck.java | 180 ------------ .../api/health/DatabaseReadinessCheck.java | 72 +++++ .../api/v1/DatabaseReadinessResource.java | 103 +++++++ .../sgv2/jsonapi/config/OperationsConfig.java | 8 +- .../metrics/TenantRequestMetricsFilter.java | 46 +-- src/main/resources/application.yaml | 2 +- .../CassandraConnectionHealthCheckTest.java | 263 ------------------ .../health/DatabaseReadinessCheckTest.java | 147 ++++++++++ .../api/v1/DatabaseReadinessResourceTest.java | 181 ++++++++++++ .../v1/SessionEvictionIntegrationTest.java | 34 ++- .../TenantRequestMetricsFilterTest.java | 40 +++ 12 files changed, 627 insertions(+), 503 deletions(-) delete mode 100644 src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java create mode 100644 src/main/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheck.java create mode 100644 src/main/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResource.java delete mode 100644 src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java create mode 100644 src/test/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheckTest.java create mode 100644 src/test/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResourceTest.java create mode 100644 src/test/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilterTest.java diff --git a/CONFIGURATION.md b/CONFIGURATION.md index be02d9ba25..70e0b43147 100644 --- a/CONFIGURATION.md +++ b/CONFIGURATION.md @@ -54,27 +54,39 @@ Other Quarkus properties that are specifically relevant for the service: | `stargate.jsonapi.operations.database-config.ddl-delay-millis` | `int` | `2000` | Delay between create table and create index to get the schema sync. | | `stargate.jsonapi.operations.vectorize-enabled` | `boolean` | `false` | Flag to enable server side vectorization. | -### Cassandra readiness - -When `stargate.jsonapi.operations.database-config.type` is `CASSANDRA`, the -`/stargate/health/ready` response includes a `cassandra-connection` check. The check obtains a -session through the application session cache and executes -`SELECT release_version FROM system.local`. - -If `stargate.jsonapi.operations.database-config.fixed-token` is configured, the readiness check -uses that token. Otherwise it connects with -`stargate.jsonapi.operations.database-config.user-name` and -`stargate.jsonapi.operations.database-config.password`, which both default to `cassandra`. These -credentials must be valid even when API clients supply different per-request credentials. The check -verifies connectivity with the configured default credentials; it does not validate every -request-specific credential. - -Readiness polling counts as session access and intentionally keeps the cached session active while -polling continues. Session acquisition and the validation query each have a five-second -timeout, so deployment probe timeouts should allow for both stages. The Helm chart defaults the -readiness probe timeout to ten seconds. - -For other database types, the check reports UP without accessing the Cassandra session cache. +### Database readiness + +`GET /v1/health/ready` is an authenticated database readiness endpoint used for both Astra and +Cassandra deployments. It uses the request's tenant, `Token` header, and `User-Agent` to obtain a +session through the normal session cache. The Data API does not store separate readiness +credentials. + +The endpoint executes `SELECT * FROM datastax_sla.check LIMIT 1` at `LOCAL_QUORUM`, using the +`table-read` driver profile for the remaining read settings. An `UP` response therefore confirms +that the coordinator can complete a read from a replicated table at local quorum. It does not +validate every tenant's credentials, write availability, or cross-region availability. + +The deployment must provide a dedicated canary tenant and credentials for this request and must +provision a `datastax_sla.check` table that the canary principal can read. Its replication factor +must be appropriate for the deployment (greater than one in a multi-node local data center) so +`LOCAL_QUORUM` requires responses from multiple replicas. Astra callers must use the canary database +hostname so the tenant and region are resolved from `Host`; Cassandra ignores the tenant portion of +`Host`. The caller must also send the exact User-Agent configured by +`stargate.jsonapi.operations.sla-user-agent`, allowing a dedicated canary session to use the shorter +SLA session TTL instead of being treated like normal client traffic. + +The check is fully asynchronous and has a five-second timeout. It returns HTTP 200 with +`{"status":"UP"}` after a successful read, HTTP 503 with `{"status":"DOWN"}` after a database +failure or timeout, and HTTP 401 when the `Token` header is missing or authentication fails. + +Kubernetes or an SLA checker must call each pod directly for this endpoint to control per-pod +readiness. An external request sent through a load balancer does not establish which pod is ready. +Restrict the endpoint to trusted probe traffic with deployment controls such as a NetworkPolicy, +mTLS, or an ingress ACL and rate limit. Kubernetes `httpGet` headers cannot reference a Secret, so +delivery of the canary token is intentionally outside the Data API configuration. Prefer an +external checker or a Secret-mounted file read by an `exec` probe; do not put the token literally in +the probe command or shell trace. The unauthenticated Quarkus health endpoints under the +non-application path do not include this database check. ## Jsonapi metering configuration diff --git a/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java b/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java deleted file mode 100644 index b5d0d2ccbc..0000000000 --- a/src/main/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheck.java +++ /dev/null @@ -1,180 +0,0 @@ -package io.stargate.sgv2.jsonapi.api.health; - -import com.datastax.oss.driver.api.core.AllNodesFailedException; -import com.datastax.oss.driver.api.core.connection.ClosedConnectionException; -import com.datastax.oss.driver.api.core.cql.SimpleStatement; -import com.google.common.annotations.VisibleForTesting; -import io.stargate.sgv2.jsonapi.api.request.UserAgent; -import io.stargate.sgv2.jsonapi.api.request.tenant.Tenant; -import io.stargate.sgv2.jsonapi.api.request.tenant.TenantFactory; -import io.stargate.sgv2.jsonapi.config.DatabaseType; -import io.stargate.sgv2.jsonapi.config.OperationsConfig; -import io.stargate.sgv2.jsonapi.service.cqldriver.CQLSessionCache; -import io.stargate.sgv2.jsonapi.service.cqldriver.CqlCredentials; -import io.stargate.sgv2.jsonapi.service.cqldriver.CqlSessionCacheSupplier; -import jakarta.enterprise.context.ApplicationScoped; -import jakarta.inject.Inject; -import java.nio.charset.StandardCharsets; -import java.time.Duration; -import java.util.Base64; -import java.util.Objects; -import org.eclipse.microprofile.health.HealthCheck; -import org.eclipse.microprofile.health.HealthCheckResponse; -import org.eclipse.microprofile.health.Readiness; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * Health check that verifies Cassandra connectivity for the Data API readiness probe. - * - *

The check obtains a session through the same session cache used by normal requests, verifies - * that it is open, and executes a lightweight query against {@code system.local}. It uses the fixed - * token when configured; otherwise it uses the configured Cassandra username and password. These - * default credentials verify base Cassandra connectivity, not the validity of every credential - * supplied on a request. - * - *

Readiness polling deliberately keeps the cached session active while polling continues. If - * requests use the same credentials, they share that session; otherwise the check maintains a - * dedicated session for the configured default credentials. - * - *

The database type is runtime configuration, so the bean remains registered for other database - * types. In those deployments it reports UP without accessing the Cassandra session cache. - */ -@Readiness -@ApplicationScoped -public class CassandraConnectionHealthCheck implements HealthCheck { - - private static final Logger LOGGER = - LoggerFactory.getLogger(CassandraConnectionHealthCheck.class); - - private static final String HEALTH_CHECK_NAME = "cassandra-connection"; - private static final String HEALTH_CHECK_QUERY = "SELECT release_version FROM system.local"; - private static final Duration HEALTH_CHECK_TIMEOUT = Duration.ofSeconds(5); - private static final UserAgent HEALTH_CHECK_USER_AGENT = new UserAgent("DataAPI-HealthCheck/1.0"); - - private final CQLSessionCache sessionCache; - private final OperationsConfig operationsConfig; - private final Duration timeout; - private final String authToken; - - @Inject - public CassandraConnectionHealthCheck( - CqlSessionCacheSupplier sessionCacheSupplier, OperationsConfig operationsConfig) { - this( - Objects.requireNonNull(sessionCacheSupplier, "sessionCacheSupplier must not be null").get(), - operationsConfig, - HEALTH_CHECK_TIMEOUT); - } - - @VisibleForTesting - CassandraConnectionHealthCheck( - CQLSessionCache sessionCache, OperationsConfig operationsConfig, Duration timeout) { - this.sessionCache = Objects.requireNonNull(sessionCache, "sessionCache must not be null"); - this.operationsConfig = - Objects.requireNonNull(operationsConfig, "operationsConfig must not be null"); - this.timeout = Objects.requireNonNull(timeout, "timeout must not be null"); - this.authToken = - operationsConfig.databaseConfig().type() == DatabaseType.CASSANDRA - ? createAuthToken(operationsConfig.databaseConfig()) - : null; - } - - @Override - public HealthCheckResponse call() { - var responseBuilder = HealthCheckResponse.named(HEALTH_CHECK_NAME); - - if (operationsConfig.databaseConfig().type() != DatabaseType.CASSANDRA) { - return responseBuilder - .up() - .withData( - "reason", - "Cassandra connectivity check is not applicable for database type " - + operationsConfig.databaseConfig().type()) - .build(); - } - - var healthCheckTenant = TenantFactory.instance().create(null); - - try { - var session = - sessionCache - .getSession(healthCheckTenant, authToken, HEALTH_CHECK_USER_AGENT) - .await() - .atMost(timeout); - - if (session.isClosed()) { - LOGGER.warn("Cassandra session is closed during health check"); - evictSession(healthCheckTenant); - return responseBuilder.down().withData("reason", "Session is closed").build(); - } - - var statement = - SimpleStatement.builder(HEALTH_CHECK_QUERY) - .setTimeout(timeout) - .setConsistencyLevel(operationsConfig.queriesConfig().consistency().reads()) - .build(); - - var resultSet = session.execute(statement); - var row = resultSet.one(); - var version = row != null ? row.getString("release_version") : "unknown"; - - LOGGER.trace("Cassandra health check passed, version: {}", version); - - return responseBuilder - .up() - .withData("cassandra_version", version) - .withData("session_name", session.getName()) - .build(); - } catch (Exception e) { - if (isUnreliableSessionFailure(e)) { - evictSession(healthCheckTenant); - } - - LOGGER.error("Cassandra health check failed", e); - return responseBuilder - .down() - .withData("error", e.getClass().getSimpleName()) - .withData("message", e.getMessage() != null ? e.getMessage() : "Unknown error") - .build(); - } - } - - private static String createAuthToken(OperationsConfig.DatabaseConfig databaseConfig) { - return databaseConfig - .fixedToken() - .orElseGet( - () -> - CqlCredentials.USERNAME_PASSWORD_TOKEN_PREFIX - + encode(databaseConfig.userName()) - + ":" - + encode(databaseConfig.password())); - } - - private static String encode(String value) { - return Base64.getEncoder() - .encodeToString( - Objects.requireNonNull(value, "Cassandra credential must not be null") - .getBytes(StandardCharsets.UTF_8)); - } - - private static boolean isUnreliableSessionFailure(Throwable throwable) { - // Session acquisition through Mutiny may wrap driver failures, so inspect the cause chain. - var current = throwable; - while (current != null) { - if (current instanceof AllNodesFailedException - || current instanceof ClosedConnectionException) { - return true; - } - current = current.getCause(); - } - return false; - } - - private void evictSession(Tenant tenant) { - try { - sessionCache.evictSession(tenant, authToken, HEALTH_CHECK_USER_AGENT); - } catch (Exception e) { - LOGGER.warn("Unable to evict the Cassandra session after a failed health check", e); - } - } -} diff --git a/src/main/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheck.java b/src/main/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheck.java new file mode 100644 index 0000000000..bbec7687d1 --- /dev/null +++ b/src/main/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheck.java @@ -0,0 +1,72 @@ +package io.stargate.sgv2.jsonapi.api.health; + +import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import com.google.common.annotations.VisibleForTesting; +import io.smallrye.mutiny.Uni; +import io.stargate.sgv2.jsonapi.api.request.RequestContext; +import io.stargate.sgv2.jsonapi.service.cqldriver.CQLSessionCache; +import io.stargate.sgv2.jsonapi.service.cqldriver.executor.CommandQueryExecutor; +import java.time.Duration; +import java.util.Objects; +import java.util.function.Supplier; + +/** + * Runs the database probe exposed at {@code GET /v1/health/ready}. + * + *

This class is constructed by the JAX-RS resource and is not a CDI bean or a MicroProfile + * health check. The caller's request context supplies the tenant, token, and User-Agent for both + * Astra and Cassandra connections. + */ +public final class DatabaseReadinessCheck { + + private static final String READINESS_QUERY = "SELECT * FROM datastax_sla.check LIMIT 1"; + private static final Duration DEFAULT_TIMEOUT = Duration.ofSeconds(5); + + private final Supplier sessionCacheSupplier; + private final SimpleStatement statement; + private final Duration timeout; + + public DatabaseReadinessCheck(Supplier sessionCacheSupplier) { + this(sessionCacheSupplier, DEFAULT_TIMEOUT); + } + + @VisibleForTesting + DatabaseReadinessCheck(CQLSessionCache sessionCache, Duration timeout) { + this(() -> sessionCache, timeout); + } + + private DatabaseReadinessCheck(Supplier sessionCacheSupplier, Duration timeout) { + this.sessionCacheSupplier = + Objects.requireNonNull(sessionCacheSupplier, "sessionCacheSupplier must not be null"); + this.timeout = Objects.requireNonNull(timeout, "timeout must not be null"); + this.statement = + SimpleStatement.builder(READINESS_QUERY) + .setConsistencyLevel(DefaultConsistencyLevel.LOCAL_QUORUM) + .setTimeout(timeout) + .build(); + } + + /** + * Executes a replicated table read at {@code LOCAL_QUORUM}, using the {@code table-read} driver + * profile for the remaining read settings. + */ + public Uni check(RequestContext requestContext) { + Objects.requireNonNull(requestContext, "requestContext must not be null"); + + return Uni.createFrom() + .deferred( + () -> { + var sessionCache = + Objects.requireNonNull( + sessionCacheSupplier.get(), "sessionCacheSupplier returned null"); + return new CommandQueryExecutor( + sessionCache, requestContext, CommandQueryExecutor.QueryTarget.TABLE) + .executeRead(statement) + .replaceWithVoid(); + }) + .ifNoItem() + .after(timeout) + .fail(); + } +} diff --git a/src/main/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResource.java b/src/main/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResource.java new file mode 100644 index 0000000000..38bafaa95f --- /dev/null +++ b/src/main/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResource.java @@ -0,0 +1,103 @@ +package io.stargate.sgv2.jsonapi.api.v1; + +import io.quarkus.security.UnauthorizedException; +import io.smallrye.mutiny.Uni; +import io.stargate.sgv2.jsonapi.api.health.DatabaseReadinessCheck; +import io.stargate.sgv2.jsonapi.api.request.RequestContext; +import io.stargate.sgv2.jsonapi.config.constants.OpenApiConstants; +import io.stargate.sgv2.jsonapi.exception.APISecurityException; +import io.stargate.sgv2.jsonapi.service.cqldriver.CqlSessionCacheSupplier; +import jakarta.inject.Inject; +import jakarta.ws.rs.GET; +import jakarta.ws.rs.Path; +import jakarta.ws.rs.Produces; +import jakarta.ws.rs.core.MediaType; +import jakarta.ws.rs.core.Response; +import org.eclipse.microprofile.openapi.annotations.Operation; +import org.eclipse.microprofile.openapi.annotations.media.Content; +import org.eclipse.microprofile.openapi.annotations.media.Schema; +import org.eclipse.microprofile.openapi.annotations.responses.APIResponse; +import org.eclipse.microprofile.openapi.annotations.responses.APIResponses; +import org.eclipse.microprofile.openapi.annotations.security.SecurityRequirement; +import org.jboss.resteasy.reactive.RestResponse; + +/** + * Authenticated database readiness endpoint registered through Quarkus JAX-RS resource discovery. + * + *

{@code GET /v1/health/ready} runs the same request-scoped probe for Astra and Cassandra. A + * successful probe returns HTTP 200; invalid credentials return HTTP 401; and a database failure or + * timeout returns HTTP 503. The existing {@code /v1/*} security policy rejects requests without a + * token before this resource is called. + */ +@Path(DatabaseReadinessResource.BASE_PATH) +@Produces(MediaType.APPLICATION_JSON) +@SecurityRequirement(name = OpenApiConstants.SecuritySchemes.TOKEN) +public class DatabaseReadinessResource { + + public static final String BASE_PATH = GeneralResource.BASE_PATH + "/health/ready"; + + private static final ReadinessResponse UP = new ReadinessResponse("UP"); + private static final ReadinessResponse DOWN = new ReadinessResponse("DOWN"); + + private final DatabaseReadinessCheck readinessCheck; + private final RequestContext requestContext; + + @Inject + public DatabaseReadinessResource( + CqlSessionCacheSupplier sessionCacheSupplier, RequestContext requestContext) { + this.readinessCheck = new DatabaseReadinessCheck(sessionCacheSupplier); + this.requestContext = requestContext; + } + + @GET + @Operation( + summary = "Check database readiness", + description = + "Uses the authenticated request tenant and token to perform a LOCAL_QUORUM read.") + @APIResponses({ + @APIResponse( + responseCode = "200", + description = "The database completed the readiness read.", + content = + @Content( + mediaType = MediaType.APPLICATION_JSON, + schema = @Schema(implementation = ReadinessResponse.class))), + @APIResponse(responseCode = "401", description = "The token is missing or invalid."), + @APIResponse( + responseCode = "503", + description = "The database read failed or timed out.", + content = + @Content( + mediaType = MediaType.APPLICATION_JSON, + schema = @Schema(implementation = ReadinessResponse.class))) + }) + public Uni> ready() { + return readinessCheck + .check(requestContext) + .map(ignored -> RestResponse.ok(UP)) + .onFailure(DatabaseReadinessResource::isUnauthorized) + .recoverWithItem( + failure -> + RestResponse.ResponseBuilder.create(Response.Status.UNAUTHORIZED, DOWN).build()) + .onFailure() + .recoverWithItem( + failure -> + RestResponse.ResponseBuilder.create(Response.Status.SERVICE_UNAVAILABLE, DOWN) + .build()); + } + + private static boolean isUnauthorized(Throwable failure) { + var current = failure; + while (current != null) { + if (current instanceof UnauthorizedException + || current instanceof APISecurityException apiException + && apiException.httpStatus == Response.Status.UNAUTHORIZED.getStatusCode()) { + return true; + } + current = current.getCause(); + } + return false; + } + + public record ReadinessResponse(String status) {} +} diff --git a/src/main/java/io/stargate/sgv2/jsonapi/config/OperationsConfig.java b/src/main/java/io/stargate/sgv2/jsonapi/config/OperationsConfig.java index 63d56f7024..981d89da89 100644 --- a/src/main/java/io/stargate/sgv2/jsonapi/config/OperationsConfig.java +++ b/src/main/java/io/stargate/sgv2/jsonapi/config/OperationsConfig.java @@ -220,16 +220,16 @@ interface DatabaseConfig { DatabaseType type(); /** - * Username used for Cassandra connections when fixedToken is configured, and by the Cassandra - * readiness check when fixedToken is not configured. + * Username when connecting to cassandra database (when type is {@link DatabaseType#CASSANDRA}) + * and fixedToken is used */ @Nullable @WithDefault("cassandra") String userName(); /** - * Password used for Cassandra connections when fixedToken is configured, and by the Cassandra - * readiness check when fixedToken is not configured. + * Password when connecting to cassandra database (when type is {@link DatabaseType#CASSANDRA}) + * and fixedToken is used */ @Nullable @WithDefault("cassandra") diff --git a/src/main/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilter.java b/src/main/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilter.java index e365b4891e..c633208225 100644 --- a/src/main/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilter.java +++ b/src/main/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilter.java @@ -21,6 +21,7 @@ import io.micrometer.core.instrument.Tag; import io.micrometer.core.instrument.Tags; import io.stargate.sgv2.jsonapi.api.request.RequestContext; +import io.stargate.sgv2.jsonapi.api.v1.DatabaseReadinessResource; import io.stargate.sgv2.jsonapi.api.v1.metrics.MetricsConfig; import jakarta.enterprise.context.ApplicationScoped; import jakarta.inject.Inject; @@ -71,31 +72,36 @@ public TenantRequestMetricsFilter( @ServerResponseFilter public void record( ContainerRequestContext requestContext, ContainerResponseContext responseContext) { - // only if enabled - if (config.enabled()) { - - // resolve tenant - Tag tenantTag = Tag.of(config.tenantTag(), this.requestContext.tenant().toString()); + if (!config.enabled() || isDatabaseReadinessRequest(requestContext)) { + return; + } - // resolve error - boolean error = responseContext.getStatus() >= 500; - Tag errorTag = error ? ExceptionMetrics.TAG_ERROR_TRUE : ExceptionMetrics.TAG_ERROR_FALSE; + // resolve tenant + Tag tenantTag = Tag.of(config.tenantTag(), this.requestContext.tenant().toString()); - // check if we need user agent as well - Tags tags = Tags.of(tenantTag, errorTag); - if (config.userAgentTagEnabled()) { - String userAgentValue = getUserAgentValue(requestContext); - tags = tags.and(Tag.of(config.userAgentTag(), userAgentValue)); - } + // resolve error + boolean error = responseContext.getStatus() >= 500; + Tag errorTag = error ? ExceptionMetrics.TAG_ERROR_TRUE : ExceptionMetrics.TAG_ERROR_FALSE; - // add http status code - if (config.statusTagEnabled()) { - tags = tags.and(Tag.of(config.statusTag(), String.valueOf(responseContext.getStatus()))); - } + // check if we need user agent as well + Tags tags = Tags.of(tenantTag, errorTag); + if (config.userAgentTagEnabled()) { + String userAgentValue = getUserAgentValue(requestContext); + tags = tags.and(Tag.of(config.userAgentTag(), userAgentValue)); + } - // record - meterRegistry.counter(config.metricName(), tags).increment(); + // add http status code + if (config.statusTagEnabled()) { + tags = tags.and(Tag.of(config.statusTag(), String.valueOf(responseContext.getStatus()))); } + + // record + meterRegistry.counter(config.metricName(), tags).increment(); + } + + private static boolean isDatabaseReadinessRequest(ContainerRequestContext requestContext) { + return DatabaseReadinessResource.BASE_PATH.equals( + requestContext.getUriInfo().getRequestUri().getPath()); } private String getUserAgentValue(ContainerRequestContext requestContext) { diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index 354232535c..1103fef1b3 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -159,7 +159,7 @@ quarkus: http-server: # ignore all non-application uris, as well as the custom set suppress-non-application-uris: true - ignore-patterns: /,/metrics,/swagger-ui.*,.*\.html + ignore-patterns: /,/metrics,/swagger-ui.*,.*\.html,/v1/health/ready # due to the https://github.com/quarkusio/quarkus/issues/24938 # we need to define uri templating on our own for now diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java deleted file mode 100644 index bbdf6e5ef2..0000000000 --- a/src/test/java/io/stargate/sgv2/jsonapi/api/health/CassandraConnectionHealthCheckTest.java +++ /dev/null @@ -1,263 +0,0 @@ -package io.stargate.sgv2.jsonapi.api.health; - -import static org.assertj.core.api.Assertions.assertThat; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Mockito.*; - -import com.datastax.oss.driver.api.core.AllNodesFailedException; -import com.datastax.oss.driver.api.core.CqlSession; -import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; -import com.datastax.oss.driver.api.core.connection.ClosedConnectionException; -import com.datastax.oss.driver.api.core.cql.ResultSet; -import com.datastax.oss.driver.api.core.cql.Row; -import com.datastax.oss.driver.api.core.cql.SimpleStatement; -import io.smallrye.mutiny.Uni; -import io.stargate.sgv2.jsonapi.TestConstants; -import io.stargate.sgv2.jsonapi.api.request.UserAgent; -import io.stargate.sgv2.jsonapi.api.request.tenant.TenantFactory; -import io.stargate.sgv2.jsonapi.config.DatabaseType; -import io.stargate.sgv2.jsonapi.config.OperationsConfig; -import io.stargate.sgv2.jsonapi.service.cqldriver.CQLSessionCache; -import io.stargate.sgv2.jsonapi.service.cqldriver.CqlCredentials; -import java.nio.charset.StandardCharsets; -import java.time.Duration; -import java.util.Base64; -import java.util.Optional; -import org.eclipse.microprofile.health.HealthCheckResponse; -import org.junit.jupiter.api.AfterEach; -import org.junit.jupiter.api.BeforeEach; -import org.junit.jupiter.api.Test; -import org.mockito.ArgumentCaptor; - -/** Tests for {@link CassandraConnectionHealthCheck}. */ -public class CassandraConnectionHealthCheckTest { - - private static final UserAgent HEALTH_CHECK_USER_AGENT = new UserAgent("DataAPI-HealthCheck/1.0"); - private static final String USER_NAME = "test-user"; - private static final String PASSWORD = "test-password"; - - private final TestConstants TEST_CONSTANTS = new TestConstants(); - - private CQLSessionCache sessionCache; - private OperationsConfig operationsConfig; - private OperationsConfig.DatabaseConfig databaseConfig; - private CqlSession session; - private CassandraConnectionHealthCheck healthCheck; - - @BeforeEach - public void setup() { - TenantFactory.reset(); - TenantFactory.initialize(DatabaseType.CASSANDRA); - - sessionCache = mock(CQLSessionCache.class); - session = mock(CqlSession.class); - - databaseConfig = mock(OperationsConfig.DatabaseConfig.class); - when(databaseConfig.type()).thenReturn(DatabaseType.CASSANDRA); - when(databaseConfig.fixedToken()).thenReturn(Optional.of(TEST_CONSTANTS.AUTH_TOKEN)); - when(databaseConfig.userName()).thenReturn(USER_NAME); - when(databaseConfig.password()).thenReturn(PASSWORD); - - var consistencyConfig = mock(OperationsConfig.QueriesConfig.ConsistencyConfig.class); - when(consistencyConfig.reads()).thenReturn(DefaultConsistencyLevel.LOCAL_QUORUM); - - var queriesConfig = mock(OperationsConfig.QueriesConfig.class); - when(queriesConfig.consistency()).thenReturn(consistencyConfig); - - operationsConfig = mock(OperationsConfig.class); - when(operationsConfig.databaseConfig()).thenReturn(databaseConfig); - when(operationsConfig.queriesConfig()).thenReturn(queriesConfig); - - healthCheck = - new CassandraConnectionHealthCheck(sessionCache, operationsConfig, Duration.ofMillis(100)); - } - - @AfterEach - public void cleanup() { - TenantFactory.reset(); - } - - @Test - public void successfulHealthCheck() { - var resultSet = mock(ResultSet.class); - var row = mock(Row.class); - - sessionRequestReturns(Uni.createFrom().item(session)); - when(session.isClosed()).thenReturn(false); - when(session.execute(any(SimpleStatement.class))).thenReturn(resultSet); - when(resultSet.one()).thenReturn(row); - when(row.getString("release_version")).thenReturn("6.9.21"); - when(session.getName()).thenReturn("SINGLE-TENANT"); - - var response = healthCheck.call(); - - assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.UP); - assertThat(response.getData()) - .hasValueSatisfying(data -> assertThat(data).containsEntry("cassandra_version", "6.9.21")); - assertThat(response.getData()) - .hasValueSatisfying( - data -> assertThat(data).containsEntry("session_name", "SINGLE-TENANT")); - - var statementCaptor = ArgumentCaptor.forClass(SimpleStatement.class); - verify(session).execute(statementCaptor.capture()); - assertThat(statementCaptor.getValue().getQuery()) - .isEqualTo("SELECT release_version FROM system.local"); - assertThat(statementCaptor.getValue().getTimeout()).isEqualTo(Duration.ofMillis(100)); - assertThat(statementCaptor.getValue().getConsistencyLevel()) - .isEqualTo(DefaultConsistencyLevel.LOCAL_QUORUM); - } - - @Test - public void sessionAcquisitionFailureReportsDown() { - sessionRequestReturns(Uni.createFrom().failure(new IllegalStateException("Cannot connect"))); - - var response = healthCheck.call(); - - assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); - assertThat(response.getData()) - .hasValueSatisfying( - data -> - assertThat(data) - .containsEntry("error", "IllegalStateException") - .containsEntry("message", "Cannot connect")); - verifySessionNotEvicted(); - } - - @Test - public void closedSessionReportsDownAndIsEvicted() { - sessionRequestReturns(Uni.createFrom().item(session)); - when(session.isClosed()).thenReturn(true); - - var response = healthCheck.call(); - - assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); - assertThat(response.getData()) - .hasValueSatisfying(data -> assertThat(data).containsEntry("reason", "Session is closed")); - verifySessionEvicted(); - verify(session, never()).execute(any(SimpleStatement.class)); - } - - @Test - public void queryFailureReportsDownAndUnreliableSessionIsEvicted() { - sessionRequestReturns(Uni.createFrom().item(session)); - when(session.isClosed()).thenReturn(false); - when(session.execute(any(SimpleStatement.class))) - .thenThrow(new ClosedConnectionException("Connection is closed")); - - var response = healthCheck.call(); - - assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); - assertThat(response.getData()) - .hasValueSatisfying( - data -> - assertThat(data) - .containsEntry("error", "ClosedConnectionException") - .containsEntry("message", "Connection is closed")); - verifySessionEvicted(); - } - - @Test - public void wrappedAllNodesFailedReportsDownAndUnreliableSessionIsEvicted() { - var allNodesFailed = mock(AllNodesFailedException.class); - - sessionRequestReturns(Uni.createFrom().item(session)); - when(session.isClosed()).thenReturn(false); - when(session.execute(any(SimpleStatement.class))) - .thenThrow(new RuntimeException("Session failed", allNodesFailed)); - - var response = healthCheck.call(); - - assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); - assertThat(response.getData()) - .hasValueSatisfying( - data -> - assertThat(data) - .containsEntry("error", "RuntimeException") - .containsEntry("message", "Session failed")); - verifySessionEvicted(); - } - - @Test - public void queryFailureDoesNotEvictReliableSession() { - sessionRequestReturns(Uni.createFrom().item(session)); - when(session.isClosed()).thenReturn(false); - when(session.execute(any(SimpleStatement.class))) - .thenThrow(new IllegalStateException("Invalid query")); - - var response = healthCheck.call(); - - assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); - verifySessionNotEvicted(); - } - - @Test - public void sessionAcquisitionTimeoutReportsDown() { - sessionRequestReturns(Uni.createFrom().nothing()); - - var response = healthCheck.call(); - - assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.DOWN); - verifySessionNotEvicted(); - verifyNoInteractions(session); - } - - @Test - public void configuredCredentialsAreUsedWhenFixedTokenIsNotSet() { - when(databaseConfig.fixedToken()).thenReturn(Optional.empty()); - healthCheck = - new CassandraConnectionHealthCheck(sessionCache, operationsConfig, Duration.ofMillis(100)); - - when(sessionCache.getSession( - eq(TEST_CONSTANTS.CASSANDRA_TENANT), any(String.class), eq(HEALTH_CHECK_USER_AGENT))) - .thenReturn(Uni.createFrom().item(session)); - when(session.isClosed()).thenReturn(true); - - healthCheck.call(); - - var expectedToken = - CqlCredentials.USERNAME_PASSWORD_TOKEN_PREFIX - + Base64.getEncoder().encodeToString(USER_NAME.getBytes(StandardCharsets.UTF_8)) - + ":" - + Base64.getEncoder().encodeToString(PASSWORD.getBytes(StandardCharsets.UTF_8)); - verify(sessionCache) - .getSession(TEST_CONSTANTS.CASSANDRA_TENANT, expectedToken, HEALTH_CHECK_USER_AGENT); - } - - @Test - public void nonCassandraDatabaseDoesNotRunConnectivityCheck() { - when(databaseConfig.type()).thenReturn(DatabaseType.ASTRA); - healthCheck = - new CassandraConnectionHealthCheck(sessionCache, operationsConfig, Duration.ofMillis(100)); - - var response = healthCheck.call(); - - assertThat(response.getStatus()).isEqualTo(HealthCheckResponse.Status.UP); - assertThat(response.getData()) - .hasValueSatisfying( - data -> - assertThat(data) - .containsEntry( - "reason", - "Cassandra connectivity check is not applicable for database type ASTRA")); - verifyNoInteractions(sessionCache); - } - - private void sessionRequestReturns(Uni sessionResult) { - when(sessionCache.getSession( - TEST_CONSTANTS.CASSANDRA_TENANT, TEST_CONSTANTS.AUTH_TOKEN, HEALTH_CHECK_USER_AGENT)) - .thenReturn(sessionResult); - } - - private void verifySessionEvicted() { - verify(sessionCache) - .evictSession( - TEST_CONSTANTS.CASSANDRA_TENANT, TEST_CONSTANTS.AUTH_TOKEN, HEALTH_CHECK_USER_AGENT); - } - - private void verifySessionNotEvicted() { - verify(sessionCache, never()) - .evictSession( - TEST_CONSTANTS.CASSANDRA_TENANT, TEST_CONSTANTS.AUTH_TOKEN, HEALTH_CHECK_USER_AGENT); - } -} diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheckTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheckTest.java new file mode 100644 index 0000000000..6893797446 --- /dev/null +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheckTest.java @@ -0,0 +1,147 @@ +package io.stargate.sgv2.jsonapi.api.health; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.same; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.DefaultConsistencyLevel; +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import io.smallrye.mutiny.TimeoutException; +import io.smallrye.mutiny.Uni; +import io.smallrye.mutiny.helpers.test.UniAssertSubscriber; +import io.stargate.sgv2.jsonapi.api.request.RequestContext; +import io.stargate.sgv2.jsonapi.api.request.UserAgent; +import io.stargate.sgv2.jsonapi.api.request.tenant.Tenant; +import io.stargate.sgv2.jsonapi.config.DatabaseType; +import io.stargate.sgv2.jsonapi.service.cqldriver.CQLSessionCache; +import java.time.Duration; +import java.util.concurrent.CompletableFuture; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +public class DatabaseReadinessCheckTest { + + private static final Duration TIMEOUT = Duration.ofMillis(100); + private static final String ASTRA_TENANT_ID = "60b5dccb-e91d-4a60-987b-7588cd8aa1e3"; + + private CQLSessionCache sessionCache; + private CqlSession session; + private RequestContext requestContext; + private DatabaseReadinessCheck readinessCheck; + + @BeforeEach + public void setup() { + sessionCache = mock(CQLSessionCache.class); + session = mock(CqlSession.class); + requestContext = + new RequestContext( + Tenant.create(DatabaseType.ASTRA, ASTRA_TENANT_ID, "us-west-2"), + "astra-token", + new UserAgent("Datastax-SLA-Checker")); + readinessCheck = new DatabaseReadinessCheck(sessionCache, TIMEOUT); + } + + @Test + public void successfulCheckUsesRequestContextAndAsyncDistributedQuery() { + var resultSet = mock(AsyncResultSet.class); + var resultFuture = new CompletableFuture(); + when(sessionCache.getSession(requestContext)).thenReturn(Uni.createFrom().item(session)); + when(session.executeAsync(any(SimpleStatement.class))).thenReturn(resultFuture); + + var subscriber = + readinessCheck + .check(requestContext) + .subscribe() + .withSubscriber(UniAssertSubscriber.create()); + + subscriber.assertSubscribed().assertNotTerminated(); + resultFuture.complete(resultSet); + subscriber.awaitItem().assertItem(null).assertCompleted(); + + verify(sessionCache).getSession(same(requestContext)); + var statementCaptor = ArgumentCaptor.forClass(SimpleStatement.class); + verify(session).executeAsync(statementCaptor.capture()); + assertThat(statementCaptor.getValue().getQuery()) + .isEqualTo("SELECT * FROM datastax_sla.check LIMIT 1"); + assertThat(statementCaptor.getValue().getTimeout()).isEqualTo(TIMEOUT); + assertThat(statementCaptor.getValue().getConsistencyLevel()) + .isEqualTo(DefaultConsistencyLevel.LOCAL_QUORUM); + assertThat(statementCaptor.getValue().getExecutionProfileName()).isEqualTo("table-read"); + verify(session, never()).execute(any(SimpleStatement.class)); + verify(session, never()).isClosed(); + verify(sessionCache, never()).evictSession(any(RequestContext.class)); + } + + @Test + public void sessionAcquisitionFailureIsPropagated() { + when(sessionCache.getSession(requestContext)) + .thenReturn(Uni.createFrom().failure(new IllegalStateException("Cannot connect"))); + + readinessCheck + .check(requestContext) + .subscribe() + .withSubscriber(UniAssertSubscriber.create()) + .awaitFailure() + .assertFailedWith(IllegalStateException.class, "Cannot connect"); + + verify(session, never()).executeAsync(any(SimpleStatement.class)); + verify(sessionCache, never()).evictSession(any(RequestContext.class)); + } + + @Test + public void asynchronousQueryFailureIsPropagated() { + var resultFuture = new CompletableFuture(); + when(sessionCache.getSession(requestContext)).thenReturn(Uni.createFrom().item(session)); + when(session.executeAsync(any(SimpleStatement.class))).thenReturn(resultFuture); + + var subscriber = + readinessCheck + .check(requestContext) + .subscribe() + .withSubscriber(UniAssertSubscriber.create()); + resultFuture.completeExceptionally(new IllegalStateException("Query failed")); + + subscriber.awaitFailure().assertFailedWith(IllegalStateException.class, "Query failed"); + verify(sessionCache, never()).evictSession(any(RequestContext.class)); + } + + @Test + public void reactiveTimeoutBoundsSessionAcquisition() { + readinessCheck = new DatabaseReadinessCheck(sessionCache, Duration.ofMillis(20)); + when(sessionCache.getSession(requestContext)).thenReturn(Uni.createFrom().nothing()); + + readinessCheck + .check(requestContext) + .subscribe() + .withSubscriber(UniAssertSubscriber.create()) + .awaitFailure() + .assertFailedWith(TimeoutException.class); + + verify(session, never()).executeAsync(any(SimpleStatement.class)); + verify(sessionCache, never()).evictSession(any(RequestContext.class)); + } + + @Test + public void reactiveTimeoutBoundsAsynchronousQuery() { + readinessCheck = new DatabaseReadinessCheck(sessionCache, Duration.ofMillis(20)); + when(sessionCache.getSession(requestContext)).thenReturn(Uni.createFrom().item(session)); + when(session.executeAsync(any(SimpleStatement.class))).thenReturn(new CompletableFuture<>()); + + readinessCheck + .check(requestContext) + .subscribe() + .withSubscriber(UniAssertSubscriber.create()) + .awaitFailure() + .assertFailedWith(TimeoutException.class); + + verify(session).executeAsync(any(SimpleStatement.class)); + verify(sessionCache, never()).evictSession(any(RequestContext.class)); + } +} diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResourceTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResourceTest.java new file mode 100644 index 0000000000..0fd4d9e8d3 --- /dev/null +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResourceTest.java @@ -0,0 +1,181 @@ +package io.stargate.sgv2.jsonapi.api.v1; + +import static io.restassured.RestAssured.given; +import static org.assertj.core.api.Assertions.assertThat; +import static org.hamcrest.Matchers.equalTo; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +import com.datastax.oss.driver.api.core.CqlSession; +import com.datastax.oss.driver.api.core.cql.AsyncResultSet; +import com.datastax.oss.driver.api.core.cql.SimpleStatement; +import io.quarkus.security.UnauthorizedException; +import io.quarkus.test.InjectMock; +import io.quarkus.test.junit.QuarkusTest; +import io.quarkus.test.junit.TestProfile; +import io.smallrye.mutiny.Uni; +import io.stargate.sgv2.jsonapi.api.request.RequestContext; +import io.stargate.sgv2.jsonapi.config.constants.HttpConstants; +import io.stargate.sgv2.jsonapi.exception.APISecurityException; +import io.stargate.sgv2.jsonapi.service.cqldriver.CQLSessionCache; +import io.stargate.sgv2.jsonapi.service.cqldriver.CqlSessionCacheSupplier; +import io.stargate.sgv2.jsonapi.testresource.NoGlobalResourcesTestProfile; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +@QuarkusTest +@TestProfile(DatabaseReadinessResourceTest.AstraProfile.class) +public class DatabaseReadinessResourceTest { + + private static final String TENANT_ID = "60b5dccb-e91d-4a60-987b-7588cd8aa1e3"; + private static final String REGION = "us-west-2"; + private static final String ASTRA_HOST = TENANT_ID + "-" + REGION + ".apps.astra.datastax.com"; + private static final String TOKEN = "astra-canary-token"; + private static final String SLA_USER_AGENT = "Datastax-SLA-Checker"; + + @InjectMock CqlSessionCacheSupplier sessionCacheSupplier; + + private CQLSessionCache sessionCache; + private CqlSession session; + + @BeforeEach + public void setup() { + sessionCache = mock(CQLSessionCache.class); + session = mock(CqlSession.class); + when(sessionCacheSupplier.get()).thenReturn(sessionCache); + } + + @Test + public void missingTokenIsRejectedBeforeDatabaseAccess() { + given() + .header("Host", ASTRA_HOST) + .header("User-Agent", SLA_USER_AGENT) + .when() + .get(DatabaseReadinessResource.BASE_PATH) + .then() + .statusCode(401); + + verifyNoInteractions(sessionCache); + } + + @Test + public void astraRequestUsesResolvedTenantTokenAndSlaUserAgent() { + var resultSet = mock(AsyncResultSet.class); + var capturedContext = new AtomicReference(); + when(sessionCache.getSession(any(RequestContext.class))) + .thenAnswer( + invocation -> { + RequestContext context = invocation.getArgument(0); + capturedContext.set( + new RequestContextSnapshot( + context.tenant().toString(), + context.tenant().region(), + context.authToken(), + context.userAgent().toString())); + return Uni.createFrom().item(session); + }); + when(session.executeAsync(any(SimpleStatement.class))) + .thenReturn(CompletableFuture.completedFuture(resultSet)); + + authenticatedRequest() + .when() + .get(DatabaseReadinessResource.BASE_PATH) + .then() + .statusCode(200) + .body("status", equalTo("UP")); + + assertThat(capturedContext.get()) + .isEqualTo(new RequestContextSnapshot(TENANT_ID, REGION, TOKEN, SLA_USER_AGENT)); + } + + @Test + public void databaseFailureReturnsServiceUnavailableWithoutDetails() { + when(sessionCache.getSession(any(RequestContext.class))) + .thenReturn(Uni.createFrom().failure(new IllegalStateException("sensitive failure"))); + + var response = + authenticatedRequest() + .when() + .get(DatabaseReadinessResource.BASE_PATH) + .then() + .statusCode(503) + .body("status", equalTo("DOWN")) + .extract() + .asString(); + + assertThat(response).doesNotContain("sensitive failure", TOKEN, TENANT_ID); + } + + @Test + public void invalidTokenFailureReturnsUnauthorizedWithoutDetails() { + when(sessionCache.getSession(any(RequestContext.class))) + .thenThrow(new UnauthorizedException("sensitive credential failure")); + + var response = + authenticatedRequest() + .when() + .get(DatabaseReadinessResource.BASE_PATH) + .then() + .statusCode(401) + .body("status", equalTo("DOWN")) + .extract() + .asString(); + + assertThat(response).doesNotContain("sensitive credential failure", TOKEN, TENANT_ID); + } + + @Test + public void databaseAuthenticationFailureReturnsUnauthorizedWithoutDetails() { + var authenticationFailure = + APISecurityException.Code.UNAUTHENTICATED_REQUEST.withPreformattedMessage( + "sensitive database authentication failure"); + when(sessionCache.getSession(any(RequestContext.class))) + .thenReturn(Uni.createFrom().failure(authenticationFailure)); + + var response = + authenticatedRequest() + .when() + .get(DatabaseReadinessResource.BASE_PATH) + .then() + .statusCode(401) + .body("status", equalTo("DOWN")) + .extract() + .asString(); + + assertThat(response) + .doesNotContain("sensitive database authentication failure", TOKEN, TENANT_ID); + } + + private io.restassured.specification.RequestSpecification authenticatedRequest() { + return given() + .header("Host", ASTRA_HOST) + .header(HttpConstants.AUTHENTICATION_TOKEN_HEADER_NAME, TOKEN) + .header("User-Agent", SLA_USER_AGENT); + } + + private record RequestContextSnapshot( + String tenantId, String region, String authToken, String userAgent) {} + + public static class AstraProfile implements NoGlobalResourcesTestProfile { + + @Override + public Map getConfigOverrides() { + return Map.of( + "stargate.jsonapi.operations.database-config.type", + "ASTRA", + "stargate.multi-tenancy.enabled", + "true", + "stargate.multi-tenancy.tenant-resolver.type", + "subdomain", + "stargate.multi-tenancy.tenant-resolver.subdomain.max-chars", + "36", + "stargate.jsonapi.operations.sla-user-agent", + SLA_USER_AGENT); + } + } +} diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java index 28ec1c2907..5afb578c28 100644 --- a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/SessionEvictionIntegrationTest.java @@ -16,6 +16,7 @@ import java.io.IOException; import java.net.ServerSocket; import java.util.Collections; +import java.util.HashMap; import java.util.Map; import org.junit.jupiter.api.Test; import org.slf4j.Logger; @@ -30,8 +31,7 @@ public class SessionEvictionIntegrationTest extends AbstractCollectionIntegratio private static final Logger LOGGER = LoggerFactory.getLogger(SessionEvictionIntegrationTest.class); - private static final String READINESS_PATH = "/stargate/health/ready"; - private static final String LIVENESS_PATH = "/stargate/health/live"; + private static final String READINESS_USER_AGENT = "DataAPI-Readiness-Test/1.0"; /** * Overridden to ensure we connect to the isolated container created for this test. @@ -52,6 +52,13 @@ protected int getCassandraCqlPort() { @Test public void testSessionEvictionOnAllNodesFailed() { + if (!executeCqlStatement( + "CREATE KEYSPACE IF NOT EXISTS datastax_sla " + + "WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 1}", + "CREATE TABLE IF NOT EXISTS datastax_sla.check (id text PRIMARY KEY)")) { + throw new AssertionError("Failed to provision the distributed readiness table"); + } + // 1. Insert and find initial data to ensure the database is healthy before the test insertDoc( """ @@ -91,7 +98,6 @@ public void testSessionEvictionOnAllNodesFailed() { "errors[0].errorCode", is(DatabaseException.Code.FAILED_TO_CONNECT_TO_DATABASE.name())); waitForReadinessStatus("DOWN", 60_000); - given().when().get(LIVENESS_PATH).then().statusCode(200).body("status", is("UP")); // 4. Restart the container to simulate recovery getDockerClient().startContainerCmd(getContainerId()).exec(); @@ -233,10 +239,7 @@ private boolean isApiReady() { } } - /** - * Polls the readiness endpoint until both the overall health and Cassandra connectivity check - * have the expected status. - */ + /** Polls the authenticated database readiness endpoint until it has the expected status. */ private void waitForReadinessStatus(String expectedStatus, long timeoutMillis) { var start = System.currentTimeMillis(); var expectedStatusCode = "UP".equals(expectedStatus) ? 200 : 503; @@ -244,15 +247,17 @@ private void waitForReadinessStatus(String expectedStatus, long timeoutMillis) { while (System.currentTimeMillis() - start < timeoutMillis) { try { - lastResponse = given().when().get(READINESS_PATH); + lastResponse = + given() + .headers(getHeaders()) + .header("User-Agent", READINESS_USER_AGENT) + .when() + .get(DatabaseReadinessResource.BASE_PATH); var jsonPath = lastResponse.jsonPath(); var overallStatus = jsonPath.getString("status"); - var cassandraStatus = - jsonPath.getString("checks.find { it.name == 'cassandra-connection' }.status"); if (lastResponse.statusCode() == expectedStatusCode - && expectedStatus.equals(overallStatus) - && expectedStatus.equals(cassandraStatus)) { + && expectedStatus.equals(overallStatus)) { return; } } catch (Exception e) { @@ -372,9 +377,10 @@ protected GenericContainer baseCassandraContainer(boolean reuse) { */ @Override public Map start() { - var props = super.start(); + var props = new HashMap<>(super.start()); + props.put("stargate.jsonapi.operations.sla-user-agent", READINESS_USER_AGENT); sessionEvictionCassandraContainer = super.getCassandraContainer(); - return props; + return Map.copyOf(props); } /** diff --git a/src/test/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilterTest.java b/src/test/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilterTest.java new file mode 100644 index 0000000000..5b9393c522 --- /dev/null +++ b/src/test/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilterTest.java @@ -0,0 +1,40 @@ +package io.stargate.sgv2.jsonapi.metrics; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +import io.micrometer.core.instrument.MeterRegistry; +import io.stargate.sgv2.jsonapi.api.request.RequestContext; +import io.stargate.sgv2.jsonapi.api.v1.DatabaseReadinessResource; +import io.stargate.sgv2.jsonapi.api.v1.metrics.MetricsConfig; +import jakarta.ws.rs.container.ContainerRequestContext; +import jakarta.ws.rs.container.ContainerResponseContext; +import jakarta.ws.rs.core.UriInfo; +import java.net.URI; +import org.junit.jupiter.api.Test; + +public class TenantRequestMetricsFilterTest { + + @Test + public void databaseReadinessIsNotCountedAsTenantTraffic() { + var meterRegistry = mock(MeterRegistry.class); + var dataApiRequestContext = mock(RequestContext.class); + var metricsConfig = mock(MetricsConfig.class); + var tenantRequestConfig = mock(MetricsConfig.TenantRequestCounterConfig.class); + when(metricsConfig.tenantRequestCounter()).thenReturn(tenantRequestConfig); + when(tenantRequestConfig.enabled()).thenReturn(true); + + var requestContext = mock(ContainerRequestContext.class); + var uriInfo = mock(UriInfo.class); + when(requestContext.getUriInfo()).thenReturn(uriInfo); + when(uriInfo.getRequestUri()) + .thenReturn(URI.create("http://localhost" + DatabaseReadinessResource.BASE_PATH)); + + var filter = + new TenantRequestMetricsFilter(meterRegistry, dataApiRequestContext, metricsConfig); + filter.record(requestContext, mock(ContainerResponseContext.class)); + + verifyNoInteractions(meterRegistry, dataApiRequestContext); + } +} From 3e7ff74751936af7268510fc4e8cc5e8aa59545d Mon Sep 17 00:00:00 2001 From: Eric Hare Date: Mon, 3 Aug 2026 20:50:06 -0700 Subject: [PATCH 5/5] fix: address review comments --- CONFIGURATION.md | 20 ++++-- .../api/health/DatabaseReadinessCheck.java | 1 + .../api/v1/DatabaseReadinessResource.java | 66 +++++++++++++---- .../metrics/TenantRequestMetricsFilter.java | 5 +- .../cqldriver/CqlSessionCacheSupplier.java | 12 +++- src/main/resources/application.yaml | 2 +- .../api/v1/DatabaseReadinessResourceTest.java | 70 +++++++++++++++++-- .../TenantRequestMetricsFilterTest.java | 11 +-- .../CqlSessionCacheSupplierTests.java | 34 +++++---- 9 files changed, 174 insertions(+), 47 deletions(-) diff --git a/CONFIGURATION.md b/CONFIGURATION.md index 70e0b43147..76b507e757 100644 --- a/CONFIGURATION.md +++ b/CONFIGURATION.md @@ -71,18 +71,26 @@ provision a `datastax_sla.check` table that the canary principal can read. Its r must be appropriate for the deployment (greater than one in a multi-node local data center) so `LOCAL_QUORUM` requires responses from multiple replicas. Astra callers must use the canary database hostname so the tenant and region are resolved from `Host`; Cassandra ignores the tenant portion of -`Host`. The caller must also send the exact User-Agent configured by -`stargate.jsonapi.operations.sla-user-agent`, allowing a dedicated canary session to use the shorter -SLA session TTL instead of being treated like normal client traffic. +`Host`. The caller must also send the full User-Agent configured by +`stargate.jsonapi.operations.sla-user-agent`. The comparison is case-insensitive. Requests with a +missing or different User-Agent are rejected before accessing the session cache, and the endpoint +fails closed when the SLA User-Agent is not configured. This ensures the canary session uses the +shorter SLA session TTL instead of being treated like normal client traffic. Do not reuse the canary +credentials for normal traffic, because using the same cached session with a non-SLA User-Agent can +extend its lifetime. The check is fully asynchronous and has a five-second timeout. It returns HTTP 200 with -`{"status":"UP"}` after a successful read, HTTP 503 with `{"status":"DOWN"}` after a database -failure or timeout, and HTTP 401 when the `Token` header is missing or authentication fails. +`{"status":"UP"}` after a successful read and HTTP 503 with `{"status":"DOWN"}` after a database +failure, timeout, or missing SLA User-Agent configuration. It returns the standard Data API error +response with HTTP 401 when the `Token` header is missing or authentication fails, and HTTP 403 when +the request User-Agent does not match the configured SLA User-Agent. Probe integrations must use the +HTTP status as the readiness contract rather than parsing the response body's `status` field alone. Kubernetes or an SLA checker must call each pod directly for this endpoint to control per-pod readiness. An external request sent through a load balancer does not establish which pod is ready. Restrict the endpoint to trusted probe traffic with deployment controls such as a NetworkPolicy, -mTLS, or an ingress ACL and rate limit. Kubernetes `httpGet` headers cannot reference a Secret, so +mTLS, or an ingress ACL and rate limit. The User-Agent check is an operational guard, not an +authentication boundary. Kubernetes `httpGet` headers cannot reference a Secret, so delivery of the canary token is intentionally outside the Data API configuration. Prefer an external checker or a Secret-mounted file read by an `exec` probe; do not put the token literally in the probe command or shell trace. The unauthenticated Quarkus health endpoints under the diff --git a/src/main/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheck.java b/src/main/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheck.java index bbec7687d1..6f4c3e0d50 100644 --- a/src/main/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheck.java +++ b/src/main/java/io/stargate/sgv2/jsonapi/api/health/DatabaseReadinessCheck.java @@ -65,6 +65,7 @@ public Uni check(RequestContext requestContext) { .executeRead(statement) .replaceWithVoid(); }) + // The statement timeout bounds driver I/O; this also bounds asynchronous session lookup. .ifNoItem() .after(timeout) .fail(); diff --git a/src/main/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResource.java b/src/main/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResource.java index 38bafaa95f..7e39b952c2 100644 --- a/src/main/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResource.java +++ b/src/main/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResource.java @@ -3,6 +3,8 @@ import io.quarkus.security.UnauthorizedException; import io.smallrye.mutiny.Uni; import io.stargate.sgv2.jsonapi.api.health.DatabaseReadinessCheck; +import io.stargate.sgv2.jsonapi.api.model.command.CommandResult; +import io.stargate.sgv2.jsonapi.api.model.command.tracing.RequestTracing; import io.stargate.sgv2.jsonapi.api.request.RequestContext; import io.stargate.sgv2.jsonapi.config.constants.OpenApiConstants; import io.stargate.sgv2.jsonapi.exception.APISecurityException; @@ -13,6 +15,8 @@ import jakarta.ws.rs.Produces; import jakarta.ws.rs.core.MediaType; import jakarta.ws.rs.core.Response; +import java.util.Collections; +import java.util.IdentityHashMap; import org.eclipse.microprofile.openapi.annotations.Operation; import org.eclipse.microprofile.openapi.annotations.media.Content; import org.eclipse.microprofile.openapi.annotations.media.Schema; @@ -25,9 +29,10 @@ * Authenticated database readiness endpoint registered through Quarkus JAX-RS resource discovery. * *

{@code GET /v1/health/ready} runs the same request-scoped probe for Astra and Cassandra. A - * successful probe returns HTTP 200; invalid credentials return HTTP 401; and a database failure or - * timeout returns HTTP 503. The existing {@code /v1/*} security policy rejects requests without a - * token before this resource is called. + * request must use the configured SLA User-Agent. A successful probe returns HTTP 200; invalid + * credentials return HTTP 401; a missing or different SLA User-Agent returns HTTP 403; and a + * database failure, timeout, or missing SLA configuration returns HTTP 503. The existing {@code + * /v1/*} security policy rejects requests without a token before this resource is called. */ @Path(DatabaseReadinessResource.BASE_PATH) @Produces(MediaType.APPLICATION_JSON) @@ -41,12 +46,14 @@ public class DatabaseReadinessResource { private final DatabaseReadinessCheck readinessCheck; private final RequestContext requestContext; + private final CqlSessionCacheSupplier sessionCacheSupplier; @Inject public DatabaseReadinessResource( CqlSessionCacheSupplier sessionCacheSupplier, RequestContext requestContext) { this.readinessCheck = new DatabaseReadinessCheck(sessionCacheSupplier); this.requestContext = requestContext; + this.sessionCacheSupplier = sessionCacheSupplier; } @GET @@ -62,33 +69,46 @@ public DatabaseReadinessResource( @Content( mediaType = MediaType.APPLICATION_JSON, schema = @Schema(implementation = ReadinessResponse.class))), - @APIResponse(responseCode = "401", description = "The token is missing or invalid."), + @APIResponse( + responseCode = "401", + description = "The token is missing or invalid.", + content = + @Content( + mediaType = MediaType.APPLICATION_JSON, + schema = @Schema(implementation = CommandResult.class))), + @APIResponse( + responseCode = "403", + description = "The request does not use the configured SLA User-Agent."), @APIResponse( responseCode = "503", - description = "The database read failed or timed out.", + description = "The SLA User-Agent is not configured, or the database check failed.", content = @Content( mediaType = MediaType.APPLICATION_JSON, schema = @Schema(implementation = ReadinessResponse.class))) }) - public Uni> ready() { + public Uni> ready() { + var configuredSlaUserAgent = sessionCacheSupplier.slaUserAgent(); + if (configuredSlaUserAgent.isEmpty()) { + return Uni.createFrom().item(response(Response.Status.SERVICE_UNAVAILABLE, DOWN)); + } + if (!configuredSlaUserAgent.get().equals(requestContext.userAgent())) { + return Uni.createFrom().item(response(Response.Status.FORBIDDEN)); + } + return readinessCheck .check(requestContext) - .map(ignored -> RestResponse.ok(UP)) + .map(ignored -> response(Response.Status.OK, UP)) .onFailure(DatabaseReadinessResource::isUnauthorized) - .recoverWithItem( - failure -> - RestResponse.ResponseBuilder.create(Response.Status.UNAUTHORIZED, DOWN).build()) + .recoverWithItem(failure -> unauthorizedResponse()) .onFailure() - .recoverWithItem( - failure -> - RestResponse.ResponseBuilder.create(Response.Status.SERVICE_UNAVAILABLE, DOWN) - .build()); + .recoverWithItem(failure -> response(Response.Status.SERVICE_UNAVAILABLE, DOWN)); } private static boolean isUnauthorized(Throwable failure) { var current = failure; - while (current != null) { + var seen = Collections.newSetFromMap(new IdentityHashMap()); + while (current != null && seen.add(current)) { if (current instanceof UnauthorizedException || current instanceof APISecurityException apiException && apiException.httpStatus == Response.Status.UNAUTHORIZED.getStatusCode()) { @@ -99,5 +119,21 @@ private static boolean isUnauthorized(Throwable failure) { return false; } + private static RestResponse unauthorizedResponse() { + var commandResult = + CommandResult.statusOnlyBuilder(RequestTracing.NO_OP) + .addThrowable(APISecurityException.Code.UNAUTHENTICATED_REQUEST.get()) + .build(); + return response(Response.Status.UNAUTHORIZED, commandResult); + } + + private static RestResponse response(Response.Status status) { + return RestResponse.ResponseBuilder.create(status).build(); + } + + private static RestResponse response(Response.Status status, Object entity) { + return RestResponse.ResponseBuilder.create(status, entity).build(); + } + public record ReadinessResponse(String status) {} } diff --git a/src/main/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilter.java b/src/main/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilter.java index c633208225..b277403ff3 100644 --- a/src/main/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilter.java +++ b/src/main/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilter.java @@ -100,8 +100,9 @@ public void record( } private static boolean isDatabaseReadinessRequest(ContainerRequestContext requestContext) { - return DatabaseReadinessResource.BASE_PATH.equals( - requestContext.getUriInfo().getRequestUri().getPath()); + var requestPath = requestContext.getUriInfo().getRequestUri().getPath(); + return DatabaseReadinessResource.BASE_PATH.equals(requestPath) + || (DatabaseReadinessResource.BASE_PATH + "/").equals(requestPath); } private String getUserAgentValue(ContainerRequestContext requestContext) { diff --git a/src/main/java/io/stargate/sgv2/jsonapi/service/cqldriver/CqlSessionCacheSupplier.java b/src/main/java/io/stargate/sgv2/jsonapi/service/cqldriver/CqlSessionCacheSupplier.java index 766c36899f..7e569d60e4 100644 --- a/src/main/java/io/stargate/sgv2/jsonapi/service/cqldriver/CqlSessionCacheSupplier.java +++ b/src/main/java/io/stargate/sgv2/jsonapi/service/cqldriver/CqlSessionCacheSupplier.java @@ -10,6 +10,7 @@ import java.time.Duration; import java.util.List; import java.util.Objects; +import java.util.Optional; import java.util.function.Supplier; import org.eclipse.microprofile.config.inject.ConfigProperty; @@ -23,6 +24,7 @@ public class CqlSessionCacheSupplier implements Supplier { private final CQLSessionCache singleton; + private final Optional slaUserAgent; @Inject public CqlSessionCacheSupplier( @@ -53,11 +55,14 @@ public CqlSessionCacheSupplier( dbConfig.cassandraPort(), () -> schemaObjectCacheSupplier.get().getSchemaChangeListener()); + slaUserAgent = + operationsConfig.slaUserAgent().filter(value -> !value.isBlank()).map(UserAgent::new); + singleton = new CQLSessionCache( dbConfig.sessionCacheMaxSize(), Duration.ofSeconds(dbConfig.sessionCacheTtlSeconds()), - operationsConfig.slaUserAgent().map(UserAgent::new).orElse(null), + slaUserAgent.orElse(null), Duration.ofSeconds(dbConfig.slaSessionCacheTtlSeconds()), credentialsFactory, sessionFactory, @@ -70,4 +75,9 @@ public CqlSessionCacheSupplier( public CQLSessionCache get() { return singleton; } + + /** Gets the configured User-Agent that selects the shorter SLA session-cache TTL. */ + public Optional slaUserAgent() { + return slaUserAgent; + } } diff --git a/src/main/resources/application.yaml b/src/main/resources/application.yaml index 1103fef1b3..a978343bbc 100644 --- a/src/main/resources/application.yaml +++ b/src/main/resources/application.yaml @@ -159,7 +159,7 @@ quarkus: http-server: # ignore all non-application uris, as well as the custom set suppress-non-application-uris: true - ignore-patterns: /,/metrics,/swagger-ui.*,.*\.html,/v1/health/ready + ignore-patterns: /,/metrics,/swagger-ui.*,.*\.html,/v1/health/ready/? # due to the https://github.com/quarkusio/quarkus/issues/24938 # we need to define uri templating on our own for now diff --git a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResourceTest.java b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResourceTest.java index 0fd4d9e8d3..c3db7999f8 100644 --- a/src/test/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResourceTest.java +++ b/src/test/java/io/stargate/sgv2/jsonapi/api/v1/DatabaseReadinessResourceTest.java @@ -17,16 +17,20 @@ import io.quarkus.test.junit.TestProfile; import io.smallrye.mutiny.Uni; import io.stargate.sgv2.jsonapi.api.request.RequestContext; +import io.stargate.sgv2.jsonapi.api.request.UserAgent; import io.stargate.sgv2.jsonapi.config.constants.HttpConstants; import io.stargate.sgv2.jsonapi.exception.APISecurityException; import io.stargate.sgv2.jsonapi.service.cqldriver.CQLSessionCache; import io.stargate.sgv2.jsonapi.service.cqldriver.CqlSessionCacheSupplier; import io.stargate.sgv2.jsonapi.testresource.NoGlobalResourcesTestProfile; import java.util.Map; +import java.util.Optional; import java.util.concurrent.CompletableFuture; import java.util.concurrent.atomic.AtomicReference; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; @QuarkusTest @TestProfile(DatabaseReadinessResourceTest.AstraProfile.class) @@ -48,6 +52,8 @@ public void setup() { sessionCache = mock(CQLSessionCache.class); session = mock(CqlSession.class); when(sessionCacheSupplier.get()).thenReturn(sessionCache); + when(sessionCacheSupplier.slaUserAgent()) + .thenReturn(Optional.of(new UserAgent(SLA_USER_AGENT))); } @Test @@ -58,13 +64,15 @@ public void missingTokenIsRejectedBeforeDatabaseAccess() { .when() .get(DatabaseReadinessResource.BASE_PATH) .then() - .statusCode(401); + .statusCode(401) + .body("errors[0].errorCode", equalTo("MISSING_AUTHENTICATION_TOKEN")); verifyNoInteractions(sessionCache); } - @Test - public void astraRequestUsesResolvedTenantTokenAndSlaUserAgent() { + @ParameterizedTest + @ValueSource(strings = {"", "/"}) + public void astraRequestUsesResolvedTenantTokenAndSlaUserAgent(String pathSuffix) { var resultSet = mock(AsyncResultSet.class); var capturedContext = new AtomicReference(); when(sessionCache.getSession(any(RequestContext.class))) @@ -84,7 +92,7 @@ public void astraRequestUsesResolvedTenantTokenAndSlaUserAgent() { authenticatedRequest() .when() - .get(DatabaseReadinessResource.BASE_PATH) + .get(DatabaseReadinessResource.BASE_PATH + pathSuffix) .then() .statusCode(200) .body("status", equalTo("UP")); @@ -93,6 +101,34 @@ public void astraRequestUsesResolvedTenantTokenAndSlaUserAgent() { .isEqualTo(new RequestContextSnapshot(TENANT_ID, REGION, TOKEN, SLA_USER_AGENT)); } + @Test + public void wrongSlaUserAgentIsRejectedBeforeDatabaseAccess() { + given() + .header("Host", ASTRA_HOST) + .header(HttpConstants.AUTHENTICATION_TOKEN_HEADER_NAME, TOKEN) + .header("User-Agent", "ordinary-client") + .when() + .get(DatabaseReadinessResource.BASE_PATH) + .then() + .statusCode(403); + + verifyNoInteractions(sessionCache); + } + + @Test + public void missingSlaUserAgentConfigurationReturnsServiceUnavailable() { + when(sessionCacheSupplier.slaUserAgent()).thenReturn(Optional.empty()); + + authenticatedRequest() + .when() + .get(DatabaseReadinessResource.BASE_PATH) + .then() + .statusCode(503) + .body("status", equalTo("DOWN")); + + verifyNoInteractions(sessionCache); + } + @Test public void databaseFailureReturnsServiceUnavailableWithoutDetails() { when(sessionCache.getSession(any(RequestContext.class))) @@ -122,7 +158,7 @@ public void invalidTokenFailureReturnsUnauthorizedWithoutDetails() { .get(DatabaseReadinessResource.BASE_PATH) .then() .statusCode(401) - .body("status", equalTo("DOWN")) + .body("errors[0].errorCode", equalTo("UNAUTHENTICATED_REQUEST")) .extract() .asString(); @@ -143,7 +179,7 @@ public void databaseAuthenticationFailureReturnsUnauthorizedWithoutDetails() { .get(DatabaseReadinessResource.BASE_PATH) .then() .statusCode(401) - .body("status", equalTo("DOWN")) + .body("errors[0].errorCode", equalTo("UNAUTHENTICATED_REQUEST")) .extract() .asString(); @@ -151,6 +187,28 @@ public void databaseAuthenticationFailureReturnsUnauthorizedWithoutDetails() { .doesNotContain("sensitive database authentication failure", TOKEN, TENANT_ID); } + @Test + public void cyclicFailureCauseReturnsServiceUnavailable() { + var firstFailure = new IllegalStateException("first sensitive failure"); + var secondFailure = new IllegalArgumentException("second sensitive failure"); + firstFailure.initCause(secondFailure); + secondFailure.initCause(firstFailure); + when(sessionCache.getSession(any(RequestContext.class))) + .thenReturn(Uni.createFrom().failure(firstFailure)); + + var response = + authenticatedRequest() + .when() + .get(DatabaseReadinessResource.BASE_PATH) + .then() + .statusCode(503) + .body("status", equalTo("DOWN")) + .extract() + .asString(); + + assertThat(response).doesNotContain("first sensitive failure", "second sensitive failure"); + } + private io.restassured.specification.RequestSpecification authenticatedRequest() { return given() .header("Host", ASTRA_HOST) diff --git a/src/test/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilterTest.java b/src/test/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilterTest.java index 5b9393c522..648e2bcaf2 100644 --- a/src/test/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilterTest.java +++ b/src/test/java/io/stargate/sgv2/jsonapi/metrics/TenantRequestMetricsFilterTest.java @@ -12,12 +12,14 @@ import jakarta.ws.rs.container.ContainerResponseContext; import jakarta.ws.rs.core.UriInfo; import java.net.URI; -import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; public class TenantRequestMetricsFilterTest { - @Test - public void databaseReadinessIsNotCountedAsTenantTraffic() { + @ParameterizedTest + @ValueSource(strings = {"", "/"}) + public void databaseReadinessIsNotCountedAsTenantTraffic(String pathSuffix) { var meterRegistry = mock(MeterRegistry.class); var dataApiRequestContext = mock(RequestContext.class); var metricsConfig = mock(MetricsConfig.class); @@ -29,7 +31,8 @@ public void databaseReadinessIsNotCountedAsTenantTraffic() { var uriInfo = mock(UriInfo.class); when(requestContext.getUriInfo()).thenReturn(uriInfo); when(uriInfo.getRequestUri()) - .thenReturn(URI.create("http://localhost" + DatabaseReadinessResource.BASE_PATH)); + .thenReturn( + URI.create("http://localhost" + DatabaseReadinessResource.BASE_PATH + pathSuffix)); var filter = new TenantRequestMetricsFilter(meterRegistry, dataApiRequestContext, metricsConfig); diff --git a/src/test/java/io/stargate/sgv2/jsonapi/service/cqldriver/CqlSessionCacheSupplierTests.java b/src/test/java/io/stargate/sgv2/jsonapi/service/cqldriver/CqlSessionCacheSupplierTests.java index de0cfca852..976344d441 100644 --- a/src/test/java/io/stargate/sgv2/jsonapi/service/cqldriver/CqlSessionCacheSupplierTests.java +++ b/src/test/java/io/stargate/sgv2/jsonapi/service/cqldriver/CqlSessionCacheSupplierTests.java @@ -7,6 +7,7 @@ import com.datastax.oss.driver.api.core.metadata.schema.SchemaChangeListener; import io.micrometer.core.instrument.simple.SimpleMeterRegistry; import io.stargate.sgv2.jsonapi.TestConstants; +import io.stargate.sgv2.jsonapi.api.request.UserAgent; import io.stargate.sgv2.jsonapi.config.DatabaseType; import io.stargate.sgv2.jsonapi.config.OperationsConfig; import io.stargate.sgv2.jsonapi.service.schema.SchemaObjectCache; @@ -24,6 +25,24 @@ public class CqlSessionCacheSupplierTests { public void testSingleton() { // Not a lot to test, just checking it always returns the same instance. + var factory = newFactory(Optional.of(TEST_CONSTANTS.SLA_USER_AGENT_NAME)); + + var sessionCache1 = factory.get(); + var sessionCache2 = factory.get(); + + assertThat(sessionCache1) + .as("Session cache should be the same instance") + .isSameAs(sessionCache2); + assertThat(factory.slaUserAgent()).contains(new UserAgent(TEST_CONSTANTS.SLA_USER_AGENT_NAME)); + } + + @Test + public void blankSlaUserAgentIsTreatedAsUnconfigured() { + assertThat(newFactory(Optional.of(" ")).slaUserAgent()).isEmpty(); + } + + private CqlSessionCacheSupplier newFactory(Optional slaUserAgent) { + var dbConfig = mock(OperationsConfig.DatabaseConfig.class); when(dbConfig.type()).thenReturn(DatabaseType.ASTRA); when(dbConfig.localDatacenter()).thenReturn("datacenter1"); @@ -35,8 +54,7 @@ public void testSingleton() { var operationsConfig = mock(OperationsConfig.class); when(operationsConfig.databaseConfig()).thenReturn(dbConfig); - when(operationsConfig.slaUserAgent()) - .thenReturn(Optional.of(TEST_CONSTANTS.SLA_USER_AGENT_NAME)); + when(operationsConfig.slaUserAgent()).thenReturn(slaUserAgent); var mockSchemaObjectCacheSupplier = mock(SchemaObjectCacheSupplier.class); var mockSchemaObjectCache = mock(SchemaObjectCache.class); @@ -44,15 +62,7 @@ public void testSingleton() { when(mockSchemaObjectCache.getSchemaChangeListener()) .thenReturn(mock(SchemaChangeListener.class)); - var factory = - new CqlSessionCacheSupplier( - "testApp", operationsConfig, new SimpleMeterRegistry(), mockSchemaObjectCacheSupplier); - - var sessionCache1 = factory.get(); - var sessionCache2 = factory.get(); - - assertThat(sessionCache1) - .as("Session cache should be the same instance") - .isSameAs(sessionCache2); + return new CqlSessionCacheSupplier( + "testApp", operationsConfig, new SimpleMeterRegistry(), mockSchemaObjectCacheSupplier); } }