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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions .fern/metadata.json
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
{
"cliVersion": "5.95.1",
"generatorName": "fernapi/fern-java-sdk",
"generatorVersion": "4.21.0",
"generatorVersion": "4.23.1",
"generatorConfig": {
"package-prefix": "com.deepgram",
"base-api-exception-class-name": "DeepgramHttpException",
Expand All @@ -12,8 +12,8 @@
"enable-wire-tests": true,
"runtime-version": true
},
"originGitCommit": "f6f5d3876c8996e2a6b6e59136bb222183a8e317",
"originGitCommit": "6c10a16dcc3b061e3030440e339f405af1ad29a9",
"originGitCommitIsDirty": true,
"invokedBy": "manual",
"sdkVersion": "0.10.2"
}
"sdkVersion": "0.10.3"
}
5 changes: 0 additions & 5 deletions .fernignore
Original file line number Diff line number Diff line change
Expand Up @@ -51,11 +51,6 @@ src/main/java/com/deepgram/resources/listen/v2/websocket/V2WebSocketClient.java
src/main/java/com/deepgram/resources/listen/v1/websocket/V1WebSocketClient.java
src/main/java/com/deepgram/resources/speak/v1/websocket/V1WebSocketClient.java

# Close is safe before connect in every generated WebSocket client. The listener is created only
# during connect(), so disconnect() must tolerate a null listener for try-with-resources clients.
# Remove these patches once Fern emits the null guard.
src/main/java/com/deepgram/resources/agent/v1/websocket/V1WebSocketClient.java

# Restores the FLUX_RENEE_EN constant that generator 4.18.0 dropped. The voice is live:
# POST /v2/speak?model=flux-renee-en returns 200 with valid audio, and the name resolves in the
# server's model registry (an invented flux-* name is rejected with INVALID_QUERY_PARAMETER), so the
Expand Down
5 changes: 3 additions & 2 deletions src/main/java/com/deepgram/AsyncDeepgramApiClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,9 @@ public AsyncVoiceAgentClient voiceAgent() {
}

