diff --git a/.fernignore b/.fernignore index 4734b1ba..cfc1da8a 100644 --- a/.fernignore +++ b/.fernignore @@ -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 @@ -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 diff --git a/AGENTS.md b/AGENTS.md index 39cb6388..2ca17e57 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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). diff --git a/src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java b/src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java index a5330688..67eaeb3b 100644 --- a/src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java +++ b/src/main/java/com/deepgram/core/ReconnectingWebSocketListener.java @@ -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; @@ -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 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; @@ -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 */ @@ -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. * diff --git a/src/main/java/com/deepgram/resources/listen/v2/websocket/V2WebSocketClient.java b/src/main/java/com/deepgram/resources/listen/v2/websocket/V2WebSocketClient.java index 445a96aa..91f4794e 100644 --- a/src/main/java/com/deepgram/resources/listen/v2/websocket/V2WebSocketClient.java +++ b/src/main/java/com/deepgram/resources/listen/v2/websocket/V2WebSocketClient.java @@ -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; @@ -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 closeStreamSocket = new AtomicReference<>(); + private volatile Consumer connectedHandler; private volatile Consumer turnInfoHandler; @@ -86,6 +90,7 @@ public V2WebSocketClient(ClientOptions clientOptions) { * @param options connection options including query parameters */ public CompletableFuture connect(V2ConnectOptions options) { + closeStreamSocket.set(null); connectionFuture = new CompletableFuture<>(); String baseUrl = clientOptions.environment().getProductionURL(); String fullPath = "/v2/listen"; @@ -181,6 +186,7 @@ public CompletableFuture connect(V2ConnectOptions options) { }) { @Override protected void onWebSocketOpen(WebSocket webSocket, Response response) { + closeStreamSocket.set(null); readyState = WebSocketReadyState.OPEN; if (onConnectedHandler != null) { onConnectedHandler.run(); @@ -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; @@ -263,7 +274,7 @@ public CompletableFuture sendMedia(ByteString message) { * @return a CompletableFuture that completes when the message is sent */ public CompletableFuture sendCloseStream(ListenV2CloseStream message) { - return sendMessage(message); + return sendMessage(message, closeStreamSocket::set); } /** @@ -391,12 +402,16 @@ private void assertSocketIsOpen() { } private CompletableFuture sendMessage(Object body) { + return sendMessage(body, null); + } + + private CompletableFuture sendMessage(Object body, Consumer onSent) { CompletableFuture 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); diff --git a/src/test/java/com/deepgram/ListenV2ServerCloseTest.java b/src/test/java/com/deepgram/ListenV2ServerCloseTest.java new file mode 100644 index 00000000..ad6fb55f --- /dev/null +++ b/src/test/java/com/deepgram/ListenV2ServerCloseTest.java @@ -0,0 +1,212 @@ +package com.deepgram; + +import static org.assertj.core.api.Assertions.assertThat; + +import com.deepgram.core.Environment; +import com.deepgram.core.ReconnectingWebSocketListener; +import com.deepgram.core.WebSocketFactory; +import com.deepgram.core.WebSocketReadyState; +import com.deepgram.resources.listen.v2.types.ListenV2CloseStream; +import com.deepgram.resources.listen.v2.websocket.V2ConnectOptions; +import com.deepgram.resources.listen.v2.websocket.V2WebSocketClient; +import com.deepgram.types.ListenV2Model; +import java.lang.reflect.Field; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import okhttp3.Request; +import okhttp3.WebSocket; +import okhttp3.WebSocketListener; +import okhttp3.mockwebserver.MockResponse; +import okhttp3.mockwebserver.MockWebServer; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +/** Regression coverage for a server-initiated streaming close. */ +class ListenV2ServerCloseTest { + private MockWebServer server; + private DeepgramClient client; + + @BeforeEach + void setUp() throws Exception { + server = new MockWebServer(); + server.start(); + String base = server.url("/").toString().replaceAll("/$", ""); + Environment env = Environment.custom() + .base(base) + .production(base) + .agent(base) + .agentRest(base) + .build(); + client = DeepgramClient.builder().apiKey("test").environment(env).build(); + } + + @AfterEach + void tearDown() throws Exception { + server.shutdown(); + } + + @Test + @DisplayName("a server close notifies the client without an explicit disconnect") + void serverInitiatedCloseNotifiesClient() throws Exception { + server.enqueue(new MockResponse().withWebSocketUpgrade(new WebSocketListener() { + @Override + public void onOpen(WebSocket webSocket, okhttp3.Response response) { + webSocket.close(1000, "stream complete"); + } + })); + V2WebSocketClient ws = client.listen().v2().v2WebSocket(); + CountDownLatch disconnected = new CountDownLatch(1); + AtomicInteger disconnectCount = new AtomicInteger(); + ws.onDisconnected(reason -> { + disconnectCount.incrementAndGet(); + disconnected.countDown(); + }); + + try { + ws.connect(V2ConnectOptions.builder() + .model(ListenV2Model.FLUX_GENERAL_EN) + .build()) + .get(5, TimeUnit.SECONDS); + + assertThat(disconnected.await(2, TimeUnit.SECONDS)).isTrue(); + assertThat(disconnectCount).hasValue(1); + assertThat(ws.getReadyState()).isEqualTo(WebSocketReadyState.CLOSED); + } finally { + ws.disconnect(); + } + } + + @Test + @DisplayName("CloseStream makes a following no-status close terminal") + void closeStreamNoStatusCloseDoesNotReconnect() throws Exception { + server.enqueue(new MockResponse().withWebSocketUpgrade(new WebSocketListener() {})); + V2WebSocketClient ws = client.listen().v2().v2WebSocket(); + CountDownLatch disconnected = new CountDownLatch(1); + ws.reconnectOptions(ReconnectingWebSocketListener.ReconnectOptions.builder() + .minReconnectionDelayMs(10) + .maxReconnectionDelayMs(10) + .build()); + ws.onDisconnected(reason -> disconnected.countDown()); + + try { + ws.connect(V2ConnectOptions.builder() + .model(ListenV2Model.FLUX_GENERAL_EN) + .build()) + .get(5, TimeUnit.SECONDS); + assertThat(server.takeRequest(5, TimeUnit.SECONDS)).isNotNull(); + + ws.sendCloseStream(ListenV2CloseStream.builder().build()).get(5, TimeUnit.SECONDS); + ReconnectingWebSocketListener listener = getListener(ws); + listener.onClosed(listener.getWebSocket(), 1005, ""); + + assertThat(disconnected.await(2, TimeUnit.SECONDS)).isTrue(); + assertThat(server.takeRequest(100, TimeUnit.MILLISECONDS)).isNull(); + } finally { + ws.disconnect(); + } + } + + @Test + @DisplayName("a rejected CloseStream does not make a no-status close terminal") + void rejectedCloseStreamDoesNotSuppressReconnect() throws Exception { + AtomicInteger connectionCount = new AtomicInteger(); + RejectingWebSocket webSocket = new RejectingWebSocket(); + WebSocketFactory factory = (request, listener) -> { + connectionCount.incrementAndGet(); + listener.onOpen(webSocket, null); + return webSocket; + }; + V2WebSocketClient ws = new V2WebSocketClient(com.deepgram.core.ClientOptions.builder() + .environment(Environment.PRODUCTION) + .webSocketFactory(factory) + .build()); + ws.reconnectOptions(ReconnectingWebSocketListener.ReconnectOptions.builder() + .minReconnectionDelayMs(10) + .maxReconnectionDelayMs(10) + .build()); + + try { + ws.connect(V2ConnectOptions.builder() + .model(ListenV2Model.FLUX_GENERAL_EN) + .build()) + .get(5, TimeUnit.SECONDS); + ws.sendCloseStream(ListenV2CloseStream.builder().build()).get(5, TimeUnit.SECONDS); + + ReconnectingWebSocketListener listener = getListener(ws); + listener.onClosed(webSocket, 1005, ""); + + Thread.sleep(100); + assertThat(connectionCount).hasValue(2); + } finally { + ws.disconnect(); + } + } + + @Test + @DisplayName("a new connection clears CloseStream terminal state") + void newConnectionClearsCloseStreamTerminalState() throws Exception { + server.enqueue(new MockResponse().withWebSocketUpgrade(new WebSocketListener() {})); + V2WebSocketClient ws = client.listen().v2().v2WebSocket(); + ws.reconnectOptions(ReconnectingWebSocketListener.ReconnectOptions.builder() + .minReconnectionDelayMs(10) + .maxReconnectionDelayMs(10) + .build()); + + try { + ws.connect(V2ConnectOptions.builder() + .model(ListenV2Model.FLUX_GENERAL_EN) + .build()) + .get(5, TimeUnit.SECONDS); + assertThat(server.takeRequest(5, TimeUnit.SECONDS)).isNotNull(); + + ws.sendCloseStream(ListenV2CloseStream.builder().build()).get(5, TimeUnit.SECONDS); + ReconnectingWebSocketListener listener = getListener(ws); + WebSocket webSocket = listener.getWebSocket(); + listener.onOpen(webSocket, null); + listener.onClosed(webSocket, 1005, ""); + + assertThat(server.takeRequest(100, TimeUnit.MILLISECONDS)).isNotNull(); + } finally { + ws.disconnect(); + } + } + + private ReconnectingWebSocketListener getListener(V2WebSocketClient ws) throws Exception { + Field field = V2WebSocketClient.class.getDeclaredField("reconnectingListener"); + field.setAccessible(true); + return (ReconnectingWebSocketListener) field.get(ws); + } + + private static final class RejectingWebSocket implements WebSocket { + @Override + public Request request() { + return new Request.Builder().url("ws://localhost/").build(); + } + + @Override + public long queueSize() { + return 0; + } + + @Override + public boolean send(String text) { + return false; + } + + @Override + public boolean send(okio.ByteString bytes) { + return false; + } + + @Override + public boolean close(int code, String reason) { + return true; + } + + @Override + public void cancel() {} + } +} diff --git a/src/test/java/com/deepgram/core/ReconnectingWebSocketListenerTest.java b/src/test/java/com/deepgram/core/ReconnectingWebSocketListenerTest.java index da8001e0..0fa0ed7f 100644 --- a/src/test/java/com/deepgram/core/ReconnectingWebSocketListenerTest.java +++ b/src/test/java/com/deepgram/core/ReconnectingWebSocketListenerTest.java @@ -39,6 +39,8 @@ public WebSocket get() { } private static final class FakeWebSocket implements WebSocket { + int closeCode = -1; + @Override public okhttp3.Request request() { return new okhttp3.Request.Builder().url("ws://localhost/").build(); @@ -61,6 +63,7 @@ public boolean send(ByteString bytes) { @Override public boolean close(int code, String reason) { + closeCode = code; return true; } @@ -149,6 +152,42 @@ void initialAttemptProceedsWhenMaxRetriesIsZero() { } } + @Nested + @DisplayName("server close handling") + class ServerCloseTests { + @Test + @DisplayName("reconnects after a no-status close without protocol context") + void noStatusCloseReconnectsWithoutProtocolContext() throws Exception { + CountingSupplier supplier = new CountingSupplier(false); + ReconnectOptions opts = ReconnectOptions.builder() + .minReconnectionDelayMs(10) + .maxReconnectionDelayMs(10) + .build(); + TestListener listener = new TestListener(opts, supplier); + + try { + listener.onClosed(new FakeWebSocket(), 1005, ""); + + Thread.sleep(100); + assertThat(supplier.calls.get()).isEqualTo(1); + } finally { + listener.disconnect(); + } + } + + @Test + @DisplayName("acknowledges a no-status close with a valid normal close code") + void noStatusCloseIsAcknowledgedWithNormalCloseCode() { + TestListener listener = new TestListener(ReconnectOptions.builder().build(), new CountingSupplier(false)); + FakeWebSocket webSocket = new FakeWebSocket(); + + listener.onClosing(webSocket, 1005, ""); + + assertThat(webSocket.closeCode).isEqualTo(1000); + listener.disconnect(); + } + } + @Nested @DisplayName("applyOptionsOverride") class ApplyOverrideTests {