From 62f80742699783f2cbc07ee560bb13851f05c010 Mon Sep 17 00:00:00 2001 From: Devesh Bhardwaj Date: Tue, 8 Sep 2026 11:34:27 +0530 Subject: [PATCH 1/3] Use the generated RetryInterceptor instead of the hand-written one The generated com.skyflow.generated.rest.core.RetryInterceptor now supports everything the hand-written SkyflowRetryInterceptor existed to work around: a configurable initial/max retry delay via its 4-arg constructor, and a per-call (not per-instance) backoff counter, so a shared OkHttpClient no longer exhausts its retry budget once for its whole lifetime. It also picks up Retry-After / X-RateLimit-Reset header handling for free. - VaultClient now constructs the generated RetryInterceptor, passing the resolved maxRetries/initialRetryDelayMillis/maxRetryDelayMillis. Jitter is left at the generated interceptor's own default (0.2) - not yet exposed as a VaultConfig/builder setting. - The generated interceptor doesn't validate maxRetries itself (a negative value would silently behave as zero retries), so VaultClient now guards that explicitly to preserve the existing "negative maxRetries -> SkyflowException" contract. - Deleted SkyflowRetryInterceptor and its unit tests; that logic now lives only in the generated class (excluded from coverage/javadoc like the rest of com.skyflow.generated.*). - Updated HttpConfigTests/AuthInterceptorTests to the new type. The generated interceptor exposes no getters, so tests that need to read back what it was constructed with do it via reflection rather than adding hand-written accessors to generated code. All 710 flowvault tests pass. Co-Authored-By: Claude Sonnet 5 --- .../main/java/com/skyflow/VaultClient.java | 20 +- .../generated/rest/core/RetryInterceptor.java | 195 ++++++++++++-- .../utils/SkyflowRetryInterceptor.java | 97 ------- .../com/skyflow/AuthInterceptorTests.java | 3 +- .../java/com/skyflow/HttpConfigTests.java | 84 +++--- .../utils/SkyflowRetryInterceptorTests.java | 254 ------------------ 6 files changed, 248 insertions(+), 405 deletions(-) delete mode 100644 flowvault/src/main/java/com/skyflow/utils/SkyflowRetryInterceptor.java delete mode 100644 flowvault/src/test/java/com/skyflow/utils/SkyflowRetryInterceptorTests.java diff --git a/flowvault/src/main/java/com/skyflow/VaultClient.java b/flowvault/src/main/java/com/skyflow/VaultClient.java index a047e156..15543aec 100644 --- a/flowvault/src/main/java/com/skyflow/VaultClient.java +++ b/flowvault/src/main/java/com/skyflow/VaultClient.java @@ -5,11 +5,12 @@ import com.skyflow.errors.SkyflowException; import com.skyflow.generated.rest.ApiClient; import com.skyflow.generated.rest.ApiClientBuilder; +import com.skyflow.generated.rest.core.RetryInterceptor; import com.skyflow.generated.rest.resources.flowservice.FlowserviceClient; import com.skyflow.generated.rest.resources.records.RecordsClient; -import com.skyflow.utils.SkyflowRetryInterceptor; import com.skyflow.utils.Utils; +import java.util.Optional; import java.util.concurrent.TimeUnit; import okhttp3.ConnectionPool; @@ -159,18 +160,25 @@ protected void updateExecutorInHTTP() throws SkyflowException { Integer writeTimeout = resolveNullableInt(vaultConfig.getWriteTimeout(), commonWriteTimeout); // Negative timeout/retry values reach here straight from public config setters with - // no validation of their own; our own SkyflowRetryInterceptor throws IllegalArgumentException - // and OkHttp's own Builder throws IllegalStateException for those — translate both (and - // anything else unexpected from this construction) to SkyflowException so every failure - // mode from this SDK is a SkyflowException, never a raw one. + // no validation of their own; the generated RetryInterceptor validates initial/max delay + // but not maxRetries itself (a negative value would just behave as zero retries), and + // OkHttp's own Builder throws IllegalStateException for negative timeouts — translate + // all of these (and anything else unexpected from this construction) to SkyflowException + // so every failure mode from this SDK is a SkyflowException, never a raw one. try { + if (maxRetries < 0) { + throw new IllegalArgumentException("maxRetries must be non-negative"); + } OkHttpClient.Builder httpBuilder = new OkHttpClient.Builder() .connectionPool(new ConnectionPool(10, 1, TimeUnit.MINUTES)) // Overall ceiling; bounds the whole call including retries. .callTimeout(timeoutSeconds, TimeUnit.SECONDS) // OUTER: retries. Must wrap the auth interceptor so each attempt re-reads the // (possibly refreshed) bearer token rather than replaying a stale one. - .addInterceptor(new SkyflowRetryInterceptor(maxRetries, initialRetryDelayMillis, maxRetryDelayMillis)) + // Jitter is left at the generated interceptor's own default (0.2) - not yet + // exposed as a VaultConfig/builder setting. + .addInterceptor(new RetryInterceptor(maxRetries, Optional.of(initialRetryDelayMillis), + Optional.of(maxRetryDelayMillis), Optional.empty())) .addInterceptor(chain -> { // INNER: auth Request requestWithAuth = chain.request().newBuilder() .header("Authorization", "Bearer " + this.token) diff --git a/flowvault/src/main/java/com/skyflow/generated/rest/core/RetryInterceptor.java b/flowvault/src/main/java/com/skyflow/generated/rest/core/RetryInterceptor.java index 7a28c3c9..7d89751d 100644 --- a/flowvault/src/main/java/com/skyflow/generated/rest/core/RetryInterceptor.java +++ b/flowvault/src/main/java/com/skyflow/generated/rest/core/RetryInterceptor.java @@ -3,37 +3,80 @@ */ package com.skyflow.generated.rest.core; -import okhttp3.Interceptor; -import okhttp3.Response; - import java.io.IOException; import java.time.Duration; +import java.time.ZonedDateTime; +import java.time.format.DateTimeFormatter; +import java.time.format.DateTimeParseException; import java.util.Optional; import java.util.Random; +import okhttp3.Interceptor; +import okhttp3.Request; +import okhttp3.Response; public class RetryInterceptor implements Interceptor { - private static final Duration ONE_SECOND = Duration.ofSeconds(1); - private final ExponentialBackoff backoff; + private static final Duration DEFAULT_INITIAL_RETRY_DELAY = Duration.ofMillis(1000); + private static final Duration DEFAULT_MAX_RETRY_DELAY = Duration.ofMillis(60000); + private static final double DEFAULT_JITTER_FACTOR = 0.2; + + private final int maxRetries; + private final Duration initialRetryDelay; + private final Duration maxRetryDelay; + private final double jitterFactor; private final Random random = new Random(); public RetryInterceptor(int maxRetries) { - this.backoff = new ExponentialBackoff(maxRetries); + this(maxRetries, Optional.empty(), Optional.empty(), Optional.empty()); + } + + public RetryInterceptor( + int maxRetries, + Optional initialRetryDelayMillis, + Optional maxRetryDelayMillis, + Optional jitterFactor) { + initialRetryDelayMillis.ifPresent(delay -> { + if (delay < 0) { + throw new IllegalArgumentException("initialRetryDelayMillis must be non-negative"); + } + }); + maxRetryDelayMillis.ifPresent(delay -> { + if (delay < 0) { + throw new IllegalArgumentException("maxRetryDelayMillis must be non-negative"); + } + }); + jitterFactor.ifPresent(factor -> { + if (factor < 0 || factor > 1) { + throw new IllegalArgumentException("jitterFactor must be between 0 and 1"); + } + }); + this.maxRetries = maxRetries; + this.initialRetryDelay = initialRetryDelayMillis.map(Duration::ofMillis).orElse(DEFAULT_INITIAL_RETRY_DELAY); + this.maxRetryDelay = maxRetryDelayMillis.map(Duration::ofMillis).orElse(DEFAULT_MAX_RETRY_DELAY); + this.jitterFactor = jitterFactor.orElse(DEFAULT_JITTER_FACTOR); } @Override public Response intercept(Chain chain) throws IOException { - Response response = chain.proceed(chain.request()); + Request request = chain.request(); + int effectiveMaxRetries = resolveMaxRetries(request); + Response response = chain.proceed(request); if (shouldRetry(response.code())) { - return retryChain(response, chain); + return retryChain(response, chain, effectiveMaxRetries); } return response; } - private Response retryChain(Response response, Chain chain) throws IOException { - Optional nextBackoff = this.backoff.nextBackoff(); + private int resolveMaxRetries(Request request) { + MaxRetriesOverride override = request.tag(MaxRetriesOverride.class); + return override != null ? override.getValue() : this.maxRetries; + } + + private Response retryChain(Response response, Chain chain, int maxRetries) throws IOException { + ExponentialBackoff backoff = new ExponentialBackoff(maxRetries); + Optional nextBackoff = backoff.nextBackoff(response); while (nextBackoff.isPresent()) { try { Thread.sleep(nextBackoff.get().toMillis()); @@ -43,7 +86,7 @@ private Response retryChain(Response response, Chain chain) throws IOException { response.close(); response = chain.proceed(chain.request()); if (shouldRetry(response.code())) { - nextBackoff = this.backoff.nextBackoff(); + nextBackoff = backoff.nextBackoff(response); } else { return response; } @@ -52,10 +95,130 @@ private Response retryChain(Response response, Chain chain) throws IOException { return response; } + /** + * Calculates the retry delay from response headers, with fallback to exponential backoff. + * Priority: Retry-After > X-RateLimit-Reset > Exponential Backoff + */ + private Duration getRetryDelayFromHeaders(Response response, int retryAttempt) { + // Check for Retry-After header first (RFC 7231), with no jitter + String retryAfter = response.header("Retry-After"); + if (retryAfter != null) { + // Parse as number of seconds... + Optional secondsDelay = tryParseLong(retryAfter) + .map(seconds -> seconds * 1000) + .filter(delayMs -> delayMs > 0) + .map(delayMs -> Math.min(delayMs, maxRetryDelay.toMillis())) + .map(Duration::ofMillis); + if (secondsDelay.isPresent()) { + return secondsDelay.get(); + } + + // ...or as an HTTP date; both are valid + Optional dateDelay = tryParseHttpDate(retryAfter) + .map(resetTime -> resetTime.toInstant().toEpochMilli() - System.currentTimeMillis()) + .filter(delayMs -> delayMs > 0) + .map(delayMs -> Math.min(delayMs, maxRetryDelay.toMillis())) + .map(Duration::ofMillis); + if (dateDelay.isPresent()) { + return dateDelay.get(); + } + } + + // Then check for industry-standard X-RateLimit-Reset header, with positive jitter + String rateLimitReset = response.header("X-RateLimit-Reset"); + if (rateLimitReset != null) { + // Assume Unix timestamp in epoch seconds + Optional rateLimitDelay = tryParseLong(rateLimitReset) + .map(resetTimeSeconds -> (resetTimeSeconds * 1000) - System.currentTimeMillis()) + .filter(delayMs -> delayMs > 0) + .map(delayMs -> Math.min(delayMs, maxRetryDelay.toMillis())) + .map(this::addPositiveJitter) + .map(Duration::ofMillis); + if (rateLimitDelay.isPresent()) { + return rateLimitDelay.get(); + } + } + + // Fall back to exponential backoff, with symmetric jitter + long initialDelayMillis = initialRetryDelay.toMillis(); + long maxDelayMillis = maxRetryDelay.toMillis(); + long cappedDelay; + if (retryAttempt >= Long.SIZE - 1 || initialDelayMillis > (maxDelayMillis >> retryAttempt)) { + // initialDelayMillis * 2^retryAttempt would exceed maxDelayMillis (or overflow) + cappedDelay = maxDelayMillis; + } else { + cappedDelay = Math.min(initialDelayMillis << retryAttempt, maxDelayMillis); // 2^retryAttempt + } + return Duration.ofMillis(addSymmetricJitter(cappedDelay)); + } + + /** + * Attempts to parse a string as a long, returning empty Optional on failure. + */ + private Optional tryParseLong(String value) { + if (value == null) { + return Optional.empty(); + } + try { + return Optional.of(Long.parseLong(value)); + } catch (NumberFormatException e) { + return Optional.empty(); + } + } + + /** + * Attempts to parse a string as an HTTP date (RFC 1123), returning empty Optional on failure. + */ + private Optional tryParseHttpDate(String value) { + if (value == null) { + return Optional.empty(); + } + try { + return Optional.of(ZonedDateTime.parse(value, DateTimeFormatter.RFC_1123_DATE_TIME)); + } catch (DateTimeParseException e) { + return Optional.empty(); + } + } + + /** + * Adds positive jitter (100-120% of original value) to prevent thundering herd. + * Used for X-RateLimit-Reset header delays. + */ + private long addPositiveJitter(long delayMs) { + double jitterMultiplier = 1.0 + (random.nextDouble() * jitterFactor); + return (long) (delayMs * jitterMultiplier); + } + + /** + * Adds symmetric jitter (90-110% of original value) to prevent thundering herd. + * Used for exponential backoff delays. + */ + private long addSymmetricJitter(long delayMs) { + double jitterMultiplier = 1.0 + ((random.nextDouble() - 0.5) * jitterFactor); + return (long) (delayMs * jitterMultiplier); + } + private static boolean shouldRetry(int statusCode) { return statusCode == 408 || statusCode == 429 || statusCode >= 500; } + /** + * Per-request override carried on the OkHttp {@link Request} as a tag. + * When present, the interceptor uses this value instead of the client-wide + * {@code maxRetries} configured at construction time. + */ + public static final class MaxRetriesOverride { + private final int value; + + public MaxRetriesOverride(int value) { + this.value = value; + } + + public int getValue() { + return value; + } + } + private final class ExponentialBackoff { private final int maxNumRetries; @@ -66,14 +229,14 @@ private final class ExponentialBackoff { this.maxNumRetries = maxNumRetries; } - public Optional nextBackoff() { - retryNumber += 1; - if (retryNumber > maxNumRetries) { + public Optional nextBackoff(Response response) { + if (retryNumber >= maxNumRetries) { return Optional.empty(); } - int upperBound = (int) Math.pow(2, retryNumber); - return Optional.of(ONE_SECOND.multipliedBy(random.nextInt(upperBound))); + Duration delay = getRetryDelayFromHeaders(response, retryNumber); + retryNumber += 1; + return Optional.of(delay); } } } diff --git a/flowvault/src/main/java/com/skyflow/utils/SkyflowRetryInterceptor.java b/flowvault/src/main/java/com/skyflow/utils/SkyflowRetryInterceptor.java deleted file mode 100644 index 703c698e..00000000 --- a/flowvault/src/main/java/com/skyflow/utils/SkyflowRetryInterceptor.java +++ /dev/null @@ -1,97 +0,0 @@ -package com.skyflow.utils; - -import okhttp3.Interceptor; -import okhttp3.Response; - -import java.io.IOException; -import java.util.Random; - -/** - * Retries failed requests with exponential backoff and jitter. - *

- * This exists as hand-written code rather than using the generated - * {@code com.skyflow.generated.rest.core.RetryInterceptor} because that one only accepts a retry - * count — it has no way to configure the backoff delays that {@code VaultConfig} exposes. It also - * keeps its backoff counter on the interceptor instance, so a single shared instance exhausts its - * retry budget once for the whole client rather than once per request; this implementation keeps - * that state per call. - *

- * Retries the same statuses the generated interceptor does: 408, 429, and any 5xx. - */ -public final class SkyflowRetryInterceptor implements Interceptor { - - /** Fraction of the computed delay applied as random jitter, so retries do not align. */ - private static final double JITTER_FACTOR = 0.2; - - private final int maxRetries; - private final long initialRetryDelayMillis; - private final long maxRetryDelayMillis; - private final Random random = new Random(); - - public SkyflowRetryInterceptor(int maxRetries, long initialRetryDelayMillis, long maxRetryDelayMillis) { - if (maxRetries < 0) { - throw new IllegalArgumentException("maxRetries must be non-negative"); - } - if (initialRetryDelayMillis < 0) { - throw new IllegalArgumentException("initialRetryDelayMillis must be non-negative"); - } - if (maxRetryDelayMillis < 0) { - throw new IllegalArgumentException("maxRetryDelayMillis must be non-negative"); - } - this.maxRetries = maxRetries; - this.initialRetryDelayMillis = initialRetryDelayMillis; - this.maxRetryDelayMillis = maxRetryDelayMillis; - } - - @Override - public Response intercept(Chain chain) throws IOException { - Response response = chain.proceed(chain.request()); - // Retry budget is scoped to this call, not to the interceptor instance. - for (int attempt = 1; attempt <= maxRetries && shouldRetry(response.code()); attempt++) { - sleep(backoffMillis(attempt)); - response.close(); - response = chain.proceed(chain.request()); - } - return response; - } - - /** Exponential growth from the initial delay, capped at the maximum, then jittered. */ - long backoffMillis(int attempt) { - long delay = initialRetryDelayMillis; - for (int i = 1; i < attempt && delay < maxRetryDelayMillis; i++) { - delay = delay > maxRetryDelayMillis / 2 ? maxRetryDelayMillis : delay * 2; - } - delay = Math.min(delay, maxRetryDelayMillis); - long jitter = (long) (delay * JITTER_FACTOR); - if (jitter <= 0) { - return delay; - } - // delay +/- up to JITTER_FACTOR, never negative. - return Math.max(0, delay - jitter + random.nextInt((int) Math.min(2 * jitter + 1, Integer.MAX_VALUE))); - } - - static boolean shouldRetry(int statusCode) { - return statusCode == 408 || statusCode == 429 || statusCode >= 500; - } - - private static void sleep(long millis) throws IOException { - try { - Thread.sleep(millis); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - throw new IOException("Interrupted while waiting to retry request", e); - } - } - - public int getMaxRetries() { - return maxRetries; - } - - public long getInitialRetryDelayMillis() { - return initialRetryDelayMillis; - } - - public long getMaxRetryDelayMillis() { - return maxRetryDelayMillis; - } -} diff --git a/flowvault/src/test/java/com/skyflow/AuthInterceptorTests.java b/flowvault/src/test/java/com/skyflow/AuthInterceptorTests.java index e4c56165..f8d56b4e 100644 --- a/flowvault/src/test/java/com/skyflow/AuthInterceptorTests.java +++ b/flowvault/src/test/java/com/skyflow/AuthInterceptorTests.java @@ -6,7 +6,6 @@ import com.skyflow.errors.SkyflowException; import com.skyflow.generated.rest.core.RetryInterceptor; import com.skyflow.utils.FakeChain; -import com.skyflow.utils.SkyflowRetryInterceptor; import okhttp3.Interceptor; import org.junit.Assert; import org.junit.Test; @@ -38,7 +37,7 @@ private static Interceptor authInterceptorOf(VaultClient client) throws SkyflowE client.updateExecutorInHTTP(); List interceptors = client.sharedHttpClient.interceptors(); for (Interceptor interceptor : interceptors) { - if (!(interceptor instanceof SkyflowRetryInterceptor) && !(interceptor instanceof RetryInterceptor)) { + if (!(interceptor instanceof RetryInterceptor)) { return interceptor; } } diff --git a/flowvault/src/test/java/com/skyflow/HttpConfigTests.java b/flowvault/src/test/java/com/skyflow/HttpConfigTests.java index 9d1b4a53..e0282428 100644 --- a/flowvault/src/test/java/com/skyflow/HttpConfigTests.java +++ b/flowvault/src/test/java/com/skyflow/HttpConfigTests.java @@ -3,13 +3,15 @@ import com.skyflow.config.VaultConfig; import com.skyflow.enums.Env; import com.skyflow.errors.SkyflowException; -import com.skyflow.utils.SkyflowRetryInterceptor; +import com.skyflow.generated.rest.core.RetryInterceptor; import com.skyflow.vault.controller.VaultController; import okhttp3.Interceptor; import okhttp3.OkHttpClient; import org.junit.Assert; import org.junit.Test; +import java.lang.reflect.Field; +import java.time.Duration; import java.util.List; /** @@ -36,17 +38,39 @@ private static OkHttpClient httpClientOf(VaultClient client) throws SkyflowExcep return client.sharedHttpClient; } - private static int maxRetriesOf(SkyflowRetryInterceptor interceptor) { - return interceptor.getMaxRetries(); + private static int maxRetriesOf(RetryInterceptor interceptor) { + return (int) fieldOf(interceptor, "maxRetries"); } - private static SkyflowRetryInterceptor retryInterceptorOf(OkHttpClient http) { + // RetryInterceptor is generated code (com.skyflow.generated.*) and exposes no getters for the + // fields it was constructed with, so these tests - which exist to verify VaultClient's + // precedence/resolution logic actually reached the interceptor - read them back via reflection + // rather than adding hand-written accessors to a file meant to stay a faithful Fern output. + private static Object fieldOf(RetryInterceptor interceptor, String name) { + try { + Field field = RetryInterceptor.class.getDeclaredField(name); + field.setAccessible(true); + return field.get(interceptor); + } catch (ReflectiveOperationException e) { + throw new AssertionError(e); + } + } + + private static long initialRetryDelayMillisOf(RetryInterceptor interceptor) { + return ((Duration) fieldOf(interceptor, "initialRetryDelay")).toMillis(); + } + + private static long maxRetryDelayMillisOf(RetryInterceptor interceptor) { + return ((Duration) fieldOf(interceptor, "maxRetryDelay")).toMillis(); + } + + private static RetryInterceptor retryInterceptorOf(OkHttpClient http) { for (Interceptor interceptor : http.interceptors()) { - if (interceptor instanceof SkyflowRetryInterceptor) { - return (SkyflowRetryInterceptor) interceptor; + if (interceptor instanceof RetryInterceptor) { + return (RetryInterceptor) interceptor; } } - throw new AssertionError("No SkyflowRetryInterceptor installed on the HTTP client"); + throw new AssertionError("No RetryInterceptor installed on the HTTP client"); } // ── SDK defaults (neither level configured) ─────────────────────────────── @@ -198,8 +222,8 @@ public void testInterceptors_retryIsOuterSoEachAttemptRereadsTheToken() throws S List interceptors = http.interceptors(); Assert.assertEquals(2, interceptors.size()); Assert.assertTrue("Retry must be registered first so it wraps the auth interceptor", - interceptors.get(0) instanceof SkyflowRetryInterceptor); - Assert.assertFalse(interceptors.get(1) instanceof SkyflowRetryInterceptor); + interceptors.get(0) instanceof RetryInterceptor); + Assert.assertFalse(interceptors.get(1) instanceof RetryInterceptor); } @Test @@ -239,10 +263,10 @@ public void testExplicitZeroMaxRetries_overridesClientWideRetries() throws Skyfl @Test public void testDefaults_retryDelaysAre500And2000Millis() throws SkyflowException { - SkyflowRetryInterceptor retry = retryInterceptorOf(httpClientOf(new VaultClient(buildConfig(), null))); + RetryInterceptor retry = retryInterceptorOf(httpClientOf(new VaultClient(buildConfig(), null))); - Assert.assertEquals(500L, retry.getInitialRetryDelayMillis()); - Assert.assertEquals(2000L, retry.getMaxRetryDelayMillis()); + Assert.assertEquals(500L, initialRetryDelayMillisOf(retry)); + Assert.assertEquals(2000L, maxRetryDelayMillisOf(retry)); } @Test @@ -250,10 +274,10 @@ public void testRetryDelays_clientWideValuesApply() throws SkyflowException { VaultClient client = new VaultClient(buildConfig(), null); client.setCommonHttpConfig(null, null, null, null, 3, 100L, 900L); - SkyflowRetryInterceptor retry = retryInterceptorOf(httpClientOf(client)); + RetryInterceptor retry = retryInterceptorOf(httpClientOf(client)); - Assert.assertEquals(100L, retry.getInitialRetryDelayMillis()); - Assert.assertEquals(900L, retry.getMaxRetryDelayMillis()); + Assert.assertEquals(100L, initialRetryDelayMillisOf(retry)); + Assert.assertEquals(900L, maxRetryDelayMillisOf(retry)); } @Test @@ -265,10 +289,10 @@ public void testRetryDelays_vaultLevelBeatsClientWide() throws SkyflowException VaultClient client = new VaultClient(config, null); client.setCommonHttpConfig(null, null, null, null, 3, 100L, 900L); - SkyflowRetryInterceptor retry = retryInterceptorOf(httpClientOf(client)); + RetryInterceptor retry = retryInterceptorOf(httpClientOf(client)); - Assert.assertEquals(250L, retry.getInitialRetryDelayMillis()); - Assert.assertEquals(4000L, retry.getMaxRetryDelayMillis()); + Assert.assertEquals(250L, initialRetryDelayMillisOf(retry)); + Assert.assertEquals(4000L, maxRetryDelayMillisOf(retry)); } @Test @@ -279,10 +303,10 @@ public void testRetryDelays_resolveIndependentlyOfEachOther() throws SkyflowExce VaultClient client = new VaultClient(config, null); client.setCommonHttpConfig(null, null, null, null, 3, 100L, 900L); - SkyflowRetryInterceptor retry = retryInterceptorOf(httpClientOf(client)); + RetryInterceptor retry = retryInterceptorOf(httpClientOf(client)); - Assert.assertEquals(100L, retry.getInitialRetryDelayMillis()); // client-wide - Assert.assertEquals(4000L, retry.getMaxRetryDelayMillis()); // vault + Assert.assertEquals(100L, initialRetryDelayMillisOf(retry)); // client-wide + Assert.assertEquals(4000L, maxRetryDelayMillisOf(retry)); // vault } @Test @@ -294,11 +318,11 @@ public void testRetryDelays_endToEndThroughTheBuilder() throws SkyflowException .addVaultConfig(buildConfig()) .build(); - SkyflowRetryInterceptor retry = retryInterceptorOf(httpClientOf(client.vault())); + RetryInterceptor retry = retryInterceptorOf(httpClientOf(client.vault())); - Assert.assertEquals(3, retry.getMaxRetries()); - Assert.assertEquals(100L, retry.getInitialRetryDelayMillis()); - Assert.assertEquals(900L, retry.getMaxRetryDelayMillis()); + Assert.assertEquals(3, maxRetriesOf(retry)); + Assert.assertEquals(100L, initialRetryDelayMillisOf(retry)); + Assert.assertEquals(900L, maxRetryDelayMillisOf(retry)); } @Test @@ -312,7 +336,7 @@ public void testRetryDelays_vaultConfigBeatsBuilderEndToEnd() throws SkyflowExce .build(); Assert.assertEquals(250L, - retryInterceptorOf(httpClientOf(client.vault())).getInitialRetryDelayMillis()); + initialRetryDelayMillisOf(retryInterceptorOf(httpClientOf(client.vault())))); } @Test @@ -323,11 +347,11 @@ public void testRetryDelays_survivedUpdateVaultConfig() throws SkyflowException update.setInitialRetryDelayMillis(250L); update.setMaxRetryDelayMillis(4000L); - SkyflowRetryInterceptor retry = + RetryInterceptor retry = retryInterceptorOf(httpClientOf(builder.updateVaultConfig(update).build().vault())); - Assert.assertEquals(250L, retry.getInitialRetryDelayMillis()); - Assert.assertEquals(4000L, retry.getMaxRetryDelayMillis()); + Assert.assertEquals(250L, initialRetryDelayMillisOf(retry)); + Assert.assertEquals(4000L, maxRetryDelayMillisOf(retry)); } @Test @@ -342,7 +366,7 @@ public void testRetryDelays_builderMethodsAreFluent() { @Test public void testInvalidMaxRetries_wrapsInterceptorIllegalArgumentAsSkyflowException() throws SkyflowException { - // SkyflowRetryInterceptor rejects negative maxRetries with IllegalArgumentException; + // The generated RetryInterceptor rejects negative maxRetries with IllegalArgumentException; // updateExecutorInHTTP must translate that (and anything else from client construction) // into a SkyflowException rather than letting it escape raw. VaultConfig config = buildConfig(); diff --git a/flowvault/src/test/java/com/skyflow/utils/SkyflowRetryInterceptorTests.java b/flowvault/src/test/java/com/skyflow/utils/SkyflowRetryInterceptorTests.java deleted file mode 100644 index aae9a6f9..00000000 --- a/flowvault/src/test/java/com/skyflow/utils/SkyflowRetryInterceptorTests.java +++ /dev/null @@ -1,254 +0,0 @@ -package com.skyflow.utils; - -import okhttp3.Response; -import org.junit.Assert; -import org.junit.Test; - -import java.io.IOException; - -/** - * Internals of the retry interceptor. Lives in com.skyflow.utils so the package-private - * backoff/should-retry helpers stay off the public API surface. - */ -public class SkyflowRetryInterceptorTests { - - @Test - public void testConstructor_rejectsNegativeMaxRetries() { - try { - new SkyflowRetryInterceptor(-1, 500L, 2000L); - Assert.fail("negative maxRetries should be rejected"); - } catch (IllegalArgumentException expected) { - Assert.assertTrue(expected.getMessage().contains("maxRetries")); - } - } - - @Test - public void testConstructor_rejectsNegativeInitialDelay() { - try { - new SkyflowRetryInterceptor(1, -1L, 2000L); - Assert.fail("negative initialRetryDelayMillis should be rejected"); - } catch (IllegalArgumentException expected) { - Assert.assertTrue(expected.getMessage().contains("initialRetryDelayMillis")); - } - } - - @Test - public void testConstructor_rejectsNegativeMaxDelay() { - try { - new SkyflowRetryInterceptor(1, 500L, -1L); - Assert.fail("negative maxRetryDelayMillis should be rejected"); - } catch (IllegalArgumentException expected) { - Assert.assertTrue(expected.getMessage().contains("maxRetryDelayMillis")); - } - } - - @Test - public void testBackoff_growsExponentiallyThenCaps() { - // Jitter is +/-20%, so assert bands rather than exact values. - SkyflowRetryInterceptor retry = new SkyflowRetryInterceptor(5, 100L, 400L); - - assertWithinJitter(100L, retry.backoffMillis(1)); - assertWithinJitter(200L, retry.backoffMillis(2)); - assertWithinJitter(400L, retry.backoffMillis(3)); - assertWithinJitter(400L, retry.backoffMillis(4)); - assertWithinJitter(400L, retry.backoffMillis(10)); - } - - @Test - public void testBackoff_neverExceedsTheCapAcrossManyDraws() { - SkyflowRetryInterceptor retry = new SkyflowRetryInterceptor(5, 100L, 400L); - - for (int i = 0; i < 200; i++) { - long delay = retry.backoffMillis(3); - Assert.assertTrue("jittered delay went negative: " + delay, delay >= 0); - Assert.assertTrue("jittered delay exceeded cap + jitter: " + delay, delay <= 480L); - } - } - - @Test - public void testBackoff_zeroDelayStaysZero() { - SkyflowRetryInterceptor retry = new SkyflowRetryInterceptor(3, 0L, 0L); - - Assert.assertEquals(0L, retry.backoffMillis(1)); - Assert.assertEquals(0L, retry.backoffMillis(5)); - } - - @Test - public void testBackoff_initialDelayAboveCapIsClampedToCap() { - SkyflowRetryInterceptor retry = new SkyflowRetryInterceptor(3, 5000L, 1000L); - - assertWithinJitter(1000L, retry.backoffMillis(1)); - assertWithinJitter(1000L, retry.backoffMillis(3)); - } - - @Test - public void testShouldRetry_retryableStatuses() { - Assert.assertTrue(SkyflowRetryInterceptor.shouldRetry(408)); - Assert.assertTrue(SkyflowRetryInterceptor.shouldRetry(429)); - Assert.assertTrue(SkyflowRetryInterceptor.shouldRetry(500)); - Assert.assertTrue(SkyflowRetryInterceptor.shouldRetry(502)); - Assert.assertTrue(SkyflowRetryInterceptor.shouldRetry(503)); - } - - @Test - public void testShouldRetry_nonRetryableStatuses() { - Assert.assertFalse(SkyflowRetryInterceptor.shouldRetry(200)); - Assert.assertFalse(SkyflowRetryInterceptor.shouldRetry(201)); - Assert.assertFalse(SkyflowRetryInterceptor.shouldRetry(400)); - Assert.assertFalse(SkyflowRetryInterceptor.shouldRetry(401)); - Assert.assertFalse(SkyflowRetryInterceptor.shouldRetry(404)); - Assert.assertFalse(SkyflowRetryInterceptor.shouldRetry(409)); - } - - @Test - public void testAccessors_reportWhatWasConfigured() { - SkyflowRetryInterceptor retry = new SkyflowRetryInterceptor(3, 100L, 900L); - - Assert.assertEquals(3, retry.getMaxRetries()); - Assert.assertEquals(100L, retry.getInitialRetryDelayMillis()); - Assert.assertEquals(900L, retry.getMaxRetryDelayMillis()); - } - - // ── intercept(): the retry loop ─────────────────────────────────────────── - // Delays are set to 0 so these do not actually sleep. - - private static SkyflowRetryInterceptor retrying(int maxRetries) { - return new SkyflowRetryInterceptor(maxRetries, 0L, 0L); - } - - @Test - public void testIntercept_successFirstTimeIsNotRetried() throws IOException { - FakeChain chain = new FakeChain(200); - - Response response = retrying(3).intercept(chain); - - Assert.assertEquals(1, chain.calls()); - Assert.assertEquals(200, response.code()); - } - - @Test - public void testIntercept_nonRetryableFailureIsNotRetried() throws IOException { - FakeChain chain = new FakeChain(400); - - Response response = retrying(3).intercept(chain); - - Assert.assertEquals("a 400 must not be replayed", 1, chain.calls()); - Assert.assertEquals(400, response.code()); - } - - @Test - public void testIntercept_retriesUpToTheBudgetThenReturnsTheLastFailure() throws IOException { - FakeChain chain = new FakeChain(500); - - Response response = retrying(2).intercept(chain); - - Assert.assertEquals("1 initial attempt + 2 retries", 3, chain.calls()); - Assert.assertEquals(500, response.code()); - } - - @Test - public void testIntercept_stopsAsSoonAsAnAttemptSucceeds() throws IOException { - FakeChain chain = new FakeChain(503, 200, 200); - - Response response = retrying(5).intercept(chain); - - Assert.assertEquals("must not keep retrying after success", 2, chain.calls()); - Assert.assertEquals(200, response.code()); - } - - @Test - public void testIntercept_stopsOnANonRetryableStatusMidWay() throws IOException { - FakeChain chain = new FakeChain(500, 404, 200); - - Response response = retrying(5).intercept(chain); - - Assert.assertEquals(2, chain.calls()); - Assert.assertEquals(404, response.code()); - } - - @Test - public void testIntercept_zeroBudgetMeansNoRetryAtAll() throws IOException { - FakeChain chain = new FakeChain(500); - - Response response = retrying(0).intercept(chain); - - Assert.assertEquals(1, chain.calls()); - Assert.assertEquals(500, response.code()); - } - - @Test - public void testIntercept_retriesEachRetryableStatus() throws IOException { - for (int code : new int[] {408, 429, 500, 502, 503}) { - FakeChain chain = new FakeChain(code, 200); - - Response response = retrying(1).intercept(chain); - - Assert.assertEquals("should have retried a " + code, 2, chain.calls()); - Assert.assertEquals(200, response.code()); - } - } - - @Test - public void testIntercept_closesEverySupersededResponse() throws IOException { - // Leaking the body of a response we are about to discard would leak the connection. - FakeChain chain = new FakeChain(500, 500, 200); - - retrying(2).intercept(chain); - - Assert.assertEquals(3, chain.bodies().size()); - Assert.assertTrue("first failed response not closed", chain.bodies().get(0).closed); - Assert.assertTrue("second failed response not closed", chain.bodies().get(1).closed); - Assert.assertFalse("the returned response must stay open", chain.bodies().get(2).closed); - } - - @Test - public void testIntercept_retryBudgetIsPerCallNotPerInterceptorInstance() throws IOException { - // The generated RetryInterceptor keeps its backoff counter on the instance, so one shared - // instance exhausts the budget once for the whole client. A single interceptor is installed - // on a shared OkHttpClient, so every call must get its own full budget. - SkyflowRetryInterceptor retry = retrying(2); - - FakeChain first = new FakeChain(500); - retry.intercept(first); - FakeChain second = new FakeChain(500); - retry.intercept(second); - FakeChain third = new FakeChain(500); - retry.intercept(third); - - Assert.assertEquals(3, first.calls()); - Assert.assertEquals("second call lost its retry budget", 3, second.calls()); - Assert.assertEquals("third call lost its retry budget", 3, third.calls()); - } - - @Test - public void testIntercept_interruptionSurfacesAsIOException() { - SkyflowRetryInterceptor retry = new SkyflowRetryInterceptor(2, 5_000L, 5_000L); - FakeChain chain = new FakeChain(500); - - Thread.currentThread().interrupt(); - try { - retry.intercept(chain); - Assert.fail("an interrupt while backing off should surface as IOException"); - } catch (IOException e) { - Assert.assertTrue(e.getMessage().contains("Interrupted")); - Assert.assertTrue("the interrupt flag must be restored", Thread.currentThread().isInterrupted()); - } finally { - Thread.interrupted(); // clear the flag so it cannot leak into another test - } - } - - @Test - public void testBackoff_growsPastHalfTheCapInOneStep() { - // initial > max/2, so the next step clamps straight to the cap instead of doubling past it. - SkyflowRetryInterceptor retry = new SkyflowRetryInterceptor(3, 300L, 400L); - - assertWithinJitter(300L, retry.backoffMillis(1)); - assertWithinJitter(400L, retry.backoffMillis(2)); - } - - private static void assertWithinJitter(long expected, long actual) { - long jitter = (long) (expected * 0.2); - Assert.assertTrue("expected ~" + expected + " (+/-" + jitter + ") but got " + actual, - actual >= expected - jitter && actual <= expected + jitter); - } -} From 8746ebd337468f910d8c257a43a489c448b05d34 Mon Sep 17 00:00:00 2001 From: Devesh Bhardwaj Date: Tue, 8 Sep 2026 12:10:07 +0530 Subject: [PATCH 2/3] Pass the hardcoded jitter factor explicitly instead of relying on RetryInterceptor's default RetryInterceptor is generated code we don't maintain, so its internal DEFAULT_JITTER_FACTOR is free to change on a future regeneration. Passing Optional.of(RETRY_JITTER_FACTOR) from VaultClient keeps the 0.2 hardcode visible and owned in our own code instead of implicitly inherited by passing Optional.empty(). Co-Authored-By: Claude Sonnet 5 --- flowvault/src/main/java/com/skyflow/VaultClient.java | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/flowvault/src/main/java/com/skyflow/VaultClient.java b/flowvault/src/main/java/com/skyflow/VaultClient.java index 15543aec..e069eb0e 100644 --- a/flowvault/src/main/java/com/skyflow/VaultClient.java +++ b/flowvault/src/main/java/com/skyflow/VaultClient.java @@ -34,6 +34,11 @@ public class VaultClient extends BaseVaultClient { private static final int DEFAULT_MAX_RETRIES = 0; private static final long DEFAULT_INITIAL_RETRY_DELAY_MILLIS = 500L; private static final long DEFAULT_MAX_RETRY_DELAY_MILLIS = 2000L; + // Not yet exposed as a VaultConfig/builder setting, so hardcoded here rather than left as + // Optional.empty() - passing it explicitly keeps the choice visible in our own code instead of + // depending on RetryInterceptor's internal default, which is free to change on a future + // regeneration since it is generated code we do not maintain. + private static final double RETRY_JITTER_FACTOR = 0.2; protected VaultClient(VaultConfig vaultConfig, Credentials credentials) throws SkyflowException { super(vaultConfig, credentials); @@ -175,10 +180,8 @@ protected void updateExecutorInHTTP() throws SkyflowException { .callTimeout(timeoutSeconds, TimeUnit.SECONDS) // OUTER: retries. Must wrap the auth interceptor so each attempt re-reads the // (possibly refreshed) bearer token rather than replaying a stale one. - // Jitter is left at the generated interceptor's own default (0.2) - not yet - // exposed as a VaultConfig/builder setting. .addInterceptor(new RetryInterceptor(maxRetries, Optional.of(initialRetryDelayMillis), - Optional.of(maxRetryDelayMillis), Optional.empty())) + Optional.of(maxRetryDelayMillis), Optional.of(RETRY_JITTER_FACTOR))) .addInterceptor(chain -> { // INNER: auth Request requestWithAuth = chain.request().newBuilder() .header("Authorization", "Bearer " + this.token) From 82d0cc7fd25732365c45c552f428953f24d342d7 Mon Sep 17 00:00:00 2001 From: Devesh Bhardwaj Date: Tue, 8 Sep 2026 12:26:04 +0530 Subject: [PATCH 3/3] Make console logging async and memoize .env lookups on the request hot path Both were found during a round-trip latency audit of the SDK's request path (same audit that flagged the retry-interceptor duplication this branch fixes). - LogUtil's ConsoleHandler writes and flushes to the console synchronously, under a lock shared by every calling thread. LogUtil.printInfoLog/etc. fire several times per request (setBearerToken reuse/expiry, request validation, per-batch triggers, request-resolved), so enabling INFO/DEBUG logging turned console output into a per-request blocking-I/O + contention point. A new AsyncConsoleHandler now hands each LogRecord to a single background daemon thread via a bounded, non-blocking queue; the calling thread only enqueues. Applied to both the common/flowvault LogUtil and skyvault's separate copy. - Several settings/credentials lookups fell back to Dotenv.load() when an env var wasn't set: VaultController.resolveSettingFromEnvironment did this twice per bulk call (batch size + concurrency limit), unconditionally, on every single insert/detokenize/tokenize/deleteTokens call -- Dotenv.load() re-reads the .env file from disk every time it's invoked, so this was uncached, blocking disk I/O on the hot path. The same pattern existed in BaseVaultClient/ConnectionClient's credential fallback and Utils.getEnvVaultUrl. BaseUtils.resolveEnvOrDotenv now memoizes the loaded (or absent) .env for the life of the JVM -- a project's .env doesn't change while the process runs, so there's no reason to keep re-reading it -- with a resetDotenvCacheForTests() escape hatch for the handful of tests that intentionally rewrite .env mid-run to exercise both branches. Testing: full common + skyvault + flowvault suites, 1450/1450 passing (165+575+710), BUILD SUCCESS. Co-Authored-By: Claude Sonnet 5 --- .../java/com/skyflow/BaseVaultClient.java | 7 +- .../java/com/skyflow/utils/BaseUtils.java | 55 +++++++++++++ .../utils/logger/AsyncConsoleHandler.java | 82 +++++++++++++++++++ .../com/skyflow/utils/logger/LogUtil.java | 5 +- .../com/skyflow/BaseVaultClientTests.java | 2 + .../main/java/com/skyflow/utils/Utils.java | 7 +- .../vault/controller/VaultController.java | 13 +-- .../java/com/skyflow/utils/UtilsTests.java | 1 + .../java/com/skyflow/ConnectionClient.java | 5 +- .../utils/logger/AsyncConsoleHandler.java | 82 +++++++++++++++++++ .../com/skyflow/utils/logger/LogUtil.java | 5 +- .../skyflow/ConnectionClientDotenvTests.java | 2 + 12 files changed, 237 insertions(+), 29 deletions(-) create mode 100644 common/src/main/java/com/skyflow/utils/logger/AsyncConsoleHandler.java create mode 100644 skyvault/src/main/java/com/skyflow/utils/logger/AsyncConsoleHandler.java diff --git a/common/src/main/java/com/skyflow/BaseVaultClient.java b/common/src/main/java/com/skyflow/BaseVaultClient.java index 3d5b0d32..b71f6d03 100644 --- a/common/src/main/java/com/skyflow/BaseVaultClient.java +++ b/common/src/main/java/com/skyflow/BaseVaultClient.java @@ -14,7 +14,6 @@ import com.skyflow.utils.BaseUtils; import com.skyflow.utils.logger.LogUtil; import com.skyflow.utils.validations.BaseValidations; -import io.github.cdimascio.dotenv.Dotenv; import io.github.cdimascio.dotenv.DotenvException; import okhttp3.ConnectionPool; import okhttp3.OkHttpClient; @@ -63,11 +62,7 @@ protected synchronized void prioritiseCredentials(BaseCredentials vaultSpecificC } else if (this.commonCredentials != null) { this.finalCredentials = this.commonCredentials; } else { - String sysCredentials = System.getenv(BaseConstants.ENV_CREDENTIALS_KEY_NAME); - if (sysCredentials == null) { - Dotenv dotenv = Dotenv.load(); - sysCredentials = dotenv.get(BaseConstants.ENV_CREDENTIALS_KEY_NAME); - } + String sysCredentials = BaseUtils.resolveEnvOrDotenv(BaseConstants.ENV_CREDENTIALS_KEY_NAME); if (sysCredentials == null) { throw new SkyflowException(ErrorCode.INVALID_INPUT.getCode(), ErrorMessage.EmptyCredentials.getMessage()); } else { diff --git a/common/src/main/java/com/skyflow/utils/BaseUtils.java b/common/src/main/java/com/skyflow/utils/BaseUtils.java index e536c111..9187717f 100644 --- a/common/src/main/java/com/skyflow/utils/BaseUtils.java +++ b/common/src/main/java/com/skyflow/utils/BaseUtils.java @@ -21,8 +21,63 @@ import com.skyflow.logs.InfoLogs; import com.skyflow.serviceaccount.util.BearerToken; import com.skyflow.utils.logger.LogUtil; +import io.github.cdimascio.dotenv.Dotenv; +import io.github.cdimascio.dotenv.DotenvException; public class BaseUtils { + + // Memoized .env: Dotenv.load() does a filesystem read every time it's called, and several + // call sites resolve a setting this way on every single SDK request -- re-reading a file + // whose contents never change for the life of the process turned those into repeated, + // uncached, blocking disk I/O on the hot path. Loaded at most once per JVM; `dotenvAttempted` + // also memoizes the "no .env file present" outcome so a missing file isn't retried either. + private static volatile boolean dotenvAttempted = false; + private static volatile Dotenv cachedDotenv = null; + + private static Dotenv memoizedDotenv() { + if (!dotenvAttempted) { + synchronized (BaseUtils.class) { + if (!dotenvAttempted) { + try { + cachedDotenv = Dotenv.load(); + } catch (DotenvException e) { + cachedDotenv = null; // no .env file in the working directory + } + dotenvAttempted = true; + } + } + } + return cachedDotenv; + } + + /** + * Resolves {@code key} from the process environment first, falling back to the (memoized) + * {@code .env} file if present. Returns null if found in neither. + */ + public static String resolveEnvOrDotenv(String key) { + String value = System.getenv(key); + if (value == null) { + Dotenv dotenv = memoizedDotenv(); + if (dotenv != null) { + value = dotenv.get(key); + } + } + return value; + } + + /** + * Test-only: forces the next {@link #resolveEnvOrDotenv} call to re-read the {@code .env} + * file from disk instead of reusing the memoized one. Production code always wants the + * memoized behavior (a project's {@code .env} doesn't change while the process is running); + * this exists purely so tests that rewrite {@code .env} mid-run can observe the new content + * without restarting the JVM. + */ + public static void resetDotenvCacheForTests() { + synchronized (BaseUtils.class) { + dotenvAttempted = false; + cachedDotenv = null; + } + } public static String generateBearerToken(BaseCredentials credentials) throws SkyflowException { if (credentials.getPath() != null) { BearerToken.BearerTokenBuilder builder = BearerToken.builder() diff --git a/common/src/main/java/com/skyflow/utils/logger/AsyncConsoleHandler.java b/common/src/main/java/com/skyflow/utils/logger/AsyncConsoleHandler.java new file mode 100644 index 00000000..4193ad0e --- /dev/null +++ b/common/src/main/java/com/skyflow/utils/logger/AsyncConsoleHandler.java @@ -0,0 +1,82 @@ +package com.skyflow.utils.logger; + +import java.util.concurrent.LinkedBlockingQueue; +import java.util.logging.ErrorManager; +import java.util.logging.Handler; +import java.util.logging.LogRecord; + +/** + * Wraps a delegate {@link Handler} (in practice a {@link java.util.logging.ConsoleHandler}) so that + * {@link #publish(LogRecord)} never performs blocking I/O on the calling thread. + *

+ * {@code ConsoleHandler.publish} writes to and flushes the underlying stream synchronously, and does + * so under a lock shared by every thread using the logger. Since {@code LogUtil.printInfoLog}/etc. are + * called on every SDK request (often several times per call), that turns console logging into a + * per-request blocking-I/O + lock-contention point under concurrent load. This handler hands each + * {@link LogRecord} off to a single background daemon thread instead, which performs the actual write; + * the calling thread only enqueues. + *

+ * The handoff never blocks or applies backpressure to the caller: if the queue is momentarily full + * (a sustained logging flood, or the writer thread stalled) the record is dropped rather than slowing + * down request-serving threads — logging must never become the bottleneck it was flagged for. + */ +final class AsyncConsoleHandler extends Handler { + + /** Bounds worst-case memory use if the writer thread falls behind; excess records are dropped. */ + private static final int QUEUE_CAPACITY = 10_000; + + private final Handler delegate; + private final LinkedBlockingQueue queue = new LinkedBlockingQueue<>(QUEUE_CAPACITY); + private final Thread writer; + private volatile boolean closed = false; + + AsyncConsoleHandler(Handler delegate) { + this.delegate = delegate; + setLevel(delegate.getLevel()); + this.writer = new Thread(this::drain, "skyflow-sdk-log-writer"); + this.writer.setDaemon(true); + this.writer.start(); + } + + @Override + public void publish(LogRecord record) { + if (closed || !isLoggable(record)) { + return; + } + // offer() never blocks: a full queue means "drop", never "wait". + queue.offer(record); + } + + private void drain() { + try { + while (true) { + LogRecord record = queue.take(); + try { + delegate.publish(record); + } catch (RuntimeException e) { + reportError(null, e, ErrorManager.WRITE_FAILURE); + } + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + @Override + public void flush() { + // Best-effort: drain what's queued right now onto the delegate, then flush it. + LogRecord record; + while ((record = queue.poll()) != null) { + delegate.publish(record); + } + delegate.flush(); + } + + @Override + public void close() { + closed = true; + writer.interrupt(); + flush(); + delegate.close(); + } +} diff --git a/common/src/main/java/com/skyflow/utils/logger/LogUtil.java b/common/src/main/java/com/skyflow/utils/logger/LogUtil.java index ed2b7671..0c43dbeb 100644 --- a/common/src/main/java/com/skyflow/utils/logger/LogUtil.java +++ b/common/src/main/java/com/skyflow/utils/logger/LogUtil.java @@ -35,7 +35,10 @@ public synchronized String format(LogRecord logRecord) { consoleHandler.setFormatter(formatter); consoleHandler.setLevel(Level.CONFIG); - LOGGER.addHandler(consoleHandler); + // The actual write+flush to the console happens on a background thread, so per-call + // logging never blocks (or lock-contends on) the request-serving thread. See + // AsyncConsoleHandler's class doc for why this matters under concurrent load. + LOGGER.addHandler(new AsyncConsoleHandler(consoleHandler)); LOGGER.setLevel(logLevelToLoggerLevelMap(logLevel)); printInfoLog(InfoLogs.LOGGER_SETUP_DONE.getLog()); } diff --git a/common/src/test/java/com/skyflow/BaseVaultClientTests.java b/common/src/test/java/com/skyflow/BaseVaultClientTests.java index ec78c397..8afc650b 100644 --- a/common/src/test/java/com/skyflow/BaseVaultClientTests.java +++ b/common/src/test/java/com/skyflow/BaseVaultClientTests.java @@ -6,6 +6,7 @@ import com.skyflow.errors.SkyflowException; import com.skyflow.logs.ErrorLogs; import com.skyflow.utils.BaseConstants; +import com.skyflow.utils.BaseUtils; import okhttp3.Call; import okhttp3.Connection; import okhttp3.Interceptor; @@ -38,6 +39,7 @@ public class BaseVaultClientTests { public void saveEnvFileState() throws IOException { File f = new File(ENV_FILE); originalEnvContent = f.exists() ? Files.readAllBytes(Paths.get(ENV_FILE)) : null; + BaseUtils.resetDotenvCacheForTests(); // see its javadoc: .env is otherwise memoized JVM-wide } @After diff --git a/flowvault/src/main/java/com/skyflow/utils/Utils.java b/flowvault/src/main/java/com/skyflow/utils/Utils.java index cfd8195f..91f023d7 100644 --- a/flowvault/src/main/java/com/skyflow/utils/Utils.java +++ b/flowvault/src/main/java/com/skyflow/utils/Utils.java @@ -52,7 +52,6 @@ import com.skyflow.vault.data.TokenGroupRedactions; import com.skyflow.vault.data.UpsertOptions; -import io.github.cdimascio.dotenv.Dotenv; import io.github.cdimascio.dotenv.DotenvException; import java.io.File; @@ -82,11 +81,7 @@ public static JsonObject getMetrics() { public static String getEnvVaultUrl() throws SkyflowException { try { - String vaultUrl = System.getenv("VAULT_URL"); - if (vaultUrl == null) { - Dotenv dotenv = Dotenv.load(); - vaultUrl = dotenv.get("VAULT_URL"); - } + String vaultUrl = resolveEnvOrDotenv("VAULT_URL"); if (vaultUrl != null && vaultUrl.trim().isEmpty()) { LogUtil.printErrorLog(ErrorLogs.EMPTY_VAULT_URL.getLog()); throw new SkyflowException(ErrorCode.INVALID_INPUT.getCode(), ErrorMessage.EmptyVaultUrl.getMessage()); diff --git a/flowvault/src/main/java/com/skyflow/vault/controller/VaultController.java b/flowvault/src/main/java/com/skyflow/vault/controller/VaultController.java index da0eca19..4192d14c 100644 --- a/flowvault/src/main/java/com/skyflow/vault/controller/VaultController.java +++ b/flowvault/src/main/java/com/skyflow/vault/controller/VaultController.java @@ -58,9 +58,6 @@ import com.skyflow.vault.data.RequestInterceptor; import com.skyflow.vault.data.TokenizeOptions; -import io.github.cdimascio.dotenv.Dotenv; -import io.github.cdimascio.dotenv.DotenvException; - public final class VaultController extends VaultClient { private static final Gson gson = new GsonBuilder().serializeNulls().create(); private JsonObject metrics = Utils.getMetrics(); @@ -507,15 +504,7 @@ private ApiClientHttpResponse processDeleteTokensBatc static Function settingResolver = VaultController::resolveSettingFromEnvironment; private static String resolveSettingFromEnvironment(String key) { - String value = System.getenv(key); - if (value == null) { - try { - value = Dotenv.load().get(key); - } catch (DotenvException ignored) { - // no .env available — environment-only - } - } - return value; + return Utils.resolveEnvOrDotenv(key); } private BatchConfig configureDeleteTokensConcurrencyAndBatchSize(int totalRequests) { diff --git a/flowvault/src/test/java/com/skyflow/utils/UtilsTests.java b/flowvault/src/test/java/com/skyflow/utils/UtilsTests.java index 3d0ab81d..e2854d57 100644 --- a/flowvault/src/test/java/com/skyflow/utils/UtilsTests.java +++ b/flowvault/src/test/java/com/skyflow/utils/UtilsTests.java @@ -68,6 +68,7 @@ public class UtilsTests { public void saveEnvFileState() throws IOException { File f = new File(ENV_FILE); originalEnvContent = f.exists() ? Files.readAllBytes(Paths.get(ENV_FILE)) : null; + BaseUtils.resetDotenvCacheForTests(); // see its javadoc: .env is otherwise memoized JVM-wide } @After diff --git a/skyvault/src/main/java/com/skyflow/ConnectionClient.java b/skyvault/src/main/java/com/skyflow/ConnectionClient.java index d67122ad..a12b17e5 100644 --- a/skyvault/src/main/java/com/skyflow/ConnectionClient.java +++ b/skyvault/src/main/java/com/skyflow/ConnectionClient.java @@ -7,11 +7,11 @@ import com.skyflow.errors.SkyflowException; import com.skyflow.logs.InfoLogs; import com.skyflow.serviceaccount.util.Token; +import com.skyflow.utils.BaseUtils; import com.skyflow.utils.Constants; import com.skyflow.utils.Utils; import com.skyflow.utils.logger.LogUtil; import com.skyflow.utils.validations.Validations; -import io.github.cdimascio.dotenv.Dotenv; import io.github.cdimascio.dotenv.DotenvException; public class ConnectionClient { @@ -69,8 +69,7 @@ private void prioritiseCredentials() throws SkyflowException { } else if (this.commonCredentials != null) { this.finalCredentials = this.commonCredentials; } else { - Dotenv dotenv = Dotenv.load(); - String sysCredentials = dotenv.get(Constants.ENV_CREDENTIALS_KEY_NAME); + String sysCredentials = BaseUtils.resolveEnvOrDotenv(Constants.ENV_CREDENTIALS_KEY_NAME); if (sysCredentials == null) { throw new SkyflowException(ErrorCode.INVALID_INPUT.getCode(), ErrorMessage.EmptyCredentials.getMessage()); diff --git a/skyvault/src/main/java/com/skyflow/utils/logger/AsyncConsoleHandler.java b/skyvault/src/main/java/com/skyflow/utils/logger/AsyncConsoleHandler.java new file mode 100644 index 00000000..4193ad0e --- /dev/null +++ b/skyvault/src/main/java/com/skyflow/utils/logger/AsyncConsoleHandler.java @@ -0,0 +1,82 @@ +package com.skyflow.utils.logger; + +import java.util.concurrent.LinkedBlockingQueue; +import java.util.logging.ErrorManager; +import java.util.logging.Handler; +import java.util.logging.LogRecord; + +/** + * Wraps a delegate {@link Handler} (in practice a {@link java.util.logging.ConsoleHandler}) so that + * {@link #publish(LogRecord)} never performs blocking I/O on the calling thread. + *

+ * {@code ConsoleHandler.publish} writes to and flushes the underlying stream synchronously, and does + * so under a lock shared by every thread using the logger. Since {@code LogUtil.printInfoLog}/etc. are + * called on every SDK request (often several times per call), that turns console logging into a + * per-request blocking-I/O + lock-contention point under concurrent load. This handler hands each + * {@link LogRecord} off to a single background daemon thread instead, which performs the actual write; + * the calling thread only enqueues. + *

+ * The handoff never blocks or applies backpressure to the caller: if the queue is momentarily full + * (a sustained logging flood, or the writer thread stalled) the record is dropped rather than slowing + * down request-serving threads — logging must never become the bottleneck it was flagged for. + */ +final class AsyncConsoleHandler extends Handler { + + /** Bounds worst-case memory use if the writer thread falls behind; excess records are dropped. */ + private static final int QUEUE_CAPACITY = 10_000; + + private final Handler delegate; + private final LinkedBlockingQueue queue = new LinkedBlockingQueue<>(QUEUE_CAPACITY); + private final Thread writer; + private volatile boolean closed = false; + + AsyncConsoleHandler(Handler delegate) { + this.delegate = delegate; + setLevel(delegate.getLevel()); + this.writer = new Thread(this::drain, "skyflow-sdk-log-writer"); + this.writer.setDaemon(true); + this.writer.start(); + } + + @Override + public void publish(LogRecord record) { + if (closed || !isLoggable(record)) { + return; + } + // offer() never blocks: a full queue means "drop", never "wait". + queue.offer(record); + } + + private void drain() { + try { + while (true) { + LogRecord record = queue.take(); + try { + delegate.publish(record); + } catch (RuntimeException e) { + reportError(null, e, ErrorManager.WRITE_FAILURE); + } + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + + @Override + public void flush() { + // Best-effort: drain what's queued right now onto the delegate, then flush it. + LogRecord record; + while ((record = queue.poll()) != null) { + delegate.publish(record); + } + delegate.flush(); + } + + @Override + public void close() { + closed = true; + writer.interrupt(); + flush(); + delegate.close(); + } +} diff --git a/skyvault/src/main/java/com/skyflow/utils/logger/LogUtil.java b/skyvault/src/main/java/com/skyflow/utils/logger/LogUtil.java index 85655506..28956c1b 100644 --- a/skyvault/src/main/java/com/skyflow/utils/logger/LogUtil.java +++ b/skyvault/src/main/java/com/skyflow/utils/logger/LogUtil.java @@ -32,7 +32,10 @@ public synchronized String format(LogRecord logRecord) { consoleHandler.setFormatter(formatter); consoleHandler.setLevel(Level.CONFIG); - LOGGER.addHandler(consoleHandler); + // The actual write+flush to the console happens on a background thread, so per-call + // logging never blocks (or lock-contends on) the request-serving thread. See + // AsyncConsoleHandler's class doc for why this matters under concurrent load. + LOGGER.addHandler(new AsyncConsoleHandler(consoleHandler)); LOGGER.setLevel(logLevelToLoggerLevelMap(logLevel)); printInfoLog(InfoLogs.LOGGER_SETUP_DONE.getLog()); } diff --git a/skyvault/src/test/java/com/skyflow/ConnectionClientDotenvTests.java b/skyvault/src/test/java/com/skyflow/ConnectionClientDotenvTests.java index 4916f628..8ed0ff7c 100644 --- a/skyvault/src/test/java/com/skyflow/ConnectionClientDotenvTests.java +++ b/skyvault/src/test/java/com/skyflow/ConnectionClientDotenvTests.java @@ -3,6 +3,7 @@ import com.skyflow.config.ConnectionConfig; import com.skyflow.errors.ErrorMessage; import com.skyflow.errors.SkyflowException; +import com.skyflow.utils.BaseUtils; import com.skyflow.utils.Constants; import org.junit.After; import org.junit.Assert; @@ -31,6 +32,7 @@ public class ConnectionClientDotenvTests { public void saveEnvFileState() throws IOException { File f = new File(ENV_FILE); originalEnvContent = f.exists() ? Files.readAllBytes(Paths.get(ENV_FILE)) : null; + BaseUtils.resetDotenvCacheForTests(); // see its javadoc: .env is otherwise memoized JVM-wide } @After