/**
* Releases resources owned by this client. See {@code ClientOptions.close()} for what is
* and is not released.
* Releases resources owned by this client: any WebSocket clients still connected through
* it are disconnected first, then the SDK-owned HTTP client is shut down. See
* {@code ClientOptions.close()} for what is and is not released.
*/
@Override
public void close() {
Expand Down
5 changes: 3 additions & 2 deletions src/main/java/com/deepgram/DeepgramApiClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -79,8 +79,9 @@ public VoiceAgentClient voiceAgent() {
}

/**
* Releases resources owned by this client. See {@code ClientOptions.close()} for what is
* and is not released.
* Releases resources owned by this client: any WebSocket clients still connected through
* it are disconnected first, then the SDK-owned HTTP client is shut down. See
* {@code ClientOptions.close()} for what is and is not released.
*/
@Override
public void close() {
Expand Down
75 changes: 71 additions & 4 deletions src/main/java/com/deepgram/core/ClientOptions.java
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,14 @@
*/
package com.deepgram.core;

import java.util.ArrayList;
import java.util.HashMap;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Supplier;
import okhttp3.OkHttpClient;

Expand All @@ -21,6 +25,10 @@ public final class ClientOptions {

private final boolean ownsHttpClient;

private final AtomicBoolean closed;

private final Set<AutoCloseable> openWebSockets;

private final int timeout;

private final int maxRetries;
Expand All @@ -41,6 +49,8 @@ private ClientOptions(
Map<String, Supplier<String>> headerSuppliers,
OkHttpClient httpClient,
boolean ownsHttpClient,
AtomicBoolean closed,
Set<AutoCloseable> openWebSockets,
int timeout,
int maxRetries,
Optional<Long> initialRetryDelayMillis,
Expand All @@ -62,6 +72,8 @@ private ClientOptions(
this.headerSuppliers = headerSuppliers;
this.httpClient = httpClient;
this.ownsHttpClient = ownsHttpClient;
this.closed = closed;
this.openWebSockets = openWebSockets;
this.timeout = timeout;
this.maxRetries = maxRetries;
this.initialRetryDelayMillis = initialRetryDelayMillis;
Expand Down Expand Up @@ -132,19 +144,62 @@ public Optional<Double> retryJitterFactor() {
* OkHttpClient itself; an OkHttpClient supplied via httpClient is left running, since the
* caller owns its lifecycle.
* <p>
* In-flight calls are not cancelled or awaited. Later asynchronous requests fail when the
* dispatcher rejects work, while synchronous requests may still run. Options derived from this
* one via {@code Builder.from(...)} share the same dispatcher and connection pool, so closing
* either releases them for both. Calling this method more than once has no further effect.
* In-flight calls are not cancelled or awaited, and any request issued after this method
* returns fails with a {@code RejectedExecutionException}. Options derived from this one via
* {@code Builder.from(...)} share the same dispatcher and connection pool, so closing either
* releases them for both. Calling this method more than once has no further effect.
* <p>
* WebSocket clients created from this client that are still connected are disconnected
* first (whether or not the OkHttpClient is owned), so they stop reconnecting before the
* dispatcher goes away. Any WebSocket connect or reconnect attempted afterwards fails with
* an {@code IllegalStateException} explaining that the client has been closed.
*/
public void close() {
if (!this.closed.compareAndSet(false, true)) {
return;
}
for (AutoCloseable webSocket : new ArrayList<>(this.openWebSockets)) {
try {
webSocket.close();
} catch (Exception e) {
// best effort; keep closing the remaining sockets and the HTTP client
}
}
this.openWebSockets.clear();
if (!this.ownsHttpClient) {
return;
}
this.httpClient.dispatcher().executorService().shutdown();
this.httpClient.connectionPool().evictAll();
}

/**
* Returns whether {@link #close()} has been called on this client.
*/
public boolean isClosed() {
return this.closed.get();
}

/**
* Tracks a connected WebSocket client so that {@link #close()} disconnects it.
*
* @throws IllegalStateException if this client has already been closed
*/
public void registerWebSocket(AutoCloseable webSocket) {
this.openWebSockets.add(webSocket);
if (this.closed.get()) {
this.openWebSockets.remove(webSocket);
throw new IllegalStateException("root client has been closed");
}
}

/**
* Stops tracking a WebSocket client that has been disconnected.
*/
public void unregisterWebSocket(AutoCloseable webSocket) {
this.openWebSockets.remove(webSocket);
}

public Optional<WebSocketFactory> webSocketFactory() {
return this.webSocketFactory;
}
Expand Down Expand Up @@ -178,6 +233,10 @@ public static class Builder {

private boolean ownsHttpClient = true;

private AtomicBoolean closed;

private Set<AutoCloseable> openWebSockets;

private Optional<LogConfig> logging = Optional.empty();

private Optional<WebSocketFactory> webSocketFactory = Optional.empty();
Expand Down Expand Up @@ -303,12 +362,18 @@ public ClientOptions build() {
this.httpClient = httpClientBuilder.build();
this.timeout = Optional.of(httpClient.callTimeoutMillis() / 1000);

AtomicBoolean closedToUse = this.closed != null ? this.closed : new AtomicBoolean(false);
Set<AutoCloseable> openWebSocketsToUse =
this.openWebSockets != null ? this.openWebSockets : ConcurrentHashMap.newKeySet();

return new ClientOptions(
environment,
headers,
headerSuppliers,
httpClient,
this.ownsHttpClient,
closedToUse,
openWebSocketsToUse,
this.timeout.get(),
this.maxRetries,
this.initialRetryDelayMillis,
Expand All @@ -327,6 +392,8 @@ public static Builder from(ClientOptions clientOptions) {
builder.timeout = Optional.of(clientOptions.timeout(null));
builder.httpClient = clientOptions.httpClient();
builder.ownsHttpClient = clientOptions.ownsHttpClient;
builder.closed = clientOptions.closed;
builder.openWebSockets = clientOptions.openWebSockets;
builder.headers.putAll(clientOptions.headers);
builder.headerSuppliers.putAll(clientOptions.headerSuppliers);
builder.maxRetries = clientOptions.maxRetries();
Expand Down
22 changes: 22 additions & 0 deletions src/main/java/com/deepgram/core/OkHttpWebSocketFactory.java
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
*/
package com.deepgram.core;

import java.util.function.BooleanSupplier;
import okhttp3.OkHttpClient;
import okhttp3.Request;
import okhttp3.WebSocket;
Expand All @@ -15,17 +16,38 @@
public final class OkHttpWebSocketFactory implements WebSocketFactory {
private final OkHttpClient okHttpClient;

private final BooleanSupplier closedCheck;

/**
* Creates a new OkHttpWebSocketFactory with the specified OkHttpClient.
*
* @param okHttpClient The OkHttpClient instance to use for creating WebSockets
*/
public OkHttpWebSocketFactory(OkHttpClient okHttpClient) {
this(okHttpClient, () -> false);
}

/**
* Creates a new OkHttpWebSocketFactory that refuses to open WebSockets once the owning
* client has been closed.
*
* @param okHttpClient The OkHttpClient instance to use for creating WebSockets
* @param closedCheck Returns true once the owning client has been closed (e.g.
* {@code clientOptions::isClosed})
*/
public OkHttpWebSocketFactory(OkHttpClient okHttpClient, BooleanSupplier closedCheck) {
this.okHttpClient = okHttpClient;
this.closedCheck = closedCheck;
}

/**
* @throws IllegalStateException if the owning client has been closed
*/
@Override
public WebSocket create(Request request, WebSocketListener listener) {
if (closedCheck.getAsBoolean()) {
throw new IllegalStateException("root client has been closed");
}
return okHttpClient.newWebSocket(request, listener);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.BooleanSupplier;
import java.util.function.Consumer;
import java.util.function.Supplier;
import okhttp3.Response;
Expand Down Expand Up @@ -49,6 +50,8 @@ public abstract class ReconnectingWebSocketListener extends WebSocketListener {

private final Supplier<? extends WebSocket> connectionSupplier;

private final BooleanSupplier closedCheck;

/**
* Creates a new reconnecting WebSocket listener.
*
Expand All @@ -57,9 +60,26 @@ public abstract class ReconnectingWebSocketListener extends WebSocketListener {
*/
public ReconnectingWebSocketListener(
ReconnectingWebSocketListener.ReconnectOptions options, Supplier<? extends WebSocket> connectionSupplier) {
this(options, connectionSupplier, () -> false);
}

/**
* Creates a new reconnecting WebSocket listener.
*
* @param options Reconnection configuration options
* @param connectionSupplier Supplier that creates new WebSocket connections
* @param closedCheck Returns true once the owning client has been closed; connect and
* reconnect attempts made after that point fail with an {@link IllegalStateException}
* instead of being retried
*/
public ReconnectingWebSocketListener(
ReconnectingWebSocketListener.ReconnectOptions options,
Supplier<? extends WebSocket> connectionSupplier,
BooleanSupplier closedCheck) {
this.activeOptions = options;
this.maxEnqueuedMessages = options.maxEnqueuedMessages;
this.connectionSupplier = connectionSupplier;
this.closedCheck = closedCheck;
}

/**
Expand All @@ -85,11 +105,18 @@ public void applyOptionsOverride(ReconnectOptions options) {
* - TimeoutException: Includes retry attempt context
* - InterruptedException: Preserves thread interruption status
* - ExecutionException: Extracts actual cause and adds context
* - Owning client closed: reports an IllegalStateException once and stops reconnecting
*/
public void connect() {
if (!connectLock.compareAndSet(false, true)) {
return;
}
if (closedCheck.getAsBoolean()) {
shouldReconnect.set(false);
connectLock.set(false);
onWebSocketFailure(null, new IllegalStateException("root client has been closed"), null);
return;
}
ReconnectOptions options = this.activeOptions;
if (retryCount.get() > options.maxRetries) {
connectLock.set(false);
Expand Down Expand Up @@ -403,8 +430,13 @@ private long getNextDelay() {
/**
* Schedules a reconnection attempt with appropriate delay.
* Increments retry count and uses exponential backoff.
* Does nothing once the owning client has been closed.
*/
private void scheduleReconnect() {
if (closedCheck.getAsBoolean()) {
shouldReconnect.set(false);
return;
}
retryCount.incrementAndGet();
long delay = getNextDelay();
reconnectExecutor.schedule(this::connect, delay, MILLISECONDS);
Expand Down
Loading
Loading