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
12 changes: 7 additions & 5 deletions .fernignore
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,10 @@ src/main/java/com/deepgram/core/ClientOptions.java
# Transport abstraction (pluggable transport for SageMaker, etc.)
src/main/java/com/deepgram/core/transport/

# Bug fixes for maxRetries(0) semantics ("connect once, don't retry") and a
# configurable connectionTimeoutMs on ReconnectOptions (was hardcoded 4000ms).
# Pull this back out once the fixes are upstreamed into the Fern generator.
# Bug fixes for maxRetries(0) semantics ("connect once, don't retry"), a
# configurable connectionTimeoutMs on ReconnectOptions (was hardcoded 4000ms),
# and server-close acknowledgement. Pull this back out once the fixes are
# upstreamed into the Fern generator.
src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java

# Forward-compat patch: Fern's generated dispatcher routes any unrecognized message
Expand All @@ -30,8 +31,9 @@ src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java
# frame is still delivered via onMessage(String)); mirrors the JS/Python SDKs. Re-apply
# after regen. Unfreeze once the generator stops treating unknown frames as errors.
# NOTE: the listen v2 and speak v2 clients below also carry the streaming query-param patches
# (multi-value serialization + the additionalProperties escape hatch) documented lower down —
# keep them frozen until those fixes are upstreamed too.
# (multi-value serialization + the additionalProperties escape hatch) documented lower down.
# Listen v2 additionally treats a no-status close as terminal only when the closing socket accepted
# CloseStream. Keep them frozen until those fixes are upstreamed too.
src/main/java/com/deepgram/resources/speak/v2/websocket/V2WebSocketClient.java
src/main/java/com/deepgram/resources/listen/v2/websocket/V2WebSocketClient.java

Expand Down
4 changes: 2 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -48,8 +48,8 @@ How to identify:
Current temporarily frozen files:

- `src/main/java/com/deepgram/core/ClientOptions.java` - preserves release-please version markers and correct SDK header constants that Fern currently overwrites; use the standard `.bak` swap/restore workflow during regen review. Since generator 4.18.0 Fern emits a `getSdkVersion()` helper reading `Package.getImplementationVersion()` instead of a literal. That *does* resolve in the published artifact (CI publishes via `mvn deploy -P release`, and `pom.xml`'s maven-jar-plugin sets `addDefaultImplementationEntries=true`, so the JAR manifest carries `Implementation-Version`), but it resolves to `null` under Gradle and in tests, where it silently falls back to a hardcoded literal the generator does not keep current. We keep the explicit literals because they are correct in every context and because `.github/release-please-config.json` already lists this file in `extra-files`, so release-please bumps it alongside `pom.xml`, `build.gradle`, and `.fern/metadata.json`. Fern also emits `User-Agent` with a `com.deepgram.` prefix while leaving `X-Fern-SDK-Name` on the `com.deepgram:` Maven-coordinate form; we keep both on the colon form.
- `src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java` - carries bug fixes for `maxRetries(0)` semantics ("connect once, don't retry") and a configurable `connectionTimeoutMs` field (was hardcoded 4000ms), plus an `applyOptionsOverride(...)` hook used by `TransportWebSocketFactory` to apply per-transport reconnect policy; pull this back out once the fixes are upstreamed into the Fern generator. Use the standard `.bak` swap/restore workflow during regen review.
- `src/main/java/com/deepgram/resources/speak/v2/websocket/V2WebSocketClient.java` and `src/main/java/com/deepgram/resources/listen/v2/websocket/V2WebSocketClient.java` - forward-compat patch (both clients). Fern's generated `handleIncomingMessage` dispatcher routes any unrecognized message type to `onError` with "Update your SDK version...", which makes a benign new server control frame look fatal to a deployed client. Patched so the unrecognized-type branch is a no-op — the raw frame is already delivered via `onMessage(String)` earlier in the method, so consumers still see it. Mirrors the JS/Python SDKs' forward-compat behavior and is regression-guarded by `src/test/java/com/deepgram/SpeakV2ForwardCompatTest.java` and `src/test/java/com/deepgram/ListenV2ForwardCompatTest.java`. These two clients also carry the streaming query-param patches described in the next entry. Use the standard `.bak` swap/restore workflow during regen review; re-apply the no-op to both after regen, and unfreeze once the generator stops treating unknown frames as errors.
- `src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java` - carries bug fixes for `maxRetries(0)` semantics ("connect once, don't retry"), a configurable `connectionTimeoutMs` field (was hardcoded 4000ms), server-initiated close acknowledgement, and an `applyOptionsOverride(...)` hook used by `TransportWebSocketFactory` to apply per-transport reconnect policy; pull this back out once the fixes are upstreamed into the Fern generator. Use the standard `.bak` swap/restore workflow during regen review.
- `src/main/java/com/deepgram/resources/speak/v2/websocket/V2WebSocketClient.java` and `src/main/java/com/deepgram/resources/listen/v2/websocket/V2WebSocketClient.java` - forward-compat patch (both clients). Fern's generated `handleIncomingMessage` dispatcher routes any unrecognized message type to `onError` with "Update your SDK version...", which makes a benign new server control frame look fatal to a deployed client. Patched so the unrecognized-type branch is a no-op — the raw frame is already delivered via `onMessage(String)` earlier in the method, so consumers still see it. Mirrors the JS/Python SDKs' forward-compat behavior and is regression-guarded by `src/test/java/com/deepgram/SpeakV2ForwardCompatTest.java` and `src/test/java/com/deepgram/ListenV2ForwardCompatTest.java`. These two clients also carry the streaming query-param patches described in the next entry; Listen v2 additionally treats a no-status close as terminal only when the closing socket accepted `sendCloseStream(...)`. Use the standard `.bak` swap/restore workflow during regen review; re-apply the no-op to both after regen, and unfreeze once the generator stops treating unknown frames as errors.
- `src/main/java/com/deepgram/resources/listen/v1/websocket/V1WebSocketClient.java` and `src/main/java/com/deepgram/resources/speak/v1/websocket/V1WebSocketClient.java` (and the v2 clients above) - streaming query-param patches on the generated `connect()` builders. Two fixes: (1) multi-value serialization — array-valued params (listen: `keyterm`, `keywords`, `replace`, `search`, `tag`, `extra`, `language_hint`; speak: `tag`) were serialized with `String.valueOf(union.get())`, collapsing a `List` into one param (`keyterm=[a, b]`) instead of repeats (`keyterm=a&keyterm=b`); (2) an `additionalProperties` escape hatch — the builder exposes `additionalProperty(key, value)` for unmodeled params (e.g. `no_delay`) but `connect()` never emitted them to the URL. Both patched to route through `QueryStringMapper(arraysAsRepeats=true)`, matching the REST path. Use the standard `.bak` swap/restore workflow during regen review; re-apply after regen and unfreeze once the generator emits array params as repeats and serializes `additionalProperties` on the WS `connect()` path (tracked as an upstream Fern request).
- Fields-less message types carrying a manual `hashCode()` patch (Fern generates `equals()` but no `hashCode()` for these, violating the Object contract): `src/main/java/com/deepgram/resources/listen/v2/types/ListenV2CloseStream.java`, `src/main/java/com/deepgram/resources/listen/v2/types/ListenV2ForceEndTurn.java`, `src/main/java/com/deepgram/resources/speak/v2/types/SpeakV2Close.java`, `src/main/java/com/deepgram/resources/speak/v2/types/SpeakV2Flush.java`, and the `AgentV1*` event types `src/main/java/com/deepgram/resources/agent/v1/types/{AgentV1ListenUpdated,AgentV1SpeakUpdated,AgentV1AgentAudioDone,AgentV1SettingsApplied,AgentV1UserStartedSpeaking,AgentV1KeepAlive,AgentV1ThinkUpdated,AgentV1PromptUpdated,AgentV1ForceEndTurn}.java`. Use the standard `.bak` swap/restore workflow during regen review; drop the patches and unfreeze all of them once the generator emits a matching equals/hashCode pair for fields-less types (tracked as an upstream Fern request).
- `src/main/java/com/deepgram/types/DeepgramModel.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 removal is a spec regression rather than a retirement, and dropping the constant would break 0.8.0 callers for nothing. Five touchpoints: the constant, the `Value` enum entry, the `visit()` case, the `valueOf()` case, and the `Visitor` method. **This file is unlike the other temporarily frozen ones — it receives frequent additive spec changes (4.18.0 alone added 25 constants), so on the next regen do NOT restore the `.bak` wholesale.** Diff the `.bak` against the newly generated file, carry forward every new voice, and re-apply only the `FLUX_RENEE_EN` touchpoints. Drop the patch and unfreeze once the spec lists the voice again (tracked as an upstream spec request).
Expand Down
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.Consumer;
import java.util.function.Supplier;
import okhttp3.Response;
import okhttp3.WebSocket;
Expand Down Expand Up @@ -202,9 +203,20 @@ public void disconnect() {
* @return true if sent immediately, false if queued or dropped
*/
public synchronized boolean send(String message) {
return send(message, null);
}

/**
* Sends a message and exposes the accepting socket to the caller. The callback is invoked
* only when the message was sent directly rather than queued or dropped.
*/
public synchronized boolean send(String message, Consumer<WebSocket> onSent) {
WebSocket ws = webSocket;
if (ws != null) {
boolean sent = ws.send(message);
if (sent && onSent != null) {
onSent.accept(ws);
}
if (!sent && messageQueue.size() < maxEnqueuedMessages) {
messageQueue.offer(message);
return false;
Expand Down Expand Up @@ -281,6 +293,17 @@ public void onMessage(WebSocket webSocket, ByteString bytes) {
onWebSocketBinaryMessage(webSocket, bytes);
}

/**
* Acknowledge a peer-initiated close so OkHttp can complete the close handshake and invoke
* {@link #onClosed(WebSocket, int, String)}.
*/
@Override
public void onClosing(WebSocket webSocket, int code, String reason) {
// 1005 is a local no-status sentinel, not a valid close frame code. Acknowledge it with
// the normal closure code so OkHttp can finish the handshake instead of throwing.
webSocket.close(code == 1005 ? 1000 : code, reason);
}

/**
* @hidden
*/
Expand Down Expand Up @@ -328,11 +351,19 @@ public void onClosed(WebSocket webSocket, int code, String reason) {
}
connectionEstablishedTime = 0L;
onWebSocketClosed(webSocket, code, reason);
if (code != 1000 && shouldReconnect.get()) {
if (shouldReconnect.get() && shouldReconnectAfterClose(webSocket, code)) {
scheduleReconnect();
}
}

/**
* Returns whether a close status should reconnect. Resource-specific listeners can override
* this when a protocol operation establishes that a particular close is terminal.
*/
protected boolean shouldReconnectAfterClose(WebSocket webSocket, int code) {
return code != 1000;
}

/**
* Calculates the next reconnection delay using exponential backoff.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import com.fasterxml.jackson.databind.ObjectMapper;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Consumer;
import okhttp3.HttpUrl;
import okhttp3.OkHttpClient;
Expand Down Expand Up @@ -61,6 +62,9 @@ public class V2WebSocketClient implements AutoCloseable {

private ReconnectingWebSocketListener reconnectingListener;

// A no-status close is terminal only after this socket accepted the Flux CloseStream frame.
private final AtomicReference<WebSocket> closeStreamSocket = new AtomicReference<>();

private volatile Consumer<ListenV2Connected> connectedHandler;

private volatile Consumer<ListenV2TurnInfo> turnInfoHandler;
Expand All @@ -86,6 +90,7 @@ public V2WebSocketClient(ClientOptions clientOptions) {
* @param options connection options including query parameters
*/
public CompletableFuture<Void> connect(V2ConnectOptions options) {
closeStreamSocket.set(null);
connectionFuture = new CompletableFuture<>();
String baseUrl = clientOptions.environment().getProductionURL();
String fullPath = "/v2/listen";
Expand Down Expand Up @@ -181,6 +186,7 @@ public CompletableFuture<Void> connect(V2ConnectOptions options) {
}) {
@Override
protected void onWebSocketOpen(WebSocket webSocket, Response response) {
closeStreamSocket.set(null);
readyState = WebSocketReadyState.OPEN;
if (onConnectedHandler != null) {
onConnectedHandler.run();
Expand Down Expand Up @@ -212,6 +218,11 @@ protected void onWebSocketClosed(WebSocket webSocket, int code, String reason) {
onDisconnectedHandler.accept(new DisconnectReason(code, reason));
}
}

@Override
protected boolean shouldReconnectAfterClose(WebSocket webSocket, int code) {
return code != 1005 || closeStreamSocket.get() != webSocket;
}
};
reconnectingListener.connect();
return connectionFuture;
Expand Down Expand Up @@ -263,7 +274,7 @@ public CompletableFuture<Void> sendMedia(ByteString message) {
* @return a CompletableFuture that completes when the message is sent
*/
public CompletableFuture<Void> sendCloseStream(ListenV2CloseStream message) {
return sendMessage(message);
return sendMessage(message, closeStreamSocket::set);
}

/**
Expand Down Expand Up @@ -391,12 +402,16 @@ private void assertSocketIsOpen() {
}

private CompletableFuture<Void> sendMessage(Object body) {
return sendMessage(body, null);
}

private CompletableFuture<Void> sendMessage(Object body, Consumer<WebSocket> onSent) {
CompletableFuture<Void> future = new CompletableFuture<>();
try {
assertSocketIsOpen();
String json = objectMapper.writeValueAsString(body);
// Use reconnecting listener's send method which handles queuing
reconnectingListener.send(json);
reconnectingListener.send(json, onSent);
future.complete(null);
} catch (IllegalStateException e) {
future.completeExceptionally(e);
Expand Down
Loading
Loading